@cryptotaxi247 / kubo / commits / 35ed8c460

change providers map and lock over to an agent based approach for managing providers

Jeromy committed Aug 18, 2014 at 20:38 UTC 35ed8c460cf61f3abc01793c10de0658ce2ab072
2 files changed +95 -48
routing/dht/dht.go
+12 -48
@@ -36,9 +36,7 @@ type IpfsDHT struct {
36 datastore ds.Datastore
37 dslock sync.Mutex
38
39 - // Map keys to peers that can provide their value
40 - providers map[u.Key][]*providerInfo
41 - providerLock sync.RWMutex
39 + providers *ProviderManager
40
41 // Signal to shutdown dht
42 shutdown chan struct{}
@@ -59,7 +57,7 @@ func NewDHT(p *peer.Peer, net swarm.Network) *IpfsDHT {
57 dht.network = net
58 dht.datastore = ds.NewMapDatastore()
59 dht.self = p
62 - dht.providers = make(map[u.Key][]*providerInfo)
60 + dht.providers = NewProviderManager()
61 dht.shutdown = make(chan struct{})
62
63 dht.routingTables = make([]*kb.RoutingTable, 3)
@@ -102,7 +100,6 @@ func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) (*peer.Peer, error) {
100 func (dht *IpfsDHT) handleMessages() {
101 u.DOut("Begin message handling routine")
102
105 - checkTimeouts := time.NewTicker(time.Minute * 5)
103 ch := dht.network.GetChan()
104 for {
105 select {
@@ -146,34 +143,18 @@ func (dht *IpfsDHT) handleMessages() {
143 dht.handlePing(mes.Peer, pmes)
144 case PBDHTMessage_DIAGNOSTIC:
145 dht.handleDiagnostic(mes.Peer, pmes)
146 + default:
147 + u.PErr("Recieved invalid message type")
148 }
149
150 case err := <-ch.Errors:
151 u.PErr("dht err: %s\n", err)
152 case <-dht.shutdown:
154 - checkTimeouts.Stop()
153 return
156 - case <-checkTimeouts.C:
157 - // Time to collect some garbage!
158 - dht.cleanExpiredProviders()
154 }
155 }
156 }
157
163 -func (dht *IpfsDHT) cleanExpiredProviders() {
164 - dht.providerLock.Lock()
165 - for k, parr := range dht.providers {
166 - var cleaned []*providerInfo
167 - for _, v := range parr {
168 - if time.Since(v.Creation) < time.Hour {
169 - cleaned = append(cleaned, v)
170 - }
171 - }
172 - dht.providers[k] = cleaned
173 - }
174 - dht.providerLock.Unlock()
175 -}
176 -
158 func (dht *IpfsDHT) putValueToNetwork(p *peer.Peer, key string, value []byte) error {
159 pmes := Message{
160 Type: PBDHTMessage_PUT_VALUE,
@@ -202,14 +183,10 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
183 resp.Value = iVal.([]byte)
184 } else if err == ds.ErrNotFound {
185 // Check if we know any providers for the requested value
205 - dht.providerLock.RLock()
206 - provs, ok := dht.providers[u.Key(pmes.GetKey())]
207 - dht.providerLock.RUnlock()
208 - if ok && len(provs) > 0 {
186 + provs := dht.providers.GetProviders(u.Key(pmes.GetKey()))
187 + if len(provs) > 0 {
188 u.DOut("handleGetValue returning %d provider[s]\n", len(provs))
210 - for _, prov := range provs {
211 - resp.Peers = append(resp.Peers, prov.Value)
212 - }
189 + resp.Peers = provs
190 resp.Success = true
191 } else {
192 // No providers?
@@ -313,9 +290,7 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *PBDHTMessage) {
290 Response: true,
291 }
292
316 - dht.providerLock.RLock()
317 - providers := dht.providers[u.Key(pmes.GetKey())]
318 - dht.providerLock.RUnlock()
293 + providers := dht.providers.GetProviders(u.Key(pmes.GetKey()))
294 if providers == nil || len(providers) == 0 {
295 level := 0
296 if len(pmes.GetValue()) > 0 {
@@ -329,9 +304,7 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *PBDHTMessage) {
304 resp.Peers = []*peer.Peer{closer}
305 }
306 } else {
332 - for _, prov := range providers {
333 - resp.Peers = append(resp.Peers, prov.Value)
334 - }
307 + resp.Peers = providers
308 resp.Success = true
309 }
310
@@ -345,9 +318,8 @@ type providerInfo struct {
318 }
319
320 func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *PBDHTMessage) {
348 - //TODO: need to implement TTLs on providers
321 key := u.Key(pmes.GetKey())
350 - dht.addProviderEntry(key, p)
322 + dht.providers.AddProvider(key, p)
323 }
324
325 // Halt stops all communications from this peer and shut down
@@ -356,14 +328,6 @@ func (dht *IpfsDHT) Halt() {
328 dht.network.Close()
329 }
330
359 -func (dht *IpfsDHT) addProviderEntry(key u.Key, p *peer.Peer) {
360 - u.DOut("Adding %s as provider for '%s'\n", p.Key().Pretty(), key)
361 - dht.providerLock.Lock()
362 - provs := dht.providers[key]
363 - dht.providers[key] = append(provs, &providerInfo{time.Now(), p})
364 - dht.providerLock.Unlock()
365 -}
366 -
331 // NOTE: not yet finished, low priority
332 func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *PBDHTMessage) {
333 seq := dht.routingTables[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
@@ -514,7 +478,7 @@ func (dht *IpfsDHT) getFromPeerList(key u.Key, timeout time.Duration,
478 u.DErr("getFromPeers error: %s\n", err)
479 continue
480 }
517 - dht.addProviderEntry(key, p)
481 + dht.providers.AddProvider(key, p)
482
483 // Make sure it was a successful get
484 if pmes.GetSuccess() && pmes.Value != nil {
@@ -656,7 +620,7 @@ func (dht *IpfsDHT) addPeerList(key u.Key, peers []*PBDHTMessage_PBPeer) []*peer
620 continue
621 }
622 }
659 - dht.addProviderEntry(key, p)
623 + dht.providers.AddProvider(key, p)
624 provArr = append(provArr, p)
625 }
626 return provArr
routing/dht/providers.go new
+83
@@ -0,0 +1,83 @@
1 +package dht
2 +
3 +import (
4 + "time"
5 +
6 + u "github.com/jbenet/go-ipfs/util"
7 + peer "github.com/jbenet/go-ipfs/peer"
8 +)
9 +
10 +type ProviderManager struct {
11 + providers map[u.Key][]*providerInfo
12 + newprovs chan *addProv
13 + getprovs chan *getProv
14 + halt chan struct{}
15 +}
16 +
17 +type addProv struct {
18 + k u.Key
19 + val *peer.Peer
20 +}
21 +
22 +type getProv struct {
23 + k u.Key
24 + resp chan []*peer.Peer
25 +}
26 +
27 +func NewProviderManager() *ProviderManager {
28 + pm := new(ProviderManager)
29 + pm.getprovs = make(chan *getProv)
30 + pm.newprovs = make(chan *addProv)
31 + pm.providers = make(map[u.Key][]*providerInfo)
32 + pm.halt = make(chan struct{})
33 + go pm.run()
34 + return pm
35 +}
36 +
37 +func (pm *ProviderManager) run() {
38 + tick := time.NewTicker(time.Hour)
39 + for {
40 + select {
41 + case np := <-pm.newprovs:
42 + pi := new(providerInfo)
43 + pi.Creation = time.Now()
44 + pi.Value = np.val
45 + arr := pm.providers[np.k]
46 + pm.providers[np.k] = append(arr, pi)
47 + case gp := <-pm.getprovs:
48 + var parr []*peer.Peer
49 + provs := pm.providers[gp.k]
50 + for _, p := range provs {
51 + parr = append(parr, p.Value)
52 + }
53 + gp.resp <- parr
54 + case <-tick.C:
55 + for k, provs := range pm.providers {
56 + var filtered []*providerInfo
57 + for _, p := range provs {
58 + if time.Now().Sub(p.Creation) < time.Hour * 24 {
59 + filtered = append(filtered, p)
60 + }
61 + }
62 + pm.providers[k] = filtered
63 + }
64 + case <-pm.halt:
65 + return
66 + }
67 + }
68 +}
69 +
70 +func (pm *ProviderManager) AddProvider(k u.Key, val *peer.Peer) {
71 + pm.newprovs <- &addProv{
72 + k: k,
73 + val: val,
74 + }
75 +}
76 +
77 +func (pm *ProviderManager) GetProviders(k u.Key) []*peer.Peer {
78 + gp := new(getProv)
79 + gp.k = k
80 + gp.resp = make(chan []*peer.Peer)
81 + pm.getprovs <- gp
82 + return <-gp.resp
83 +}