@cryptotaxi247 / kubo / commits / 0773ab384

fix cleanup of empty provider sets

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

Jeromy committed Jun 3, 2016 at 13:52 UTC 0773ab3840809cafb6614d4e4aec1f4a12eb2faf
2 files changed +53 -4
routing/dht/providers.go
+15 -4
@@ -11,6 +11,9 @@ import (
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
@@ -23,6 +26,8 @@ type ProviderManager struct {
26 getprovs chan *getProv
27 period time.Duration
28 proc goprocess.Process
29 +
30 + cleanupInterval time.Duration
31 }
32
33 type providerSet struct {
@@ -48,13 +53,14 @@ func NewProviderManager(ctx context.Context, local peer.ID) *ProviderManager {
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() {
57 - tick := time.NewTicker(time.Hour)
63 + tick := time.NewTicker(pm.cleanupInterval)
64 for {
65 select {
66 case np := <-pm.newprovs:
@@ -85,16 +91,21 @@ func (pm *ProviderManager) run() {
91 lc <- keys
92
93 case <-tick.C:
88 - for _, provs := range pm.providers {
94 + for k, provs := range pm.providers {
95 var filtered []peer.ID
96 for p, t := range provs.set {
91 - if time.Now().Sub(t) > time.Hour*24 {
97 + if time.Now().Sub(t) > ProvideValidity {
98 delete(provs.set, p)
99 } else {
100 filtered = append(filtered, p)
101 }
102 }
97 - provs.providers = filtered
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():
routing/dht/providers_test.go
+38
@@ -2,6 +2,7 @@ package dht
2
3 import (
4 "testing"
5 + "time"
6
7 key "github.com/ipfs/go-ipfs/blocks/key"
8 peer "gx/ipfs/QmQGwpJy9P4yXZySmqkZEXCmbBpJUb8xntCv8Ca4taZwDC/go-libp2p-peer"
@@ -21,3 +22,40 @@ func TestProviderManager(t *testing.T) {
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 +}