@cryptotaxi247 / kubo / commits / 762faa742

rewrite findpeer and other dht tweaks

Jeromy committed Oct 24, 2014 at 18:32 UTC 762faa742144a7d2854930c1cc24ef5e4f40e110
4 files changed +63 -83
routing/dht/dht.go
-7
@@ -384,18 +384,11 @@ func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.Peer, id peer.ID,
384 return dht.sendRequest(ctx, p, pmes)
385 }
386
387 -func (dht *IpfsDHT) printTables() {
388 - for _, route := range dht.routingTables {
389 - route.Print()
390 - }
391 -}
392 -
387 func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.Peer, key u.Key, level int) (*Message, error) {
388 pmes := newMessage(Message_GET_PROVIDERS, string(key), level)
389 return dht.sendRequest(ctx, p, pmes)
390 }
391
398 -// TODO: Could be done async
392 func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []peer.Peer {
393 var provArr []peer.Peer
394 for _, prov := range peers {
routing/dht/dht_test.go
+7 -1
@@ -283,7 +283,13 @@ func TestProvidesAsync(t *testing.T) {
283 ctxT, _ := context.WithTimeout(ctx, time.Millisecond*300)
284 provs := dhts[0].FindProvidersAsync(ctxT, u.Key("hello"), 5)
285 select {
286 - case p := <-provs:
286 + case p, ok := <-provs:
287 + if !ok {
288 + t.Fatal("Provider channel was closed...")
289 + }
290 + if p == nil {
291 + t.Fatal("Got back nil provider!")
292 + }
293 if !p.ID().Equal(dhts[3].self.ID()) {
294 t.Fatalf("got a provider, but not the right one. %s", p)
295 }
routing/dht/routing.go
+56 -74
@@ -61,7 +61,7 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
61 closest := dht.routingTables[routeLevel].NearestPeers(kb.ConvertKey(key), PoolSize)
62 if closest == nil || len(closest) == 0 {
63 log.Warning("Got no peers back from routing table!")
64 - return nil, nil
64 + return nil, kb.ErrLookupFailure
65 }
66
67 // setup the Query
@@ -152,22 +152,32 @@ func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int
152 return peerOut
153 }
154
155 -//TODO: this function could also be done asynchronously
155 func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*Message_Peer, ps *peerSet, count int, out chan peer.Peer) {
156 + done := make(chan struct{})
157 for _, pbp := range peers {
158 + go func(mp *Message_Peer) {
159 + defer func() { done <- struct{}{} }()
160 + // construct new peer
161 + p, err := dht.ensureConnectedToPeer(mp)
162 + if err != nil {
163 + log.Error("%s", err)
164 + return
165 + }
166 + if p == nil {
167 + log.Error("Got nil peer from ensureConnectedToPeer")
168 + return
169 + }
170
159 - // construct new peer
160 - p, err := dht.ensureConnectedToPeer(pbp)
161 - if err != nil {
162 - continue
163 - }
164 -
165 - dht.providers.AddProvider(k, p)
166 - if ps.AddIfSmallerThan(p, count) {
167 - out <- p
168 - } else if ps.Size() >= count {
169 - return
170 - }
171 + dht.providers.AddProvider(k, p)
172 + if ps.AddIfSmallerThan(p, count) {
173 + out <- p
174 + } else if ps.Size() >= count {
175 + return
176 + }
177 + }(pbp)
178 + }
179 + for _ = range peers {
180 + <-done
181 }
182 }
183
@@ -182,92 +192,64 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (peer.Peer, error)
192 }
193
194 routeLevel := 0
185 - p = dht.routingTables[routeLevel].NearestPeer(kb.ConvertPeerID(id))
186 - if p == nil {
187 - return nil, nil
188 - }
189 - if p.ID().Equal(id) {
190 - return p, nil
195 + closest := dht.routingTables[routeLevel].NearestPeers(kb.ConvertPeerID(id), AlphaValue)
196 + if closest == nil || len(closest) == 0 {
197 + return nil, kb.ErrLookupFailure
198 }
199
193 - for routeLevel < len(dht.routingTables) {
194 - pmes, err := dht.findPeerSingle(ctx, p, id, routeLevel)
195 - plist := pmes.GetCloserPeers()
196 - if plist == nil || len(plist) == 0 {
197 - routeLevel++
198 - continue
200 + // Sanity...
201 + for _, p := range closest {
202 + if p.ID().Equal(id) {
203 + log.Error("Found target peer in list of closest peers...")
204 + return p, nil
205 }
200 - found := plist[0]
201 -
202 - nxtPeer, err := dht.ensureConnectedToPeer(found)
203 - if err != nil {
204 - routeLevel++
205 - continue
206 - }
207 -
208 - if nxtPeer.ID().Equal(id) {
209 - return nxtPeer, nil
210 - }
211 -
212 - p = nxtPeer
213 - }
214 - return nil, u.ErrNotFound
215 -}
216 -
217 -func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (peer.Peer, error) {
218 -
219 - // Check if were already connected to them
220 - p, _ := dht.FindLocal(id)
221 - if p != nil {
222 - return p, nil
223 - }
224 -
225 - // get the peers we need to announce to
226 - routeLevel := 0
227 - peers := dht.routingTables[routeLevel].NearestPeers(kb.ConvertPeerID(id), AlphaValue)
228 - if len(peers) == 0 {
229 - return nil, nil
206 }
207
232 - // setup query function
208 + // setup the Query
209 query := newQuery(u.Key(id), dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
210 +
211 pmes, err := dht.findPeerSingle(ctx, p, id, routeLevel)
212 if err != nil {
236 - log.Error("%s getPeer error: %v", dht.self, err)
213 return nil, err
214 }
215
240 - plist := pmes.GetCloserPeers()
241 - if len(plist) == 0 {
242 - routeLevel++
243 - }
244 -
245 - nxtprs := make([]peer.Peer, len(plist))
246 - for i, fp := range plist {
247 - nxtp, err := dht.peerFromInfo(fp)
216 + closer := pmes.GetCloserPeers()
217 + var clpeers []peer.Peer
218 + for _, pbp := range closer {
219 + np, err := dht.getPeer(peer.ID(pbp.GetId()))
220 if err != nil {
249 - log.Error("%s findPeer error: %v", dht.self, err)
221 + log.Warning("Received invalid peer from query")
222 continue
223 }
252 -
253 - if nxtp.ID().Equal(id) {
254 - return &dhtQueryResult{peer: nxtp, success: true}, nil
224 + ma, err := pbp.Address()
225 + if err != nil {
226 + log.Warning("Received peer with bad or missing address.")
227 + continue
228 }
256 -
257 - nxtprs[i] = nxtp
229 + np.AddAddress(ma)
230 + if pbp.GetId() == string(id) {
231 + return &dhtQueryResult{
232 + peer: np,
233 + success: true,
234 + }, nil
235 + }
236 + clpeers = append(clpeers, np)
237 }
238
260 - return &dhtQueryResult{closerPeers: nxtprs}, nil
239 + return &dhtQueryResult{closerPeers: clpeers}, nil
240 })
241
263 - result, err := query.Run(ctx, peers)
242 + // run it!
243 + result, err := query.Run(ctx, closest)
244 if err != nil {
245 return nil, err
246 }
247
248 + log.Debug("FindPeer %v %v", id, result.success)
249 if result.peer == nil {
250 return nil, u.ErrNotFound
251 }
252 +
253 return result.peer, nil
254 }
255
routing/kbucket/table.go
-1
@@ -88,7 +88,6 @@ func (rt *RoutingTable) nextBucket() peer.Peer {
88 newBucket := bucket.Split(len(rt.Buckets)-1, rt.local)
89 rt.Buckets = append(rt.Buckets, newBucket)
90 if newBucket.len() > rt.bucketsize {
91 - // TODO: This is a very rare and annoying case
91 return rt.nextBucket()
92 }
93