test(bitswap) send entire wantlist to peers
fix(bitswap) pass go vet fixes #97 https://github.com/jbenet/go-ipfs/issues/97
Brian Tiger Chow committed
Sep 21, 2014 at 21:39 UTC
faee10effee1ad1ab79004532978e0b58426d359
2 files changed
+93
-27
exchange/bitswap/bitswap.go
+61
-9
@@ -2,6 +2,7 @@ package bitswap
2
3
import (
4
"errors"
5
+ "sync"
6
7
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
@@ -28,6 +29,9 @@ func NetMessageSession(parent context.Context, p *peer.Peer, s bsnet.NetMessageS
29
strategy: strategy.New(),
30
routing: directory,
31
sender: networkAdapter,
32
+ wantlist: WantList{
33
+ data: make(map[u.Key]struct{}),
34
+ },
35
}
36
networkAdapter.SetDelegate(bs)
37
@@ -53,6 +57,39 @@ type bitswap struct {
57
// interact with partners.
58
// TODO(brian): save the strategy's state to the datastore
59
strategy strategy.Strategy
60
+
61
+ wantlist WantList
62
+}
63
+
64
+type WantList struct {
65
+ lock sync.RWMutex
66
+ data map[u.Key]struct{}
67
+}
68
+
69
+func (wl *WantList) Add(k u.Key) {
70
+ u.DOut("Adding %v to Wantlist\n", k.Pretty())
71
+ wl.lock.Lock()
72
+ defer wl.lock.Unlock()
73
+
74
+ wl.data[k] = struct{}{}
75
+}
76
+
77
+func (wl *WantList) Remove(k u.Key) {
78
+ u.DOut("Removing %v from Wantlist\n", k.Pretty())
79
+ wl.lock.Lock()
80
+ defer wl.lock.Unlock()
81
+
82
+ delete(wl.data, k)
83
+}
84
+
85
+func (wl *WantList) Keys() []u.Key {
86
+ wl.lock.RLock()
87
+ defer wl.lock.RUnlock()
88
+ keys := make([]u.Key, 0)
89
+ for k, _ := range wl.data {
90
+ keys = append(keys, k)
91
+ }
92
+ return keys
93
}
94
95
// GetBlock attempts to retrieve a particular block from peers within the
@@ -60,9 +97,10 @@ type bitswap struct {
97
//
98
// TODO ensure only one active request per key
99
func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error) {
100
+ u.DOut("Get Block %v\n", k.Pretty())
101
102
ctx, cancelFunc := context.WithCancel(parent)
65
- // TODO add to wantlist
103
+ bs.wantlist.Add(k)
104
promise := bs.notifications.Subscribe(ctx, k)
105
106
const maxProviders = 20
@@ -70,6 +108,9 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
108
109
go func() {
110
message := bsmsg.New()
111
+ for _, wanted := range bs.wantlist.Keys() {
112
+ message.AppendWanted(wanted)
113
+ }
114
message.AppendWanted(k)
115
for iiiii := range peersToQuery {
116
// u.DOut("bitswap got peersToQuery: %s\n", iiiii)
@@ -94,6 +135,7 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
135
select {
136
case block := <-promise:
137
cancelFunc()
138
+ bs.wantlist.Remove(k)
139
// TODO remove from wantlist
140
return &block, nil
141
case <-parent.Done():
@@ -104,6 +146,8 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
146
// HasBlock announces the existance of a block to bitswap, potentially sending
147
// it to peers (Partners) whose WantLists include it.
148
func (bs *bitswap) HasBlock(ctx context.Context, blk blocks.Block) error {
149
+ u.DOut("Has Block %v\n", blk.Key().Pretty())
150
+ bs.wantlist.Remove(blk.Key())
151
bs.sendToPeersThatWant(ctx, blk)
152
return bs.routing.Provide(ctx, blk.Key())
153
}
@@ -111,6 +155,7 @@ func (bs *bitswap) HasBlock(ctx context.Context, blk blocks.Block) error {
155
// TODO(brian): handle errors
156
func (bs *bitswap) ReceiveMessage(ctx context.Context, p *peer.Peer, incoming bsmsg.BitSwapMessage) (
157
*peer.Peer, bsmsg.BitSwapMessage, error) {
158
+ u.DOut("ReceiveMessage from %v\n", p.Key().Pretty())
159
160
if p == nil {
161
return nil, nil, errors.New("Received nil Peer")
@@ -132,19 +177,21 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p *peer.Peer, incoming bs
177
}(block)
178
}
179
180
+ message := bsmsg.New()
181
+ for _, wanted := range bs.wantlist.Keys() {
182
+ message.AppendWanted(wanted)
183
+ }
184
for _, key := range incoming.Wantlist() {
185
if bs.strategy.ShouldSendBlockToPeer(key, p) {
137
- block, errBlockNotFound := bs.blockstore.Get(key)
138
- if errBlockNotFound != nil {
139
- return nil, nil, errBlockNotFound
186
+ if block, errBlockNotFound := bs.blockstore.Get(key); errBlockNotFound != nil {
187
+ continue
188
+ } else {
189
+ message.AppendBlock(*block)
190
}
141
- message := bsmsg.New()
142
- message.AppendBlock(*block)
143
- defer bs.strategy.MessageSent(p, message)
144
- return p, message, nil
191
}
192
}
147
- return nil, nil, nil
193
+ defer bs.strategy.MessageSent(p, message)
194
+ return p, message, nil
195
}
196
197
// send strives to ensure that accounting is always performed when a message is
@@ -155,11 +202,16 @@ func (bs *bitswap) send(ctx context.Context, p *peer.Peer, m bsmsg.BitSwapMessag
202
}
203
204
func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block blocks.Block) {
205
+ u.DOut("Sending %v to peers that want it\n", block.Key().Pretty())
206
for _, p := range bs.strategy.Peers() {
207
if bs.strategy.BlockIsWantedByPeer(block.Key(), p) {
208
+ u.DOut("%v wants %v\n", p.Key().Pretty(), block.Key().Pretty())
209
if bs.strategy.ShouldSendBlockToPeer(block.Key(), p) {
210
message := bsmsg.New()
211
message.AppendBlock(block)
212
+ for _, wanted := range bs.wantlist.Keys() {
213
+ message.AppendWanted(wanted)
214
+ }
215
go bs.send(ctx, p, message)
216
}
217
}
exchange/bitswap/bitswap_test.go
+32
-18
@@ -16,6 +16,7 @@ import (
16
strategy "github.com/jbenet/go-ipfs/exchange/bitswap/strategy"
17
tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
18
peer "github.com/jbenet/go-ipfs/peer"
19
+ util "github.com/jbenet/go-ipfs/util"
20
testutil "github.com/jbenet/go-ipfs/util/testutil"
21
)
22
@@ -145,7 +146,10 @@ func getOrFail(bitswap instance, b *blocks.Block, t *testing.T, wg *sync.WaitGro
146
wg.Done()
147
}
148
149
+// TODO simplify this test. get to the _essence_!
150
func TestSendToWantingPeer(t *testing.T) {
151
+ util.Debug = true
152
+
153
net := tn.VirtualNetwork()
154
rs := tn.VirtualRoutingServer()
155
sg := NewSessionGenerator(net, rs)
@@ -155,48 +159,55 @@ func TestSendToWantingPeer(t *testing.T) {
159
w := sg.Next()
160
o := sg.Next()
161
162
+ t.Logf("Session %v\n", me.peer.Key().Pretty())
163
+ t.Logf("Session %v\n", w.peer.Key().Pretty())
164
+ t.Logf("Session %v\n", o.peer.Key().Pretty())
165
+
166
alpha := bg.Next()
167
160
- const timeout = 100 * time.Millisecond
161
- const wait = 100 * time.Millisecond
168
+ const timeout = 1 * time.Millisecond // FIXME don't depend on time
169
163
- t.Log("Peer |w| attempts to get a file |alpha|. NB: alpha not available")
170
+ t.Logf("Peer %v attempts to get %v. NB: not available\n", w.peer.Key().Pretty(), alpha.Key().Pretty())
171
ctx, _ := context.WithTimeout(context.Background(), timeout)
172
_, err := w.exchange.Block(ctx, alpha.Key())
173
if err == nil {
167
- t.Error("Expected alpha to NOT be available")
174
+ t.Fatalf("Expected %v to NOT be available", alpha.Key().Pretty())
175
}
169
- time.Sleep(wait)
176
171
- t.Log("Peer |w| announces availability of a file |beta|")
177
beta := bg.Next()
178
+ t.Logf("Peer %v announes availability of %v\n", w.peer.Key().Pretty(), beta.Key().Pretty())
179
ctx, _ = context.WithTimeout(context.Background(), timeout)
180
+ if err := w.blockstore.Put(beta); err != nil {
181
+ t.Fatal(err)
182
+ }
183
w.exchange.HasBlock(ctx, beta)
175
- time.Sleep(wait)
184
177
- t.Log("I request and get |beta| from |w|. In the message, I receive |w|'s wants [alpha]")
178
- t.Log("I don't have alpha, but I keep it on my wantlist.")
185
+ t.Logf("%v gets %v from %v and discovers it wants %v\n", me.peer.Key().Pretty(), beta.Key().Pretty(), w.peer.Key().Pretty(), alpha.Key().Pretty())
186
ctx, _ = context.WithTimeout(context.Background(), timeout)
180
- me.exchange.Block(ctx, beta.Key())
181
- time.Sleep(wait)
187
+ if _, err := me.exchange.Block(ctx, beta.Key()); err != nil {
188
+ t.Fatal(err)
189
+ }
190
183
- t.Log("Peer |o| announces the availability of |alpha|")
191
+ t.Logf("%v announces availability of %v\n", o.peer.Key().Pretty(), alpha.Key().Pretty())
192
ctx, _ = context.WithTimeout(context.Background(), timeout)
193
+ if err := o.blockstore.Put(alpha); err != nil {
194
+ t.Fatal(err)
195
+ }
196
o.exchange.HasBlock(ctx, alpha)
186
- time.Sleep(wait)
197
188
- t.Log("I request |alpha| for myself.")
198
+ t.Logf("%v requests %v\n", me.peer.Key().Pretty(), alpha.Key().Pretty())
199
ctx, _ = context.WithTimeout(context.Background(), timeout)
190
- me.exchange.Block(ctx, alpha.Key())
191
- time.Sleep(wait)
200
+ if _, err := me.exchange.Block(ctx, alpha.Key()); err != nil {
201
+ t.Fatal(err)
202
+ }
203
193
- t.Log("After receiving |f| from |o|, I send it to the wanting peer |w|")
204
+ t.Logf("%v should now have %v\n", w.peer.Key().Pretty(), alpha.Key().Pretty())
205
block, err := w.blockstore.Get(alpha.Key())
206
if err != nil {
207
t.Fatal("Should not have received an error")
208
}
209
if block.Key() != alpha.Key() {
199
- t.Error("Expected to receive alpha from me")
210
+ t.Fatal("Expected to receive alpha from me")
211
}
212
}
213
@@ -278,6 +289,9 @@ func session(net tn.Network, rs tn.RoutingServer, id peer.ID) instance {
289
strategy: strategy.New(),
290
routing: htc,
291
sender: adapter,
292
+ wantlist: WantList{
293
+ data: make(map[util.Key]struct{}),
294
+ },
295
}
296
adapter.SetDelegate(bs)
297
return instance{