@cryptotaxi247 / kubo / commits / 8f91069e5

Write providers to disk to avoid memory leaks

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

Jeromy committed May 16, 2016 at 17:01 UTC 8f91069e50c1c755d43a905bb086d332fea01000
5 files changed +501 -229
routing/dht/dht.go
+4 -3
@@ -12,6 +12,7 @@ import (
12 key "github.com/ipfs/go-ipfs/blocks/key"
13 routing "github.com/ipfs/go-ipfs/routing"
14 pb "github.com/ipfs/go-ipfs/routing/dht/pb"
15 + providers "github.com/ipfs/go-ipfs/routing/dht/providers"
16 kb "github.com/ipfs/go-ipfs/routing/kbucket"
17 record "github.com/ipfs/go-ipfs/routing/record"
18
@@ -48,7 +49,7 @@ type IpfsDHT struct {
49 datastore ds.Datastore // Local data
50
51 routingTable *kb.RoutingTable // Array of routing tables for differently distanced nodes
51 - providers *ProviderManager
52 + providers *providers.ProviderManager
53
54 birth time.Time // When this peer started up
55 diaglock sync.Mutex // lock to make diagnostics work better
@@ -84,8 +85,8 @@ func NewDHT(ctx context.Context, h host.Host, dstore ds.Datastore) *IpfsDHT {
85 dht.ctx = ctx
86
87 h.SetStreamHandler(ProtocolDHT, dht.handleNewStream)
87 - dht.providers = NewProviderManager(dht.ctx, dht.self)
88 - dht.proc.AddChild(dht.providers.proc)
88 + dht.providers = providers.NewProviderManager(dht.ctx, dht.self, dstore)
89 + dht.proc.AddChild(dht.providers.Process())
90 goprocessctx.CloseAfterContext(dht.proc, ctx)
91
92 dht.routingTable = kb.NewRoutingTable(20, kb.ConvertPeerID(dht.self), time.Minute, dht.peerstore)
routing/dht/providers.go deleted
-165
@@ -1,165 +0,0 @@
1 -package dht
2 -
3 -import (
4 - "time"
5 -
6 - key "github.com/ipfs/go-ipfs/blocks/key"
7 - goprocess "gx/ipfs/QmQopLATEYMNg7dVqZRNDfeE2S1yKy8zrRh5xnYiuqeZBn/goprocess"
8 - goprocessctx "gx/ipfs/QmQopLATEYMNg7dVqZRNDfeE2S1yKy8zrRh5xnYiuqeZBn/goprocess/context"
9 - peer "gx/ipfs/QmRBqJF7hb8ZSpRcMwUt8hNhydWcxGEhtk81HKq6oUwKvs/go-libp2p-peer"
10 -
11 - context "gx/ipfs/QmZy2y8t9zQH2a1b8q2ZSLKp17ATuJoCNxxyMFG5qFExpt/go-net/context"
12 -)
13 -
14 -var ProvideValidity = time.Hour * 24
15 -var defaultCleanupInterval = time.Hour
16 -
17 -type ProviderManager struct {
18 - // all non channel fields are meant to be accessed only within
19 - // the run method
20 - providers map[key.Key]*providerSet
21 - local map[key.Key]struct{}
22 - lpeer peer.ID
23 -
24 - getlocal chan chan []key.Key
25 - newprovs chan *addProv
26 - getprovs chan *getProv
27 - period time.Duration
28 - proc goprocess.Process
29 -
30 - cleanupInterval time.Duration
31 -}
32 -
33 -type providerSet struct {
34 - providers []peer.ID
35 - set map[peer.ID]time.Time
36 -}
37 -
38 -type addProv struct {
39 - k key.Key
40 - val peer.ID
41 -}
42 -
43 -type getProv struct {
44 - k key.Key
45 - resp chan []peer.ID
46 -}
47 -
48 -func NewProviderManager(ctx context.Context, local peer.ID) *ProviderManager {
49 - pm := new(ProviderManager)
50 - pm.getprovs = make(chan *getProv)
51 - pm.newprovs = make(chan *addProv)
52 - pm.providers = make(map[key.Key]*providerSet)
53 - pm.getlocal = make(chan chan []key.Key)
54 - pm.local = make(map[key.Key]struct{})
55 - pm.proc = goprocessctx.WithContext(ctx)
56 - pm.cleanupInterval = defaultCleanupInterval
57 - pm.proc.Go(func(p goprocess.Process) { pm.run() })
58 -
59 - return pm
60 -}
61 -
62 -func (pm *ProviderManager) run() {
63 - tick := time.NewTicker(pm.cleanupInterval)
64 - for {
65 - select {
66 - case np := <-pm.newprovs:
67 - if np.val == pm.lpeer {
68 - pm.local[np.k] = struct{}{}
69 - }
70 - provs, ok := pm.providers[np.k]
71 - if !ok {
72 - provs = newProviderSet()
73 - pm.providers[np.k] = provs
74 - }
75 - provs.Add(np.val)
76 -
77 - case gp := <-pm.getprovs:
78 - var parr []peer.ID
79 - provs, ok := pm.providers[gp.k]
80 - if ok {
81 - parr = provs.providers
82 - }
83 -
84 - gp.resp <- parr
85 -
86 - case lc := <-pm.getlocal:
87 - var keys []key.Key
88 - for k := range pm.local {
89 - keys = append(keys, k)
90 - }
91 - lc <- keys
92 -
93 - case <-tick.C:
94 - for k, provs := range pm.providers {
95 - var filtered []peer.ID
96 - for p, t := range provs.set {
97 - if time.Now().Sub(t) > ProvideValidity {
98 - delete(provs.set, p)
99 - } else {
100 - filtered = append(filtered, p)
101 - }
102 - }
103 -
104 - if len(filtered) > 0 {
105 - provs.providers = filtered
106 - } else {
107 - delete(pm.providers, k)
108 - }
109 - }
110 -
111 - case <-pm.proc.Closing():
112 - return
113 - }
114 - }
115 -}
116 -
117 -func (pm *ProviderManager) AddProvider(ctx context.Context, k key.Key, val peer.ID) {
118 - prov := &addProv{
119 - k: k,
120 - val: val,
121 - }
122 - select {
123 - case pm.newprovs <- prov:
124 - case <-ctx.Done():
125 - }
126 -}
127 -
128 -func (pm *ProviderManager) GetProviders(ctx context.Context, k key.Key) []peer.ID {
129 - gp := &getProv{
130 - k: k,
131 - resp: make(chan []peer.ID, 1), // buffered to prevent sender from blocking
132 - }
133 - select {
134 - case <-ctx.Done():
135 - return nil
136 - case pm.getprovs <- gp:
137 - }
138 - select {
139 - case <-ctx.Done():
140 - return nil
141 - case peers := <-gp.resp:
142 - return peers
143 - }
144 -}
145 -
146 -func (pm *ProviderManager) GetLocal() []key.Key {
147 - resp := make(chan []key.Key)
148 - pm.getlocal <- resp
149 - return <-resp
150 -}
151 -
152 -func newProviderSet() *providerSet {
153 - return &providerSet{
154 - set: make(map[peer.ID]time.Time),
155 - }
156 -}
157 -
158 -func (ps *providerSet) Add(p peer.ID) {
159 - _, found := ps.set[p]
160 - if !found {
161 - ps.providers = append(ps.providers, p)
162 - }
163 -
164 - ps.set[p] = time.Now()
165 -}
routing/dht/providers/providers.go new
+363
@@ -0,0 +1,363 @@
1 +package providers
2 +
3 +import (
4 + "encoding/base32"
5 + "encoding/binary"
6 + "fmt"
7 + "strings"
8 + "time"
9 +
10 + logging "gx/ipfs/QmNQynaz7qfriSUJkiEZUrm2Wen1u3Kj9goZzWtrPyu7XR/go-log"
11 + goprocess "gx/ipfs/QmQopLATEYMNg7dVqZRNDfeE2S1yKy8zrRh5xnYiuqeZBn/goprocess"
12 + goprocessctx "gx/ipfs/QmQopLATEYMNg7dVqZRNDfeE2S1yKy8zrRh5xnYiuqeZBn/goprocess/context"
13 + peer "gx/ipfs/QmRBqJF7hb8ZSpRcMwUt8hNhydWcxGEhtk81HKq6oUwKvs/go-libp2p-peer"
14 + lru "gx/ipfs/QmVYxfoJQiZijTgPNHCHgHELvQpbsJNTg6Crmc3dQkj3yy/golang-lru"
15 + ds "gx/ipfs/QmZ6A6P6AMo8SR3jXAwzTuSU6B9R2Y4eqW2yW9VvfUayDN/go-datastore"
16 + dsq "gx/ipfs/QmZ6A6P6AMo8SR3jXAwzTuSU6B9R2Y4eqW2yW9VvfUayDN/go-datastore/query"
17 +
18 + key "github.com/ipfs/go-ipfs/blocks/key"
19 +
20 + context "gx/ipfs/QmZy2y8t9zQH2a1b8q2ZSLKp17ATuJoCNxxyMFG5qFExpt/go-net/context"
21 +)
22 +
23 +var log = logging.Logger("providers")
24 +
25 +var lruCacheSize = 256
26 +var ProvideValidity = time.Hour * 24
27 +var defaultCleanupInterval = time.Hour
28 +
29 +type ProviderManager struct {
30 + // all non channel fields are meant to be accessed only within
31 + // the run method
32 + providers *lru.Cache
33 + local map[key.Key]struct{}
34 + lpeer peer.ID
35 + dstore ds.Datastore
36 +
37 + getlocal chan chan []key.Key
38 + newprovs chan *addProv
39 + getprovs chan *getProv
40 + period time.Duration
41 + proc goprocess.Process
42 +
43 + cleanupInterval time.Duration
44 +}
45 +
46 +type providerSet struct {
47 + providers []peer.ID
48 + set map[peer.ID]time.Time
49 +}
50 +
51 +type addProv struct {
52 + k key.Key
53 + val peer.ID
54 +}
55 +
56 +type getProv struct {
57 + k key.Key
58 + resp chan []peer.ID
59 +}
60 +
61 +func NewProviderManager(ctx context.Context, local peer.ID, dstore ds.Datastore) *ProviderManager {
62 + pm := new(ProviderManager)
63 + pm.getprovs = make(chan *getProv)
64 + pm.newprovs = make(chan *addProv)
65 + pm.dstore = dstore
66 + cache, err := lru.New(lruCacheSize)
67 + if err != nil {
68 + panic(err) //only happens if negative value is passed to lru constructor
69 + }
70 + pm.providers = cache
71 +
72 + pm.getlocal = make(chan chan []key.Key)
73 + pm.local = make(map[key.Key]struct{})
74 + pm.proc = goprocessctx.WithContext(ctx)
75 + pm.cleanupInterval = defaultCleanupInterval
76 + pm.proc.Go(func(p goprocess.Process) { pm.run() })
77 +
78 + return pm
79 +}
80 +
81 +const providersKeyPrefix = "/providers/"
82 +
83 +func mkProvKey(k key.Key) ds.Key {
84 + return ds.NewKey(providersKeyPrefix + base32.StdEncoding.EncodeToString([]byte(k)))
85 +}
86 +
87 +func (pm *ProviderManager) Process() goprocess.Process {
88 + return pm.proc
89 +}
90 +
91 +func (pm *ProviderManager) providersForKey(k key.Key) ([]peer.ID, error) {
92 + pset, err := pm.getProvSet(k)
93 + if err != nil {
94 + return nil, err
95 + }
96 + return pset.providers, nil
97 +}
98 +
99 +func (pm *ProviderManager) getProvSet(k key.Key) (*providerSet, error) {
100 + cached, ok := pm.providers.Get(k)
101 + if ok {
102 + return cached.(*providerSet), nil
103 + }
104 +
105 + pset, err := loadProvSet(pm.dstore, k)
106 + if err != nil {
107 + return nil, err
108 + }
109 +
110 + if len(pset.providers) > 0 {
111 + pm.providers.Add(k, pset)
112 + }
113 +
114 + return pset, nil
115 +}
116 +
117 +func loadProvSet(dstore ds.Datastore, k key.Key) (*providerSet, error) {
118 + res, err := dstore.Query(dsq.Query{Prefix: mkProvKey(k).String()})
119 + if err != nil {
120 + return nil, err
121 + }
122 +
123 + out := newProviderSet()
124 + for e := range res.Next() {
125 + if e.Error != nil {
126 + log.Error("got an error: ", err)
127 + continue
128 + }
129 + parts := strings.Split(e.Key, "/")
130 + if len(parts) != 4 {
131 + log.Warning("incorrectly formatted key: ", e.Key)
132 + continue
133 + }
134 +
135 + decstr, err := base32.StdEncoding.DecodeString(parts[len(parts)-1])
136 + if err != nil {
137 + log.Error("base32 decoding error: ", err)
138 + continue
139 + }
140 +
141 + pid := peer.ID(decstr)
142 +
143 + t, err := readTimeValue(e.Value)
144 + if err != nil {
145 + log.Warning("parsing providers record from disk: ", err)
146 + continue
147 + }
148 +
149 + out.setVal(pid, t)
150 + }
151 +
152 + return out, nil
153 +}
154 +
155 +func readTimeValue(i interface{}) (time.Time, error) {
156 + data, ok := i.([]byte)
157 + if !ok {
158 + return time.Time{}, fmt.Errorf("data was not a []byte")
159 + }
160 +
161 + nsec, _ := binary.Varint(data)
162 +
163 + return time.Unix(0, nsec), nil
164 +}
165 +
166 +func (pm *ProviderManager) addProv(k key.Key, p peer.ID) error {
167 + iprovs, ok := pm.providers.Get(k)
168 + if !ok {
169 + iprovs = newProviderSet()
170 + pm.providers.Add(k, iprovs)
171 + }
172 + provs := iprovs.(*providerSet)
173 + now := time.Now()
174 + provs.setVal(p, now)
175 +
176 + return writeProviderEntry(pm.dstore, k, p, now)
177 +}
178 +
179 +func writeProviderEntry(dstore ds.Datastore, k key.Key, p peer.ID, t time.Time) error {
180 + dsk := mkProvKey(k).ChildString(base32.StdEncoding.EncodeToString([]byte(p)))
181 +
182 + buf := make([]byte, 16)
183 + n := binary.PutVarint(buf, t.UnixNano())
184 +
185 + return dstore.Put(dsk, buf[:n])
186 +}
187 +
188 +func (pm *ProviderManager) deleteProvSet(k key.Key) error {
189 + pm.providers.Remove(k)
190 +
191 + res, err := pm.dstore.Query(dsq.Query{
192 + KeysOnly: true,
193 + Prefix: mkProvKey(k).String(),
194 + })
195 +
196 + entries, err := res.Rest()
197 + if err != nil {
198 + return err
199 + }
200 +
201 + for _, e := range entries {
202 + err := pm.dstore.Delete(ds.NewKey(e.Key))
203 + if err != nil {
204 + log.Error("deleting provider set: ", err)
205 + }
206 + }
207 + return nil
208 +}
209 +
210 +func (pm *ProviderManager) getAllProvKeys() ([]key.Key, error) {
211 + res, err := pm.dstore.Query(dsq.Query{
212 + KeysOnly: true,
213 + Prefix: providersKeyPrefix,
214 + })
215 +
216 + if err != nil {
217 + return nil, err
218 + }
219 +
220 + entries, err := res.Rest()
221 + if err != nil {
222 + return nil, err
223 + }
224 +
225 + out := make([]key.Key, 0, len(entries))
226 + seen := make(map[key.Key]struct{})
227 + for _, e := range entries {
228 + parts := strings.Split(e.Key, "/")
229 + if len(parts) != 4 {
230 + log.Warning("incorrectly formatted provider entry in datastore")
231 + continue
232 + }
233 + decoded, err := base32.StdEncoding.DecodeString(parts[2])
234 + if err != nil {
235 + log.Warning("error decoding base32 provider key")
236 + continue
237 + }
238 +
239 + k := key.Key(decoded)
240 + if _, ok := seen[k]; !ok {
241 + out = append(out, key.Key(decoded))
242 + seen[k] = struct{}{}
243 + }
244 + }
245 +
246 + return out, nil
247 +}
248 +
249 +func (pm *ProviderManager) run() {
250 + tick := time.NewTicker(pm.cleanupInterval)
251 + for {
252 + select {
253 + case np := <-pm.newprovs:
254 + if np.val == pm.lpeer {
255 + pm.local[np.k] = struct{}{}
256 + }
257 + err := pm.addProv(np.k, np.val)
258 + if err != nil {
259 + log.Error("error adding new providers: ", err)
260 + }
261 + case gp := <-pm.getprovs:
262 + provs, err := pm.providersForKey(gp.k)
263 + if err != nil && err != ds.ErrNotFound {
264 + log.Error("error reading providers: ", err)
265 + }
266 +
267 + gp.resp <- provs
268 + case lc := <-pm.getlocal:
269 + var keys []key.Key
270 + for k := range pm.local {
271 + keys = append(keys, k)
272 + }
273 + lc <- keys
274 +
275 + case <-tick.C:
276 + keys, err := pm.getAllProvKeys()
277 + if err != nil {
278 + log.Error("Error loading provider keys: ", err)
279 + continue
280 + }
281 + for _, k := range keys {
282 + provs, err := pm.getProvSet(k)
283 + if err != nil {
284 + log.Error("error loading known provset: ", err)
285 + continue
286 + }
287 + var filtered []peer.ID
288 + for p, t := range provs.set {
289 + if time.Now().Sub(t) > ProvideValidity {
290 + delete(provs.set, p)
291 + } else {
292 + filtered = append(filtered, p)
293 + }
294 + }
295 +
296 + if len(filtered) > 0 {
297 + provs.providers = filtered
298 + } else {
299 + err := pm.deleteProvSet(k)
300 + if err != nil {
301 + log.Error("error deleting provider set: ", err)
302 + }
303 + }
304 + }
305 + case <-pm.proc.Closing():
306 + return
307 + }
308 + }
309 +}
310 +
311 +func (pm *ProviderManager) AddProvider(ctx context.Context, k key.Key, val peer.ID) {
312 + prov := &addProv{
313 + k: k,
314 + val: val,
315 + }
316 + select {
317 + case pm.newprovs <- prov:
318 + case <-ctx.Done():
319 + }
320 +}
321 +
322 +func (pm *ProviderManager) GetProviders(ctx context.Context, k key.Key) []peer.ID {
323 + gp := &getProv{
324 + k: k,
325 + resp: make(chan []peer.ID, 1), // buffered to prevent sender from blocking
326 + }
327 + select {
328 + case <-ctx.Done():
329 + return nil
330 + case pm.getprovs <- gp:
331 + }
332 + select {
333 + case <-ctx.Done():
334 + return nil
335 + case peers := <-gp.resp:
336 + return peers
337 + }
338 +}
339 +
340 +func (pm *ProviderManager) GetLocal() []key.Key {
341 + resp := make(chan []key.Key)
342 + pm.getlocal <- resp
343 + return <-resp
344 +}
345 +
346 +func newProviderSet() *providerSet {
347 + return &providerSet{
348 + set: make(map[peer.ID]time.Time),
349 + }
350 +}
351 +
352 +func (ps *providerSet) Add(p peer.ID) {
353 + ps.setVal(p, time.Now())
354 +}
355 +
356 +func (ps *providerSet) setVal(p peer.ID, t time.Time) {
357 + _, found := ps.set[p]
358 + if !found {
359 + ps.providers = append(ps.providers, p)
360 + }
361 +
362 + ps.set[p] = t
363 +}
routing/dht/providers/providers_test.go new
+134
@@ -0,0 +1,134 @@
1 +package providers
2 +
3 +import (
4 + "fmt"
5 + "testing"
6 + "time"
7 +
8 + key "github.com/ipfs/go-ipfs/blocks/key"
9 + peer "gx/ipfs/QmRBqJF7hb8ZSpRcMwUt8hNhydWcxGEhtk81HKq6oUwKvs/go-libp2p-peer"
10 + ds "gx/ipfs/QmZ6A6P6AMo8SR3jXAwzTuSU6B9R2Y4eqW2yW9VvfUayDN/go-datastore"
11 +
12 + context "gx/ipfs/QmZy2y8t9zQH2a1b8q2ZSLKp17ATuJoCNxxyMFG5qFExpt/go-net/context"
13 +)
14 +
15 +func TestProviderManager(t *testing.T) {
16 + ctx := context.Background()
17 + mid := peer.ID("testing")
18 + p := NewProviderManager(ctx, mid, ds.NewMapDatastore())
19 + a := key.Key("test")
20 + p.AddProvider(ctx, a, peer.ID("testingprovider"))
21 + resp := p.GetProviders(ctx, a)
22 + if len(resp) != 1 {
23 + t.Fatal("Could not retrieve provider.")
24 + }
25 + p.proc.Close()
26 +}
27 +
28 +func TestProvidersDatastore(t *testing.T) {
29 + old := lruCacheSize
30 + lruCacheSize = 10
31 + defer func() { lruCacheSize = old }()
32 +
33 + ctx := context.Background()
34 + mid := peer.ID("testing")
35 + p := NewProviderManager(ctx, mid, ds.NewMapDatastore())
36 + defer p.proc.Close()
37 +
38 + friend := peer.ID("friend")
39 + var keys []key.Key
40 + for i := 0; i < 100; i++ {
41 + k := key.Key(fmt.Sprint(i))
42 + keys = append(keys, k)
43 + p.AddProvider(ctx, k, friend)
44 + }
45 +
46 + for _, k := range keys {
47 + resp := p.GetProviders(ctx, k)
48 + if len(resp) != 1 {
49 + t.Fatal("Could not retrieve provider.")
50 + }
51 + if resp[0] != friend {
52 + t.Fatal("expected provider to be 'friend'")
53 + }
54 + }
55 +}
56 +
57 +func TestProvidersSerialization(t *testing.T) {
58 + dstore := ds.NewMapDatastore()
59 +
60 + k := key.Key("my key!")
61 + p := peer.ID("my peer")
62 + pt := time.Now()
63 +
64 + err := writeProviderEntry(dstore, k, p, pt)
65 + if err != nil {
66 + t.Fatal(err)
67 + }
68 +
69 + pset, err := loadProvSet(dstore, k)
70 + if err != nil {
71 + t.Fatal(err)
72 + }
73 +
74 + lt, ok := pset.set[p]
75 + if !ok {
76 + t.Fatal("failed to load set correctly")
77 + }
78 +
79 + if pt != lt {
80 + t.Fatal("time wasnt serialized correctly")
81 + }
82 +}
83 +
84 +func TestProvidesExpire(t *testing.T) {
85 + pval := ProvideValidity
86 + cleanup := defaultCleanupInterval
87 + ProvideValidity = time.Second / 2
88 + defaultCleanupInterval = time.Second / 2
89 + defer func() {
90 + ProvideValidity = pval
91 + defaultCleanupInterval = cleanup
92 + }()
93 +
94 + ctx := context.Background()
95 + mid := peer.ID("testing")
96 + p := NewProviderManager(ctx, mid, ds.NewMapDatastore())
97 +
98 + peers := []peer.ID{"a", "b"}
99 + var keys []key.Key
100 + for i := 0; i < 10; i++ {
101 + k := key.Key(i)
102 + keys = append(keys, k)
103 + p.AddProvider(ctx, k, peers[0])
104 + p.AddProvider(ctx, k, peers[1])
105 + }
106 +
107 + for i := 0; i < 10; i++ {
108 + out := p.GetProviders(ctx, keys[i])
109 + if len(out) != 2 {
110 + t.Fatal("expected providers to still be there")
111 + }
112 + }
113 +
114 + time.Sleep(time.Second)
115 + for i := 0; i < 10; i++ {
116 + out := p.GetProviders(ctx, keys[i])
117 + if len(out) > 2 {
118 + t.Fatal("expected providers to be cleaned up")
119 + }
120 + }
121 +
122 + if p.providers.Len() != 0 {
123 + t.Fatal("providers map not cleaned up")
124 + }
125 +
126 + allprovs, err := p.getAllProvKeys()
127 + if err != nil {
128 + t.Fatal(err)
129 + }
130 +
131 + if len(allprovs) != 0 {
132 + t.Fatal("expected everything to be cleaned out of the datastore")
133 + }
134 +}
routing/dht/providers_test.go deleted
-61
@@ -1,61 +0,0 @@
1 -package dht
2 -
3 -import (
4 - "testing"
5 - "time"
6 -
7 - key "github.com/ipfs/go-ipfs/blocks/key"
8 - peer "gx/ipfs/QmRBqJF7hb8ZSpRcMwUt8hNhydWcxGEhtk81HKq6oUwKvs/go-libp2p-peer"
9 -
10 - context "gx/ipfs/QmZy2y8t9zQH2a1b8q2ZSLKp17ATuJoCNxxyMFG5qFExpt/go-net/context"
11 -)
12 -
13 -func TestProviderManager(t *testing.T) {
14 - ctx := context.Background()
15 - mid := peer.ID("testing")
16 - p := NewProviderManager(ctx, mid)
17 - a := key.Key("test")
18 - p.AddProvider(ctx, a, peer.ID("testingprovider"))
19 - resp := p.GetProviders(ctx, a)
20 - if len(resp) != 1 {
21 - t.Fatal("Could not retrieve provider.")
22 - }
23 - p.proc.Close()
24 -}
25 -
26 -func TestProvidesExpire(t *testing.T) {
27 - ProvideValidity = time.Second
28 - defaultCleanupInterval = time.Second
29 -
30 - ctx := context.Background()
31 - mid := peer.ID("testing")
32 - p := NewProviderManager(ctx, mid)
33 -
34 - peers := []peer.ID{"a", "b"}
35 - var keys []key.Key
36 - for i := 0; i < 10; i++ {
37 - k := key.Key(i)
38 - keys = append(keys, k)
39 - p.AddProvider(ctx, k, peers[0])
40 - p.AddProvider(ctx, k, peers[1])
41 - }
42 -
43 - for i := 0; i < 10; i++ {
44 - out := p.GetProviders(ctx, keys[i])
45 - if len(out) != 2 {
46 - t.Fatal("expected providers to still be there")
47 - }
48 - }
49 -
50 - time.Sleep(time.Second * 3)
51 - for i := 0; i < 10; i++ {
52 - out := p.GetProviders(ctx, keys[i])
53 - if len(out) > 2 {
54 - t.Fatal("expected providers to be cleaned up")
55 - }
56 - }
57 -
58 - if len(p.providers) != 0 {
59 - t.Fatal("providers map not cleaned up")
60 - }
61 -}