p2p/peer/addr: manager with expirations
Juan Batiz-Benet committed
Feb 2, 2015 at 10:33 UTC
8b98c8fbdc219ab6639a62f90ba99dc60aea374c
2 files changed
+334
p2p/peer/addr/addr_manager.go
new
+152
@@ -0,0 +1,152 @@
1
+// package addr provides useful address utilities for p2p
2
+// applications. It buys into the multi-transport addressing
3
+// scheme Multiaddr, and uses it to build its own p2p addressing.
4
+// All Addrs must have an associated peer.ID.
5
+package addr
6
+
7
+import (
8
+ "sync"
9
+ "time"
10
+
11
+ ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
12
+
13
+ peer "github.com/jbenet/go-ipfs/p2p/peer"
14
+)
15
+
16
+type expiringAddr struct {
17
+ Addr ma.Multiaddr
18
+ TTL time.Time
19
+}
20
+
21
+func (e *expiringAddr) ExpiredBy(t time.Time) bool {
22
+ return t.After(e.TTL)
23
+}
24
+
25
+type addrSet map[string]expiringAddr
26
+
27
+// Manager manages addresses.
28
+// The zero-value is ready to be used.
29
+type Manager struct {
30
+ addrmu sync.Mutex // guards addrs
31
+ addrs map[peer.ID]addrSet
32
+}
33
+
34
+// ensures the Manager is initialized.
35
+// So we can use the zero value.
36
+func (mgr *Manager) init() {
37
+ if mgr.addrs == nil {
38
+ mgr.addrs = make(map[peer.ID]addrSet)
39
+ }
40
+}
41
+
42
+// AddAddr calls AddAddrs(p, []ma.Multiaddr{addr}, ttl)
43
+func (mgr *Manager) AddAddr(p peer.ID, addr ma.Multiaddr, ttl time.Duration) {
44
+ mgr.AddAddrs(p, []ma.Multiaddr{addr}, ttl)
45
+}
46
+
47
+// AddAddrs gives Manager addresses to use, with a given ttl
48
+// (time-to-live), after which the address is no longer valid.
49
+// If the manager has a longer TTL, the operation is a no-op for that address
50
+func (mgr *Manager) AddAddrs(p peer.ID, addrs []ma.Multiaddr, ttl time.Duration) {
51
+ mgr.addrmu.Lock()
52
+ defer mgr.addrmu.Unlock()
53
+
54
+ // if ttl is zero, exit. nothing to do.
55
+ if ttl <= 0 {
56
+ return
57
+ }
58
+
59
+ // so zero value can be used
60
+ mgr.init()
61
+
62
+ amap, found := mgr.addrs[p]
63
+ if !found {
64
+ amap = make(addrSet)
65
+ mgr.addrs[p] = amap
66
+ }
67
+
68
+ // only expand ttls
69
+ exp := time.Now().Add(ttl)
70
+ for _, addr := range addrs {
71
+ addrstr := addr.String()
72
+ a, found := amap[addrstr]
73
+ if !found || exp.After(a.TTL) {
74
+ amap[addrstr] = expiringAddr{Addr: addr, TTL: exp}
75
+ }
76
+ }
77
+}
78
+
79
+// SetAddr calls mgr.SetAddrs(p, addr, ttl)
80
+func (mgr *Manager) SetAddr(p peer.ID, addr ma.Multiaddr, ttl time.Duration) {
81
+ mgr.SetAddrs(p, []ma.Multiaddr{addr}, ttl)
82
+}
83
+
84
+// SetAddrs sets the ttl on addresses. This clears any TTL there previously.
85
+// This is used when we receive the best estimate of the validity of an address.
86
+func (mgr *Manager) SetAddrs(p peer.ID, addrs []ma.Multiaddr, ttl time.Duration) {
87
+ mgr.addrmu.Lock()
88
+ defer mgr.addrmu.Unlock()
89
+
90
+ // so zero value can be used
91
+ mgr.init()
92
+
93
+ amap, found := mgr.addrs[p]
94
+ if !found {
95
+ amap = make(addrSet)
96
+ mgr.addrs[p] = amap
97
+ }
98
+
99
+ exp := time.Now().Add(ttl)
100
+ for _, addr := range addrs {
101
+ // re-set all of them for new ttl.
102
+ addrs := addr.String()
103
+
104
+ if ttl > 0 {
105
+ amap[addrs] = expiringAddr{Addr: addr, TTL: exp}
106
+ } else {
107
+ delete(amap, addrs)
108
+ }
109
+ }
110
+}
111
+
112
+// Addresses returns all known (and valid) addresses for a given peer.
113
+func (mgr *Manager) Addrs(p peer.ID) []ma.Multiaddr {
114
+ mgr.addrmu.Lock()
115
+ defer mgr.addrmu.Unlock()
116
+
117
+ // not initialized? nothing to give.
118
+ if mgr.addrs == nil {
119
+ return nil
120
+ }
121
+
122
+ maddrs, found := mgr.addrs[p]
123
+ if !found {
124
+ return nil
125
+ }
126
+
127
+ now := time.Now()
128
+ good := make([]ma.Multiaddr, 0, len(maddrs))
129
+ var expired []string
130
+ for s, m := range maddrs {
131
+ if m.ExpiredBy(now) {
132
+ expired = append(expired, s)
133
+ } else {
134
+ good = append(good, m.Addr)
135
+ }
136
+ }
137
+
138
+ // clean up the expired ones.
139
+ for _, s := range expired {
140
+ delete(maddrs, s)
141
+ }
142
+ return good
143
+}
144
+
145
+// ClearAddresses removes all previously stored addresses
146
+func (mgr *Manager) ClearAddrs(p peer.ID) {
147
+ mgr.addrmu.Lock()
148
+ defer mgr.addrmu.Unlock()
149
+ mgr.init()
150
+
151
+ mgr.addrs[p] = make(addrSet) // clear what was there before
152
+}
p2p/peer/addr/addr_manager_test.go
new
+182
@@ -0,0 +1,182 @@
1
+package addr
2
+
3
+import (
4
+ "testing"
5
+ "time"
6
+
7
+ ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
8
+
9
+ peer "github.com/jbenet/go-ipfs/p2p/peer"
10
+)
11
+
12
+func IDS(t *testing.T, ids string) peer.ID {
13
+ id, err := peer.IDB58Decode(ids)
14
+ if err != nil {
15
+ t.Fatal(err)
16
+ }
17
+ return id
18
+}
19
+
20
+func MA(t *testing.T, m string) ma.Multiaddr {
21
+ maddr, err := ma.NewMultiaddr(m)
22
+ if err != nil {
23
+ t.Fatal(err)
24
+ }
25
+ return maddr
26
+}
27
+
28
+func testHas(t *testing.T, exp, act []ma.Multiaddr) {
29
+ if len(exp) != len(act) {
30
+ t.Fatal("lengths not the same")
31
+ }
32
+
33
+ for _, a := range exp {
34
+ found := false
35
+
36
+ for _, b := range act {
37
+ if a.Equal(b) {
38
+ found = true
39
+ break
40
+ }
41
+ }
42
+
43
+ if !found {
44
+ t.Fatal("expected address %s not found", a)
45
+ }
46
+ }
47
+}
48
+
49
+func TestAddresses(t *testing.T) {
50
+
51
+ id1 := IDS(t, "QmcNstKuwBBoVTpSCSDrwzjgrRcaYXK833Psuz2EMHwyQN")
52
+ id2 := IDS(t, "QmRmPL3FDZKE3Qiwv1RosLdwdvbvg17b2hB39QPScgWKKZ")
53
+ id3 := IDS(t, "QmPhi7vBsChP7sjRoZGgg7bcKqF6MmCcQwvRbDte8aJ6Kn")
54
+ id4 := IDS(t, "QmPhi7vBsChP7sjRoZGgg7bcKqF6MmCcQwvRbDte8aJ5Kn")
55
+ id5 := IDS(t, "QmPhi7vBsChP7sjRoZGgg7bcKqF6MmCcQwvRbDte8aJ5Km")
56
+
57
+ ma11 := MA(t, "/ip4/1.2.3.1/tcp/1111")
58
+ ma21 := MA(t, "/ip4/2.2.3.2/tcp/1111")
59
+ ma22 := MA(t, "/ip4/2.2.3.2/tcp/2222")
60
+ ma31 := MA(t, "/ip4/3.2.3.3/tcp/1111")
61
+ ma32 := MA(t, "/ip4/3.2.3.3/tcp/2222")
62
+ ma33 := MA(t, "/ip4/3.2.3.3/tcp/3333")
63
+ ma41 := MA(t, "/ip4/4.2.3.3/tcp/1111")
64
+ ma42 := MA(t, "/ip4/4.2.3.3/tcp/2222")
65
+ ma43 := MA(t, "/ip4/4.2.3.3/tcp/3333")
66
+ ma44 := MA(t, "/ip4/4.2.3.3/tcp/4444")
67
+ ma51 := MA(t, "/ip4/5.2.3.3/tcp/1111")
68
+ ma52 := MA(t, "/ip4/5.2.3.3/tcp/2222")
69
+ ma53 := MA(t, "/ip4/5.2.3.3/tcp/3333")
70
+ ma54 := MA(t, "/ip4/5.2.3.3/tcp/4444")
71
+ ma55 := MA(t, "/ip4/5.2.3.3/tcp/5555")
72
+
73
+ ttl := time.Hour
74
+ m := Manager{}
75
+ m.AddAddr(id1, ma11, ttl)
76
+
77
+ m.AddAddrs(id2, []ma.Multiaddr{ma21, ma22}, ttl)
78
+ m.AddAddrs(id2, []ma.Multiaddr{ma21, ma22}, ttl) // idempotency
79
+
80
+ m.AddAddr(id3, ma31, ttl)
81
+ m.AddAddr(id3, ma32, ttl)
82
+ m.AddAddr(id3, ma33, ttl)
83
+ m.AddAddr(id3, ma33, ttl) // idempotency
84
+ m.AddAddr(id3, ma33, ttl)
85
+
86
+ m.AddAddrs(id4, []ma.Multiaddr{ma41, ma42, ma43, ma44}, ttl) // multiple
87
+
88
+ m.AddAddrs(id5, []ma.Multiaddr{ma21, ma22}, ttl) // clearing
89
+ m.AddAddrs(id5, []ma.Multiaddr{ma41, ma42, ma43, ma44}, ttl) // clearing
90
+ m.ClearAddrs(id5)
91
+ m.AddAddrs(id5, []ma.Multiaddr{ma51, ma52, ma53, ma54, ma55}, ttl) // clearing
92
+
93
+ // test the Addresses return value
94
+ testHas(t, []ma.Multiaddr{ma11}, m.Addrs(id1))
95
+ testHas(t, []ma.Multiaddr{ma21, ma22}, m.Addrs(id2))
96
+ testHas(t, []ma.Multiaddr{ma31, ma32, ma33}, m.Addrs(id3))
97
+ testHas(t, []ma.Multiaddr{ma41, ma42, ma43, ma44}, m.Addrs(id4))
98
+ testHas(t, []ma.Multiaddr{ma51, ma52, ma53, ma54, ma55}, m.Addrs(id5))
99
+
100
+}
101
+
102
+func TestAddressesExpire(t *testing.T) {
103
+
104
+ id1 := IDS(t, "QmcNstKuwBBoVTpSCSDrwzjgrRcaYXK833Psuz2EMHwyQN")
105
+ id2 := IDS(t, "QmcNstKuwBBoVTpSCSDrwzjgrRcaYXK833Psuz2EMHwyQM")
106
+ ma11 := MA(t, "/ip4/1.2.3.1/tcp/1111")
107
+ ma12 := MA(t, "/ip4/2.2.3.2/tcp/2222")
108
+ ma13 := MA(t, "/ip4/3.2.3.3/tcp/3333")
109
+ ma24 := MA(t, "/ip4/4.2.3.3/tcp/4444")
110
+ ma25 := MA(t, "/ip4/5.2.3.3/tcp/5555")
111
+
112
+ m := Manager{}
113
+ m.AddAddr(id1, ma11, time.Hour)
114
+ m.AddAddr(id1, ma12, time.Hour)
115
+ m.AddAddr(id1, ma13, time.Hour)
116
+ m.AddAddr(id2, ma24, time.Hour)
117
+ m.AddAddr(id2, ma25, time.Hour)
118
+
119
+ testHas(t, []ma.Multiaddr{ma11, ma12, ma13}, m.Addrs(id1))
120
+ testHas(t, []ma.Multiaddr{ma24, ma25}, m.Addrs(id2))
121
+
122
+ m.SetAddr(id1, ma11, 2*time.Hour)
123
+ m.SetAddr(id1, ma12, 2*time.Hour)
124
+ m.SetAddr(id1, ma13, 2*time.Hour)
125
+ m.SetAddr(id2, ma24, 2*time.Hour)
126
+ m.SetAddr(id2, ma25, 2*time.Hour)
127
+
128
+ testHas(t, []ma.Multiaddr{ma11, ma12, ma13}, m.Addrs(id1))
129
+ testHas(t, []ma.Multiaddr{ma24, ma25}, m.Addrs(id2))
130
+
131
+ m.SetAddr(id1, ma11, time.Millisecond)
132
+ <-time.After(time.Millisecond)
133
+ testHas(t, []ma.Multiaddr{ma12, ma13}, m.Addrs(id1))
134
+ testHas(t, []ma.Multiaddr{ma24, ma25}, m.Addrs(id2))
135
+
136
+ m.SetAddr(id1, ma13, time.Millisecond)
137
+ <-time.After(time.Millisecond)
138
+ testHas(t, []ma.Multiaddr{ma12}, m.Addrs(id1))
139
+ testHas(t, []ma.Multiaddr{ma24, ma25}, m.Addrs(id2))
140
+
141
+ m.SetAddr(id2, ma24, time.Millisecond)
142
+ <-time.After(time.Millisecond)
143
+ testHas(t, []ma.Multiaddr{ma12}, m.Addrs(id1))
144
+ testHas(t, []ma.Multiaddr{ma25}, m.Addrs(id2))
145
+
146
+ m.SetAddr(id2, ma25, time.Millisecond)
147
+ <-time.After(time.Millisecond)
148
+ testHas(t, []ma.Multiaddr{ma12}, m.Addrs(id1))
149
+ testHas(t, nil, m.Addrs(id2))
150
+
151
+ m.SetAddr(id1, ma12, time.Millisecond)
152
+ <-time.After(time.Millisecond)
153
+ testHas(t, nil, m.Addrs(id1))
154
+ testHas(t, nil, m.Addrs(id2))
155
+}
156
+
157
+func TestClearWorks(t *testing.T) {
158
+
159
+ id1 := IDS(t, "QmcNstKuwBBoVTpSCSDrwzjgrRcaYXK833Psuz2EMHwyQN")
160
+ id2 := IDS(t, "QmcNstKuwBBoVTpSCSDrwzjgrRcaYXK833Psuz2EMHwyQM")
161
+ ma11 := MA(t, "/ip4/1.2.3.1/tcp/1111")
162
+ ma12 := MA(t, "/ip4/2.2.3.2/tcp/2222")
163
+ ma13 := MA(t, "/ip4/3.2.3.3/tcp/3333")
164
+ ma24 := MA(t, "/ip4/4.2.3.3/tcp/4444")
165
+ ma25 := MA(t, "/ip4/5.2.3.3/tcp/5555")
166
+
167
+ m := Manager{}
168
+ m.AddAddr(id1, ma11, time.Hour)
169
+ m.AddAddr(id1, ma12, time.Hour)
170
+ m.AddAddr(id1, ma13, time.Hour)
171
+ m.AddAddr(id2, ma24, time.Hour)
172
+ m.AddAddr(id2, ma25, time.Hour)
173
+
174
+ testHas(t, []ma.Multiaddr{ma11, ma12, ma13}, m.Addrs(id1))
175
+ testHas(t, []ma.Multiaddr{ma24, ma25}, m.Addrs(id2))
176
+
177
+ m.ClearAddrs(id1)
178
+ m.ClearAddrs(id2)
179
+
180
+ testHas(t, nil, m.Addrs(id1))
181
+ testHas(t, nil, m.Addrs(id2))
182
+}