@cryptotaxi247 / kubo / commits / 44c874611

a few small changes to make the dht more efficient

License: MIT Signed-off-by: Jeromy <why@ipfs.io>

Jeromy committed Jun 11, 2016 at 17:08 UTC 44c874611400f6663c3bf4a1368acf17e70989a8
3 files changed +25 -31
routing/dht/dht.go
-28
@@ -117,34 +117,6 @@ func (dht *IpfsDHT) putValueToPeer(ctx context.Context, p peer.ID,
117 return nil
118 }
119
120 -// putProvider sends a message to peer 'p' saying that the local node
121 -// can provide the value of 'key'
122 -func (dht *IpfsDHT) putProvider(ctx context.Context, p peer.ID, skey string) error {
123 -
124 - // add self as the provider
125 - pi := pstore.PeerInfo{
126 - ID: dht.self,
127 - Addrs: dht.host.Addrs(),
128 - }
129 -
130 - // // only share WAN-friendly addresses ??
131 - // pi.Addrs = addrutil.WANShareableAddrs(pi.Addrs)
132 - if len(pi.Addrs) < 1 {
133 - // log.Infof("%s putProvider: %s for %s error: no wan-friendly addresses", dht.self, p, key.Key(key), pi.Addrs)
134 - return fmt.Errorf("no known addresses for self. cannot put provider.")
135 - }
136 -
137 - pmes := pb.NewMessage(pb.Message_ADD_PROVIDER, skey, 0)
138 - pmes.ProviderPeers = pb.RawPeerInfosToPBPeers([]pstore.PeerInfo{pi})
139 - err := dht.sendMessage(ctx, p, pmes)
140 - if err != nil {
141 - return err
142 - }
143 -
144 - log.Debugf("%s putProvider: %s for %s (%s)", dht.self, p, key.Key(skey), pi.Addrs)
145 - return nil
146 -}
147 -
120 var errInvalidRecord = errors.New("received invalid record")
121
122 // getValueOrPeers queries a particular peer p for the value for
routing/dht/routing.go
+23 -1
@@ -2,6 +2,7 @@ package dht
2
3 import (
4 "bytes"
5 + "fmt"
6 "sync"
7 "time"
8
@@ -243,13 +244,18 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key key.Key) error {
244 return err
245 }
246
247 + mes, err := dht.makeProvRecord(key)
248 + if err != nil {
249 + return err
250 + }
251 +
252 wg := sync.WaitGroup{}
253 for p := range peers {
254 wg.Add(1)
255 go func(p peer.ID) {
256 defer wg.Done()
257 log.Debugf("putProvider(%s, %s)", key, p)
252 - err := dht.putProvider(ctx, p, string(key))
258 + err := dht.sendMessage(ctx, p, mes)
259 if err != nil {
260 log.Debug(err)
261 }
@@ -258,6 +264,22 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key key.Key) error {
264 wg.Wait()
265 return nil
266 }
267 +func (dht *IpfsDHT) makeProvRecord(skey key.Key) (*pb.Message, error) {
268 + pi := pstore.PeerInfo{
269 + ID: dht.self,
270 + Addrs: dht.host.Addrs(),
271 + }
272 +
273 + // // only share WAN-friendly addresses ??
274 + // pi.Addrs = addrutil.WANShareableAddrs(pi.Addrs)
275 + if len(pi.Addrs) < 1 {
276 + return nil, fmt.Errorf("no known addresses for self. cannot put provider.")
277 + }
278 +
279 + pmes := pb.NewMessage(pb.Message_ADD_PROVIDER, string(skey), 0)
280 + pmes.ProviderPeers = pb.RawPeerInfosToPBPeers([]pstore.PeerInfo{pi})
281 + return pmes, nil
282 +}
283
284 // FindProviders searches until the context expires.
285 func (dht *IpfsDHT) FindProviders(ctx context.Context, key key.Key) ([]pstore.PeerInfo, error) {
routing/kbucket/table.go
+2 -2
@@ -48,11 +48,11 @@ func NewRoutingTable(bucketsize int, localID ID, latency time.Duration, m pstore
48 // Update adds or moves the given peer to the front of its respective bucket
49 // If a peer gets removed from a bucket, it is returned
50 func (rt *RoutingTable) Update(p peer.ID) {
51 - rt.tabLock.Lock()
52 - defer rt.tabLock.Unlock()
51 peerID := ConvertPeerID(p)
52 cpl := commonPrefixLen(peerID, rt.local)
53
54 + rt.tabLock.Lock()
55 + defer rt.tabLock.Unlock()
56 bucketID := cpl
57 if bucketID >= len(rt.Buckets) {
58 bucketID = len(rt.Buckets) - 1