@cryptotaxi247 / kubo / commits / 198aa1959

it's not a queue yet but it's okay to name it as such

License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>

Brian Tiger Chow committed Dec 16, 2014 at 21:36 UTC 198aa1959ae721ec646e0f0ad317525703a69068
2 files changed +15 -15
exchange/bitswap/strategy/ledgermanager.go
+9 -9
@@ -26,9 +26,9 @@ type LedgerManager struct {
26 lock sync.RWMutex
27 ledgerMap ledgerMap
28 bs bstore.Blockstore
29 - // FIXME tasklist isn't threadsafe nor is it protected by a mutex. consider
30 - // a way to avoid sharing the tasklist between the worker and the receiver
31 - tasklist *taskList
29 + // FIXME taskqueue isn't threadsafe nor is it protected by a mutex. consider
30 + // a way to avoid sharing the taskqueue between the worker and the receiver
31 + taskqueue *taskQueue
32 outbox chan Envelope
33 workSignal chan struct{}
34 }
@@ -37,7 +37,7 @@ func NewLedgerManager(ctx context.Context, bs bstore.Blockstore) *LedgerManager
37 lm := &LedgerManager{
38 ledgerMap: make(ledgerMap),
39 bs: bs,
40 - tasklist: newTaskList(),
40 + taskqueue: newTaskQueue(),
41 outbox: make(chan Envelope, 4), // TODO extract constant
42 workSignal: make(chan struct{}),
43 }
@@ -47,7 +47,7 @@ func NewLedgerManager(ctx context.Context, bs bstore.Blockstore) *LedgerManager
47
48 func (lm *LedgerManager) taskWorker(ctx context.Context) {
49 for {
50 - nextTask := lm.tasklist.Pop()
50 + nextTask := lm.taskqueue.Pop()
51 if nextTask == nil {
52 // No tasks in the list?
53 // Wait until there are!
@@ -124,11 +124,11 @@ func (lm *LedgerManager) MessageReceived(p peer.Peer, m bsmsg.BitSwapMessage) er
124 for _, e := range m.Wantlist() {
125 if e.Cancel {
126 l.CancelWant(e.Key)
127 - lm.tasklist.Cancel(e.Key, p)
127 + lm.taskqueue.Cancel(e.Key, p)
128 } else {
129 l.Wants(e.Key, e.Priority)
130 newWorkExists = true
131 - lm.tasklist.Push(e.Key, e.Priority, p)
131 + lm.taskqueue.Push(e.Key, e.Priority, p)
132 }
133 }
134
@@ -138,7 +138,7 @@ func (lm *LedgerManager) MessageReceived(p peer.Peer, m bsmsg.BitSwapMessage) er
138 for _, l := range lm.ledgerMap {
139 if l.WantListContains(block.Key()) {
140 newWorkExists = true
141 - lm.tasklist.Push(block.Key(), 1, l.Partner)
141 + lm.taskqueue.Push(block.Key(), 1, l.Partner)
142 }
143 }
144 }
@@ -159,7 +159,7 @@ func (lm *LedgerManager) MessageSent(p peer.Peer, m bsmsg.BitSwapMessage) error
159 for _, block := range m.Blocks() {
160 l.SentBytes(len(block.Data))
161 l.wantList.Remove(block.Key())
162 - lm.tasklist.Cancel(block.Key(), p)
162 + lm.taskqueue.Cancel(block.Key(), p)
163 }
164
165 return nil
exchange/bitswap/strategy/taskqueue.go renamed
+6 -6
@@ -8,13 +8,13 @@ import (
8 // TODO: at some point, the strategy needs to plug in here
9 // to help decide how to sort tasks (on add) and how to select
10 // tasks (on getnext). For now, we are assuming a dumb/nice strategy.
11 -type taskList struct {
11 +type taskQueue struct {
12 tasks []*Task
13 taskmap map[string]*Task
14 }
15
16 -func newTaskList() *taskList {
17 - return &taskList{
16 +func newTaskQueue() *taskQueue {
17 + return &taskQueue{
18 taskmap: make(map[string]*Task),
19 }
20 }
@@ -27,7 +27,7 @@ type Task struct {
27
28 // Push currently adds a new task to the end of the list
29 // TODO: make this into a priority queue
30 -func (tl *taskList) Push(block u.Key, priority int, to peer.Peer) {
30 +func (tl *taskQueue) Push(block u.Key, priority int, to peer.Peer) {
31 if task, ok := tl.taskmap[taskKey(to, block)]; ok {
32 // TODO: when priority queue is implemented,
33 // rearrange this Task
@@ -44,7 +44,7 @@ func (tl *taskList) Push(block u.Key, priority int, to peer.Peer) {
44 }
45
46 // Pop 'pops' the next task to be performed. Returns nil no task exists.
47 -func (tl *taskList) Pop() *Task {
47 +func (tl *taskQueue) Pop() *Task {
48 var out *Task
49 for len(tl.tasks) > 0 {
50 // TODO: instead of zero, use exponential distribution
@@ -63,7 +63,7 @@ func (tl *taskList) Pop() *Task {
63 }
64
65 // Cancel lazily cancels the sending of a block to a given peer
66 -func (tl *taskList) Cancel(k u.Key, p peer.Peer) {
66 +func (tl *taskQueue) Cancel(k u.Key, p peer.Peer) {
67 t, ok := tl.taskmap[taskKey(p, k)]
68 if ok {
69 t.theirPriority = -1