@cryptotaxi247 / kubo / commits / 57e7dd7b8

extracted ledgerset from strategy, cleaned up a few comments from the PR

Jeromy committed Dec 10, 2014 at 07:57 UTC 57e7dd7b8bc2a952ec34e5c3e576af6195b28847
9 files changed +208 -301
exchange/bitswap/bitswap.go
+20 -41
@@ -27,11 +27,14 @@ var log = eventlog.Logger("bitswap")
27 // TODO: if a 'non-nice' strategy is implemented, consider increasing this value
28 const maxProvidersPerRequest = 3
29
30 -const providerRequestTimeout = time.Second * 10
31 -const hasBlockTimeout = time.Second * 15
30 +var providerRequestTimeout = time.Second * 10
31 +var hasBlockTimeout = time.Second * 15
32 +var rebroadcastDelay = time.Second * 10
33
34 const roundTime = time.Second / 2
35
36 +var bandwidthPerRound = 500000
37 +
38 // New initializes a BitSwap instance that communicates over the
39 // provided BitSwapNetwork. This function registers the returned instance as
40 // the network delegate.
@@ -53,13 +56,14 @@ func New(parent context.Context, p peer.Peer, network bsnet.BitSwapNetwork, rout
56 cancelFunc: cancelFunc,
57 notifications: notif,
58 strategy: strategy.New(nice),
59 + ledgerset: strategy.NewLedgerSet(),
60 routing: routing,
61 sender: network,
58 - wantlist: wl.NewWantlist(),
62 + wantlist: wl.New(),
63 batchRequests: make(chan []u.Key, 32),
64 }
65 network.SetDelegate(bs)
62 - go bs.loop(ctx)
66 + go bs.clientWorker(ctx)
67 go bs.roundWorker(ctx)
68
69 return bs
@@ -85,11 +89,11 @@ type bitswap struct {
89 // have more than a single block in the set
90 batchRequests chan []u.Key
91
88 - // strategy listens to network traffic and makes decisions about how to
89 - // interact with partners.
90 - // TODO(brian): save the strategy's state to the datastore
92 + // strategy makes decisions about how to interact with partners.
93 strategy strategy.Strategy
94
95 + ledgerset *strategy.LedgerSet
96 +
97 wantlist *wl.Wantlist
98
99 // cancelFunc signals cancellation to the bitswap event loop
@@ -159,10 +163,6 @@ func (bs *bitswap) HasBlock(ctx context.Context, blk *blocks.Block) error {
163 bs.wantlist.Remove(blk.Key())
164 bs.notifications.Publish(blk)
165 child, _ := context.WithTimeout(ctx, hasBlockTimeout)
162 - if err := bs.sendToPeersThatWant(child, blk); err != nil {
163 - return err
164 - }
165 - child, _ = context.WithTimeout(ctx, hasBlockTimeout)
166 return bs.routing.Provide(child, blk.Key())
167 }
168
@@ -194,7 +194,7 @@ func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) e
194 // FIXME ensure accounting is handled correctly when
195 // communication fails. May require slightly different API to
196 // get better guarantees. May need shared sequence numbers.
197 - bs.strategy.MessageSent(p, message)
197 + bs.ledgerset.MessageSent(p, message)
198 }(peerToQuery)
199 }
200 return nil
@@ -220,17 +220,16 @@ func (bs *bitswap) sendWantlistToProviders(ctx context.Context, wantlist *wl.Wan
220
221 func (bs *bitswap) roundWorker(ctx context.Context) {
222 roundTicker := time.NewTicker(roundTime)
223 - bandwidthPerRound := 500000
223 for {
224 select {
225 case <-ctx.Done():
226 return
227 case <-roundTicker.C:
229 - alloc, err := bs.strategy.GetAllocation(bandwidthPerRound, bs.blockstore)
228 + alloc, err := bs.strategy.GetTasks(bandwidthPerRound, bs.ledgerset, bs.blockstore)
229 if err != nil {
230 log.Critical("%s", err)
231 }
233 - //log.Errorf("Allocation: %v", alloc)
232 + log.Error(alloc)
233 bs.processStrategyAllocation(ctx, alloc)
234 }
235 }
@@ -241,9 +240,6 @@ func (bs *bitswap) processStrategyAllocation(ctx context.Context, alloc []*strat
240 for _, block := range t.Blocks {
241 message := bsmsg.New()
242 message.AddBlock(block)
244 - for _, wanted := range bs.wantlist.Entries() {
245 - message.AddEntry(wanted.Value, wanted.Priority, false)
246 - }
243 if err := bs.send(ctx, t.Peer, message); err != nil {
244 log.Errorf("Message Send Failed: %s", err)
245 }
@@ -252,11 +248,11 @@ func (bs *bitswap) processStrategyAllocation(ctx context.Context, alloc []*strat
248 }
249
250 // TODO ensure only one active request per key
255 -func (bs *bitswap) loop(parent context.Context) {
251 +func (bs *bitswap) clientWorker(parent context.Context) {
252
253 ctx, cancel := context.WithCancel(parent)
254
259 - broadcastSignal := time.NewTicker(bs.strategy.GetRebroadcastDelay())
255 + broadcastSignal := time.NewTicker(rebroadcastDelay)
256 defer func() {
257 cancel() // signal to derived async functions
258 broadcastSignal.Stop()
@@ -317,13 +313,14 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
313
314 // This call records changes to wantlists, blocks received,
315 // and number of bytes transfered.
320 - bs.strategy.MessageReceived(p, incoming)
316 + bs.ledgerset.MessageReceived(p, incoming)
317 // TODO: this is bad, and could be easily abused.
318 // Should only track *useful* messages in ledger
319
320 var blkeys []u.Key
321 for _, block := range incoming.Blocks() {
322 blkeys = append(blkeys, block.Key())
323 + log.Errorf("Got block: %s", block)
324 if err := bs.HasBlock(ctx, block); err != nil {
325 log.Error(err)
326 }
@@ -342,7 +339,7 @@ func (bs *bitswap) cancelBlocks(ctx context.Context, bkeys []u.Key) {
339 for _, k := range bkeys {
340 message.AddEntry(k, 0, true)
341 }
345 - for _, p := range bs.strategy.Peers() {
342 + for _, p := range bs.ledgerset.Peers() {
343 err := bs.send(ctx, p, message)
344 if err != nil {
345 log.Errorf("Error sending message: %s", err)
@@ -362,25 +359,7 @@ func (bs *bitswap) send(ctx context.Context, p peer.Peer, m bsmsg.BitSwapMessage
359 if err := bs.sender.SendMessage(ctx, p, m); err != nil {
360 return err
361 }
365 - return bs.strategy.MessageSent(p, m)
366 -}
367 -
368 -func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block *blocks.Block) error {
369 - for _, p := range bs.strategy.Peers() {
370 - if bs.strategy.BlockIsWantedByPeer(block.Key(), p) {
371 - if bs.strategy.ShouldSendBlockToPeer(block.Key(), p) {
372 - message := bsmsg.New()
373 - message.AddBlock(block)
374 - for _, wanted := range bs.wantlist.Entries() {
375 - message.AddEntry(wanted.Value, wanted.Priority, false)
376 - }
377 - if err := bs.send(ctx, p, message); err != nil {
378 - return err
379 - }
380 - }
381 - }
382 - }
383 - return nil
362 + return bs.ledgerset.MessageSent(p, m)
363 }
364
365 func (bs *bitswap) Close() error {
exchange/bitswap/bitswap_test.go
+24 -40
@@ -206,60 +206,44 @@ func TestSendToWantingPeer(t *testing.T) {
206 defer sg.Stop()
207 bg := blocksutil.NewBlockGenerator()
208
209 - me := sg.Next()
210 - w := sg.Next()
211 - o := sg.Next()
209 + oldVal := rebroadcastDelay
210 + rebroadcastDelay = time.Second / 2
211 + defer func() { rebroadcastDelay = oldVal }()
212
213 - t.Logf("Session %v\n", me.Peer)
214 - t.Logf("Session %v\n", w.Peer)
215 - t.Logf("Session %v\n", o.Peer)
213 + peerA := sg.Next()
214 + peerB := sg.Next()
215
217 - alpha := bg.Next()
218 -
219 - const timeout = 1000 * time.Millisecond // FIXME don't depend on time
216 + t.Logf("Session %v\n", peerA.Peer)
217 + t.Logf("Session %v\n", peerB.Peer)
218
221 - t.Logf("Peer %v attempts to get %v. NB: not available\n", w.Peer, alpha.Key())
222 - ctx, _ := context.WithTimeout(context.Background(), timeout)
223 - _, err := w.Exchange.GetBlock(ctx, alpha.Key())
224 - if err == nil {
225 - t.Fatalf("Expected %v to NOT be available", alpha.Key())
226 - }
219 + timeout := time.Second
220 + waitTime := time.Second * 5
221
228 - beta := bg.Next()
229 - t.Logf("Peer %v announes availability of %v\n", w.Peer, beta.Key())
230 - ctx, _ = context.WithTimeout(context.Background(), timeout)
231 - if err := w.Blockstore().Put(beta); err != nil {
222 + alpha := bg.Next()
223 + // peerA requests and waits for block alpha
224 + ctx, _ := context.WithTimeout(context.TODO(), waitTime)
225 + alphaPromise, err := peerA.Exchange.GetBlocks(ctx, []u.Key{alpha.Key()})
226 + if err != nil {
227 t.Fatal(err)
228 }
234 - w.Exchange.HasBlock(ctx, beta)
229
236 - t.Logf("%v gets %v from %v and discovers it wants %v\n", me.Peer, beta.Key(), w.Peer, alpha.Key())
237 - ctx, _ = context.WithTimeout(context.Background(), timeout)
238 - if _, err := me.Exchange.GetBlock(ctx, beta.Key()); err != nil {
230 + // peerB announces to the network that he has block alpha
231 + ctx, _ = context.WithTimeout(context.TODO(), timeout)
232 + err = peerB.Exchange.HasBlock(ctx, alpha)
233 + if err != nil {
234 t.Fatal(err)
235 }
236
242 - t.Logf("%v announces availability of %v\n", o.Peer, alpha.Key())
243 - ctx, _ = context.WithTimeout(context.Background(), timeout)
244 - if err := o.Blockstore().Put(alpha); err != nil {
245 - t.Fatal(err)
237 + // At some point, peerA should get alpha (or timeout)
238 + blkrecvd, ok := <-alphaPromise
239 + if !ok {
240 + t.Fatal("context timed out and broke promise channel!")
241 }
247 - o.Exchange.HasBlock(ctx, alpha)
242
249 - t.Logf("%v requests %v\n", me.Peer, alpha.Key())
250 - ctx, _ = context.WithTimeout(context.Background(), timeout)
251 - if _, err := me.Exchange.GetBlock(ctx, alpha.Key()); err != nil {
252 - t.Fatal(err)
243 + if blkrecvd.Key() != alpha.Key() {
244 + t.Fatal("Wrong block!")
245 }
246
255 - t.Logf("%v should now have %v\n", w.Peer, alpha.Key())
256 - block, err := w.Blockstore().Get(alpha.Key())
257 - if err != nil {
258 - t.Fatalf("Should not have received an error: %s", err)
259 - }
260 - if block.Key() != alpha.Key() {
261 - t.Fatal("Expected to receive alpha from me")
262 - }
247 }
248
249 func TestBasicBitswap(t *testing.T) {
exchange/bitswap/message/message.go
+2 -2
@@ -24,13 +24,13 @@ type BitSwapMessage interface {
24 Blocks() []*blocks.Block
25
26 // AddEntry adds an entry to the Wantlist.
27 - AddEntry(u.Key, int, bool)
27 + AddEntry(key u.Key, priority int, cancel bool)
28
29 // Sets whether or not the contained wantlist represents the entire wantlist
30 // true = full wantlist
31 // false = wantlist 'patch'
32 // default: true
33 - SetFull(bool)
33 + SetFull(isFull bool)
34
35 Full() bool
36
exchange/bitswap/strategy/interface.go
+1 -32
@@ -1,43 +1,12 @@
1 package strategy
2
3 import (
4 - "time"
5 -
4 bstore "github.com/jbenet/go-ipfs/blocks/blockstore"
7 - bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
8 - peer "github.com/jbenet/go-ipfs/peer"
9 - u "github.com/jbenet/go-ipfs/util"
5 )
6
7 type Strategy interface {
13 - // Returns a slice of Peers with whom the local node has active sessions
14 - Peers() []peer.Peer
15 -
16 - // BlockIsWantedByPeer returns true if peer wants the block given by this
17 - // key
18 - BlockIsWantedByPeer(u.Key, peer.Peer) bool
19 -
20 - // ShouldSendTo(Peer) decides whether to send data to this Peer
21 - ShouldSendBlockToPeer(u.Key, peer.Peer) bool
22 -
8 // Seed initializes the decider to a deterministic state
9 Seed(int64)
10
26 - // MessageReceived records receipt of message for accounting purposes
27 - MessageReceived(peer.Peer, bsmsg.BitSwapMessage) error
28 -
29 - // MessageSent records sending of message for accounting purposes
30 - MessageSent(peer.Peer, bsmsg.BitSwapMessage) error
31 -
32 - NumBytesSentTo(peer.Peer) uint64
33 -
34 - NumBytesReceivedFrom(peer.Peer) uint64
35 -
36 - BlockSentToPeer(u.Key, peer.Peer)
37 -
38 - GetAllocation(int, bstore.Blockstore) ([]*Task, error)
39 -
40 - // Values determining bitswap behavioural patterns
41 - GetBatchSize() int
42 - GetRebroadcastDelay() time.Duration
11 + GetTasks(bandwidth int, ledgers *LedgerSet, bs bstore.Blockstore) ([]*Task, error)
12 }
exchange/bitswap/strategy/ledger.go
+2 -9
@@ -12,10 +12,9 @@ import (
12 // access/lookups.
13 type keySet map[u.Key]struct{}
14
15 -func newLedger(p peer.Peer, strategy strategyFunc) *ledger {
15 +func newLedger(p peer.Peer) *ledger {
16 return &ledger{
17 - wantList: wl.NewWantlist(),
18 - Strategy: strategy,
17 + wantList: wl.New(),
18 Partner: p,
19 sentToPeer: make(map[u.Key]time.Time),
20 }
@@ -45,12 +44,6 @@ type ledger struct {
44 // sentToPeer is a set of keys to ensure we dont send duplicate blocks
45 // to a given peer
46 sentToPeer map[u.Key]time.Time
48 -
49 - Strategy strategyFunc
50 -}
51 -
52 -func (l *ledger) ShouldSend() bool {
53 - return l.Strategy(l)
47 }
48
49 func (l *ledger) SentBytes(n int) {
exchange/bitswap/strategy/ledgerset.go new
+125
@@ -0,0 +1,125 @@
1 +package strategy
2 +
3 +import (
4 + "sync"
5 +
6 + bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
7 + wl "github.com/jbenet/go-ipfs/exchange/bitswap/wantlist"
8 + peer "github.com/jbenet/go-ipfs/peer"
9 + u "github.com/jbenet/go-ipfs/util"
10 +)
11 +
12 +// LedgerMap lists Ledgers by their Partner key.
13 +type ledgerMap map[peerKey]*ledger
14 +
15 +// FIXME share this externally
16 +type peerKey u.Key
17 +
18 +type LedgerSet struct {
19 + lock sync.RWMutex
20 + ledgerMap ledgerMap
21 +}
22 +
23 +func NewLedgerSet() *LedgerSet {
24 + return &LedgerSet{
25 + ledgerMap: make(ledgerMap),
26 + }
27 +}
28 +
29 +// Returns a slice of Peers with whom the local node has active sessions
30 +func (ls *LedgerSet) Peers() []peer.Peer {
31 + ls.lock.RLock()
32 + defer ls.lock.RUnlock()
33 +
34 + response := make([]peer.Peer, 0)
35 + for _, ledger := range ls.ledgerMap {
36 + response = append(response, ledger.Partner)
37 + }
38 + return response
39 +}
40 +
41 +// BlockIsWantedByPeer returns true if peer wants the block given by this
42 +// key
43 +func (ls *LedgerSet) BlockIsWantedByPeer(k u.Key, p peer.Peer) bool {
44 + ls.lock.RLock()
45 + defer ls.lock.RUnlock()
46 +
47 + ledger := ls.ledger(p)
48 + return ledger.WantListContains(k)
49 +}
50 +
51 +// MessageReceived performs book-keeping. Returns error if passed invalid
52 +// arguments.
53 +func (ls *LedgerSet) MessageReceived(p peer.Peer, m bsmsg.BitSwapMessage) error {
54 + ls.lock.Lock()
55 + defer ls.lock.Unlock()
56 +
57 + // TODO find a more elegant way to handle this check
58 + /*
59 + if p == nil {
60 + return errors.New("Strategy received nil peer")
61 + }
62 + if m == nil {
63 + return errors.New("Strategy received nil message")
64 + }
65 + */
66 + l := ls.ledger(p)
67 + if m.Full() {
68 + l.wantList = wl.New()
69 + }
70 + for _, e := range m.Wantlist() {
71 + if e.Cancel {
72 + l.CancelWant(e.Key)
73 + } else {
74 + l.Wants(e.Key, e.Priority)
75 + }
76 + }
77 + for _, block := range m.Blocks() {
78 + // FIXME extract blocks.NumBytes(block) or block.NumBytes() method
79 + l.ReceivedBytes(len(block.Data))
80 + }
81 + return nil
82 +}
83 +
84 +// TODO add contents of m.WantList() to my local wantlist? NB: could introduce
85 +// race conditions where I send a message, but MessageSent gets handled after
86 +// MessageReceived. The information in the local wantlist could become
87 +// inconsistent. Would need to ensure that Sends and acknowledgement of the
88 +// send happen atomically
89 +
90 +func (ls *LedgerSet) MessageSent(p peer.Peer, m bsmsg.BitSwapMessage) error {
91 + ls.lock.Lock()
92 + defer ls.lock.Unlock()
93 +
94 + l := ls.ledger(p)
95 + for _, block := range m.Blocks() {
96 + l.SentBytes(len(block.Data))
97 + l.wantList.Remove(block.Key())
98 + }
99 +
100 + return nil
101 +}
102 +
103 +func (ls *LedgerSet) NumBytesSentTo(p peer.Peer) uint64 {
104 + ls.lock.RLock()
105 + defer ls.lock.RUnlock()
106 +
107 + return ls.ledger(p).Accounting.BytesSent
108 +}
109 +
110 +func (ls *LedgerSet) NumBytesReceivedFrom(p peer.Peer) uint64 {
111 + ls.lock.RLock()
112 + defer ls.lock.RUnlock()
113 +
114 + return ls.ledger(p).Accounting.BytesRecv
115 +}
116 +
117 +// ledger lazily instantiates a ledger
118 +func (ls *LedgerSet) ledger(p peer.Peer) *ledger {
119 + l, ok := ls.ledgerMap[peerKey(p.Key())]
120 + if !ok {
121 + l = newLedger(p)
122 + ls.ledgerMap[peerKey(p.Key())] = l
123 + }
124 + return l
125 +}
exchange/bitswap/strategy/ledgerset_test.go renamed
+26 -25
@@ -10,21 +10,22 @@ import (
10 testutil "github.com/jbenet/go-ipfs/util/testutil"
11 )
12
13 -type peerAndStrategist struct {
13 +type peerAndLedgerset struct {
14 peer.Peer
15 - Strategy
15 + ls *LedgerSet
16 }
17
18 -func newPeerAndStrategist(idStr string) peerAndStrategist {
19 - return peerAndStrategist{
20 - Peer: testutil.NewPeerWithIDString(idStr),
21 - Strategy: New(true),
18 +func newPeerAndLedgerset(idStr string) peerAndLedgerset {
19 + return peerAndLedgerset{
20 + Peer: testutil.NewPeerWithIDString(idStr),
21 + //Strategy: New(true),
22 + ls: NewLedgerSet(),
23 }
24 }
25
26 func TestConsistentAccounting(t *testing.T) {
26 - sender := newPeerAndStrategist("Ernie")
27 - receiver := newPeerAndStrategist("Bert")
27 + sender := newPeerAndLedgerset("Ernie")
28 + receiver := newPeerAndLedgerset("Bert")
29
30 // Send messages from Ernie to Bert
31 for i := 0; i < 1000; i++ {
@@ -33,69 +34,69 @@ func TestConsistentAccounting(t *testing.T) {
34 content := []string{"this", "is", "message", "i"}
35 m.AddBlock(blocks.NewBlock([]byte(strings.Join(content, " "))))
36
36 - sender.MessageSent(receiver.Peer, m)
37 - receiver.MessageReceived(sender.Peer, m)
37 + sender.ls.MessageSent(receiver.Peer, m)
38 + receiver.ls.MessageReceived(sender.Peer, m)
39 }
40
41 // Ensure sender records the change
41 - if sender.NumBytesSentTo(receiver.Peer) == 0 {
42 + if sender.ls.NumBytesSentTo(receiver.Peer) == 0 {
43 t.Fatal("Sent bytes were not recorded")
44 }
45
46 // Ensure sender and receiver have the same values
46 - if sender.NumBytesSentTo(receiver.Peer) != receiver.NumBytesReceivedFrom(sender.Peer) {
47 + if sender.ls.NumBytesSentTo(receiver.Peer) != receiver.ls.NumBytesReceivedFrom(sender.Peer) {
48 t.Fatal("Inconsistent book-keeping. Strategies don't agree")
49 }
50
51 // Ensure sender didn't record receving anything. And that the receiver
52 // didn't record sending anything
52 - if receiver.NumBytesSentTo(sender.Peer) != 0 || sender.NumBytesReceivedFrom(receiver.Peer) != 0 {
53 + if receiver.ls.NumBytesSentTo(sender.Peer) != 0 || sender.ls.NumBytesReceivedFrom(receiver.Peer) != 0 {
54 t.Fatal("Bert didn't send bytes to Ernie")
55 }
56 }
57
58 func TestBlockRecordedAsWantedAfterMessageReceived(t *testing.T) {
58 - beggar := newPeerAndStrategist("can't be chooser")
59 - chooser := newPeerAndStrategist("chooses JIF")
59 + beggar := newPeerAndLedgerset("can't be chooser")
60 + chooser := newPeerAndLedgerset("chooses JIF")
61
62 block := blocks.NewBlock([]byte("data wanted by beggar"))
63
64 messageFromBeggarToChooser := message.New()
65 messageFromBeggarToChooser.AddEntry(block.Key(), 1, false)
66
66 - chooser.MessageReceived(beggar.Peer, messageFromBeggarToChooser)
67 + chooser.ls.MessageReceived(beggar.Peer, messageFromBeggarToChooser)
68 // for this test, doesn't matter if you record that beggar sent
69
69 - if !chooser.BlockIsWantedByPeer(block.Key(), beggar.Peer) {
70 + if !chooser.ls.BlockIsWantedByPeer(block.Key(), beggar.Peer) {
71 t.Fatal("chooser failed to record that beggar wants block")
72 }
73 }
74
75 func TestPeerIsAddedToPeersWhenMessageReceivedOrSent(t *testing.T) {
76
76 - sanfrancisco := newPeerAndStrategist("sf")
77 - seattle := newPeerAndStrategist("sea")
77 + sanfrancisco := newPeerAndLedgerset("sf")
78 + seattle := newPeerAndLedgerset("sea")
79
80 m := message.New()
81
81 - sanfrancisco.MessageSent(seattle.Peer, m)
82 - seattle.MessageReceived(sanfrancisco.Peer, m)
82 + sanfrancisco.ls.MessageSent(seattle.Peer, m)
83 + seattle.ls.MessageReceived(sanfrancisco.Peer, m)
84
85 if seattle.Peer.Key() == sanfrancisco.Peer.Key() {
86 t.Fatal("Sanity Check: Peers have same Key!")
87 }
88
88 - if !peerIsPartner(seattle.Peer, sanfrancisco.Strategy) {
89 + if !peerIsPartner(seattle.Peer, sanfrancisco.ls) {
90 t.Fatal("Peer wasn't added as a Partner")
91 }
92
92 - if !peerIsPartner(sanfrancisco.Peer, seattle.Strategy) {
93 + if !peerIsPartner(sanfrancisco.Peer, seattle.ls) {
94 t.Fatal("Peer wasn't added as a Partner")
95 }
96 }
97
97 -func peerIsPartner(p peer.Peer, s Strategy) bool {
98 - for _, partner := range s.Peers() {
98 +func peerIsPartner(p peer.Peer, ls *LedgerSet) bool {
99 + for _, partner := range ls.Peers() {
100 if partner.Key() == p.Key() {
101 return true
102 }
exchange/bitswap/strategy/strategy.go
+7 -151
@@ -1,20 +1,13 @@
1 package strategy
2
3 import (
4 - "errors"
5 - "sync"
6 - "time"
7 -
4 blocks "github.com/jbenet/go-ipfs/blocks"
5 bstore "github.com/jbenet/go-ipfs/blocks/blockstore"
10 - bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
6 wl "github.com/jbenet/go-ipfs/exchange/bitswap/wantlist"
7 peer "github.com/jbenet/go-ipfs/peer"
8 u "github.com/jbenet/go-ipfs/util"
9 )
10
16 -const resendTimeoutPeriod = time.Minute
17 -
11 var log = u.Logger("strategy")
12
13 // TODO niceness should be on a per-peer basis. Use-case: Certain peers are
@@ -28,81 +21,37 @@ func New(nice bool) Strategy {
21 stratFunc = standardStrategy
22 }
23 return &strategist{
31 - ledgerMap: ledgerMap{},
24 strategyFunc: stratFunc,
25 }
26 }
27
28 type strategist struct {
37 - lock sync.RWMutex
38 - ledgerMap
29 strategyFunc
30 }
31
42 -// LedgerMap lists Ledgers by their Partner key.
43 -type ledgerMap map[peerKey]*ledger
44 -
45 -// FIXME share this externally
46 -type peerKey u.Key
47 -
48 -// Peers returns a list of peers
49 -func (s *strategist) Peers() []peer.Peer {
50 - s.lock.RLock()
51 - defer s.lock.RUnlock()
52 -
53 - response := make([]peer.Peer, 0)
54 - for _, ledger := range s.ledgerMap {
55 - response = append(response, ledger.Partner)
56 - }
57 - return response
58 -}
59 -
60 -func (s *strategist) BlockIsWantedByPeer(k u.Key, p peer.Peer) bool {
61 - s.lock.RLock()
62 - defer s.lock.RUnlock()
63 -
64 - ledger := s.ledger(p)
65 - return ledger.WantListContains(k)
66 -}
67 -
68 -func (s *strategist) ShouldSendBlockToPeer(k u.Key, p peer.Peer) bool {
69 - s.lock.RLock()
70 - defer s.lock.RUnlock()
71 -
72 - ledger := s.ledger(p)
73 -
74 - // Dont resend blocks within a certain time period
75 - t, ok := ledger.sentToPeer[k]
76 - if ok && t.Add(resendTimeoutPeriod).After(time.Now()) {
77 - return false
78 - }
79 -
80 - return ledger.ShouldSend()
81 -}
82 -
32 type Task struct {
33 Peer peer.Peer
34 Blocks []*blocks.Block
35 }
36
88 -func (s *strategist) GetAllocation(bandwidth int, bs bstore.Blockstore) ([]*Task, error) {
37 +func (s *strategist) GetTasks(bandwidth int, ledgers *LedgerSet, bs bstore.Blockstore) ([]*Task, error) {
38 var tasks []*Task
39
91 - s.lock.RLock()
92 - defer s.lock.RUnlock()
40 + ledgers.lock.RLock()
41 var partners []peer.Peer
94 - for _, ledger := range s.ledgerMap {
95 - if ledger.ShouldSend() {
42 + for _, ledger := range ledgers.ledgerMap {
43 + if s.strategyFunc(ledger) {
44 partners = append(partners, ledger.Partner)
45 }
46 }
47 + ledgers.lock.RUnlock()
48 if len(partners) == 0 {
49 return nil, nil
50 }
51
52 bandwidthPerPeer := bandwidth / len(partners)
53 for _, p := range partners {
105 - blksForPeer, err := s.getSendableBlocks(s.ledger(p).wantList, bs, bandwidthPerPeer)
54 + blksForPeer, err := s.getSendableBlocks(ledgers.ledger(p).wantList, bs, bandwidthPerPeer)
55 if err != nil {
56 return nil, err
57 }
@@ -134,100 +83,7 @@ func (s *strategist) getSendableBlocks(wantlist *wl.Wantlist, bs bstore.Blocksto
83 return outblocks, nil
84 }
85
137 -func (s *strategist) BlockSentToPeer(k u.Key, p peer.Peer) {
138 - s.lock.Lock()
139 - defer s.lock.Unlock()
140 -
141 - ledger := s.ledger(p)
142 - ledger.sentToPeer[k] = time.Now()
143 -}
144 -
86 +func test() {}
87 func (s *strategist) Seed(int64) {
146 - s.lock.Lock()
147 - defer s.lock.Unlock()
148 -
88 // TODO
89 }
151 -
152 -// MessageReceived performs book-keeping. Returns error if passed invalid
153 -// arguments.
154 -func (s *strategist) MessageReceived(p peer.Peer, m bsmsg.BitSwapMessage) error {
155 - s.lock.Lock()
156 - defer s.lock.Unlock()
157 -
158 - // TODO find a more elegant way to handle this check
159 - if p == nil {
160 - return errors.New("Strategy received nil peer")
161 - }
162 - if m == nil {
163 - return errors.New("Strategy received nil message")
164 - }
165 - l := s.ledger(p)
166 - if m.Full() {
167 - l.wantList = wl.NewWantlist()
168 - }
169 - for _, e := range m.Wantlist() {
170 - if e.Cancel {
171 - l.CancelWant(e.Key)
172 - } else {
173 - l.Wants(e.Key, e.Priority)
174 - }
175 - }
176 - for _, block := range m.Blocks() {
177 - // FIXME extract blocks.NumBytes(block) or block.NumBytes() method
178 - l.ReceivedBytes(len(block.Data))
179 - }
180 - return nil
181 -}
182 -
183 -// TODO add contents of m.WantList() to my local wantlist? NB: could introduce
184 -// race conditions where I send a message, but MessageSent gets handled after
185 -// MessageReceived. The information in the local wantlist could become
186 -// inconsistent. Would need to ensure that Sends and acknowledgement of the
187 -// send happen atomically
188 -
189 -func (s *strategist) MessageSent(p peer.Peer, m bsmsg.BitSwapMessage) error {
190 - s.lock.Lock()
191 - defer s.lock.Unlock()
192 -
193 - l := s.ledger(p)
194 - for _, block := range m.Blocks() {
195 - l.SentBytes(len(block.Data))
196 - }
197 -
198 - // TODO remove these blocks from peer's want list
199 -
200 - return nil
201 -}
202 -
203 -func (s *strategist) NumBytesSentTo(p peer.Peer) uint64 {
204 - s.lock.RLock()
205 - defer s.lock.RUnlock()
206 -
207 - return s.ledger(p).Accounting.BytesSent
208 -}
209 -
210 -func (s *strategist) NumBytesReceivedFrom(p peer.Peer) uint64 {
211 - s.lock.RLock()
212 - defer s.lock.RUnlock()
213 -
214 - return s.ledger(p).Accounting.BytesRecv
215 -}
216 -
217 -// ledger lazily instantiates a ledger
218 -func (s *strategist) ledger(p peer.Peer) *ledger {
219 - l, ok := s.ledgerMap[peerKey(p.Key())]
220 - if !ok {
221 - l = newLedger(p, s.strategyFunc)
222 - s.ledgerMap[peerKey(p.Key())] = l
223 - }
224 - return l
225 -}
226 -
227 -func (s *strategist) GetBatchSize() int {
228 - return 10
229 -}
230 -
231 -func (s *strategist) GetRebroadcastDelay() time.Duration {
232 - return time.Second * 10
233 -}
exchange/bitswap/wantlist/wantlist.go
+1 -1
@@ -9,7 +9,7 @@ type Wantlist struct {
9 set map[u.Key]*Entry
10 }
11
12 -func NewWantlist() *Wantlist {
12 +func New() *Wantlist {
13 return &Wantlist{
14 set: make(map[u.Key]*Entry),
15 }