fix up FindProvidersAsync
Jeromy committed
Oct 11, 2014 at 10:43 UTC
c77ed6d2aa7a39935d8a549c9cdde2bb1c204a1b
4 files changed
+20
-80
routing/dht/dht.go
+3
-3
@@ -368,8 +368,8 @@ func (dht *IpfsDHT) Update(p *peer.Peer) {
368
// after some deadline of inactivity.
369
}
370
371
-// Find looks for a peer with a given ID connected to this dht and returns the peer and the table it was found in.
372
-func (dht *IpfsDHT) Find(id peer.ID) (*peer.Peer, *kb.RoutingTable) {
371
+// FindLocal looks for a peer with a given ID connected to this dht and returns the peer and the table it was found in.
372
+func (dht *IpfsDHT) FindLocal(id peer.ID) (*peer.Peer, *kb.RoutingTable) {
373
for _, table := range dht.routingTables {
374
p := table.Find(id)
375
if p != nil {
@@ -465,7 +465,7 @@ func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (*peer.Peer, error) {
465
466
p, _ := dht.peerstore.Get(id)
467
if p == nil {
468
- p, _ = dht.Find(id)
468
+ p, _ = dht.FindLocal(id)
469
if p != nil {
470
panic("somehow peer not getting into peerstore")
471
}
routing/dht/dht_test.go
+9
-6
@@ -227,13 +227,16 @@ func TestProvides(t *testing.T) {
227
time.Sleep(time.Millisecond * 60)
228
229
ctxT, _ := context.WithTimeout(context.Background(), time.Second)
230
- provs, err := dhts[0].FindProviders(ctxT, u.Key("hello"))
231
- if err != nil {
232
- t.Fatal(err)
233
- }
230
+ provchan := dhts[0].FindProvidersAsync(ctxT, u.Key("hello"), 1)
231
235
- if len(provs) != 1 {
236
- t.Fatal("Didnt get back providers")
232
+ after := time.After(time.Second)
233
+ select {
234
+ case prov := <-provchan:
235
+ if prov == nil {
236
+ t.Fatal("Got back nil provider")
237
+ }
238
+ case <-after:
239
+ t.Fatal("Did not get a provider back.")
240
}
241
}
242
routing/dht/routing.go
+8
-67
@@ -3,6 +3,7 @@ package dht
3
import (
4
"bytes"
5
"encoding/json"
6
+ "sync"
7
8
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9
@@ -117,26 +118,7 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key u.Key) error {
118
return nil
119
}
120
120
-// NB: not actually async. Used to keep the interface consistent while the
121
-// actual async method, FindProvidersAsync2 is under construction
121
func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int) <-chan *peer.Peer {
123
- ch := make(chan *peer.Peer)
124
- providers, err := dht.FindProviders(ctx, key)
125
- if err != nil {
126
- close(ch)
127
- return ch
128
- }
129
- go func() {
130
- defer close(ch)
131
- for _, p := range providers {
132
- ch <- p
133
- }
134
- }()
135
- return ch
136
-}
137
-
138
-// FIXME: there's a bug here!
139
-func (dht *IpfsDHT) FindProvidersAsync2(ctx context.Context, key u.Key, count int) <-chan *peer.Peer {
122
peerOut := make(chan *peer.Peer, count)
123
go func() {
124
ps := newPeerSet()
@@ -151,9 +133,12 @@ func (dht *IpfsDHT) FindProvidersAsync2(ctx context.Context, key u.Key, count in
133
}
134
}
135
136
+ wg := new(sync.WaitGroup)
137
peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(key), AlphaValue)
138
for _, pp := range peers {
139
+ wg.Add(1)
140
go func(p *peer.Peer) {
141
+ defer wg.Done()
142
pmes, err := dht.findProvidersSingle(ctx, p, key, 0)
143
if err != nil {
144
log.Error("%s", err)
@@ -162,7 +147,8 @@ func (dht *IpfsDHT) FindProvidersAsync2(ctx context.Context, key u.Key, count in
147
dht.addPeerListAsync(key, pmes.GetProviderPeers(), ps, count, peerOut)
148
}(pp)
149
}
165
-
150
+ wg.Wait()
151
+ close(peerOut)
152
}()
153
return peerOut
154
}
@@ -186,61 +172,16 @@ func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*Message_Peer, ps *peerSet
172
}
173
}
174
189
-// FindProviders searches for peers who can provide the value for given key.
190
-func (dht *IpfsDHT) FindProviders(ctx context.Context, key u.Key) ([]*peer.Peer, error) {
191
- // get closest peer
192
- log.Debug("Find providers for: '%s'", key)
193
- p := dht.routingTables[0].NearestPeer(kb.ConvertKey(key))
194
- if p == nil {
195
- log.Warning("Got no nearest peer for find providers: '%s'", key)
196
- return nil, nil
197
- }
198
-
199
- for level := 0; level < len(dht.routingTables); {
200
-
201
- // attempt retrieving providers
202
- pmes, err := dht.findProvidersSingle(ctx, p, key, level)
203
- if err != nil {
204
- return nil, err
205
- }
206
-
207
- // handle providers
208
- provs := pmes.GetProviderPeers()
209
- if provs != nil {
210
- log.Debug("Got providers back from findProviders call!")
211
- return dht.addProviders(key, provs), nil
212
- }
213
-
214
- log.Debug("Didnt get providers, just closer peers.")
215
- closer := pmes.GetCloserPeers()
216
- if len(closer) == 0 {
217
- level++
218
- continue
219
- }
220
-
221
- np, err := dht.peerFromInfo(closer[0])
222
- if err != nil {
223
- log.Debug("no peerFromInfo")
224
- level++
225
- continue
226
- }
227
- p = np
228
- }
229
- return nil, u.ErrNotFound
230
-}
231
-
175
// Find specific Peer
233
-
176
// FindPeer searches for a peer with given ID.
177
func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (*peer.Peer, error) {
178
179
// Check if were already connected to them
238
- p, _ := dht.Find(id)
180
+ p, _ := dht.FindLocal(id)
181
if p != nil {
182
return p, nil
183
}
184
243
- // @whyrusleeping why is this here? doesn't the dht.Find above cover it?
185
routeLevel := 0
186
p = dht.routingTables[routeLevel].NearestPeer(kb.ConvertPeerID(id))
187
if p == nil {
@@ -277,7 +218,7 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (*peer.Peer, error
218
func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (*peer.Peer, error) {
219
220
// Check if were already connected to them
280
- p, _ := dht.Find(id)
221
+ p, _ := dht.FindLocal(id)
222
if p != nil {
223
return p, nil
224
}
routing/routing.go
-4
@@ -26,11 +26,7 @@ type IpfsRouting interface {
26
// Announce that this node can provide value for given key
27
Provide(context.Context, u.Key) error
28
29
- // FindProviders searches for peers who can provide the value for given key.
30
- FindProviders(context.Context, u.Key) ([]*peer.Peer, error)
31
-
29
// Find specific Peer
33
-
30
// FindPeer searches for a peer with given ID.
31
FindPeer(context.Context, peer.ID) (*peer.Peer, error)
32
}