fix(bitswap) implement, test concrete strategist
Brian Tiger Chow committed
Sep 18, 2014 at 14:22 UTC
be8e08675d2b175b652ef1421d1040563c23a824
4 files changed
+113
-28
bitswap/bitswap.go
+1
-1
@@ -55,7 +55,7 @@ func NewSession(parent context.Context, s bsnet.NetworkService, p *peer.Peer, d
55
bs := &bitswap{
56
blockstore: blockstore.NewBlockstore(d),
57
notifications: notifications.New(),
58
- strategist: strategy.New(d),
58
+ strategist: strategy.New(),
59
peer: p,
60
routing: directory,
61
sender: bsnet.NewNetworkAdapter(s, &receiver),
bitswap/strategy/ledger.go
+7
-3
@@ -12,6 +12,13 @@ import (
12
// access/lookups.
13
type keySet map[u.Key]struct{}
14
15
+func newLedger(p *peer.Peer, strategy strategyFunc) *ledger {
16
+ return &ledger{
17
+ Strategy: strategy,
18
+ Partner: p,
19
+ }
20
+}
21
+
22
// ledger stores the data exchange relationship between two peers.
23
type ledger struct {
24
lock sync.RWMutex
@@ -37,9 +44,6 @@ type ledger struct {
44
Strategy strategyFunc
45
}
46
40
-// LedgerMap lists Ledgers by their Partner key.
41
-type ledgerMap map[u.Key]*ledger
42
-
47
func (l *ledger) ShouldSend() bool {
48
l.lock.Lock()
49
defer l.lock.Unlock()
bitswap/strategy/strategy.go
+53
-24
@@ -3,56 +3,85 @@ package strategy
3
import (
4
"errors"
5
6
- ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
6
bsmsg "github.com/jbenet/go-ipfs/bitswap/message"
7
"github.com/jbenet/go-ipfs/peer"
8
u "github.com/jbenet/go-ipfs/util"
9
)
10
11
// TODO declare thread-safe datastore
13
-func New(d ds.Datastore) Strategist {
12
+func New() Strategist {
13
return &strategist{
15
- datastore: d,
16
- peers: ledgerMap{},
14
+ ledgerMap: ledgerMap{},
15
+ strategyFunc: yesManStrategy,
16
}
17
}
18
19
type strategist struct {
21
- datastore ds.Datastore // FIXME(brian): enforce thread-safe datastore
22
-
23
- peers ledgerMap
20
+ ledgerMap
21
+ strategyFunc
22
}
23
26
-// Peers returns a list of this instance is connected to
24
+// LedgerMap lists Ledgers by their Partner key.
25
+type ledgerMap map[peerKey]*ledger
26
+
27
+// FIXME share this externally
28
+type peerKey u.Key
29
+
30
+// Peers returns a list of peers
31
func (s *strategist) Peers() []*peer.Peer {
28
- response := make([]*peer.Peer, 0) // TODO
32
+ response := make([]*peer.Peer, 0)
33
+ for _, ledger := range s.ledgerMap {
34
+ response = append(response, ledger.Partner)
35
+ }
36
return response
37
}
38
32
-func (s *strategist) IsWantedByPeer(u.Key, *peer.Peer) bool {
33
- return true // TODO
39
+func (s *strategist) IsWantedByPeer(k u.Key, p *peer.Peer) bool {
40
+ ledger := s.ledger(p)
41
+ return ledger.WantListContains(k)
42
}
43
36
-func (s *strategist) ShouldSendToPeer(u.Key, *peer.Peer) bool {
37
- return true // TODO
44
+func (s *strategist) ShouldSendToPeer(k u.Key, p *peer.Peer) bool {
45
+ ledger := s.ledger(p)
46
+ return ledger.ShouldSend()
47
}
48
49
func (s *strategist) Seed(int64) {
50
// TODO
51
}
52
44
-func (s *strategist) MessageReceived(*peer.Peer, bsmsg.BitSwapMessage) error {
45
- // TODO add peer to partners if doesn't already exist.
46
- // TODO initialize ledger for peer if doesn't already exist
47
- // TODO get wantlist from message and update contents in local wantlist for peer
48
- // TODO acknowledge receipt of blocks and do accounting in ledger
53
+func (s *strategist) MessageReceived(p *peer.Peer, m bsmsg.BitSwapMessage) error {
54
+ l := s.ledger(p)
55
+ for _, key := range m.Wantlist() {
56
+ l.Wants(key)
57
+ }
58
+ for _, block := range m.Blocks() {
59
+ // FIXME extract blocks.NumBytes(block) or block.NumBytes() method
60
+ l.ReceivedBytes(len(block.Data))
61
+ }
62
return errors.New("TODO")
63
}
64
52
-func (s *strategist) MessageSent(*peer.Peer, bsmsg.BitSwapMessage) error {
53
- // TODO add peer to partners if doesn't already exist.
54
- // TODO initialize ledger for peer if doesn't already exist
55
- // TODO add block to my wantlist
56
- // TODO acknowledge receipt of blocks and do accounting in ledger
57
- return errors.New("TODO")
65
+// TODO add contents of m.WantList() to my local wantlist? NB: could introduce
66
+// race conditions where I send a message, but MessageSent gets handled after
67
+// MessageReceived. The information in the local wantlist could become
68
+// inconsistent. Would need to ensure that Sends and acknowledgement of the
69
+// send happen atomically
70
+
71
+func (s *strategist) MessageSent(p *peer.Peer, m bsmsg.BitSwapMessage) error {
72
+ l := s.ledger(p)
73
+ for _, block := range m.Blocks() {
74
+ l.SentBytes(len(block.Data))
75
+ }
76
+ return nil
77
+}
78
+
79
+// ledger lazily instantiates a ledger
80
+func (s *strategist) ledger(p *peer.Peer) *ledger {
81
+ l, ok := s.ledgerMap[peerKey(p.Key())]
82
+ if !ok {
83
+ l = newLedger(p, s.strategyFunc)
84
+ s.ledgerMap[peerKey(p.Key())] = l
85
+ }
86
+ return l
87
}
bitswap/strategy/strategy_test.go
new
+52
@@ -0,0 +1,52 @@
1
+package strategy
2
+
3
+import (
4
+ "testing"
5
+
6
+ message "github.com/jbenet/go-ipfs/bitswap/message"
7
+ "github.com/jbenet/go-ipfs/peer"
8
+)
9
+
10
+type peerAndStrategist struct {
11
+ *peer.Peer
12
+ Strategist
13
+}
14
+
15
+func newPeerAndStrategist(idStr string) peerAndStrategist {
16
+ return peerAndStrategist{
17
+ Peer: &peer.Peer{ID: peer.ID(idStr)},
18
+ Strategist: New(),
19
+ }
20
+}
21
+
22
+func TestPeerIsAddedToPeersWhenMessageReceivedOrSent(t *testing.T) {
23
+
24
+ sanfrancisco := newPeerAndStrategist("sf")
25
+ seattle := newPeerAndStrategist("sea")
26
+
27
+ m := message.New()
28
+
29
+ sanfrancisco.MessageSent(seattle.Peer, m)
30
+ seattle.MessageReceived(sanfrancisco.Peer, m)
31
+
32
+ if seattle.Peer.Key() == sanfrancisco.Peer.Key() {
33
+ t.Fatal("Sanity Check: Peers have same Key!")
34
+ }
35
+
36
+ if !peerIsPartner(seattle.Peer, sanfrancisco.Strategist) {
37
+ t.Fatal("Peer wasn't added as a Partner")
38
+ }
39
+
40
+ if !peerIsPartner(sanfrancisco.Peer, seattle.Strategist) {
41
+ t.Fatal("Peer wasn't added as a Partner")
42
+ }
43
+}
44
+
45
+func peerIsPartner(p *peer.Peer, s Strategist) bool {
46
+ for _, partner := range s.Peers() {
47
+ if partner.Key() == p.Key() {
48
+ return true
49
+ }
50
+ }
51
+ return false
52
+}