swarm rewrite, doesnt yet work (tests)
Juan Batiz-Benet committed
Sep 13, 2014 at 18:30 UTC
0ac4a2ba9342ff08a590e47003ca221cfaab320b
5 files changed
+595
-18
net/conn/conn.go
+36
-17
@@ -1,4 +1,4 @@
1
-package swarm
1
+package conn
2
3
import (
4
"fmt"
@@ -32,6 +32,32 @@ type Conn struct {
32
// Map maps Keys (Peer.IDs) to Connections.
33
type Map map[u.Key]*Conn
34
35
+// NewConn constructs a new connection
36
+func NewConn(peer *peer.Peer, addr *ma.Multiaddr, nconn net.Conn) (*Conn, error) {
37
+ conn := &Conn{
38
+ Peer: peer,
39
+ Addr: addr,
40
+ Conn: nconn,
41
+ }
42
+
43
+ if err := conn.newChans(); err != nil {
44
+ return nil, err
45
+ }
46
+
47
+ return conn, nil
48
+}
49
+
50
+// NewNetConn constructs a new connection with given net.Conn
51
+func NewNetConn(nconn net.Conn) (*Conn, error) {
52
+
53
+ addr, err := ma.FromNetAddr(nconn.RemoteAddr())
54
+ if err != nil {
55
+ return nil, err
56
+ }
57
+
58
+ return NewConn(new(peer.Peer), addr, nconn)
59
+}
60
+
61
// Dial connects to a particular peer, over a given network
62
// Example: Dial("udp", peer)
63
func Dial(network string, peer *peer.Peer) (*Conn, error) {
@@ -50,18 +76,11 @@ func Dial(network string, peer *peer.Peer) (*Conn, error) {
76
return nil, err
77
}
78
53
- conn := &Conn{
54
- Peer: peer,
55
- Addr: addr,
56
- Conn: nconn,
57
- }
58
-
59
- newConnChans(conn)
60
- return conn, nil
79
+ return NewConn(peer, addr, nconn)
80
}
81
82
// Construct new channels for given Conn.
64
-func newConnChans(c *Conn) error {
83
+func (c *Conn) newChans() error {
84
if c.Outgoing != nil || c.Incoming != nil {
85
return fmt.Errorf("Conn already initialized")
86
}
@@ -77,18 +96,18 @@ func newConnChans(c *Conn) error {
96
}
97
98
// Close closes the connection, and associated channels.
80
-func (s *Conn) Close() error {
99
+func (c *Conn) Close() error {
100
u.DOut("Closing Conn.\n")
82
- if s.Conn == nil {
101
+ if c.Conn == nil {
102
return fmt.Errorf("Already closed") // already closed
103
}
104
105
// closing net connection
87
- err := s.Conn.Close()
88
- s.Conn = nil
106
+ err := c.Conn.Close()
107
+ c.Conn = nil
108
// closing channels
90
- s.Incoming.Close()
91
- s.Outgoing.Close()
92
- s.Closed <- true
109
+ c.Incoming.Close()
110
+ c.Outgoing.Close()
111
+ c.Closed <- true
112
return err
113
}
net/conn/conn_test.go
+1
-1
@@ -1,4 +1,4 @@
1
-package swarm
1
+package conn
2
3
import (
4
"net"
net/swarm/conn.go
new
+69
@@ -0,0 +1,69 @@
1
+package swarm
2
+
3
+import (
4
+ "errors"
5
+ "fmt"
6
+ "net"
7
+
8
+ ident "github.com/jbenet/go-ipfs/identify"
9
+ conn "github.com/jbenet/go-ipfs/net/conn"
10
+
11
+ u "github.com/jbenet/go-ipfs/util"
12
+)
13
+
14
+// Handle getting ID from this peer, handshake, and adding it into the map
15
+func (s *Swarm) handleIncomingConn(nconn net.Conn) {
16
+
17
+ c, err := conn.NewNetConn(nconn)
18
+ if err != nil {
19
+ s.errChan <- err
20
+ return
21
+ }
22
+
23
+ //TODO(jbenet) the peer might potentially already be in the global PeerBook.
24
+ // maybe use the handshake to populate peer.
25
+ c.Peer.AddAddress(c.Addr)
26
+
27
+ // Setup the new connection
28
+ err = s.connSetup(c)
29
+ if err != nil && err != ErrAlreadyOpen {
30
+ s.errChan <- err
31
+ c.Close()
32
+ }
33
+}
34
+
35
+// connSetup adds the passed in connection to its peerMap and starts
36
+// the fanIn routine for that connection
37
+func (s *Swarm) connSetup(c *conn.Conn) error {
38
+ if c == nil {
39
+ return errors.New("Tried to start nil connection.")
40
+ }
41
+
42
+ u.DOut("Starting connection: %s\n", c.Peer.Key().Pretty())
43
+
44
+ // handshake
45
+ if err := s.connHandshake(c); err != nil {
46
+ return fmt.Errorf("Conn handshake error: %v", err)
47
+ }
48
+
49
+ // add to conns
50
+ s.connsLock.Lock()
51
+ if _, ok := s.conns[c.Peer.Key()]; ok {
52
+ s.connsLock.Unlock()
53
+ return ErrAlreadyOpen
54
+ }
55
+ s.conns[c.Peer.Key()] = c
56
+ s.connsLock.Unlock()
57
+
58
+ // kick off reader goroutine
59
+ go s.fanIn(c)
60
+ return nil
61
+}
62
+
63
+// connHandshake runs the handshake with the remote connection.
64
+func (s *Swarm) connHandshake(c *conn.Conn) error {
65
+
66
+ //TODO(jbenet) this Handshake stuff should be moved elsewhere.
67
+ // needs cleanup. needs context. use msg.Pipe.
68
+ return ident.Handshake(s.local, c.Peer, c.Incoming.MsgChan, c.Outgoing.MsgChan)
69
+}
net/swarm/swarm.go
new
+335
@@ -0,0 +1,335 @@
1
+package swarm
2
+
3
+import (
4
+ "errors"
5
+ "fmt"
6
+ "net"
7
+ "sync"
8
+
9
+ conn "github.com/jbenet/go-ipfs/net/conn"
10
+ msg "github.com/jbenet/go-ipfs/net/message"
11
+ peer "github.com/jbenet/go-ipfs/peer"
12
+ u "github.com/jbenet/go-ipfs/util"
13
+
14
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
15
+ ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
16
+)
17
+
18
+// ErrAlreadyOpen signals that a connection to a peer is already open.
19
+var ErrAlreadyOpen = errors.New("Error: Connection to this peer already open.")
20
+
21
+// ListenErr contains a set of errors mapping to each of the swarms addresses.
22
+// Used to return multiple errors, as in listen.
23
+type ListenErr struct {
24
+ Errors []error
25
+}
26
+
27
+func (e *ListenErr) Error() string {
28
+ if e == nil {
29
+ return "<nil error>"
30
+ }
31
+ var out string
32
+ for i, v := range e.Errors {
33
+ if v != nil {
34
+ out += fmt.Sprintf("%d: %s\n", i, v)
35
+ }
36
+ }
37
+ return out
38
+}
39
+
40
+// Swarm is a connection muxer, allowing connections to other peers to
41
+// be opened and closed, while still using the same Chan for all
42
+// communication. The Chan sends/receives Messages, which note the
43
+// destination or source Peer.
44
+type Swarm struct {
45
+
46
+ // local is the peer this swarm represents
47
+ local *peer.Peer
48
+
49
+ // Swarm includes a Pipe object.
50
+ *msg.Pipe
51
+
52
+ // errChan is the channel of errors.
53
+ errChan chan error
54
+
55
+ // conns are the open connections the swarm is handling.
56
+ conns conn.Map
57
+ connsLock sync.RWMutex
58
+
59
+ // listeners for each network address
60
+ listeners []net.Listener
61
+
62
+ // cancel is an internal function used to stop the Swarm's processing.
63
+ cancel context.CancelFunc
64
+ ctx context.Context
65
+}
66
+
67
+// NewSwarm constructs a Swarm, with a Chan.
68
+func NewSwarm(ctx context.Context, local *peer.Peer) (*Swarm, error) {
69
+ s := &Swarm{
70
+ Pipe: msg.NewPipe(10),
71
+ conns: conn.Map{},
72
+ local: local,
73
+ errChan: make(chan error, 100),
74
+ }
75
+
76
+ s.ctx, s.cancel = context.WithCancel(ctx)
77
+ go s.fanOut()
78
+ return s, s.listen()
79
+}
80
+
81
+// Open listeners for each network the swarm should listen on
82
+func (s *Swarm) listen() error {
83
+ hasErr := false
84
+ retErr := &ListenErr{
85
+ Errors: make([]error, len(s.local.Addresses)),
86
+ }
87
+
88
+ // listen on every address
89
+ for i, addr := range s.local.Addresses {
90
+ err := s.connListen(addr)
91
+ if err != nil {
92
+ hasErr = true
93
+ retErr.Errors[i] = err
94
+ u.PErr("Failed to listen on: %s [%s]", addr, err)
95
+ }
96
+ }
97
+
98
+ if hasErr {
99
+ return retErr
100
+ }
101
+ return nil
102
+}
103
+
104
+// Listen for new connections on the given multiaddr
105
+func (s *Swarm) connListen(maddr *ma.Multiaddr) error {
106
+ netstr, addr, err := maddr.DialArgs()
107
+ if err != nil {
108
+ return err
109
+ }
110
+
111
+ list, err := net.Listen(netstr, addr)
112
+ if err != nil {
113
+ return err
114
+ }
115
+
116
+ // NOTE: this may require a lock around it later. currently, only run on setup
117
+ s.listeners = append(s.listeners, list)
118
+
119
+ // Accept and handle new connections on this listener until it errors
120
+ go func() {
121
+ for {
122
+ nconn, err := list.Accept()
123
+ if err != nil {
124
+ e := fmt.Errorf("Failed to accept connection: %s - %s [%s]",
125
+ netstr, addr, err)
126
+ s.errChan <- e
127
+
128
+ // if cancel is nil, we're closed.
129
+ if s.cancel == nil {
130
+ return
131
+ }
132
+ } else {
133
+ go s.handleIncomingConn(nconn)
134
+ }
135
+ }
136
+ }()
137
+
138
+ return nil
139
+}
140
+
141
+// Close stops a swarm.
142
+func (s *Swarm) Close() error {
143
+ if s.cancel == nil {
144
+ return errors.New("Swarm already closed.")
145
+ }
146
+
147
+ // issue cancel for the context
148
+ s.cancel()
149
+
150
+ // set cancel to nil to prevent calling Close again, and signal to Listeners
151
+ s.cancel = nil
152
+
153
+ // close listeners
154
+ for _, list := range s.listeners {
155
+ list.Close()
156
+ }
157
+ return nil
158
+}
159
+
160
+// Dial connects to a peer.
161
+//
162
+// The idea is that the client of Swarm does not need to know what network
163
+// the connection will happen over. Swarm can use whichever it choses.
164
+// This allows us to use various transport protocols, do NAT traversal/relay,
165
+// etc. to achive connection.
166
+//
167
+// For now, Dial uses only TCP. This will be extended.
168
+func (s *Swarm) Dial(peer *peer.Peer) (*conn.Conn, error) {
169
+ if peer.ID.Equal(s.local.ID) {
170
+ return nil, errors.New("Attempted connection to self!")
171
+ }
172
+
173
+ k := peer.Key()
174
+
175
+ // check if we already have an open connection first
176
+ s.connsLock.RLock()
177
+ c, found := s.conns[k]
178
+ s.connsLock.RUnlock()
179
+ if found {
180
+ return c, nil
181
+ }
182
+
183
+ // open connection to peer
184
+ c, err := conn.Dial("tcp", peer)
185
+ if err != nil {
186
+ return nil, err
187
+ }
188
+
189
+ if err := s.connSetup(c); err != nil {
190
+ c.Close()
191
+ return nil, err
192
+ }
193
+
194
+ return c, nil
195
+}
196
+
197
+// DialAddr is for connecting to a peer when you know their addr but not their ID.
198
+// Should only be used when sure that not connected to peer in question
199
+// TODO(jbenet) merge with Dial? need way to patch back.
200
+func (s *Swarm) DialAddr(addr *ma.Multiaddr) (*conn.Conn, error) {
201
+ if addr == nil {
202
+ return nil, errors.New("addr must be a non-nil Multiaddr")
203
+ }
204
+
205
+ npeer := new(peer.Peer)
206
+ npeer.AddAddress(addr)
207
+
208
+ c, err := conn.Dial("tcp", npeer)
209
+ if err != nil {
210
+ return nil, err
211
+ }
212
+
213
+ if err := s.connSetup(c); err != nil {
214
+ c.Close()
215
+ return nil, err
216
+ }
217
+
218
+ return c, err
219
+}
220
+
221
+// Handles the unwrapping + sending of messages to the right connection.
222
+func (s *Swarm) fanOut() {
223
+ for {
224
+ select {
225
+ case <-s.ctx.Done():
226
+ return // told to close.
227
+
228
+ case msg, ok := <-s.Outgoing:
229
+ if !ok {
230
+ return
231
+ }
232
+
233
+ s.connsLock.RLock()
234
+ conn, found := s.conns[msg.Peer.Key()]
235
+ s.connsLock.RUnlock()
236
+
237
+ if !found {
238
+ e := fmt.Errorf("Sent msg to peer without open conn: %v",
239
+ msg.Peer)
240
+ s.errChan <- e
241
+ continue
242
+ }
243
+
244
+ // queue it in the connection's buffer
245
+ conn.Outgoing.MsgChan <- msg.Data
246
+ }
247
+ }
248
+}
249
+
250
+// Handles the receiving + wrapping of messages, per conn.
251
+// Consider using reflect.Select with one goroutine instead of n.
252
+func (s *Swarm) fanIn(c *conn.Conn) {
253
+ for {
254
+ select {
255
+ case <-s.ctx.Done():
256
+ // close Conn.
257
+ c.Close()
258
+ goto out
259
+
260
+ case <-c.Closed:
261
+ goto out
262
+
263
+ case data, ok := <-c.Incoming.MsgChan:
264
+ if !ok {
265
+ e := fmt.Errorf("Error retrieving from conn: %v", c.Peer.Key().Pretty())
266
+ s.errChan <- e
267
+ goto out
268
+ }
269
+
270
+ msg := &msg.Message{Peer: c.Peer, Data: data}
271
+ s.Incoming <- msg
272
+ }
273
+ }
274
+
275
+out:
276
+ s.connsLock.Lock()
277
+ delete(s.conns, c.Peer.Key())
278
+ s.connsLock.Unlock()
279
+}
280
+
281
+// GetPeer returns the peer in the swarm with given key id.
282
+func (s *Swarm) GetPeer(key u.Key) *peer.Peer {
283
+ s.connsLock.RLock()
284
+ conn, found := s.conns[key]
285
+ s.connsLock.RUnlock()
286
+
287
+ if !found {
288
+ return nil
289
+ }
290
+ return conn.Peer
291
+}
292
+
293
+// GetConnection will check if we are already connected to the peer in question
294
+// and only open a new connection if we arent already
295
+func (s *Swarm) GetConnection(id peer.ID, addr *ma.Multiaddr) (*peer.Peer, error) {
296
+ p := &peer.Peer{
297
+ ID: id,
298
+ Addresses: []*ma.Multiaddr{addr},
299
+ }
300
+
301
+ c, err := s.Dial(p)
302
+ if err != nil {
303
+ return nil, err
304
+ }
305
+
306
+ return c.Peer, nil
307
+}
308
+
309
+// CloseConnection removes a given peer from swarm + closes the connection
310
+func (s *Swarm) CloseConnection(p *peer.Peer) error {
311
+ s.connsLock.RLock()
312
+ conn, found := s.conns[u.Key(p.ID)]
313
+ s.connsLock.RUnlock()
314
+ if !found {
315
+ return u.ErrNotFound
316
+ }
317
+
318
+ s.connsLock.Lock()
319
+ delete(s.conns, u.Key(p.ID))
320
+ s.connsLock.Unlock()
321
+
322
+ return conn.Close()
323
+}
324
+
325
+func (s *Swarm) Error(e error) {
326
+ s.errChan <- e
327
+}
328
+
329
+// GetErrChan returns the errors chan.
330
+func (s *Swarm) GetErrChan() chan error {
331
+ return s.errChan
332
+}
333
+
334
+// Temporary to ensure that the Swarm always matches the Network interface as we are changing it
335
+// var _ Network = &Swarm{}
net/swarm/swarm_test.go
new
+154
@@ -0,0 +1,154 @@
1
+package swarm
2
+
3
+import (
4
+ "fmt"
5
+ "net"
6
+ "testing"
7
+
8
+ msg "github.com/jbenet/go-ipfs/net/message"
9
+ peer "github.com/jbenet/go-ipfs/peer"
10
+ u "github.com/jbenet/go-ipfs/util"
11
+
12
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
13
+ msgio "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-msgio"
14
+ ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
15
+ mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
16
+)
17
+
18
+func pingListen(t *testing.T, listener *net.TCPListener, peer *peer.Peer) {
19
+ for {
20
+ c, err := listener.Accept()
21
+ if err == nil {
22
+ fmt.Println("accepted")
23
+ go pong(t, c, peer)
24
+ }
25
+ }
26
+}
27
+
28
+func pong(t *testing.T, c net.Conn, peer *peer.Peer) {
29
+ mrw := msgio.NewReadWriter(c)
30
+ for {
31
+ data := make([]byte, 1024)
32
+ n, err := mrw.ReadMsg(data)
33
+ if err != nil {
34
+ fmt.Printf("error %v\n", err)
35
+ return
36
+ }
37
+ d := string(data[:n])
38
+ if d != "ping" {
39
+ t.Errorf("error: didn't receive ping: '%v'\n", d)
40
+ return
41
+ }
42
+ err = mrw.WriteMsg([]byte("pong"))
43
+ if err != nil {
44
+ fmt.Printf("error %v\n", err)
45
+ return
46
+ }
47
+ }
48
+}
49
+
50
+func setupPeer(id string, addr string) (*peer.Peer, error) {
51
+ tcp, err := ma.NewMultiaddr(addr)
52
+ if err != nil {
53
+ return nil, err
54
+ }
55
+
56
+ mh, err := mh.FromHexString(id)
57
+ if err != nil {
58
+ return nil, err
59
+ }
60
+
61
+ p := &peer.Peer{ID: peer.ID(mh)}
62
+ p.AddAddress(tcp)
63
+ return p, nil
64
+}
65
+
66
+func TestSwarm(t *testing.T) {
67
+
68
+ local, err := setupPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a30",
69
+ "/ip4/127.0.0.1/tcp/1234")
70
+ if err != nil {
71
+ t.Fatal("error setting up peer", err)
72
+ }
73
+
74
+ swarm, err := NewSwarm(context.Background(), local)
75
+ if err != nil {
76
+ t.Error(err)
77
+ }
78
+ var peers []*peer.Peer
79
+ var listeners []net.Listener
80
+ peerNames := map[string]string{
81
+ "11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a30": "/ip4/127.0.0.1/tcp/1234",
82
+ "11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31": "/ip4/127.0.0.1/tcp/2345",
83
+ "11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a32": "/ip4/127.0.0.1/tcp/3456",
84
+ "11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a33": "/ip4/127.0.0.1/tcp/4567",
85
+ }
86
+
87
+ for k, n := range peerNames {
88
+ peer, err := setupPeer(k, n)
89
+ if err != nil {
90
+ t.Fatal("error setting up peer", err)
91
+ }
92
+ a := peer.NetAddress("tcp")
93
+ if a == nil {
94
+ t.Fatal("error setting up peer (addr is nil)", peer)
95
+ }
96
+ n, h, err := a.DialArgs()
97
+ if err != nil {
98
+ t.Fatal("error getting dial args from addr")
99
+ }
100
+ listener, err := net.Listen(n, h)
101
+ if err != nil {
102
+ t.Fatal("error setting up listener", err)
103
+ }
104
+ go pingListen(t, listener.(*net.TCPListener), peer)
105
+
106
+ u.POut("wat?\n")
107
+ _, err = swarm.Dial(peer)
108
+ if err != nil {
109
+ t.Fatal("error swarm dialing to peer", err)
110
+ }
111
+
112
+ u.POut("wut?\n")
113
+ // ok done, add it.
114
+ peers = append(peers, peer)
115
+ listeners = append(listeners, listener)
116
+ }
117
+
118
+ MsgNum := 1000
119
+ for k := 0; k < MsgNum; k++ {
120
+ for _, p := range peers {
121
+ swarm.Outgoing <- &msg.Message{Peer: p, Data: []byte("ping")}
122
+ u.POut("sending ping to %v\n", p)
123
+ }
124
+ }
125
+
126
+ got := map[u.Key]int{}
127
+
128
+ for k := 0; k < (MsgNum * len(peers)); k++ {
129
+ u.POut("listening for ping...")
130
+ msg := <-swarm.Incoming
131
+ if string(msg.Data) != "pong" {
132
+ t.Error("unexpected conn output", msg.Data)
133
+ }
134
+
135
+ n, _ := got[msg.Peer.Key()]
136
+ got[msg.Peer.Key()] = n + 1
137
+ }
138
+
139
+ if len(peers) != len(got) {
140
+ t.Error("got less messages than sent")
141
+ }
142
+
143
+ for p, n := range got {
144
+ if n != MsgNum {
145
+ t.Error("peer did not get all msgs", p, n, "/", MsgNum)
146
+ }
147
+ }
148
+
149
+ fmt.Println("closing")
150
+ swarm.Close()
151
+ for _, listener := range listeners {
152
+ listener.(*net.TCPListener).Close()
153
+ }
154
+}