@cryptotaxi247 / kubo / commits / 8cdf56686

remove multilayered routing table from the DHT (for now)

Jeromy committed Dec 11, 2014 at 06:08 UTC 8cdf5668654ff60a478616bec3ad33bacf01bcd1
3 files changed +32 -62
routing/dht/dht.go
+20 -43
@@ -37,7 +37,7 @@ const doPinging = false
37 type IpfsDHT struct {
38 // Array of routing tables for differently distanced nodes
39 // NOTE: (currently, only a single table is used)
40 - routingTables []*kb.RoutingTable
40 + routingTable *kb.RoutingTable
41
42 // the network services we need
43 dialer inet.Dialer
@@ -80,10 +80,7 @@ func NewDHT(ctx context.Context, p peer.Peer, ps peer.Peerstore, dialer inet.Dia
80 dht.providers = NewProviderManager(dht.Context(), p.ID())
81 dht.AddCloserChild(dht.providers)
82
83 - dht.routingTables = make([]*kb.RoutingTable, 3)
84 - dht.routingTables[0] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Millisecond*1000)
85 - dht.routingTables[1] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Millisecond*1000)
86 - dht.routingTables[2] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Hour)
83 + dht.routingTable = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Minute)
84 dht.birth = time.Now()
85
86 dht.Validators = make(map[string]ValidatorFunc)
@@ -243,9 +240,9 @@ func (dht *IpfsDHT) putProvider(ctx context.Context, p peer.Peer, key string) er
240 }
241
242 func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p peer.Peer,
246 - key u.Key, level int) ([]byte, []peer.Peer, error) {
243 + key u.Key) ([]byte, []peer.Peer, error) {
244
248 - pmes, err := dht.getValueSingle(ctx, p, key, level)
245 + pmes, err := dht.getValueSingle(ctx, p, key)
246 if err != nil {
247 return nil, nil, err
248 }
@@ -265,7 +262,7 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p peer.Peer,
262
263 // TODO decide on providers. This probably shouldn't be happening.
264 if prv := pmes.GetProviderPeers(); prv != nil && len(prv) > 0 {
268 - val, err := dht.getFromPeerList(ctx, key, prv, level)
265 + val, err := dht.getFromPeerList(ctx, key, prv)
266 if err != nil {
267 return nil, nil, err
268 }
@@ -292,9 +289,9 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p peer.Peer,
289
290 // getValueSingle simply performs the get value RPC with the given parameters
291 func (dht *IpfsDHT) getValueSingle(ctx context.Context, p peer.Peer,
295 - key u.Key, level int) (*pb.Message, error) {
292 + key u.Key) (*pb.Message, error) {
293
297 - pmes := pb.NewMessage(pb.Message_GET_VALUE, string(key), level)
294 + pmes := pb.NewMessage(pb.Message_GET_VALUE, string(key), 0)
295 return dht.sendRequest(ctx, p, pmes)
296 }
297
@@ -303,7 +300,7 @@ func (dht *IpfsDHT) getValueSingle(ctx context.Context, p peer.Peer,
300 // one to get the value from? Or just connect to one at a time until we get a
301 // successful connection and request the value from it?
302 func (dht *IpfsDHT) getFromPeerList(ctx context.Context, key u.Key,
306 - peerlist []*pb.Message_Peer, level int) ([]byte, error) {
303 + peerlist []*pb.Message_Peer) ([]byte, error) {
304
305 for _, pinfo := range peerlist {
306 p, err := dht.ensureConnectedToPeer(ctx, pinfo)
@@ -312,7 +309,7 @@ func (dht *IpfsDHT) getFromPeerList(ctx context.Context, key u.Key,
309 continue
310 }
311
315 - pmes, err := dht.getValueSingle(ctx, p, key, level)
312 + pmes, err := dht.getValueSingle(ctx, p, key)
313 if err != nil {
314 log.Errorf("getFromPeers error: %s\n", err)
315 continue
@@ -379,47 +376,30 @@ func (dht *IpfsDHT) putLocal(key u.Key, value []byte) error {
376 return dht.datastore.Put(key.DsKey(), data)
377 }
378
382 -// Update signals to all routingTables to Update their last-seen status
379 +// Update signals the routingTable to Update its last-seen status
380 // on the given peer.
381 func (dht *IpfsDHT) Update(ctx context.Context, p peer.Peer) {
382 log.Event(ctx, "updatePeer", p)
386 - removedCount := 0
387 - for _, route := range dht.routingTables {
388 - removed := route.Update(p)
389 - // Only close the connection if no tables refer to this peer
390 - if removed != nil {
391 - removedCount++
392 - }
393 - }
394 -
395 - // Only close the connection if no tables refer to this peer
396 - // if removedCount == len(dht.routingTables) {
397 - // dht.network.ClosePeer(p)
398 - // }
399 - // ACTUALLY, no, let's not just close the connection. it may be connected
400 - // due to other things. it seems that we just need connection timeouts
401 - // after some deadline of inactivity.
383 + dht.routingTable.Update(p)
384 }
385
386 // FindLocal looks for a peer with a given ID connected to this dht and returns the peer and the table it was found in.
387 func (dht *IpfsDHT) FindLocal(id peer.ID) (peer.Peer, *kb.RoutingTable) {
406 - for _, table := range dht.routingTables {
407 - p := table.Find(id)
408 - if p != nil {
409 - return p, table
410 - }
388 + p := dht.routingTable.Find(id)
389 + if p != nil {
390 + return p, dht.routingTable
391 }
392 return nil, nil
393 }
394
395 // findPeerSingle asks peer 'p' if they know where the peer with id 'id' is
416 -func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.Peer, id peer.ID, level int) (*pb.Message, error) {
417 - pmes := pb.NewMessage(pb.Message_FIND_NODE, string(id), level)
396 +func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.Peer, id peer.ID) (*pb.Message, error) {
397 + pmes := pb.NewMessage(pb.Message_FIND_NODE, string(id), 0)
398 return dht.sendRequest(ctx, p, pmes)
399 }
400
421 -func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.Peer, key u.Key, level int) (*pb.Message, error) {
422 - pmes := pb.NewMessage(pb.Message_GET_PROVIDERS, string(key), level)
401 +func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.Peer, key u.Key) (*pb.Message, error) {
402 + pmes := pb.NewMessage(pb.Message_GET_PROVIDERS, string(key), 0)
403 return dht.sendRequest(ctx, p, pmes)
404 }
405
@@ -446,11 +426,8 @@ func (dht *IpfsDHT) addProviders(key u.Key, pbps []*pb.Message_Peer) []peer.Peer
426
427 // nearestPeersToQuery returns the routing tables closest peers.
428 func (dht *IpfsDHT) nearestPeersToQuery(pmes *pb.Message, count int) []peer.Peer {
449 - level := pmes.GetClusterLevel()
450 - cluster := dht.routingTables[level]
451 -
429 key := u.Key(pmes.GetKey())
453 - closer := cluster.NearestPeers(kb.ConvertKey(key), count)
430 + closer := dht.routingTable.NearestPeers(kb.ConvertKey(key), count)
431 return closer
432 }
433
@@ -537,7 +514,7 @@ func (dht *IpfsDHT) PingRoutine(t time.Duration) {
514 case <-tick:
515 id := make([]byte, 16)
516 rand.Read(id)
540 - peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(u.Key(id)), 5)
517 + peers := dht.routingTable.NearestPeers(kb.ConvertKey(u.Key(id)), 5)
518 for _, p := range peers {
519 ctx, _ := context.WithTimeout(dht.Context(), time.Second*5)
520 err := dht.Ping(ctx, p)
routing/dht/diag.go
+1 -1
@@ -36,7 +36,7 @@ func (dht *IpfsDHT) getDiagInfo() *diagInfo {
36 di.LifeSpan = time.Since(dht.birth)
37 di.Keys = nil // Currently no way to query datastore
38
39 - for _, p := range dht.routingTables[0].ListPeers() {
39 + for _, p := range dht.routingTable.ListPeers() {
40 d := connDiagInfo{p.GetLatency(), p.ID()}
41 di.Connections = append(di.Connections, d)
42 }
routing/dht/routing.go
+11 -18
@@ -38,11 +38,7 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
38 return err
39 }
40
41 - var peers []peer.Peer
42 - for _, route := range dht.routingTables {
43 - npeers := route.NearestPeers(kb.ConvertKey(key), KValue)
44 - peers = append(peers, npeers...)
45 - }
41 + peers := dht.routingTable.NearestPeers(kb.ConvertKey(key), KValue)
42
43 query := newQuery(key, dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
44 log.Debugf("%s PutValue qry part %v", dht.self, p)
@@ -71,9 +67,8 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
67 return val, nil
68 }
69
74 - // get closest peers in the routing tables
75 - routeLevel := 0
76 - closest := dht.routingTables[routeLevel].NearestPeers(kb.ConvertKey(key), PoolSize)
70 + // get closest peers in the routing table
71 + closest := dht.routingTable.NearestPeers(kb.ConvertKey(key), PoolSize)
72 if closest == nil || len(closest) == 0 {
73 log.Warning("Got no peers back from routing table!")
74 return nil, kb.ErrLookupFailure
@@ -82,7 +77,7 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
77 // setup the Query
78 query := newQuery(key, dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
79
85 - val, peers, err := dht.getValueOrPeers(ctx, p, key, routeLevel)
80 + val, peers, err := dht.getValueOrPeers(ctx, p, key)
81 if err != nil {
82 return nil, err
83 }
@@ -116,7 +111,7 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
111 func (dht *IpfsDHT) Provide(ctx context.Context, key u.Key) error {
112
113 dht.providers.AddProvider(key, dht.self)
119 - peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(key), PoolSize)
114 + peers := dht.routingTable.NearestPeers(kb.ConvertKey(key), PoolSize)
115 if len(peers) == 0 {
116 return nil
117 }
@@ -166,7 +161,7 @@ func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, co
161 // setup the Query
162 query := newQuery(key, dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
163
169 - pmes, err := dht.findProvidersSingle(ctx, p, key, 0)
164 + pmes, err := dht.findProvidersSingle(ctx, p, key)
165 if err != nil {
166 return nil, err
167 }
@@ -205,7 +200,7 @@ func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, co
200 return &dhtQueryResult{closerPeers: clpeers}, nil
201 })
202
208 - peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(key), AlphaValue)
203 + peers := dht.routingTable.NearestPeers(kb.ConvertKey(key), AlphaValue)
204 _, err := query.Run(ctx, peers)
205 if err != nil {
206 log.Errorf("FindProviders Query error: %s", err)
@@ -253,8 +248,7 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (peer.Peer, error)
248 return p, nil
249 }
250
256 - routeLevel := 0
257 - closest := dht.routingTables[routeLevel].NearestPeers(kb.ConvertPeerID(id), AlphaValue)
251 + closest := dht.routingTable.NearestPeers(kb.ConvertPeerID(id), AlphaValue)
252 if closest == nil || len(closest) == 0 {
253 return nil, kb.ErrLookupFailure
254 }
@@ -270,7 +264,7 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (peer.Peer, error)
264 // setup the Query
265 query := newQuery(u.Key(id), dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
266
273 - pmes, err := dht.findPeerSingle(ctx, p, id, routeLevel)
267 + pmes, err := dht.findPeerSingle(ctx, p, id)
268 if err != nil {
269 return nil, err
270 }
@@ -316,8 +310,7 @@ func (dht *IpfsDHT) FindPeersConnectedToPeer(ctx context.Context, id peer.ID) (<
310 peerchan := make(chan peer.Peer, asyncQueryBuffer)
311 peersSeen := map[string]peer.Peer{}
312
319 - routeLevel := 0
320 - closest := dht.routingTables[routeLevel].NearestPeers(kb.ConvertPeerID(id), AlphaValue)
313 + closest := dht.routingTable.NearestPeers(kb.ConvertPeerID(id), AlphaValue)
314 if closest == nil || len(closest) == 0 {
315 return nil, kb.ErrLookupFailure
316 }
@@ -325,7 +318,7 @@ func (dht *IpfsDHT) FindPeersConnectedToPeer(ctx context.Context, id peer.ID) (<
318 // setup the Query
319 query := newQuery(u.Key(id), dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
320
328 - pmes, err := dht.findPeerSingle(ctx, p, id, routeLevel)
321 + pmes, err := dht.findPeerSingle(ctx, p, id)
322 if err != nil {
323 return nil, err
324 }