refactor: change Tasks to Outbox
notice that moving the blockstore fetch into the manager removes the weird error handling case. License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>
Brian Tiger Chow committed
Dec 16, 2014 at 20:27 UTC
2ea8ed81ac5d4cd0f99a6cbf3a6f33e43fbf4cf4
2 files changed
+21
-18
exchange/bitswap/bitswap.go
+2
-12
@@ -239,18 +239,8 @@ func (bs *bitswap) taskWorker(ctx context.Context) {
239
select {
240
case <-ctx.Done():
241
return
242
- case task := <-bs.ledgermanager.GetTaskChan():
243
- block, err := bs.blockstore.Get(task.Key)
244
- if err != nil {
245
- log.Errorf("Expected to have block %s, but it was not found!", task.Key)
246
- continue
247
- }
248
-
249
- message := bsmsg.New()
250
- message.AddBlock(block)
251
- // TODO: maybe add keys from our wantlist?
252
-
253
- bs.send(ctx, task.Target, message)
242
+ case envelope := <-bs.ledgermanager.Outbox():
243
+ bs.send(ctx, envelope.Peer, envelope.Message)
244
}
245
}
246
}
exchange/bitswap/strategy/ledgermanager.go
+19
-6
@@ -20,12 +20,17 @@ type ledgerMap map[peerKey]*ledger
20
// FIXME share this externally
21
type peerKey u.Key
22
23
+type Envelope struct {
24
+ Peer peer.Peer
25
+ Message bsmsg.BitSwapMessage
26
+}
27
+
28
type LedgerManager struct {
29
lock sync.RWMutex
30
ledgerMap ledgerMap
31
bs bstore.Blockstore
32
tasklist *TaskList
28
- taskOut chan *Task
33
+ outbox chan Envelope
34
workSignal chan struct{}
35
}
36
@@ -34,7 +39,7 @@ func NewLedgerManager(bs bstore.Blockstore, ctx context.Context) *LedgerManager
39
ledgerMap: make(ledgerMap),
40
bs: bs,
41
tasklist: NewTaskList(),
37
- taskOut: make(chan *Task, 4),
42
+ outbox: make(chan Envelope, 4), // TODO extract constant
43
workSignal: make(chan struct{}),
44
}
45
go lm.taskWorker(ctx)
@@ -54,17 +59,25 @@ func (lm *LedgerManager) taskWorker(ctx context.Context) {
59
}
60
continue
61
}
57
-
62
+ block, err := lm.bs.Get(nextTask.Key)
63
+ if err != nil {
64
+ continue // TODO maybe return an error
65
+ }
66
+ // construct message here so we can make decisions about any additional
67
+ // information we may want to include at this time.
68
+ m := bsmsg.New()
69
+ m.AddBlock(block)
70
+ // TODO: maybe add keys from our wantlist?
71
select {
72
case <-ctx.Done():
73
return
61
- case lm.taskOut <- nextTask:
74
+ case lm.outbox <- Envelope{Peer: nextTask.Target, Message: m}:
75
}
76
}
77
}
78
66
-func (lm *LedgerManager) GetTaskChan() <-chan *Task {
67
- return lm.taskOut
79
+func (lm *LedgerManager) Outbox() <-chan Envelope {
80
+ return lm.outbox
81
}
82
83
// Returns a slice of Peers with whom the local node has active sessions