@cryptotaxi247 / kubo / commits / 5bd0b9546

rename to strategy.LedgerManager to decision.Engine

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

Brian Tiger Chow committed Dec 16, 2014 at 22:52 UTC 5bd0b9546288947d2d400827d11310e71af73375
7 files changed +81 -86
exchange/bitswap/bitswap.go
+8 -13
@@ -7,14 +7,13 @@ import (
7 "time"
8
9 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
10 -
10 blocks "github.com/jbenet/go-ipfs/blocks"
11 blockstore "github.com/jbenet/go-ipfs/blocks/blockstore"
12 exchange "github.com/jbenet/go-ipfs/exchange"
13 + decision "github.com/jbenet/go-ipfs/exchange/bitswap/decision"
14 bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
15 bsnet "github.com/jbenet/go-ipfs/exchange/bitswap/network"
16 notifications "github.com/jbenet/go-ipfs/exchange/bitswap/notifications"
17 - strategy "github.com/jbenet/go-ipfs/exchange/bitswap/strategy"
17 wantlist "github.com/jbenet/go-ipfs/exchange/bitswap/wantlist"
18 peer "github.com/jbenet/go-ipfs/peer"
19 u "github.com/jbenet/go-ipfs/util"
@@ -56,7 +55,7 @@ func New(parent context.Context, p peer.Peer, network bsnet.BitSwapNetwork, rout
55 blockstore: bstore,
56 cancelFunc: cancelFunc,
57 notifications: notif,
59 - ledgermanager: strategy.NewLedgerManager(ctx, bstore),
58 + engine: decision.NewEngine(ctx, bstore),
59 routing: routing,
60 sender: network,
61 wantlist: wantlist.NewThreadSafe(),
@@ -89,11 +88,7 @@ type bitswap struct {
88 // have more than a single block in the set
89 batchRequests chan []u.Key
90
92 - // strategy makes decisions about how to interact with partners.
93 - // TODO: strategy commented out until we have a use for it again
94 - //strategy strategy.Strategy
95 -
96 - ledgermanager *strategy.LedgerManager
91 + engine *decision.Engine
92
93 wantlist *wantlist.ThreadSafe
94
@@ -196,7 +191,7 @@ func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) e
191 // FIXME ensure accounting is handled correctly when
192 // communication fails. May require slightly different API to
193 // get better guarantees. May need shared sequence numbers.
199 - bs.ledgermanager.MessageSent(p, message)
194 + bs.engine.MessageSent(p, message)
195 }(peerToQuery)
196 }
197 wg.Wait()
@@ -239,7 +234,7 @@ func (bs *bitswap) taskWorker(ctx context.Context) {
234 select {
235 case <-ctx.Done():
236 return
242 - case envelope := <-bs.ledgermanager.Outbox():
237 + case envelope := <-bs.engine.Outbox():
238 bs.send(ctx, envelope.Peer, envelope.Message)
239 }
240 }
@@ -305,7 +300,7 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
300
301 // This call records changes to wantlists, blocks received,
302 // and number of bytes transfered.
308 - bs.ledgermanager.MessageReceived(p, incoming)
303 + bs.engine.MessageReceived(p, incoming)
304 // TODO: this is bad, and could be easily abused.
305 // Should only track *useful* messages in ledger
306
@@ -334,7 +329,7 @@ func (bs *bitswap) cancelBlocks(ctx context.Context, bkeys []u.Key) {
329 for _, k := range bkeys {
330 message.Cancel(k)
331 }
337 - for _, p := range bs.ledgermanager.Peers() {
332 + for _, p := range bs.engine.Peers() {
333 err := bs.send(ctx, p, message)
334 if err != nil {
335 log.Errorf("Error sending message: %s", err)
@@ -354,7 +349,7 @@ func (bs *bitswap) send(ctx context.Context, p peer.Peer, m bsmsg.BitSwapMessage
349 if err := bs.sender.SendMessage(ctx, p, m); err != nil {
350 return err
351 }
357 - return bs.ledgermanager.MessageSent(p, m)
352 + return bs.engine.MessageSent(p, m)
353 }
354
355 func (bs *bitswap) Close() error {
exchange/bitswap/decision/engine.go renamed
+51 -51
@@ -1,4 +1,4 @@
1 -package strategy
1 +package decision
2
3 import (
4 "sync"
@@ -11,7 +11,7 @@ import (
11 u "github.com/jbenet/go-ipfs/util"
12 )
13
14 -var log = u.Logger("strategy")
14 +var log = u.Logger("engine")
15
16 // Envelope contains a message for a Peer
17 type Envelope struct {
@@ -21,7 +21,7 @@ type Envelope struct {
21 Message bsmsg.BitSwapMessage
22 }
23
24 -type LedgerManager struct {
24 +type Engine struct {
25 // FIXME taskqueue isn't threadsafe nor is it protected by a mutex. consider
26 // a way to avoid sharing the taskqueue between the worker and the receiver
27 taskqueue *taskQueue
@@ -37,32 +37,32 @@ type LedgerManager struct {
37 ledgerMap map[u.Key]*ledger
38 }
39
40 -func NewLedgerManager(ctx context.Context, bs bstore.Blockstore) *LedgerManager {
41 - lm := &LedgerManager{
40 +func NewEngine(ctx context.Context, bs bstore.Blockstore) *Engine {
41 + e := &Engine{
42 ledgerMap: make(map[u.Key]*ledger),
43 bs: bs,
44 taskqueue: newTaskQueue(),
45 outbox: make(chan Envelope, 4), // TODO extract constant
46 workSignal: make(chan struct{}),
47 }
48 - go lm.taskWorker(ctx)
49 - return lm
48 + go e.taskWorker(ctx)
49 + return e
50 }
51
52 -func (lm *LedgerManager) taskWorker(ctx context.Context) {
52 +func (e *Engine) taskWorker(ctx context.Context) {
53 for {
54 - nextTask := lm.taskqueue.Pop()
54 + nextTask := e.taskqueue.Pop()
55 if nextTask == nil {
56 // No tasks in the list?
57 // Wait until there are!
58 select {
59 case <-ctx.Done():
60 return
61 - case <-lm.workSignal:
61 + case <-e.workSignal:
62 }
63 continue
64 }
65 - block, err := lm.bs.Get(nextTask.Entry.Key)
65 + block, err := e.bs.Get(nextTask.Entry.Key)
66 if err != nil {
67 continue // TODO maybe return an error
68 }
@@ -74,22 +74,22 @@ func (lm *LedgerManager) taskWorker(ctx context.Context) {
74 select {
75 case <-ctx.Done():
76 return
77 - case lm.outbox <- Envelope{Peer: nextTask.Target, Message: m}:
77 + case e.outbox <- Envelope{Peer: nextTask.Target, Message: m}:
78 }
79 }
80 }
81
82 -func (lm *LedgerManager) Outbox() <-chan Envelope {
83 - return lm.outbox
82 +func (e *Engine) Outbox() <-chan Envelope {
83 + return e.outbox
84 }
85
86 // Returns a slice of Peers with whom the local node has active sessions
87 -func (lm *LedgerManager) Peers() []peer.Peer {
88 - lm.lock.RLock()
89 - defer lm.lock.RUnlock()
87 +func (e *Engine) Peers() []peer.Peer {
88 + e.lock.RLock()
89 + defer e.lock.RUnlock()
90
91 response := make([]peer.Peer, 0)
92 - for _, ledger := range lm.ledgerMap {
92 + for _, ledger := range e.ledgerMap {
93 response = append(response, ledger.Partner)
94 }
95 return response
@@ -97,52 +97,52 @@ func (lm *LedgerManager) Peers() []peer.Peer {
97
98 // BlockIsWantedByPeer returns true if peer wants the block given by this
99 // key
100 -func (lm *LedgerManager) BlockIsWantedByPeer(k u.Key, p peer.Peer) bool {
101 - lm.lock.RLock()
102 - defer lm.lock.RUnlock()
100 +func (e *Engine) BlockIsWantedByPeer(k u.Key, p peer.Peer) bool {
101 + e.lock.RLock()
102 + defer e.lock.RUnlock()
103
104 - ledger := lm.findOrCreate(p)
104 + ledger := e.findOrCreate(p)
105 return ledger.WantListContains(k)
106 }
107
108 // MessageReceived performs book-keeping. Returns error if passed invalid
109 // arguments.
110 -func (lm *LedgerManager) MessageReceived(p peer.Peer, m bsmsg.BitSwapMessage) error {
110 +func (e *Engine) MessageReceived(p peer.Peer, m bsmsg.BitSwapMessage) error {
111 newWorkExists := false
112 defer func() {
113 if newWorkExists {
114 // Signal task generation to restart (if stopped!)
115 select {
116 - case lm.workSignal <- struct{}{}:
116 + case e.workSignal <- struct{}{}:
117 default:
118 }
119 }
120 }()
121 - lm.lock.Lock()
122 - defer lm.lock.Unlock()
121 + e.lock.Lock()
122 + defer e.lock.Unlock()
123
124 - l := lm.findOrCreate(p)
124 + l := e.findOrCreate(p)
125 if m.Full() {
126 l.wantList = wl.New()
127 }
128 - for _, e := range m.Wantlist() {
129 - if e.Cancel {
130 - l.CancelWant(e.Key)
131 - lm.taskqueue.Remove(e.Key, p)
128 + for _, entry := range m.Wantlist() {
129 + if entry.Cancel {
130 + l.CancelWant(entry.Key)
131 + e.taskqueue.Remove(entry.Key, p)
132 } else {
133 - l.Wants(e.Key, e.Priority)
133 + l.Wants(entry.Key, entry.Priority)
134 newWorkExists = true
135 - lm.taskqueue.Push(e.Key, e.Priority, p)
135 + e.taskqueue.Push(entry.Key, entry.Priority, p)
136 }
137 }
138
139 for _, block := range m.Blocks() {
140 // FIXME extract blocks.NumBytes(block) or block.NumBytes() method
141 l.ReceivedBytes(len(block.Data))
142 - for _, l := range lm.ledgerMap {
142 + for _, l := range e.ledgerMap {
143 if l.WantListContains(block.Key()) {
144 newWorkExists = true
145 - lm.taskqueue.Push(block.Key(), 1, l.Partner)
145 + e.taskqueue.Push(block.Key(), 1, l.Partner)
146 }
147 }
148 }
@@ -155,40 +155,40 @@ func (lm *LedgerManager) MessageReceived(p peer.Peer, m bsmsg.BitSwapMessage) er
155 // inconsistent. Would need to ensure that Sends and acknowledgement of the
156 // send happen atomically
157
158 -func (lm *LedgerManager) MessageSent(p peer.Peer, m bsmsg.BitSwapMessage) error {
159 - lm.lock.Lock()
160 - defer lm.lock.Unlock()
158 +func (e *Engine) MessageSent(p peer.Peer, m bsmsg.BitSwapMessage) error {
159 + e.lock.Lock()
160 + defer e.lock.Unlock()
161
162 - l := lm.findOrCreate(p)
162 + l := e.findOrCreate(p)
163 for _, block := range m.Blocks() {
164 l.SentBytes(len(block.Data))
165 l.wantList.Remove(block.Key())
166 - lm.taskqueue.Remove(block.Key(), p)
166 + e.taskqueue.Remove(block.Key(), p)
167 }
168
169 return nil
170 }
171
172 -func (lm *LedgerManager) NumBytesSentTo(p peer.Peer) uint64 {
173 - lm.lock.RLock()
174 - defer lm.lock.RUnlock()
172 +func (e *Engine) NumBytesSentTo(p peer.Peer) uint64 {
173 + e.lock.RLock()
174 + defer e.lock.RUnlock()
175
176 - return lm.findOrCreate(p).Accounting.BytesSent
176 + return e.findOrCreate(p).Accounting.BytesSent
177 }
178
179 -func (lm *LedgerManager) NumBytesReceivedFrom(p peer.Peer) uint64 {
180 - lm.lock.RLock()
181 - defer lm.lock.RUnlock()
179 +func (e *Engine) NumBytesReceivedFrom(p peer.Peer) uint64 {
180 + e.lock.RLock()
181 + defer e.lock.RUnlock()
182
183 - return lm.findOrCreate(p).Accounting.BytesRecv
183 + return e.findOrCreate(p).Accounting.BytesRecv
184 }
185
186 // ledger lazily instantiates a ledger
187 -func (lm *LedgerManager) findOrCreate(p peer.Peer) *ledger {
188 - l, ok := lm.ledgerMap[p.Key()]
187 +func (e *Engine) findOrCreate(p peer.Peer) *ledger {
188 + l, ok := e.ledgerMap[p.Key()]
189 if !ok {
190 l = newLedger(p)
191 - lm.ledgerMap[p.Key()] = l
191 + e.ledgerMap[p.Key()] = l
192 }
193 return l
194 }
exchange/bitswap/decision/engine_test.go renamed
+19 -19
@@ -1,4 +1,4 @@
1 -package strategy
1 +package decision
2
3 import (
4 "strings"
@@ -14,16 +14,16 @@ import (
14 testutil "github.com/jbenet/go-ipfs/util/testutil"
15 )
16
17 -type peerAndLedgermanager struct {
17 +type peerAndEngine struct {
18 peer.Peer
19 - ls *LedgerManager
19 + Engine *Engine
20 }
21
22 -func newPeerAndLedgermanager(idStr string) peerAndLedgermanager {
23 - return peerAndLedgermanager{
22 +func newPeerAndLedgermanager(idStr string) peerAndEngine {
23 + return peerAndEngine{
24 Peer: testutil.NewPeerWithIDString(idStr),
25 //Strategy: New(true),
26 - ls: NewLedgerManager(context.TODO(),
26 + Engine: NewEngine(context.TODO(),
27 blockstore.NewBlockstore(sync.MutexWrap(ds.NewMapDatastore()))),
28 }
29 }
@@ -39,23 +39,23 @@ func TestConsistentAccounting(t *testing.T) {
39 content := []string{"this", "is", "message", "i"}
40 m.AddBlock(blocks.NewBlock([]byte(strings.Join(content, " "))))
41
42 - sender.ls.MessageSent(receiver.Peer, m)
43 - receiver.ls.MessageReceived(sender.Peer, m)
42 + sender.Engine.MessageSent(receiver.Peer, m)
43 + receiver.Engine.MessageReceived(sender.Peer, m)
44 }
45
46 // Ensure sender records the change
47 - if sender.ls.NumBytesSentTo(receiver.Peer) == 0 {
47 + if sender.Engine.NumBytesSentTo(receiver.Peer) == 0 {
48 t.Fatal("Sent bytes were not recorded")
49 }
50
51 // Ensure sender and receiver have the same values
52 - if sender.ls.NumBytesSentTo(receiver.Peer) != receiver.ls.NumBytesReceivedFrom(sender.Peer) {
52 + if sender.Engine.NumBytesSentTo(receiver.Peer) != receiver.Engine.NumBytesReceivedFrom(sender.Peer) {
53 t.Fatal("Inconsistent book-keeping. Strategies don't agree")
54 }
55
56 // Ensure sender didn't record receving anything. And that the receiver
57 // didn't record sending anything
58 - if receiver.ls.NumBytesSentTo(sender.Peer) != 0 || sender.ls.NumBytesReceivedFrom(receiver.Peer) != 0 {
58 + if receiver.Engine.NumBytesSentTo(sender.Peer) != 0 || sender.Engine.NumBytesReceivedFrom(receiver.Peer) != 0 {
59 t.Fatal("Bert didn't send bytes to Ernie")
60 }
61 }
@@ -69,10 +69,10 @@ func TestBlockRecordedAsWantedAfterMessageReceived(t *testing.T) {
69 messageFromBeggarToChooser := message.New()
70 messageFromBeggarToChooser.AddEntry(block.Key(), 1)
71
72 - chooser.ls.MessageReceived(beggar.Peer, messageFromBeggarToChooser)
72 + chooser.Engine.MessageReceived(beggar.Peer, messageFromBeggarToChooser)
73 // for this test, doesn't matter if you record that beggar sent
74
75 - if !chooser.ls.BlockIsWantedByPeer(block.Key(), beggar.Peer) {
75 + if !chooser.Engine.BlockIsWantedByPeer(block.Key(), beggar.Peer) {
76 t.Fatal("chooser failed to record that beggar wants block")
77 }
78 }
@@ -84,24 +84,24 @@ func TestPeerIsAddedToPeersWhenMessageReceivedOrSent(t *testing.T) {
84
85 m := message.New()
86
87 - sanfrancisco.ls.MessageSent(seattle.Peer, m)
88 - seattle.ls.MessageReceived(sanfrancisco.Peer, m)
87 + sanfrancisco.Engine.MessageSent(seattle.Peer, m)
88 + seattle.Engine.MessageReceived(sanfrancisco.Peer, m)
89
90 if seattle.Peer.Key() == sanfrancisco.Peer.Key() {
91 t.Fatal("Sanity Check: Peers have same Key!")
92 }
93
94 - if !peerIsPartner(seattle.Peer, sanfrancisco.ls) {
94 + if !peerIsPartner(seattle.Peer, sanfrancisco.Engine) {
95 t.Fatal("Peer wasn't added as a Partner")
96 }
97
98 - if !peerIsPartner(sanfrancisco.Peer, seattle.ls) {
98 + if !peerIsPartner(sanfrancisco.Peer, seattle.Engine) {
99 t.Fatal("Peer wasn't added as a Partner")
100 }
101 }
102
103 -func peerIsPartner(p peer.Peer, ls *LedgerManager) bool {
104 - for _, partner := range ls.Peers() {
103 +func peerIsPartner(p peer.Peer, e *Engine) bool {
104 + for _, partner := range e.Peers() {
105 if partner.Key() == p.Key() {
106 return true
107 }
exchange/bitswap/decision/ledger.go renamed
+1 -1
@@ -1,4 +1,4 @@
1 -package strategy
1 +package decision
2
3 import (
4 "time"
exchange/bitswap/decision/ledger_test.go new
+1
@@ -0,0 +1 @@
1 +package decision
exchange/bitswap/decision/taskqueue.go renamed
+1 -1
@@ -1,4 +1,4 @@
1 -package strategy
1 +package decision
2
3 import (
4 wantlist "github.com/jbenet/go-ipfs/exchange/bitswap/wantlist"
exchange/bitswap/strategy/ledger_test.go deleted
-1
@@ -1 +0,0 @@
1 -package strategy