@cryptotaxi247 / kubo / commits / 85273daaa

allow peers to realize that they are actually a provider for a value

Jeromy committed Sep 4, 2014 at 20:32 UTC 85273daaa505e4f876ce951513c869c04499a83f
4 files changed +68 -2
routing/dht/dht.go
+25 -1
@@ -2,6 +2,7 @@ package dht
2
3 import (
4 "bytes"
5 + "crypto/rand"
6 "fmt"
7 "sync"
8 "time"
@@ -59,7 +60,7 @@ func NewDHT(p *peer.Peer, net swarm.Network, dstore ds.Datastore) *IpfsDHT {
60 dht.netChan = net.GetChannel(swarm.PBWrapper_DHT_MESSAGE)
61 dht.datastore = dstore
62 dht.self = p
62 - dht.providers = NewProviderManager()
63 + dht.providers = NewProviderManager(p.ID)
64 dht.shutdown = make(chan struct{})
65
66 dht.routingTables = make([]*kb.RoutingTable, 3)
@@ -293,7 +294,15 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *PBDHTMessage) {
294 Response: true,
295 }
296
297 + has, err := dht.datastore.Has(ds.NewKey(pmes.GetKey()))
298 + if err != nil {
299 + dht.netChan.Errors <- err
300 + }
301 +
302 providers := dht.providers.GetProviders(u.Key(pmes.GetKey()))
303 + if has {
304 + providers = append(providers, dht.self)
305 + }
306 if providers == nil || len(providers) == 0 {
307 level := 0
308 if len(pmes.GetValue()) > 0 {
@@ -637,3 +646,18 @@ func (dht *IpfsDHT) peerFromInfo(pbp *PBDHTMessage_PBPeer) (*peer.Peer, error) {
646
647 return dht.network.GetConnection(peer.ID(pbp.GetId()), maddr)
648 }
649 +
650 +func (dht *IpfsDHT) loadProvidableKeys() error {
651 + kl := dht.datastore.KeyList()
652 + for _, k := range kl {
653 + dht.providers.AddProvider(u.Key(k.Bytes()), dht.self)
654 + }
655 + return nil
656 +}
657 +
658 +// Builds up list of peers by requesting random peer IDs
659 +func (dht *IpfsDHT) Bootstrap() {
660 + id := make([]byte, 16)
661 + rand.Read(id)
662 + dht.FindPeer(peer.ID(id), time.Second*10)
663 +}
routing/dht/providers.go
+22 -1
@@ -9,9 +9,13 @@ import (
9
10 type ProviderManager struct {
11 providers map[u.Key][]*providerInfo
12 + local map[u.Key]struct{}
13 + lpeer peer.ID
14 + getlocal chan chan []u.Key
15 newprovs chan *addProv
16 getprovs chan *getProv
17 halt chan struct{}
18 + period time.Duration
19 }
20
21 type addProv struct {
@@ -24,11 +28,13 @@ type getProv struct {
28 resp chan []*peer.Peer
29 }
30
27 -func NewProviderManager() *ProviderManager {
31 +func NewProviderManager(local peer.ID) *ProviderManager {
32 pm := new(ProviderManager)
33 pm.getprovs = make(chan *getProv)
34 pm.newprovs = make(chan *addProv)
35 pm.providers = make(map[u.Key][]*providerInfo)
36 + pm.getlocal = make(chan chan []u.Key)
37 + pm.local = make(map[u.Key]struct{})
38 pm.halt = make(chan struct{})
39 go pm.run()
40 return pm
@@ -39,6 +45,9 @@ func (pm *ProviderManager) run() {
45 for {
46 select {
47 case np := <-pm.newprovs:
48 + if np.val.ID.Equal(pm.lpeer) {
49 + pm.local[np.k] = struct{}{}
50 + }
51 pi := new(providerInfo)
52 pi.Creation = time.Now()
53 pi.Value = np.val
@@ -51,6 +60,12 @@ func (pm *ProviderManager) run() {
60 parr = append(parr, p.Value)
61 }
62 gp.resp <- parr
63 + case lc := <-pm.getlocal:
64 + var keys []u.Key
65 + for k, _ := range pm.local {
66 + keys = append(keys, k)
67 + }
68 + lc <- keys
69 case <-tick.C:
70 for k, provs := range pm.providers {
71 var filtered []*providerInfo
@@ -82,6 +97,12 @@ func (pm *ProviderManager) GetProviders(k u.Key) []*peer.Peer {
97 return <-gp.resp
98 }
99
100 +func (pm *ProviderManager) GetLocal() []u.Key {
101 + resp := make(chan []u.Key)
102 + pm.getlocal <- resp
103 + return <-resp
104 +}
105 +
106 func (pm *ProviderManager) Halt() {
107 pm.halt <- struct{}{}
108 }
routing/dht/providers_test.go new
+20
@@ -0,0 +1,20 @@
1 +package dht
2 +
3 +import (
4 + "testing"
5 +
6 + "github.com/jbenet/go-ipfs/peer"
7 + u "github.com/jbenet/go-ipfs/util"
8 +)
9 +
10 +func TestProviderManager(t *testing.T) {
11 + mid := peer.ID("testing")
12 + p := NewProviderManager(mid)
13 + a := u.Key("test")
14 + p.AddProvider(a, &peer.Peer{})
15 + resp := p.GetProviders(a)
16 + if len(resp) != 1 {
17 + t.Fatal("Could not retrieve provider.")
18 + }
19 + p.Halt()
20 +}
routing/dht/routing.go
+1
@@ -166,6 +166,7 @@ func (dht *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
166
167 // Provide makes this node announce that it can provide a value for the given key
168 func (dht *IpfsDHT) Provide(key u.Key) error {
169 + dht.providers.AddProvider(key, dht.self)
170 peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(key), PoolSize)
171 if len(peers) == 0 {
172 return kb.ErrLookupFailure