@cryptotaxi247 / kubo / commits / ac42cbe9f

mock2

Juan Batiz-Benet committed Dec 17, 2014 at 06:11 UTC ac42cbe9f955961ad716e69c2398c8104d3db1c5
7 files changed +1027
net/mock2/interface.go new
+67
@@ -0,0 +1,67 @@
1 +// Package mocknet provides a mock net.Network to test with.
2 +//
3 +// - a Mocknet has many inet.Networks
4 +// - a Mocknet has many Links
5 +// - a Link joins two inet.Networks
6 +// - inet.Conns and inet.Streams are created by inet.Networks
7 +package mocknet
8 +
9 +import (
10 + "time"
11 +
12 + inet "github.com/jbenet/go-ipfs/net"
13 + peer "github.com/jbenet/go-ipfs/peer"
14 +)
15 +
16 +type Mocknet interface {
17 + GenPeer() (inet.Network, error)
18 + AddPeer(peer.ID) (inet.Network, error)
19 +
20 + // retrieve things
21 + Peer(peer.ID) peer.Peer
22 + Peers() []peer.Peer
23 + Net(peer.ID) inet.Network
24 + Nets() []inet.Network
25 + LinksBetweenPeers(a, b peer.Peer) []Link
26 + LinksBetweenNets(a, b inet.Network) []Link
27 +
28 + // Links are the **ability to connect**.
29 + // think of Links as the physical medium.
30 + // For p1 and p2 to connect, a link must exist between them.
31 + // (this makes it possible to test dial failures, and
32 + // things like relaying traffic)
33 + LinkPeers(peer.Peer, peer.Peer) (Link, error)
34 + LinkNets(inet.Network, inet.Network) (Link, error)
35 + Unlink(Link) error
36 + UnlinkPeers(peer.Peer, peer.Peer) error
37 + UnlinkNets(inet.Network, inet.Network) error
38 +
39 + // LinkDefaults are the default options that govern links
40 + // if they do not have thier own option set.
41 + SetLinkDefaults(LinkOptions)
42 + LinkDefaults() LinkOptions
43 +
44 + // 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).
47 + ConnectPeers(peer.Peer, peer.Peer) error
48 + ConnectNets(inet.Network, inet.Network) error
49 + DisconnectPeers(peer.Peer, peer.Peer) error
50 + DisconnectNets(inet.Network, inet.Network) error
51 +}
52 +
53 +type LinkOptions struct {
54 + Latency time.Duration
55 + Bandwidth int // in bytes-per-second
56 + // we can make these values distributions down the road.
57 +}
58 +
59 +type Link interface {
60 + Networks() []inet.Network
61 + Peers() []peer.Peer
62 +
63 + SetOptions(LinkOptions)
64 + Options() LinkOptions
65 +
66 + // Metrics() Metrics
67 +}
net/mock2/mock.go new
+67
@@ -0,0 +1,67 @@
1 +package mocknet
2 +
3 +import (
4 + eventlog "github.com/jbenet/go-ipfs/util/eventlog"
5 +
6 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7 +)
8 +
9 +var log = eventlog.Logger("mocknet")
10 +
11 +// WithNPeers constructs a Mocknet with N peers.
12 +func WithNPeers(ctx context.Context, n int) (Mocknet, error) {
13 + m := New(ctx)
14 + for i := 0; i < n; i++ {
15 + if _, err := m.GenPeer(); err != nil {
16 + return nil, err
17 + }
18 + }
19 + return m, nil
20 +}
21 +
22 +// FullMeshLinked constructs a Mocknet with full mesh of Links.
23 +// This means that all the peers **can** connect to each other
24 +// (not that they already are connected. you can use m.ConnectAll())
25 +func FullMeshLinked(ctx context.Context, n int) (Mocknet, error) {
26 + m, err := WithNPeers(ctx, n)
27 + if err != nil {
28 + return nil, err
29 + }
30 +
31 + nets := m.Nets()
32 + for _, n1 := range nets {
33 + for _, n2 := range nets {
34 + // yes, even self.
35 + if _, err := m.LinkNets(n1, n2); err != nil {
36 + return nil, err
37 + }
38 + }
39 + }
40 +
41 + return m, nil
42 +}
43 +
44 +// FullMeshConnected constructs a Mocknet with full mesh of Connections.
45 +// This means that all the peers have dialed and are ready to talk to
46 +// each other.
47 +func FullMeshConnected(ctx context.Context, n int) (Mocknet, error) {
48 + m, err := FullMeshLinked(ctx, n)
49 + if err != nil {
50 + return nil, err
51 + }
52 +
53 + nets := m.Nets()
54 + for _, n1 := range nets {
55 + for _, n2 := range nets {
56 + if n1 == n2 {
57 + continue
58 + }
59 +
60 + if err := m.ConnectNets(n1, n2); err != nil {
61 + return nil, err
62 + }
63 + }
64 + }
65 +
66 + return m, nil
67 +}
net/mock2/mock_conn.go new
+106
@@ -0,0 +1,106 @@
1 +package mocknet
2 +
3 +import (
4 + "container/list"
5 + "sync"
6 +
7 + inet "github.com/jbenet/go-ipfs/net"
8 + peer "github.com/jbenet/go-ipfs/peer"
9 +
10 + ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
11 +)
12 +
13 +// conn represents one side's perspective of a
14 +// live connection between two peers.
15 +// it goes over a particular link.
16 +type conn struct {
17 + local peer.Peer
18 + remote peer.Peer
19 + net *peernet
20 + link *link
21 + rconn *conn // counterpart
22 + streams list.List
23 +
24 + sync.RWMutex
25 +}
26 +
27 +func (c *conn) Close() error {
28 + for _, s := range c.allStreams() {
29 + s.Close()
30 + }
31 + c.net.removeConn(c)
32 + return nil
33 +}
34 +
35 +func (c *conn) addStream(s *stream) {
36 + c.Lock()
37 + s.conn = c
38 + c.streams.PushBack(s)
39 + c.Unlock()
40 +}
41 +
42 +func (c *conn) removeStream(s *stream) {
43 + c.Lock()
44 + defer c.Unlock()
45 + for e := c.streams.Front(); e != nil; e = e.Next() {
46 + if s == e.Value {
47 + c.streams.Remove(e)
48 + return
49 + }
50 + }
51 +}
52 +
53 +func (c *conn) allStreams() []inet.Stream {
54 + c.RLock()
55 + defer c.RUnlock()
56 +
57 + strs := make([]inet.Stream, 0, c.streams.Len())
58 + for e := c.streams.Front(); e != nil; e = e.Next() {
59 + s := e.Value.(*stream)
60 + strs = append(strs, s)
61 + }
62 + return strs
63 +}
64 +
65 +func (c *conn) remoteOpenedStream(s *stream) {
66 + c.addStream(s)
67 + c.net.handleNewStream(s)
68 +}
69 +
70 +func (c *conn) openStream() *stream {
71 + sl, sr := c.link.newStreamPair()
72 + c.addStream(sl)
73 + c.rconn.remoteOpenedStream(sr)
74 + return sl
75 +}
76 +
77 +func (c *conn) NewStreamWithProtocol(pr inet.ProtocolID, p peer.Peer) (inet.Stream, error) {
78 + log.Debugf("Conn.NewStreamWithProtocol: %s --> %s", c.local, p)
79 +
80 + s := c.openStream()
81 + if err := inet.WriteProtocolHeader(pr, s); err != nil {
82 + s.Close()
83 + return nil, err
84 + }
85 + return s, nil
86 +}
87 +
88 +// LocalMultiaddr is the Multiaddr on this side
89 +func (c *conn) LocalMultiaddr() ma.Multiaddr {
90 + return nil
91 +}
92 +
93 +// LocalPeer is the Peer on our side of the connection
94 +func (c *conn) LocalPeer() peer.Peer {
95 + return c.local
96 +}
97 +
98 +// RemoteMultiaddr is the Multiaddr on the remote side
99 +func (c *conn) RemoteMultiaddr() ma.Multiaddr {
100 + return nil
101 +}
102 +
103 +// RemotePeer is the Peer on the remote side
104 +func (c *conn) RemotePeer() peer.Peer {
105 + return c.remote
106 +}
net/mock2/mock_link.go new
+86
@@ -0,0 +1,86 @@
1 +package mocknet
2 +
3 +import (
4 + "fmt"
5 + "io"
6 + "sync"
7 +
8 + inet "github.com/jbenet/go-ipfs/net"
9 + peer "github.com/jbenet/go-ipfs/peer"
10 +)
11 +
12 +// link implements mocknet.Link
13 +// and, for simplicity, inet.Conn
14 +type link struct {
15 + mock *mocknet
16 + nets []*peernet
17 + opts LinkOptions
18 +
19 + sync.RWMutex
20 +}
21 +
22 +func newLink(mn *mocknet) *link {
23 + return &link{mock: mn, opts: mn.linkDefaults}
24 +}
25 +
26 +func (l *link) newConnPair() (*conn, *conn) {
27 + l.RLock()
28 + defer l.RUnlock()
29 +
30 + mkconn := func(n *peernet, rid peer.ID) *conn {
31 + c := &conn{net: n, link: l}
32 + c.local = n.peer
33 +
34 + r, err := n.ps.FindOrCreate(rid)
35 + if err != nil {
36 + panic(fmt.Errorf("error creating peer: %s", err))
37 + }
38 + c.remote = r
39 + return c
40 + }
41 +
42 + c1 := mkconn(l.nets[0], l.nets[1].peer.ID())
43 + c2 := mkconn(l.nets[1], l.nets[0].peer.ID())
44 + c1.rconn = c2
45 + c2.rconn = c1
46 + return c1, c2
47 +}
48 +
49 +func (l *link) newStreamPair() (*stream, *stream) {
50 + r1, w1 := io.Pipe()
51 + r2, w2 := io.Pipe()
52 +
53 + s1 := &stream{Reader: r1, Writer: w2}
54 + s2 := &stream{Reader: r2, Writer: w1}
55 + return s1, s2
56 +}
57 +
58 +func (l *link) Networks() []inet.Network {
59 + l.RLock()
60 + defer l.RUnlock()
61 +
62 + cp := make([]inet.Network, len(l.nets))
63 + for i, n := range l.nets {
64 + cp[i] = n
65 + }
66 + return cp
67 +}
68 +
69 +func (l *link) Peers() []peer.Peer {
70 + l.RLock()
71 + defer l.RUnlock()
72 +
73 + cp := make([]peer.Peer, len(l.nets))
74 + for i, n := range l.nets {
75 + cp[i] = n.peer
76 + }
77 + return cp
78 +}
79 +
80 +func (l *link) SetOptions(o LinkOptions) {
81 + l.opts = o
82 +}
83 +
84 +func (l *link) Options() LinkOptions {
85 + return l.opts
86 +}
net/mock2/mock_net.go new
+284
@@ -0,0 +1,284 @@
1 +package mocknet
2 +
3 +import (
4 + "fmt"
5 + "sync"
6 +
7 + inet "github.com/jbenet/go-ipfs/net"
8 + peer "github.com/jbenet/go-ipfs/peer"
9 + testutil "github.com/jbenet/go-ipfs/util/testutil"
10 +
11 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
12 + ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
13 +)
14 +
15 +type peerID string
16 +
17 +// mocknet implements mocknet.Mocknet
18 +type mocknet struct {
19 + // must map on peer.ID (instead of peer.Peer) because
20 + // each inet.Network has different peerstore
21 + nets map[peerID]*peernet
22 +
23 + // links make it possible to connect two peers.
24 + // think of links as the physical medium.
25 + // usually only one, but there could be multiple
26 + // **links are shared between peers**
27 + links map[peerID]map[peerID]map[*link]struct{}
28 +
29 + linkDefaults LinkOptions
30 +
31 + cg ctxgroup.ContextGroup // for Context closing
32 + sync.RWMutex
33 +}
34 +
35 +func New(ctx context.Context) Mocknet {
36 + return &mocknet{
37 + nets: map[peerID]*peernet{},
38 + links: map[peerID]map[peerID]map[*link]struct{}{},
39 + }
40 +}
41 +
42 +func (mn *mocknet) GenPeer() (inet.Network, error) {
43 + p, err := testutil.PeerWithNewKeys()
44 + if err != nil {
45 + return nil, err
46 + }
47 +
48 + n, err := mn.AddPeer(p.ID())
49 + if err != nil {
50 + return nil, err
51 + }
52 +
53 + // copy over keys
54 + if err := n.LocalPeer().Update(p); err != nil {
55 + return nil, err
56 + }
57 +
58 + return n, nil
59 +}
60 +
61 +func (mn *mocknet) AddPeer(p peer.ID) (inet.Network, error) {
62 + n, err := newPeernet(mn.cg.Context(), mn, p)
63 + if err != nil {
64 + return nil, err
65 + }
66 +
67 + mn.cg.AddChildGroup(n.cg)
68 +
69 + mn.Lock()
70 + mn.nets[pid(n.peer)] = n
71 + mn.Unlock()
72 + return n, nil
73 +}
74 +
75 +func (mn *mocknet) Peer(pid peer.ID) peer.Peer {
76 + mn.RLock()
77 + defer mn.RUnlock()
78 +
79 + for _, n := range mn.nets {
80 + if n.peer.ID().Equal(pid) {
81 + return n.peer
82 + }
83 + }
84 + return nil
85 +}
86 +
87 +func (mn *mocknet) Peers() []peer.Peer {
88 + mn.RLock()
89 + defer mn.RUnlock()
90 +
91 + cp := make([]peer.Peer, 0, len(mn.nets))
92 + for _, n := range mn.nets {
93 + cp = append(cp, n.peer)
94 + }
95 + return cp
96 +}
97 +
98 +func (mn *mocknet) Net(pid peer.ID) inet.Network {
99 + mn.RLock()
100 + defer mn.RUnlock()
101 +
102 + for _, n := range mn.nets {
103 + if n.peer.ID().Equal(pid) {
104 + return n
105 + }
106 + }
107 + return nil
108 +}
109 +
110 +func (mn *mocknet) Nets() []inet.Network {
111 + mn.RLock()
112 + defer mn.RUnlock()
113 +
114 + cp := make([]inet.Network, 0, len(mn.nets))
115 + for _, n := range mn.nets {
116 + cp = append(cp, n)
117 + }
118 + return cp
119 +}
120 +
121 +func (mn *mocknet) LinkAll() error {
122 + nets := mn.Nets()
123 + for _, n1 := range nets {
124 + for _, n2 := range nets {
125 + if _, err := mn.LinkNets(n1, n2); err != nil {
126 + return err
127 + }
128 + }
129 + }
130 + return nil
131 +}
132 +
133 +func (mn *mocknet) LinkPeers(p1, p2 peer.Peer) (Link, error) {
134 + mn.RLock()
135 + n1 := mn.nets[pid(p1)]
136 + n2 := mn.nets[pid(p2)]
137 + mn.RUnlock()
138 +
139 + if n1 == nil {
140 + return nil, fmt.Errorf("network for p1 not in mocknet")
141 + }
142 +
143 + if n2 == nil {
144 + return nil, fmt.Errorf("network for p2 not in mocknet")
145 + }
146 +
147 + return mn.LinkNets(n1, n2)
148 +}
149 +
150 +func (mn *mocknet) validate(n inet.Network) (*peernet, error) {
151 + // WARNING: assumes locks acquired
152 +
153 + nr, ok := n.(*peernet)
154 + if !ok {
155 + return nil, fmt.Errorf("Network not supported (use mock package nets only)")
156 + }
157 +
158 + if _, found := mn.nets[pid(nr.peer)]; !found {
159 + return nil, fmt.Errorf("Network not on mocknet. is it from another mocknet?")
160 + }
161 +
162 + return nr, nil
163 +}
164 +
165 +func (mn *mocknet) LinkNets(n1, n2 inet.Network) (Link, error) {
166 + mn.Lock()
167 + defer mn.Unlock()
168 +
169 + if _, err := mn.validate(n1); err != nil {
170 + return nil, err
171 + }
172 +
173 + if _, err := mn.validate(n2); err != nil {
174 + return nil, err
175 + }
176 +
177 + l := newLink(mn)
178 + mn.addLink(l)
179 + return l, nil
180 +}
181 +
182 +func (mn *mocknet) Unlink(l2 Link) error {
183 +
184 + l, ok := l2.(*link)
185 + if !ok {
186 + return fmt.Errorf("only links from mocknet are supported")
187 + }
188 +
189 + mn.removeLink(l)
190 + return nil
191 +}
192 +
193 +func (mn *mocknet) UnlinkPeers(p1, p2 peer.Peer) error {
194 + ls := mn.LinksBetweenPeers(p1, p2)
195 + if ls == nil {
196 + return fmt.Errorf("no link between p1 and p2")
197 + }
198 +
199 + for _, l := range ls {
200 + if err := mn.Unlink(l); err != nil {
201 + return err
202 + }
203 + }
204 + return nil
205 +}
206 +
207 +func (mn *mocknet) UnlinkNets(n1, n2 inet.Network) error {
208 + return mn.DisconnectPeers(n1.LocalPeer(), n2.LocalPeer())
209 +}
210 +
211 +func (mn *mocknet) addLink(l *link) {
212 + mn.Lock()
213 + defer mn.Unlock()
214 +
215 + n1, n2 := l.nets[0], l.nets[1]
216 + mn.links[pid(n1.peer)][pid(n2.peer)][l] = struct{}{}
217 + mn.links[pid(n2.peer)][pid(n1.peer)][l] = struct{}{}
218 +}
219 +
220 +func (mn *mocknet) removeLink(l *link) {
221 + mn.Lock()
222 + defer mn.Unlock()
223 +
224 + n1, n2 := l.nets[0], l.nets[1]
225 + delete(mn.links[pid(n1.peer)][pid(n2.peer)], l)
226 + delete(mn.links[pid(n2.peer)][pid(n1.peer)], l)
227 +}
228 +
229 +func (mn *mocknet) ConnectAll(p1, p2 peer.Peer) (Link, error) {
230 + panic("nyi")
231 +}
232 +
233 +func (mn *mocknet) ConnectPeers(a, b peer.Peer) error {
234 + panic("nyi")
235 +}
236 +
237 +func (mn *mocknet) ConnectNets(inet.Network, inet.Network) error {
238 + panic("nyi")
239 +}
240 +
241 +func (mn *mocknet) DisconnectPeers(p1, p2 peer.Peer) error {
242 + panic("nyi")
243 +}
244 +
245 +func (mn *mocknet) DisconnectNets(n1, n2 inet.Network) error {
246 + return mn.DisconnectPeers(n1.LocalPeer(), n2.LocalPeer())
247 +}
248 +
249 +func (mn *mocknet) LinksBetweenPeers(p1, p2 peer.Peer) []Link {
250 + mn.RLock()
251 + defer mn.RUnlock()
252 +
253 + ls1, found := mn.links[pid(p1)]
254 + if !found {
255 + return nil
256 + }
257 +
258 + ls2, found := ls1[pid(p2)]
259 + if !found {
260 + return nil
261 + }
262 +
263 + cp := make([]Link, 0, len(ls2))
264 + for l := range ls2 {
265 + cp = append(cp, l)
266 + }
267 + return cp
268 +}
269 +
270 +func (mn *mocknet) LinksBetweenNets(n1, n2 inet.Network) []Link {
271 + return mn.LinksBetweenPeers(n1.LocalPeer(), n2.LocalPeer())
272 +}
273 +
274 +func (mn *mocknet) SetLinkDefaults(o LinkOptions) {
275 + mn.Lock()
276 + mn.linkDefaults = o
277 + mn.Unlock()
278 +}
279 +
280 +func (mn *mocknet) LinkDefaults() LinkOptions {
281 + mn.RLock()
282 + defer mn.RUnlock()
283 + return mn.linkDefaults
284 +}
net/mock2/mock_peernet.go new
+388
@@ -0,0 +1,388 @@
1 +package mocknet
2 +
3 +import (
4 + "container/list"
5 + "fmt"
6 + "math/rand"
7 + "sync"
8 +
9 + inet "github.com/jbenet/go-ipfs/net"
10 + peer "github.com/jbenet/go-ipfs/peer"
11 +
12 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
13 + ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
14 + ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
15 +)
16 +
17 +// peernet implements inet.Network
18 +type peernet struct {
19 + mocknet *mocknet // parent
20 +
21 + peer peer.Peer
22 + ps peer.Peerstore
23 +
24 + // conns are actual live connections between peers.
25 + // many conns could run over each link.
26 + // **conns are NOT shared between peers**
27 + connsByPeer map[peerID]list.List
28 + connsByLink map[*link]list.List
29 +
30 + // needed to implement inet.Network
31 + mux inet.Mux
32 +
33 + cg ctxgroup.ContextGroup
34 + sync.RWMutex
35 +}
36 +
37 +// newPeernet constructs a new peernet
38 +func newPeernet(ctx context.Context, m *mocknet, id peer.ID) (*peernet, error) {
39 +
40 + // create our own entirely, so that peers dont get shuffled across
41 + // network divides. dont share peers.
42 + ps := peer.NewPeerstore()
43 + p, err := ps.FindOrCreate(id)
44 + if err != nil {
45 + return nil, err
46 + }
47 +
48 + n := &peernet{
49 + mocknet: m,
50 + peer: p,
51 + ps: ps,
52 + mux: inet.Mux{Handlers: inet.StreamHandlerMap{}},
53 + cg: ctxgroup.WithContext(ctx),
54 +
55 + connsByPeer: map[peerID]list.List{},
56 + connsByLink: map[*link]list.List{},
57 + }
58 +
59 + n.cg.SetTeardown(n.teardown)
60 + return n, nil
61 +}
62 +
63 +func (pn *peernet) teardown() error {
64 +
65 + // close the connections
66 + for _, c := range pn.allConns() {
67 + c.Close()
68 + }
69 + return nil
70 +}
71 +
72 +// allConns returns all the connections between this peer and others
73 +func (pn *peernet) allConns() []*conn {
74 + pn.RLock()
75 + var cs []*conn
76 + for _, csl := range pn.connsByPeer {
77 + for e := csl.Front(); e != nil; e = e.Next() {
78 + c := e.Value.(*conn)
79 + cs = append(cs, c)
80 + }
81 + }
82 + pn.RUnlock()
83 + return cs
84 +}
85 +
86 +// Close calls the ContextCloser func
87 +func (pn *peernet) Close() error {
88 + return pn.cg.Close()
89 +}
90 +
91 +func (pn *peernet) String() string {
92 + return fmt.Sprintf("<mock.peernet %s - %d conns>", pn.peer, len(pn.allConns()))
93 +}
94 +
95 +// handleNewStream is an internal function to trigger the muxer handler
96 +func (pn *peernet) handleNewStream(s inet.Stream) {
97 + go pn.mux.Handle(s)
98 +}
99 +
100 +// DialPeer attempts to establish a connection to a given peer.
101 +// Respects the context.
102 +func (pn *peernet) DialPeer(ctx context.Context, p peer.Peer) error {
103 + return pn.connect(p)
104 +}
105 +
106 +func (pn *peernet) connect(p peer.Peer) error {
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)
110 + if err != nil {
111 + return err
112 + }
113 +
114 + // first, check if we already have live connections
115 + pn.RLock()
116 + cs, found := pn.connsByPeer[pid(p)]
117 + ncs := cs.Len()
118 + pn.RUnlock()
119 + if found && ncs > 0 {
120 + return nil
121 + }
122 +
123 + // ok, must create a new connection. we need a link
124 + links := pn.mocknet.LinksBetweenPeers(pn.peer, p)
125 + if len(links) < 1 {
126 + return fmt.Errorf("cannot connect to peer %s", p)
127 + }
128 +
129 + // if many links found, how do we select? for now, randomly...
130 + // this would be an interesting place to test logic that can measure
131 + // links (network interfaces) and select properly
132 + l := links[rand.Intn(len(links))]
133 +
134 + // create a new connection with link
135 + pn.openConn(p, l.(*link))
136 + return nil
137 +}
138 +
139 +func (pn *peernet) openConn(r peer.Peer, l *link) *conn {
140 + lc, rc := l.newConnPair()
141 + pn.addConn(lc)
142 + rc.net.remoteOpenedConn(rc)
143 + return lc
144 +}
145 +
146 +func (pn *peernet) remoteOpenedConn(c *conn) {
147 + pn.addConn(c)
148 +}
149 +
150 +// addConn constructs and adds a connection
151 +// to given remote peer over given link
152 +func (pn *peernet) addConn(c *conn) {
153 + pn.Lock()
154 + cs, found := pn.connsByPeer[pid(c.RemotePeer())]
155 + if !found {
156 + cs = list.List{}
157 + pn.connsByPeer[pid(c.RemotePeer())] = cs
158 + }
159 + cs.PushBack(c)
160 +
161 + cs, found = pn.connsByLink[c.link]
162 + if !found {
163 + cs = list.List{}
164 + pn.connsByLink[c.link] = cs
165 + }
166 + cs.PushBack(c)
167 + pn.Unlock()
168 +}
169 +
170 +// removeConn removes a given conn
171 +func (pn *peernet) removeConn(c *conn) {
172 + pn.Lock()
173 + defer pn.Unlock()
174 +
175 + cs, found := pn.connsByLink[c.link]
176 + if !found {
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 + }
186 +
187 + cs, found = pn.connsByPeer[pid(c.remote)]
188 + if !found {
189 + panic("attempting to remove a conn that doesnt exist")
190 + }
191 +
192 + for e := cs.Front(); e != nil; e = e.Next() {
193 + if c == e.Value {
194 + cs.Remove(e)
195 + break
196 + }
197 + }
198 +}
199 +
200 +// CtxGroup returns the network's ContextGroup
201 +func (pn *peernet) CtxGroup() ctxgroup.ContextGroup {
202 + return pn.cg
203 +}
204 +
205 +// LocalPeer the network's LocalPeer
206 +func (pn *peernet) LocalPeer() peer.Peer {
207 + return pn.peer
208 +}
209 +
210 +// Peers returns the connected peers
211 +func (pn *peernet) Peers() []peer.Peer {
212 + pn.RLock()
213 + defer pn.RUnlock()
214 +
215 + peers := make([]peer.Peer, 0, len(pn.connsByPeer))
216 + for _, cs := range pn.connsByPeer {
217 + if cs.Len() == 0 {
218 + panic("found empty connection list. not removed properly...")
219 + }
220 +
221 + c := cs.Front().Value.(*conn)
222 + peers = append(peers, c.remote)
223 + }
224 + return peers
225 +}
226 +
227 +// Conns returns all the connections of this peer
228 +func (pn *peernet) Conns() []inet.Conn {
229 + pn.RLock()
230 + defer pn.RUnlock()
231 +
232 + out := make([]inet.Conn, 0, len(pn.connsByPeer))
233 + for _, cs := range pn.connsByPeer {
234 + for e := cs.Front(); e != nil; e = e.Next() {
235 + c := e.Value.(*conn)
236 + out = append(out, c)
237 + }
238 + }
239 + return out
240 +}
241 +
242 +func (pn *peernet) ConnsToPeer(p peer.Peer) []inet.Conn {
243 + pn.RLock()
244 + defer pn.RUnlock()
245 +
246 + cs, found := pn.connsByPeer[pid(p)]
247 + if !found {
248 + return nil
249 + }
250 + if cs.Len() == 0 {
251 + panic("found empty connection list. not removed properly...")
252 + }
253 +
254 + var cs2 []inet.Conn
255 + for e := cs.Front(); e != nil; e = e.Next() {
256 + c := e.Value.(*conn)
257 + cs2 = append(cs2, c)
258 + }
259 + return cs2
260 +}
261 +
262 +// ClosePeer connections to peer
263 +func (pn *peernet) ClosePeer(p peer.Peer) error {
264 + pn.RLock()
265 + cs, found := pn.connsByPeer[pid(p)]
266 + pn.RUnlock()
267 + if !found {
268 + return nil
269 + }
270 +
271 + for e := cs.Front(); e != nil; e = e.Next() {
272 + c := e.Value.(*conn)
273 + pn.closeConn(c)
274 + }
275 + return nil
276 +}
277 +
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 +
312 +// BandwidthTotals returns the total amount of bandwidth transferred
313 +func (pn *peernet) BandwidthTotals() (in uint64, out uint64) {
314 + // need to implement this. probably best to do it in swarm this time.
315 + // need a "metrics" object
316 + return 0, 0
317 +}
318 +
319 +// ListenAddresses returns a list of addresses at which this network listens.
320 +func (pn *peernet) ListenAddresses() []ma.Multiaddr {
321 + return []ma.Multiaddr{}
322 +}
323 +
324 +// InterfaceListenAddresses returns a list of addresses at which this network
325 +// listens. It expands "any interface" addresses (/ip4/0.0.0.0, /ip6/::) to
326 +// use the known local interfaces.
327 +func (pn *peernet) InterfaceListenAddresses() ([]ma.Multiaddr, error) {
328 + return []ma.Multiaddr{}, nil
329 +}
330 +
331 +// Connectedness returns a state signaling connection capabilities
332 +// For now only returns Connecter || NotConnected. Expand into more later.
333 +func (pn *peernet) Connectedness(p peer.Peer) inet.Connectedness {
334 + pn.Lock()
335 + defer pn.Unlock()
336 +
337 + 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 +
343 + return inet.Connected
344 + }
345 + return inet.NotConnected
346 +}
347 +
348 +// NewStream returns a new stream to given peer p.
349 +// If there is no connection to p, attempts to create one.
350 +// If ProtocolID is "", writes no header.
351 +func (pn *peernet) NewStream(pr inet.ProtocolID, p peer.Peer) (inet.Stream, error) {
352 + pn.Lock()
353 + defer pn.Unlock()
354 +
355 + cs, found := pn.connsByPeer[pid(p)]
356 + if !found {
357 + return nil, fmt.Errorf("no connection to peer")
358 + }
359 +
360 + // if many conns are found, how do we select? for now, randomly...
361 + // this would be an interesting place to test logic that can measure
362 + // links (network interfaces) and select properly
363 + c := randomListElem(&cs).Value.(*conn)
364 +
365 + return c.NewStreamWithProtocol(pr, p)
366 +}
367 +
368 +// SetHandler sets the protocol handler on the Network's Muxer.
369 +// This operation is threadsafe.
370 +func (pn *peernet) SetHandler(p inet.ProtocolID, h inet.StreamHandler) {
371 + pn.mux.SetHandler(p, h)
372 +}
373 +
374 +func pid(p peer.Peer) peerID {
375 + return peerID(p.ID())
376 +}
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_stream.go new
+29
@@ -0,0 +1,29 @@
1 +package mocknet
2 +
3 +import (
4 + "io"
5 +
6 + inet "github.com/jbenet/go-ipfs/net"
7 +)
8 +
9 +// stream implements inet.Stream
10 +type stream struct {
11 + io.Reader
12 + io.Writer
13 + conn *conn
14 +}
15 +
16 +func (s *stream) Close() error {
17 + s.conn.removeStream(s)
18 + if r, ok := (s.Reader).(io.Closer); ok {
19 + r.Close()
20 + }
21 + if w, ok := (s.Writer).(io.Closer); ok {
22 + return w.Close()
23 + }
24 + return nil
25 +}
26 +
27 +func (s *stream) Conn() inet.Conn {
28 + return s.conn
29 +}