mock2: removed list to fix bugs
container/lists suck
Juan Batiz-Benet committed
Dec 17, 2014 at 07:49 UTC
e6a7179a8aa978ec184dc483db4420db41282e59
4 files changed
+57
-102
net/mock2/interface.go
+4
-2
@@ -7,6 +7,7 @@
7
package mocknet
8
9
import (
10
+ "io"
11
"time"
12
13
inet "github.com/jbenet/go-ipfs/net"
@@ -25,6 +26,8 @@ type Mocknet interface {
26
LinksBetweenPeers(a, b peer.Peer) []Link
27
LinksBetweenNets(a, b inet.Network) []Link
28
29
+ PrintLinkMap(io.Writer)
30
+
31
// Links are the **ability to connect**.
32
// think of Links as the physical medium.
33
// For p1 and p2 to connect, a link must exist between them.
@@ -42,8 +45,7 @@ type Mocknet interface {
45
LinkDefaults() LinkOptions
46
47
// Connections are the usual. Connecting means Dialing.
45
- // For convenience, if no link exists, Connect will add one.
46
- // (this is because Connect is called manually by tests).
48
+ // **to succeed, peers must be linked beforehand**
49
ConnectPeers(peer.Peer, peer.Peer) error
50
ConnectNets(inet.Network, inet.Network) error
51
DisconnectPeers(peer.Peer, peer.Peer) error
net/mock2/mock_net.go
+17
-1
@@ -2,6 +2,7 @@ package mocknet
2
3
import (
4
"fmt"
5
+ "io"
6
"sync"
7
8
inet "github.com/jbenet/go-ipfs/net"
@@ -166,7 +167,7 @@ func (mn *mocknet) validate(n inet.Network) (*peernet, error) {
167
func (mn *mocknet) LinkNets(n1, n2 inet.Network) (Link, error) {
168
mn.RLock()
169
n1r, err1 := mn.validate(n1)
169
- n2r, err2 := mn.validate(n1)
170
+ n2r, err2 := mn.validate(n2)
171
ld := mn.linkDefaults
172
mn.RUnlock()
173
@@ -214,6 +215,7 @@ func (mn *mocknet) UnlinkNets(n1, n2 inet.Network) error {
215
216
// get from the links map. and lazily contruct.
217
func (mn *mocknet) linksMapGet(p1, p2 peer.Peer) *map[*link]struct{} {
218
+
219
l1, found := mn.links[pid(p1)]
220
if !found {
221
mn.links[pid(p1)] = map[peerID]map[*link]struct{}{}
@@ -307,3 +309,17 @@ func (mn *mocknet) LinkDefaults() LinkOptions {
309
defer mn.RUnlock()
310
return mn.linkDefaults
311
}
312
+
313
+func (mn *mocknet) PrintLinkMap(w io.Writer) {
314
+ mn.RLock()
315
+ defer mn.RUnlock()
316
+
317
+ fmt.Fprintf(w, "Mocknet link map:\n")
318
+ for p1, lm := range mn.links {
319
+ fmt.Fprintf(w, "\t%s linked to:\n", peer.ID(p1))
320
+ for p2, l := range lm {
321
+ fmt.Fprintf(w, "\t\t%s (%d links)\n", peer.ID(p2), len(l))
322
+ }
323
+ }
324
+ fmt.Fprintf(w, "\n")
325
+}
net/mock2/mock_peernet.go
+35
-98
@@ -1,7 +1,6 @@
1
package mocknet
2
3
import (
4
- "container/list"
4
"fmt"
5
"math/rand"
6
"sync"
@@ -24,8 +23,8 @@ type peernet struct {
23
// conns are actual live connections between peers.
24
// many conns could run over each link.
25
// **conns are NOT shared between peers**
27
- connsByPeer map[peerID]list.List
28
- connsByLink map[*link]list.List
26
+ connsByPeer map[peerID]map[*conn]struct{}
27
+ connsByLink map[*link]map[*conn]struct{}
28
29
// needed to implement inet.Network
30
mux inet.Mux
@@ -52,8 +51,8 @@ func newPeernet(ctx context.Context, m *mocknet, id peer.ID) (*peernet, error) {
51
mux: inet.Mux{Handlers: inet.StreamHandlerMap{}},
52
cg: ctxgroup.WithContext(ctx),
53
55
- connsByPeer: map[peerID]list.List{},
56
- connsByLink: map[*link]list.List{},
54
+ connsByPeer: map[peerID]map[*conn]struct{}{},
55
+ connsByLink: map[*link]map[*conn]struct{}{},
56
}
57
58
n.cg.SetTeardown(n.teardown)
@@ -74,8 +73,7 @@ func (pn *peernet) allConns() []*conn {
73
pn.RLock()
74
var cs []*conn
75
for _, csl := range pn.connsByPeer {
77
- for e := csl.Front(); e != nil; e = e.Next() {
78
- c := e.Value.(*conn)
76
+ for c := range csl {
77
cs = append(cs, c)
78
}
79
}
@@ -104,6 +102,8 @@ func (pn *peernet) DialPeer(ctx context.Context, p peer.Peer) error {
102
}
103
104
func (pn *peernet) connect(p peer.Peer) error {
105
+ log.Debugf("%s dialing %s", pn.peer, p)
106
+
107
// cannot trust the peer we get. typical for tests to give us
108
// a peer from some other peerstore...
109
p, err := pn.ps.Add(p)
@@ -114,16 +114,15 @@ func (pn *peernet) connect(p peer.Peer) error {
114
// first, check if we already have live connections
115
pn.RLock()
116
cs, found := pn.connsByPeer[pid(p)]
117
- ncs := cs.Len()
117
pn.RUnlock()
119
- if found && ncs > 0 {
118
+ if found && len(cs) > 0 {
119
return nil
120
}
121
122
// ok, must create a new connection. we need a link
123
links := pn.mocknet.LinksBetweenPeers(pn.peer, p)
124
if len(links) < 1 {
126
- return fmt.Errorf("cannot connect to peer %s", p)
125
+ return fmt.Errorf("%s cannot connect to %s", pn.peer, p)
126
}
127
128
// if many links found, how do we select? for now, randomly...
@@ -131,6 +130,7 @@ func (pn *peernet) connect(p peer.Peer) error {
130
// links (network interfaces) and select properly
131
l := links[rand.Intn(len(links))]
132
133
+ log.Debugf("%s dialing %s openingConn", pn.peer, p)
134
// create a new connection with link
135
pn.openConn(p, l.(*link))
136
return nil
@@ -153,17 +153,17 @@ func (pn *peernet) addConn(c *conn) {
153
pn.Lock()
154
cs, found := pn.connsByPeer[pid(c.RemotePeer())]
155
if !found {
156
- cs = list.List{}
156
+ cs = map[*conn]struct{}{}
157
pn.connsByPeer[pid(c.RemotePeer())] = cs
158
}
159
- cs.PushBack(c)
159
+ pn.connsByPeer[pid(c.RemotePeer())][c] = struct{}{}
160
161
cs, found = pn.connsByLink[c.link]
162
if !found {
163
- cs = list.List{}
163
+ cs = map[*conn]struct{}{}
164
pn.connsByLink[c.link] = cs
165
}
166
- cs.PushBack(c)
166
+ pn.connsByLink[c.link][c] = struct{}{}
167
pn.Unlock()
168
}
169
@@ -173,28 +173,16 @@ func (pn *peernet) removeConn(c *conn) {
173
defer pn.Unlock()
174
175
cs, found := pn.connsByLink[c.link]
176
- if !found {
176
+ if !found || len(cs) < 1 {
177
panic("attempting to remove a conn that doesnt exist")
178
}
179
-
180
- for e := cs.Front(); e != nil; e = e.Next() {
181
- if c == e.Value {
182
- cs.Remove(e)
183
- break
184
- }
185
- }
179
+ delete(cs, c)
180
181
cs, found = pn.connsByPeer[pid(c.remote)]
182
if !found {
183
panic("attempting to remove a conn that doesnt exist")
184
}
191
-
192
- for e := cs.Front(); e != nil; e = e.Next() {
193
- if c == e.Value {
194
- cs.Remove(e)
195
- break
196
- }
197
- }
185
+ delete(cs, c)
186
}
187
188
// CtxGroup returns the network's ContextGroup
@@ -214,12 +202,10 @@ func (pn *peernet) Peers() []peer.Peer {
202
203
peers := make([]peer.Peer, 0, len(pn.connsByPeer))
204
for _, cs := range pn.connsByPeer {
217
- if cs.Len() == 0 {
218
- panic("found empty connection list. not removed properly...")
205
+ for c := range cs {
206
+ peers = append(peers, c.remote)
207
+ break
208
}
220
-
221
- c := cs.Front().Value.(*conn)
222
- peers = append(peers, c.remote)
209
}
210
return peers
211
}
@@ -231,8 +217,7 @@ func (pn *peernet) Conns() []inet.Conn {
217
218
out := make([]inet.Conn, 0, len(pn.connsByPeer))
219
for _, cs := range pn.connsByPeer {
234
- for e := cs.Front(); e != nil; e = e.Next() {
235
- c := e.Value.(*conn)
220
+ for c := range cs {
221
out = append(out, c)
222
}
223
}
@@ -244,16 +229,12 @@ func (pn *peernet) ConnsToPeer(p peer.Peer) []inet.Conn {
229
defer pn.RUnlock()
230
231
cs, found := pn.connsByPeer[pid(p)]
247
- if !found {
232
+ if !found || len(cs) == 0 {
233
return nil
234
}
250
- if cs.Len() == 0 {
251
- panic("found empty connection list. not removed properly...")
252
- }
235
236
var cs2 []inet.Conn
255
- for e := cs.Front(); e != nil; e = e.Next() {
256
- c := e.Value.(*conn)
237
+ for c := range cs {
238
cs2 = append(cs2, c)
239
}
240
return cs2
@@ -268,47 +249,12 @@ func (pn *peernet) ClosePeer(p peer.Peer) error {
249
return nil
250
}
251
271
- for e := cs.Front(); e != nil; e = e.Next() {
272
- c := e.Value.(*conn)
273
- pn.closeConn(c)
252
+ for c := range cs {
253
+ c.Close()
254
}
255
return nil
256
}
257
278
-func (pn *peernet) closeConn(c *conn) {
279
- pn.Lock()
280
- defer pn.Unlock()
281
-
282
- // remove it from connsByPeer
283
- cs, found := pn.connsByPeer[pid(c.remote)]
284
- if !found {
285
- panic("attempted to close connection that doesnt exist! (peer)")
286
- }
287
-
288
- for e := cs.Front(); e != nil; e = e.Next() {
289
- if c == e.Value.(*conn) {
290
- cs.Remove(e)
291
- }
292
- }
293
- if cs.Len() == 0 {
294
- delete(pn.connsByPeer, pid(c.remote))
295
- }
296
-
297
- // remove it from connsByLink
298
- cs, found = pn.connsByLink[c.link]
299
- if !found {
300
- panic("attempted to close connection that doesnt exist! (link)")
301
- }
302
- for e := cs.Front(); e != nil; e = e.Next() {
303
- if c == e.Value.(*conn) {
304
- cs.Remove(e)
305
- }
306
- }
307
- if cs.Len() == 0 {
308
- delete(pn.connsByLink, c.link)
309
- }
310
-}
311
-
258
// BandwidthTotals returns the total amount of bandwidth transferred
259
func (pn *peernet) BandwidthTotals() (in uint64, out uint64) {
260
// need to implement this. probably best to do it in swarm this time.
@@ -335,11 +281,7 @@ func (pn *peernet) Connectedness(p peer.Peer) inet.Connectedness {
281
defer pn.Unlock()
282
283
cs, found := pn.connsByPeer[pid(p)]
338
- if found {
339
- if cs.Len() == 0 {
340
- panic("found empty connection list. not removed properly...")
341
- }
342
-
284
+ if found && len(cs) > 0 {
285
return inet.Connected
286
}
287
return inet.NotConnected
@@ -353,14 +295,21 @@ func (pn *peernet) NewStream(pr inet.ProtocolID, p peer.Peer) (inet.Stream, erro
295
defer pn.Unlock()
296
297
cs, found := pn.connsByPeer[pid(p)]
356
- if !found {
298
+ if !found || len(cs) < 1 {
299
return nil, fmt.Errorf("no connection to peer")
300
}
301
302
// if many conns are found, how do we select? for now, randomly...
303
// this would be an interesting place to test logic that can measure
304
// links (network interfaces) and select properly
363
- c := randomListElem(&cs).Value.(*conn)
305
+ n := rand.Intn(len(cs))
306
+ var c *conn
307
+ for c = range cs {
308
+ if n == 0 {
309
+ break
310
+ }
311
+ n--
312
+ }
313
314
return c.NewStreamWithProtocol(pr, p)
315
}
@@ -374,15 +323,3 @@ func (pn *peernet) SetHandler(p inet.ProtocolID, h inet.StreamHandler) {
323
func pid(p peer.Peer) peerID {
324
return peerID(p.ID())
325
}
377
-
378
-func randomListElem(l *list.List) *list.Element {
379
- n := rand.Intn(l.Len())
380
- for e := l.Front(); e != nil; e = e.Next() {
381
- if n == 0 {
382
- return e
383
- }
384
- n--
385
- }
386
-
387
- panic("unreachable")
388
-}
net/mock2/mock_test.go
+1
-1
@@ -172,7 +172,7 @@ func TestStreamsStress(t *testing.T) {
172
to := rand.Intn(len(nets))
173
p := rand.Intn(3)
174
proto := protos[p]
175
- log.Debug("%d (%s) %d (%s) %d (%s)", from, nets[from], to, nets[to], p, protos[p])
175
+ // log.Debug("%d (%s) %d (%s) %d (%s)", from, nets[from], to, nets[to], p, protos[p])
176
s, err := nets[from].NewStream(protos[p], nets[to].LocalPeer())
177
if err != nil {
178
panic(err)