@cryptotaxi247 / kubo / commits / 5b6a5e807

implement bitswap roundWorker

make vendor

Jeromy committed Dec 10, 2014 at 02:03 UTC 5b6a5e807f43775ec18dbcc7fa3eafaf71c88678
11 files changed +379 -115
Godeps/Godeps.json
+1 -1
@@ -1,6 +1,6 @@
1 {
2 "ImportPath": "github.com/jbenet/go-ipfs",
3 - "GoVersion": "go1.3",
3 + "GoVersion": "devel +ffe33f1f1f17 Tue Nov 25 15:41:33 2014 +1100",
4 "Packages": [
5 "./..."
6 ],
exchange/bitswap/bitswap.go
+69 -35
@@ -15,6 +15,7 @@ import (
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"
18 + wl "github.com/jbenet/go-ipfs/exchange/bitswap/wantlist"
19 peer "github.com/jbenet/go-ipfs/peer"
20 u "github.com/jbenet/go-ipfs/util"
21 eventlog "github.com/jbenet/go-ipfs/util/eventlog"
@@ -29,6 +30,8 @@ const maxProvidersPerRequest = 3
30 const providerRequestTimeout = time.Second * 10
31 const hasBlockTimeout = time.Second * 15
32
33 +const roundTime = time.Second / 2
34 +
35 // New initializes a BitSwap instance that communicates over the
36 // provided BitSwapNetwork. This function registers the returned instance as
37 // the network delegate.
@@ -41,6 +44,7 @@ func New(parent context.Context, p peer.Peer, network bsnet.BitSwapNetwork, rout
44 notif := notifications.New()
45 go func() {
46 <-ctx.Done()
47 + cancelFunc()
48 notif.Shutdown()
49 }()
50
@@ -51,11 +55,12 @@ func New(parent context.Context, p peer.Peer, network bsnet.BitSwapNetwork, rout
55 strategy: strategy.New(nice),
56 routing: routing,
57 sender: network,
54 - wantlist: u.NewKeySet(),
58 + wantlist: wl.NewWantlist(),
59 batchRequests: make(chan []u.Key, 32),
60 }
61 network.SetDelegate(bs)
62 go bs.loop(ctx)
63 + go bs.roundWorker(ctx)
64
65 return bs
66 }
@@ -85,7 +90,7 @@ type bitswap struct {
90 // TODO(brian): save the strategy's state to the datastore
91 strategy strategy.Strategy
92
88 - wantlist u.KeySet
93 + wantlist *wl.Wantlist
94
95 // cancelFunc signals cancellation to the bitswap event loop
96 cancelFunc func()
@@ -166,8 +171,8 @@ func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) e
171 panic("Cant send wantlist to nil peerchan")
172 }
173 message := bsmsg.New()
169 - for _, wanted := range bs.wantlist.Keys() {
170 - message.AddWanted(wanted)
174 + for _, wanted := range bs.wantlist.Entries() {
175 + message.AddEntry(wanted.Value, wanted.Priority, false)
176 }
177 for peerToQuery := range peers {
178 log.Debug("sending query to: %s", peerToQuery)
@@ -195,9 +200,9 @@ func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) e
200 return nil
201 }
202
198 -func (bs *bitswap) sendWantlistToProviders(ctx context.Context, ks []u.Key) {
203 +func (bs *bitswap) sendWantlistToProviders(ctx context.Context, wantlist *wl.Wantlist) {
204 wg := sync.WaitGroup{}
200 - for _, k := range ks {
205 + for _, e := range wantlist.Entries() {
206 wg.Add(1)
207 go func(k u.Key) {
208 child, _ := context.WithTimeout(ctx, providerRequestTimeout)
@@ -208,11 +213,44 @@ func (bs *bitswap) sendWantlistToProviders(ctx context.Context, ks []u.Key) {
213 log.Errorf("error sending wantlist: %s", err)
214 }
215 wg.Done()
211 - }(k)
216 + }(e.Value)
217 }
218 wg.Wait()
219 }
220
221 +func (bs *bitswap) roundWorker(ctx context.Context) {
222 + roundTicker := time.NewTicker(roundTime)
223 + bandwidthPerRound := 500000
224 + for {
225 + select {
226 + case <-ctx.Done():
227 + return
228 + case <-roundTicker.C:
229 + alloc, err := bs.strategy.GetAllocation(bandwidthPerRound, bs.blockstore)
230 + if err != nil {
231 + log.Critical("%s", err)
232 + }
233 + //log.Errorf("Allocation: %v", alloc)
234 + bs.processStrategyAllocation(ctx, alloc)
235 + }
236 + }
237 +}
238 +
239 +func (bs *bitswap) processStrategyAllocation(ctx context.Context, alloc []*strategy.Task) {
240 + for _, t := range alloc {
241 + for _, block := range t.Blocks {
242 + message := bsmsg.New()
243 + message.AddBlock(block)
244 + for _, wanted := range bs.wantlist.Entries() {
245 + message.AddEntry(wanted.Value, wanted.Priority, false)
246 + }
247 + if err := bs.send(ctx, t.Peer, message); err != nil {
248 + log.Errorf("Message Send Failed: %s", err)
249 + }
250 + }
251 + }
252 +}
253 +
254 // TODO ensure only one active request per key
255 func (bs *bitswap) loop(parent context.Context) {
256
@@ -228,7 +266,7 @@ func (bs *bitswap) loop(parent context.Context) {
266 select {
267 case <-broadcastSignal.C:
268 // Resend unfulfilled wantlist keys
231 - bs.sendWantlistToProviders(ctx, bs.wantlist.Keys())
269 + bs.sendWantlistToProviders(ctx, bs.wantlist)
270 case ks := <-bs.batchRequests:
271 // TODO: implement batching on len(ks) > X for some X
272 // i.e. if given 20 keys, fetch first five, then next
@@ -239,7 +277,7 @@ func (bs *bitswap) loop(parent context.Context) {
277 continue
278 }
279 for _, k := range ks {
242 - bs.wantlist.Add(k)
280 + bs.wantlist.Add(k, 1)
281 }
282 // NB: send want list to providers for the first peer in this list.
283 // the assumption is made that the providers of the first key in
@@ -277,45 +315,41 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
315 return nil, nil
316 }
317
280 - // Record message bytes in ledger
281 - // TODO: this is bad, and could be easily abused.
282 - // Should only track *useful* messages in ledger
318 // This call records changes to wantlists, blocks received,
319 // and number of bytes transfered.
320 bs.strategy.MessageReceived(p, incoming)
321 + // TODO: this is bad, and could be easily abused.
322 + // Should only track *useful* messages in ledger
323
324 + var blkeys []u.Key
325 for _, block := range incoming.Blocks() {
326 + blkeys = append(blkeys, block.Key())
327 if err := bs.HasBlock(ctx, block); err != nil {
328 log.Error(err)
329 }
330 }
292 -
293 - for _, key := range incoming.Wantlist() {
294 - if bs.strategy.ShouldSendBlockToPeer(key, p) {
295 - if block, errBlockNotFound := bs.blockstore.Get(key); errBlockNotFound != nil {
296 - continue
297 - } else {
298 - // Create a separate message to send this block in
299 - blkmsg := bsmsg.New()
300 -
301 - // TODO: only send this the first time
302 - // no sense in sending our wantlist to the
303 - // same peer multiple times
304 - for _, k := range bs.wantlist.Keys() {
305 - blkmsg.AddWanted(k)
306 - }
307 -
308 - blkmsg.AddBlock(block)
309 - bs.send(ctx, p, blkmsg)
310 - bs.strategy.BlockSentToPeer(block.Key(), p)
311 - }
312 - }
331 + if len(blkeys) > 0 {
332 + bs.cancelBlocks(ctx, blkeys)
333 }
334
335 // TODO: consider changing this function to not return anything
336 return nil, nil
337 }
338
339 +func (bs *bitswap) cancelBlocks(ctx context.Context, bkeys []u.Key) {
340 + message := bsmsg.New()
341 + message.SetFull(false)
342 + for _, k := range bkeys {
343 + message.AddEntry(k, 0, true)
344 + }
345 + for _, p := range bs.strategy.Peers() {
346 + err := bs.send(ctx, p, message)
347 + if err != nil {
348 + log.Errorf("Error sending message: %s", err)
349 + }
350 + }
351 +}
352 +
353 func (bs *bitswap) ReceiveError(err error) {
354 log.Errorf("Bitswap ReceiveError: %s", err)
355 // TODO log the network error
@@ -337,8 +371,8 @@ func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block *blocks.Block)
371 if bs.strategy.ShouldSendBlockToPeer(block.Key(), p) {
372 message := bsmsg.New()
373 message.AddBlock(block)
340 - for _, wanted := range bs.wantlist.Keys() {
341 - message.AddWanted(wanted)
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
exchange/bitswap/bitswap_test.go
+55 -10
@@ -11,6 +11,7 @@ import (
11 blocksutil "github.com/jbenet/go-ipfs/blocks/blocksutil"
12 tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
13 mockrouting "github.com/jbenet/go-ipfs/routing/mock"
14 + u "github.com/jbenet/go-ipfs/util"
15 delay "github.com/jbenet/go-ipfs/util/delay"
16 testutil "github.com/jbenet/go-ipfs/util/testutil"
17 )
@@ -25,6 +26,7 @@ func TestClose(t *testing.T) {
26 vnet := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
27 rout := mockrouting.NewServer()
28 sesgen := NewSessionGenerator(vnet, rout)
29 + defer sesgen.Stop()
30 bgen := blocksutil.NewBlockGenerator()
31
32 block := bgen.Next()
@@ -39,6 +41,7 @@ func TestGetBlockTimeout(t *testing.T) {
41 net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
42 rs := mockrouting.NewServer()
43 g := NewSessionGenerator(net, rs)
44 + defer g.Stop()
45
46 self := g.Next()
47
@@ -56,11 +59,13 @@ func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
59 net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
60 rs := mockrouting.NewServer()
61 g := NewSessionGenerator(net, rs)
62 + defer g.Stop()
63
64 block := blocks.NewBlock([]byte("block"))
65 rs.Client(testutil.NewPeerWithIDString("testing")).Provide(context.Background(), block.Key()) // but not on network
66
67 solo := g.Next()
68 + defer solo.Exchange.Close()
69
70 ctx, _ := context.WithTimeout(context.Background(), time.Nanosecond)
71 _, err := solo.Exchange.GetBlock(ctx, block.Key())
@@ -78,8 +83,10 @@ func TestGetBlockFromPeerAfterPeerAnnounces(t *testing.T) {
83 rs := mockrouting.NewServer()
84 block := blocks.NewBlock([]byte("block"))
85 g := NewSessionGenerator(net, rs)
86 + defer g.Stop()
87
88 hasBlock := g.Next()
89 + defer hasBlock.Exchange.Close()
90
91 if err := hasBlock.Blockstore().Put(block); err != nil {
92 t.Fatal(err)
@@ -89,6 +96,7 @@ func TestGetBlockFromPeerAfterPeerAnnounces(t *testing.T) {
96 }
97
98 wantsBlock := g.Next()
99 + defer wantsBlock.Exchange.Close()
100
101 ctx, _ := context.WithTimeout(context.Background(), time.Second)
102 received, err := wantsBlock.Exchange.GetBlock(ctx, block.Key())
@@ -107,7 +115,7 @@ func TestLargeSwarm(t *testing.T) {
115 t.SkipNow()
116 }
117 t.Parallel()
110 - numInstances := 5
118 + numInstances := 500
119 numBlocks := 2
120 PerformDistributionTest(t, numInstances, numBlocks)
121 }
@@ -129,6 +137,7 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
137 net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
138 rs := mockrouting.NewServer()
139 sg := NewSessionGenerator(net, rs)
140 + defer sg.Stop()
141 bg := blocksutil.NewBlockGenerator()
142
143 t.Log("Test a few nodes trying to get one file with a lot of blocks")
@@ -138,24 +147,29 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
147
148 t.Log("Give the blocks to the first instance")
149
150 + var blkeys []u.Key
151 first := instances[0]
152 for _, b := range blocks {
153 first.Blockstore().Put(b)
154 + blkeys = append(blkeys, b.Key())
155 first.Exchange.HasBlock(context.Background(), b)
156 rs.Client(first.Peer).Provide(context.Background(), b.Key())
157 }
158
159 t.Log("Distribute!")
160
150 - var wg sync.WaitGroup
151 -
161 + wg := sync.WaitGroup{}
162 for _, inst := range instances {
153 - for _, b := range blocks {
154 - wg.Add(1)
155 - // NB: executing getOrFail concurrently puts tremendous pressure on
156 - // the goroutine scheduler
157 - getOrFail(inst, b, t, &wg)
158 - }
163 + wg.Add(1)
164 + go func(inst Instance) {
165 + defer wg.Done()
166 + outch, err := inst.Exchange.GetBlocks(context.TODO(), blkeys)
167 + if err != nil {
168 + t.Fatal(err)
169 + }
170 + for _ = range outch {
171 + }
172 + }(inst)
173 }
174 wg.Wait()
175
@@ -189,6 +203,7 @@ func TestSendToWantingPeer(t *testing.T) {
203 net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
204 rs := mockrouting.NewServer()
205 sg := NewSessionGenerator(net, rs)
206 + defer sg.Stop()
207 bg := blocksutil.NewBlockGenerator()
208
209 me := sg.Next()
@@ -201,7 +216,7 @@ func TestSendToWantingPeer(t *testing.T) {
216
217 alpha := bg.Next()
218
204 - const timeout = 100 * time.Millisecond // FIXME don't depend on time
219 + const timeout = 1000 * time.Millisecond // FIXME don't depend on time
220
221 t.Logf("Peer %v attempts to get %v. NB: not available\n", w.Peer, alpha.Key())
222 ctx, _ := context.WithTimeout(context.Background(), timeout)
@@ -246,3 +261,33 @@ func TestSendToWantingPeer(t *testing.T) {
261 t.Fatal("Expected to receive alpha from me")
262 }
263 }
264 +
265 +func TestBasicBitswap(t *testing.T) {
266 + net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
267 + rs := mockrouting.NewServer()
268 + sg := NewSessionGenerator(net, rs)
269 + bg := blocksutil.NewBlockGenerator()
270 +
271 + t.Log("Test a few nodes trying to get one file with a lot of blocks")
272 +
273 + instances := sg.Instances(2)
274 + blocks := bg.Blocks(1)
275 + err := instances[0].Exchange.HasBlock(context.TODO(), blocks[0])
276 + if err != nil {
277 + t.Fatal(err)
278 + }
279 +
280 + ctx, _ := context.WithTimeout(context.TODO(), time.Second*5)
281 + blk, err := instances[1].Exchange.GetBlock(ctx, blocks[0].Key())
282 + if err != nil {
283 + t.Fatal(err)
284 + }
285 +
286 + t.Log(blk)
287 + for _, inst := range instances {
288 + err := inst.Exchange.Close()
289 + if err != nil {
290 + t.Fatal(err)
291 + }
292 + }
293 +}
exchange/bitswap/message/internal/pb/message.pb.go
+60 -4
@@ -21,16 +21,16 @@ var _ = proto.Marshal
21 var _ = math.Inf
22
23 type Message struct {
24 - Wantlist []string `protobuf:"bytes,1,rep,name=wantlist" json:"wantlist,omitempty"`
25 - Blocks [][]byte `protobuf:"bytes,2,rep,name=blocks" json:"blocks,omitempty"`
26 - XXX_unrecognized []byte `json:"-"`
24 + Wantlist *Message_Wantlist `protobuf:"bytes,1,opt,name=wantlist" json:"wantlist,omitempty"`
25 + Blocks [][]byte `protobuf:"bytes,2,rep,name=blocks" json:"blocks,omitempty"`
26 + XXX_unrecognized []byte `json:"-"`
27 }
28
29 func (m *Message) Reset() { *m = Message{} }
30 func (m *Message) String() string { return proto.CompactTextString(m) }
31 func (*Message) ProtoMessage() {}
32
33 -func (m *Message) GetWantlist() []string {
33 +func (m *Message) GetWantlist() *Message_Wantlist {
34 if m != nil {
35 return m.Wantlist
36 }
@@ -44,5 +44,61 @@ func (m *Message) GetBlocks() [][]byte {
44 return nil
45 }
46
47 +type Message_Wantlist struct {
48 + Entries []*Message_Wantlist_Entry `protobuf:"bytes,1,rep,name=entries" json:"entries,omitempty"`
49 + Full *bool `protobuf:"varint,2,opt,name=full" json:"full,omitempty"`
50 + XXX_unrecognized []byte `json:"-"`
51 +}
52 +
53 +func (m *Message_Wantlist) Reset() { *m = Message_Wantlist{} }
54 +func (m *Message_Wantlist) String() string { return proto.CompactTextString(m) }
55 +func (*Message_Wantlist) ProtoMessage() {}
56 +
57 +func (m *Message_Wantlist) GetEntries() []*Message_Wantlist_Entry {
58 + if m != nil {
59 + return m.Entries
60 + }
61 + return nil
62 +}
63 +
64 +func (m *Message_Wantlist) GetFull() bool {
65 + if m != nil && m.Full != nil {
66 + return *m.Full
67 + }
68 + return false
69 +}
70 +
71 +type Message_Wantlist_Entry struct {
72 + Block *string `protobuf:"bytes,1,opt,name=block" json:"block,omitempty"`
73 + Priority *int32 `protobuf:"varint,2,opt,name=priority" json:"priority,omitempty"`
74 + Cancel *bool `protobuf:"varint,3,opt,name=cancel" json:"cancel,omitempty"`
75 + XXX_unrecognized []byte `json:"-"`
76 +}
77 +
78 +func (m *Message_Wantlist_Entry) Reset() { *m = Message_Wantlist_Entry{} }
79 +func (m *Message_Wantlist_Entry) String() string { return proto.CompactTextString(m) }
80 +func (*Message_Wantlist_Entry) ProtoMessage() {}
81 +
82 +func (m *Message_Wantlist_Entry) GetBlock() string {
83 + if m != nil && m.Block != nil {
84 + return *m.Block
85 + }
86 + return ""
87 +}
88 +
89 +func (m *Message_Wantlist_Entry) GetPriority() int32 {
90 + if m != nil && m.Priority != nil {
91 + return *m.Priority
92 + }
93 + return 0
94 +}
95 +
96 +func (m *Message_Wantlist_Entry) GetCancel() bool {
97 + if m != nil && m.Cancel != nil {
98 + return *m.Cancel
99 + }
100 + return false
101 +}
102 +
103 func init() {
104 }
exchange/bitswap/message/internal/pb/message.proto
+15 -2
@@ -1,6 +1,19 @@
1 package bitswap.message.pb;
2
3 message Message {
4 - repeated string wantlist = 1;
5 - repeated bytes blocks = 2;
4 +
5 + message Wantlist {
6 +
7 + message Entry {
8 + optional string block = 1; // the block key
9 + optional int32 priority = 2; // the priority (normalized). default to 1
10 + optional bool cancel = 3; // whether this revokes an entry
11 + }
12 +
13 + repeated Entry entries = 1; // a list of wantlist entries
14 + optional bool full = 2; // whether this is the full wantlist. default to false
15 + }
16 +
17 + optional Wantlist wantlist = 1;
18 + repeated bytes blocks = 2;
19 }
exchange/bitswap/message/message.go
+61 -31
@@ -9,6 +9,7 @@ import (
9 u "github.com/jbenet/go-ipfs/util"
10
11 ggio "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/gogoprotobuf/io"
12 + proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
13 )
14
15 // TODO move message.go into the bitswap package
@@ -17,21 +18,21 @@ import (
18 type BitSwapMessage interface {
19 // Wantlist returns a slice of unique keys that represent data wanted by
20 // the sender.
20 - Wantlist() []u.Key
21 + Wantlist() []*Entry
22
23 // Blocks returns a slice of unique blocks
24 Blocks() []*blocks.Block
25
25 - // AddWanted adds the key to the Wantlist.
26 - //
27 - // Insertion order determines priority. That is, earlier insertions are
28 - // deemed higher priority than keys inserted later.
29 - //
30 - // t = 0, msg.AddWanted(A)
31 - // t = 1, msg.AddWanted(B)
32 - //
33 - // implies Priority(A) > Priority(B)
34 - AddWanted(u.Key)
26 + // AddEntry adds an entry to the Wantlist.
27 + AddEntry(u.Key, int, 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)
34 +
35 + Full() bool
36
37 AddBlock(*blocks.Block)
38 Exportable
@@ -43,23 +44,30 @@ type Exportable interface {
44 }
45
46 type impl struct {
46 - existsInWantlist map[u.Key]struct{} // map to detect duplicates
47 - wantlist []u.Key // slice to preserve ordering
48 - blocks map[u.Key]*blocks.Block // map to detect duplicates
47 + full bool
48 + wantlist map[u.Key]*Entry
49 + blocks map[u.Key]*blocks.Block // map to detect duplicates
50 }
51
52 func New() BitSwapMessage {
53 return &impl{
53 - blocks: make(map[u.Key]*blocks.Block),
54 - existsInWantlist: make(map[u.Key]struct{}),
55 - wantlist: make([]u.Key, 0),
54 + blocks: make(map[u.Key]*blocks.Block),
55 + wantlist: make(map[u.Key]*Entry),
56 + full: true,
57 }
58 }
59
60 +type Entry struct {
61 + Key u.Key
62 + Priority int
63 + Cancel bool
64 +}
65 +
66 func newMessageFromProto(pbm pb.Message) BitSwapMessage {
67 m := New()
61 - for _, s := range pbm.GetWantlist() {
62 - m.AddWanted(u.Key(s))
68 + m.SetFull(pbm.GetWantlist().GetFull())
69 + for _, e := range pbm.GetWantlist().GetEntries() {
70 + m.AddEntry(u.Key(e.GetBlock()), int(e.GetPriority()), e.GetCancel())
71 }
72 for _, d := range pbm.GetBlocks() {
73 b := blocks.NewBlock(d)
@@ -68,8 +76,20 @@ func newMessageFromProto(pbm pb.Message) BitSwapMessage {
76 return m
77 }
78
71 -func (m *impl) Wantlist() []u.Key {
72 - return m.wantlist
79 +func (m *impl) SetFull(full bool) {
80 + m.full = full
81 +}
82 +
83 +func (m *impl) Full() bool {
84 + return m.full
85 +}
86 +
87 +func (m *impl) Wantlist() []*Entry {
88 + var out []*Entry
89 + for _, e := range m.wantlist {
90 + out = append(out, e)
91 + }
92 + return out
93 }
94
95 func (m *impl) Blocks() []*blocks.Block {
@@ -80,13 +100,18 @@ func (m *impl) Blocks() []*blocks.Block {
100 return bs
101 }
102
83 -func (m *impl) AddWanted(k u.Key) {
84 - _, exists := m.existsInWantlist[k]
103 +func (m *impl) AddEntry(k u.Key, priority int, cancel bool) {
104 + e, exists := m.wantlist[k]
105 if exists {
86 - return
106 + e.Priority = priority
107 + e.Cancel = cancel
108 + } else {
109 + m.wantlist[k] = &Entry{
110 + Key: k,
111 + Priority: priority,
112 + Cancel: cancel,
113 + }
114 }
88 - m.existsInWantlist[k] = struct{}{}
89 - m.wantlist = append(m.wantlist, k)
115 }
116
117 func (m *impl) AddBlock(b *blocks.Block) {
@@ -106,14 +131,19 @@ func FromNet(r io.Reader) (BitSwapMessage, error) {
131 }
132
133 func (m *impl) ToProto() *pb.Message {
109 - pb := new(pb.Message)
110 - for _, k := range m.Wantlist() {
111 - pb.Wantlist = append(pb.Wantlist, string(k))
134 + pbm := new(pb.Message)
135 + pbm.Wantlist = new(pb.Message_Wantlist)
136 + for _, e := range m.wantlist {
137 + pbm.Wantlist.Entries = append(pbm.Wantlist.Entries, &pb.Message_Wantlist_Entry{
138 + Block: proto.String(string(e.Key)),
139 + Priority: proto.Int32(int32(e.Priority)),
140 + Cancel: &e.Cancel,
141 + })
142 }
143 for _, b := range m.Blocks() {
114 - pb.Blocks = append(pb.Blocks, b.Data)
144 + pbm.Blocks = append(pbm.Blocks, b.Data)
145 }
116 - return pb
146 + return pbm
147 }
148
149 func (m *impl) ToNet(w io.Writer) error {
exchange/bitswap/message/message_test.go
+37 -22
@@ -4,6 +4,8 @@ import (
4 "bytes"
5 "testing"
6
7 + proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
8 +
9 blocks "github.com/jbenet/go-ipfs/blocks"
10 pb "github.com/jbenet/go-ipfs/exchange/bitswap/message/internal/pb"
11 u "github.com/jbenet/go-ipfs/util"
@@ -12,22 +14,26 @@ import (
14 func TestAppendWanted(t *testing.T) {
15 const str = "foo"
16 m := New()
15 - m.AddWanted(u.Key(str))
17 + m.AddEntry(u.Key(str), 1, false)
18
17 - if !contains(m.ToProto().GetWantlist(), str) {
19 + if !wantlistContains(m.ToProto().GetWantlist(), str) {
20 t.Fail()
21 }
22 + m.ToProto().GetWantlist().GetEntries()
23 }
24
25 func TestNewMessageFromProto(t *testing.T) {
26 const str = "a_key"
27 protoMessage := new(pb.Message)
25 - protoMessage.Wantlist = []string{string(str)}
26 - if !contains(protoMessage.Wantlist, str) {
28 + protoMessage.Wantlist = new(pb.Message_Wantlist)
29 + protoMessage.Wantlist.Entries = []*pb.Message_Wantlist_Entry{
30 + &pb.Message_Wantlist_Entry{Block: proto.String(str)},
31 + }
32 + if !wantlistContains(protoMessage.Wantlist, str) {
33 t.Fail()
34 }
35 m := newMessageFromProto(*protoMessage)
30 - if !contains(m.ToProto().GetWantlist(), str) {
36 + if !wantlistContains(m.ToProto().GetWantlist(), str) {
37 t.Fail()
38 }
39 }
@@ -57,7 +63,7 @@ func TestWantlist(t *testing.T) {
63 keystrs := []string{"foo", "bar", "baz", "bat"}
64 m := New()
65 for _, s := range keystrs {
60 - m.AddWanted(u.Key(s))
66 + m.AddEntry(u.Key(s), 1, false)
67 }
68 exported := m.Wantlist()
69
@@ -65,12 +71,12 @@ func TestWantlist(t *testing.T) {
71 present := false
72 for _, s := range keystrs {
73
68 - if s == string(k) {
74 + if s == string(k.Key) {
75 present = true
76 }
77 }
78 if !present {
73 - t.Logf("%v isn't in original list", string(k))
79 + t.Logf("%v isn't in original list", k.Key)
80 t.Fail()
81 }
82 }
@@ -80,19 +86,19 @@ func TestCopyProtoByValue(t *testing.T) {
86 const str = "foo"
87 m := New()
88 protoBeforeAppend := m.ToProto()
83 - m.AddWanted(u.Key(str))
84 - if contains(protoBeforeAppend.GetWantlist(), str) {
89 + m.AddEntry(u.Key(str), 1, false)
90 + if wantlistContains(protoBeforeAppend.GetWantlist(), str) {
91 t.Fail()
92 }
93 }
94
95 func TestToNetFromNetPreservesWantList(t *testing.T) {
96 original := New()
91 - original.AddWanted(u.Key("M"))
92 - original.AddWanted(u.Key("B"))
93 - original.AddWanted(u.Key("D"))
94 - original.AddWanted(u.Key("T"))
95 - original.AddWanted(u.Key("F"))
97 + original.AddEntry(u.Key("M"), 1, false)
98 + original.AddEntry(u.Key("B"), 1, false)
99 + original.AddEntry(u.Key("D"), 1, false)
100 + original.AddEntry(u.Key("T"), 1, false)
101 + original.AddEntry(u.Key("F"), 1, false)
102
103 var buf bytes.Buffer
104 if err := original.ToNet(&buf); err != nil {
@@ -106,11 +112,11 @@ func TestToNetFromNetPreservesWantList(t *testing.T) {
112
113 keys := make(map[u.Key]bool)
114 for _, k := range copied.Wantlist() {
109 - keys[k] = true
115 + keys[k.Key] = true
116 }
117
118 for _, k := range original.Wantlist() {
113 - if _, ok := keys[k]; !ok {
119 + if _, ok := keys[k.Key]; !ok {
120 t.Fatalf("Key Missing: \"%v\"", k)
121 }
122 }
@@ -146,9 +152,18 @@ func TestToAndFromNetMessage(t *testing.T) {
152 }
153 }
154
149 -func contains(s []string, x string) bool {
150 - for _, a := range s {
151 - if a == x {
155 +func wantlistContains(wantlist *pb.Message_Wantlist, x string) bool {
156 + for _, e := range wantlist.GetEntries() {
157 + if e.GetBlock() == x {
158 + return true
159 + }
160 + }
161 + return false
162 +}
163 +
164 +func contains(strs []string, x string) bool {
165 + for _, s := range strs {
166 + if s == x {
167 return true
168 }
169 }
@@ -159,8 +174,8 @@ func TestDuplicates(t *testing.T) {
174 b := blocks.NewBlock([]byte("foo"))
175 msg := New()
176
162 - msg.AddWanted(b.Key())
163 - msg.AddWanted(b.Key())
177 + msg.AddEntry(b.Key(), 1, false)
178 + msg.AddEntry(b.Key(), 1, false)
179 if len(msg.Wantlist()) != 1 {
180 t.Fatal("Duplicate in BitSwapMessage")
181 }
exchange/bitswap/strategy/interface.go
+3
@@ -3,6 +3,7 @@ package strategy
3 import (
4 "time"
5
6 + 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"
@@ -34,6 +35,8 @@ type Strategy interface {
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
exchange/bitswap/strategy/ledger.go
+10 -6
@@ -3,6 +3,7 @@ package strategy
3 import (
4 "time"
5
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 )
@@ -13,7 +14,7 @@ type keySet map[u.Key]struct{}
14
15 func newLedger(p peer.Peer, strategy strategyFunc) *ledger {
16 return &ledger{
16 - wantList: keySet{},
17 + wantList: wl.NewWantlist(),
18 Strategy: strategy,
19 Partner: p,
20 sentToPeer: make(map[u.Key]time.Time),
@@ -39,7 +40,7 @@ type ledger struct {
40 exchangeCount uint64
41
42 // wantList is a (bounded, small) set of keys that Partner desires.
42 - wantList keySet
43 + wantList *wl.Wantlist
44
45 // sentToPeer is a set of keys to ensure we dont send duplicate blocks
46 // to a given peer
@@ -65,14 +66,17 @@ func (l *ledger) ReceivedBytes(n int) {
66 }
67
68 // TODO: this needs to be different. We need timeouts.
68 -func (l *ledger) Wants(k u.Key) {
69 +func (l *ledger) Wants(k u.Key, priority int) {
70 log.Debugf("peer %s wants %s", l.Partner, k)
70 - l.wantList[k] = struct{}{}
71 + l.wantList.Add(k, priority)
72 +}
73 +
74 +func (l *ledger) CancelWant(k u.Key) {
75 + l.wantList.Remove(k)
76 }
77
78 func (l *ledger) WantListContains(k u.Key) bool {
74 - _, ok := l.wantList[k]
75 - return ok
79 + return l.wantList.Contains(k)
80 }
81
82 func (l *ledger) ExchangeCount() uint64 {
exchange/bitswap/strategy/strategy.go
+67 -3
@@ -5,7 +5,10 @@ import (
5 "sync"
6 "time"
7
8 + blocks "github.com/jbenet/go-ipfs/blocks"
9 + bstore "github.com/jbenet/go-ipfs/blocks/blockstore"
10 bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
11 + wl "github.com/jbenet/go-ipfs/exchange/bitswap/wantlist"
12 peer "github.com/jbenet/go-ipfs/peer"
13 u "github.com/jbenet/go-ipfs/util"
14 )
@@ -77,6 +80,60 @@ func (s *strategist) ShouldSendBlockToPeer(k u.Key, p peer.Peer) bool {
80 return ledger.ShouldSend()
81 }
82
83 +type Task struct {
84 + Peer peer.Peer
85 + Blocks []*blocks.Block
86 +}
87 +
88 +func (s *strategist) GetAllocation(bandwidth int, bs bstore.Blockstore) ([]*Task, error) {
89 + var tasks []*Task
90 +
91 + s.lock.RLock()
92 + defer s.lock.RUnlock()
93 + var partners []peer.Peer
94 + for _, ledger := range s.ledgerMap {
95 + if ledger.ShouldSend() {
96 + partners = append(partners, ledger.Partner)
97 + }
98 + }
99 + if len(partners) == 0 {
100 + return nil, nil
101 + }
102 +
103 + bandwidthPerPeer := bandwidth / len(partners)
104 + for _, p := range partners {
105 + blksForPeer, err := s.getSendableBlocks(s.ledger(p).wantList, bs, bandwidthPerPeer)
106 + if err != nil {
107 + return nil, err
108 + }
109 + tasks = append(tasks, &Task{
110 + Peer: p,
111 + Blocks: blksForPeer,
112 + })
113 + }
114 +
115 + return tasks, nil
116 +}
117 +
118 +func (s *strategist) getSendableBlocks(wantlist *wl.Wantlist, bs bstore.Blockstore, bw int) ([]*blocks.Block, error) {
119 + var outblocks []*blocks.Block
120 + for _, e := range wantlist.Entries() {
121 + block, err := bs.Get(e.Value)
122 + if err == u.ErrNotFound {
123 + continue
124 + }
125 + if err != nil {
126 + return nil, err
127 + }
128 + outblocks = append(outblocks, block)
129 + bw -= len(block.Data)
130 + if bw <= 0 {
131 + break
132 + }
133 + }
134 + return outblocks, nil
135 +}
136 +
137 func (s *strategist) BlockSentToPeer(k u.Key, p peer.Peer) {
138 s.lock.Lock()
139 defer s.lock.Unlock()
@@ -106,8 +163,15 @@ func (s *strategist) MessageReceived(p peer.Peer, m bsmsg.BitSwapMessage) error
163 return errors.New("Strategy received nil message")
164 }
165 l := s.ledger(p)
109 - for _, key := range m.Wantlist() {
110 - l.Wants(key)
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
@@ -165,5 +229,5 @@ func (s *strategist) GetBatchSize() int {
229 }
230
231 func (s *strategist) GetRebroadcastDelay() time.Duration {
168 - return time.Second * 5
232 + return time.Second * 10
233 }
exchange/bitswap/strategy/strategy_test.go
+1 -1
@@ -61,7 +61,7 @@ func TestBlockRecordedAsWantedAfterMessageReceived(t *testing.T) {
61 block := blocks.NewBlock([]byte("data wanted by beggar"))
62
63 messageFromBeggarToChooser := message.New()
64 - messageFromBeggarToChooser.AddWanted(block.Key())
64 + messageFromBeggarToChooser.AddEntry(block.Key(), 1, false)
65
66 chooser.MessageReceived(beggar.Peer, messageFromBeggarToChooser)
67 // for this test, doesn't matter if you record that beggar sent