new swarm with io and router
Juan Batiz-Benet committed
Dec 13, 2014 at 04:59 UTC
d94593a9551f58774ac8c137d51bef25c1a6919f
9 files changed
+642
-415
net/message/message.go
+34
-33
@@ -1,11 +1,45 @@
1
package message
2
3
import (
4
+ "errors"
5
+
6
peer "github.com/jbenet/go-ipfs/peer"
7
8
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9
proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
10
+ router "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-router"
11
)
12
13
+// ErrInvalidPayload is an error used in the router.HandlePacket implementations
14
+var ErrInvalidPayload = errors.New("invalid packet: non-[]byte payload")
15
+
16
+// Packet is used inside the network package to represent a message
17
+// flowing across the subsystems (Conn, Swarm, Mux, Service).
18
+// implements router.Packet
19
+type Packet struct {
20
+ Src router.Address // peer.ID or service string
21
+ Dst router.Address // peer.ID or service string
22
+ Data []byte // raw data
23
+ Context context.Context // context of the Packet.
24
+}
25
+
26
+func (p *Packet) Destination() router.Address {
27
+ return p.Dst
28
+}
29
+
30
+func (p *Packet) Payload() interface{} {
31
+ return p.Data
32
+}
33
+
34
+func (p *Packet) Response(data []byte) Packet {
35
+ return Packet{
36
+ Src: p.Dst,
37
+ Dst: p.Src,
38
+ Data: data,
39
+ Context: p.Context,
40
+ }
41
+}
42
+
43
// NetMessage is the interface for the message
44
type NetMessage interface {
45
Peer() peer.Peer
@@ -54,36 +88,3 @@ func FromObject(p peer.Peer, data proto.Message) (NetMessage, error) {
88
}
89
return New(p, bytes), nil
90
}
57
-
58
-// Pipe objects represent a bi-directional message channel.
59
-type Pipe struct {
60
- Incoming chan NetMessage
61
- Outgoing chan NetMessage
62
-}
63
-
64
-// NewPipe constructs a pipe with channels of a given buffer size.
65
-func NewPipe(bufsize int) *Pipe {
66
- return &Pipe{
67
- Incoming: make(chan NetMessage, bufsize),
68
- Outgoing: make(chan NetMessage, bufsize),
69
- }
70
-}
71
-
72
-// ConnectTo connects this pipe to another, using a context for termination.
73
-func (p *Pipe) ConnectTo(p2 *Pipe) {
74
- connectChans(p.Outgoing, p2.Outgoing)
75
- connectChans(p2.Incoming, p.Incoming)
76
-}
77
-
78
-func connectChans(a, b chan NetMessage) {
79
- go func() {
80
- for {
81
- m, more := <-a
82
- if !more {
83
- close(b)
84
- return
85
- }
86
- b <- m
87
- }
88
- }()
89
-}
net/swarm/addrs.go
+5
-1
@@ -95,7 +95,11 @@ func addrInList(addr ma.Multiaddr, list []ma.Multiaddr) bool {
95
96
// checkNATWarning checks if our observed addresses differ. if so,
97
// informs the user that certain things might not work yet
98
-func (s *Swarm) checkNATWarning(observed ma.Multiaddr) {
98
+func (s *Swarm) checkNATWarning(observed ma.Multiaddr, expected ma.Multiaddr) {
99
+ if observed.Equal(expected) {
100
+ return
101
+ }
102
+
103
listen, err := s.InterfaceListenAddresses()
104
if err != nil {
105
log.Errorf("Error retrieving swarm.InterfaceListenAddresses: %s", err)
net/swarm/conn.go
deleted
-275
@@ -1,275 +0,0 @@
1
-package swarm
2
-
3
-import (
4
- "errors"
5
- "fmt"
6
-
7
- conn "github.com/jbenet/go-ipfs/net/conn"
8
- msg "github.com/jbenet/go-ipfs/net/message"
9
- peer "github.com/jbenet/go-ipfs/peer"
10
-
11
- context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
12
- ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
13
-)
14
-
15
-// Open listeners for each network the swarm should listen on
16
-func (s *Swarm) listen(addrs []ma.Multiaddr) error {
17
- hasErr := false
18
- retErr := &ListenErr{
19
- Errors: make([]error, len(addrs)),
20
- }
21
-
22
- // listen on every address
23
- for i, addr := range addrs {
24
- err := s.connListen(addr)
25
- if err != nil {
26
- hasErr = true
27
- retErr.Errors[i] = err
28
- log.Errorf("Failed to listen on: %s - %s", addr, err)
29
- }
30
- }
31
-
32
- if hasErr {
33
- return retErr
34
- }
35
- return nil
36
-}
37
-
38
-// Listen for new connections on the given multiaddr
39
-func (s *Swarm) connListen(maddr ma.Multiaddr) error {
40
-
41
- resolved, err := resolveUnspecifiedAddresses([]ma.Multiaddr{maddr})
42
- if err != nil {
43
- return err
44
- }
45
-
46
- list, err := conn.Listen(s.Context(), maddr, s.local, s.peers)
47
- if err != nil {
48
- return err
49
- }
50
-
51
- // add resolved local addresses to peer
52
- for _, addr := range resolved {
53
- s.local.AddAddress(addr)
54
- }
55
-
56
- // make sure port can be reused. TOOD this doesn't work...
57
- // if err := setSocketReuse(list); err != nil {
58
- // return err
59
- // }
60
-
61
- // NOTE: this may require a lock around it later. currently, only run on setup
62
- s.listeners = append(s.listeners, list)
63
-
64
- // Accept and handle new connections on this listener until it errors
65
- // this listener is a child.
66
- s.Children().Add(1)
67
- go func() {
68
- defer s.Children().Done()
69
-
70
- for {
71
- select {
72
- case <-s.Closing():
73
- return
74
-
75
- case conn := <-list.Accept():
76
- // handler also a child.
77
- s.Children().Add(1)
78
- go s.handleIncomingConn(conn)
79
- }
80
- }
81
- }()
82
-
83
- return nil
84
-}
85
-
86
-// Handle getting ID from this peer, handshake, and adding it into the map
87
-func (s *Swarm) handleIncomingConn(nconn conn.Conn) {
88
- // this handler is a child. added by caller.
89
- defer s.Children().Done()
90
-
91
- // Setup the new connection
92
- _, err := s.connSetup(nconn)
93
- if err != nil && err != ErrAlreadyOpen {
94
- s.errChan <- err
95
- nconn.Close()
96
- }
97
-}
98
-
99
-// peerMultiConn returns the MultiConn responsible for handling this peer.
100
-// if there is none, it creates one and returns it. Note that timeouts
101
-// and connection teardowns will remove it.
102
-func (s *Swarm) peerMultiConn(p peer.Peer) (*conn.MultiConn, error) {
103
-
104
- s.connsLock.Lock()
105
- mc, found := s.conns[p.Key()]
106
- if found {
107
- s.connsLock.Unlock()
108
- return mc, nil
109
- }
110
-
111
- // multiconn doesn't exist, make a new one.
112
- mc, err := conn.NewMultiConn(s.Context(), s.local, p, nil)
113
- if err != nil {
114
- s.connsLock.Unlock()
115
- log.Errorf("error creating multiconn: %s", err)
116
- return nil, err
117
- }
118
- s.conns[p.Key()] = mc
119
- s.connsLock.Unlock()
120
-
121
- // kick off reader goroutine
122
- s.Children().Add(1)
123
- mc.Children().Add(1) // child of Conn as well.
124
- go s.fanInSingle(mc)
125
- return mc, nil
126
-}
127
-
128
-// connSetup takes a new connection, performs the IPFS handshake (handshake3)
129
-// and then adds it to the appropriate MultiConn.
130
-func (s *Swarm) connSetup(c conn.Conn) (conn.Conn, error) {
131
- if c == nil {
132
- return nil, errors.New("Tried to start nil connection.")
133
- }
134
-
135
- log.Event(context.TODO(), "connSetupBegin", c.LocalPeer(), c.RemotePeer())
136
-
137
- // add address of connection to Peer. Maybe it should happen in connSecure.
138
- // NOT adding this address here, because the incoming address in TCP
139
- // is an EPHEMERAL address, and not the address we want to keep around.
140
- // addresses should be figured out through the DHT.
141
- // c.Remote.AddAddress(c.Conn.RemoteMultiaddr())
142
-
143
- // handshake3
144
- ctxT, _ := context.WithTimeout(c.Context(), conn.HandshakeTimeout)
145
- h3result, err := conn.Handshake3(ctxT, c)
146
- if err != nil {
147
- c.Close()
148
- return nil, fmt.Errorf("Handshake3 failed: %s", err)
149
- }
150
-
151
- // check for nats. you know, just in case.
152
- if h3result.LocalObservedAddress != nil {
153
- s.checkNATWarning(h3result.LocalObservedAddress)
154
- } else {
155
- log.Warningf("Received nil observed address from %s", c.RemotePeer())
156
- }
157
-
158
- // add to conns
159
- mc, err := s.peerMultiConn(c.RemotePeer())
160
- if err != nil {
161
- c.Close()
162
- return nil, err
163
- }
164
- mc.Add(c)
165
- log.Event(context.TODO(), "connSetupSuccess", c.LocalPeer(), c.RemotePeer())
166
- return c, nil
167
-}
168
-
169
-// Handles the unwrapping + sending of messages to the right connection.
170
-func (s *Swarm) fanOut() {
171
- defer s.Children().Done()
172
-
173
- i := 0
174
- for {
175
- select {
176
- case <-s.Closing():
177
- return // told to close.
178
-
179
- case msg, ok := <-s.Outgoing:
180
- if !ok {
181
- log.Infof("%s outgoing channel closed", s)
182
- return
183
- }
184
- if len(msg.Data()) >= conn.MaxMessageSize {
185
- log.Criticalf("Attempted to send message bigger than max size. (%d)", len(msg.Data()))
186
- }
187
-
188
- s.connsLock.RLock()
189
- c, found := s.conns[msg.Peer().Key()]
190
- s.connsLock.RUnlock()
191
-
192
- if !found {
193
- e := fmt.Errorf("Sent msg to peer without open conn: %v", msg.Peer())
194
- s.errChan <- e
195
- log.Error(e)
196
- continue
197
- }
198
-
199
- i++
200
- log.Debugf("%s sent message to %s (%d)", s.local, msg.Peer(), i)
201
- log.Event(context.TODO(), "sendMessage", s.local, msg)
202
- // queue it in the connection's buffer
203
- if err := c.WriteMsg(msg.Data()); err != nil {
204
- log.Infof("%s connection failed to write: %s", c, err)
205
- continue
206
- }
207
- }
208
- }
209
-}
210
-
211
-// Handles the receiving + wrapping of messages, per conn.
212
-// Consider using reflect.Select with one goroutine instead of n.
213
-func (s *Swarm) fanInSingle(c conn.Conn) {
214
- // cleanup all data associated with this child Connection.
215
- defer func() {
216
- // remove it from the map.
217
- s.connsLock.Lock()
218
- delete(s.conns, c.RemotePeer().Key())
219
- s.connsLock.Unlock()
220
-
221
- s.Children().Done()
222
- c.Children().Done() // child of Conn as well.
223
- }()
224
-
225
- // use readChan to be able to listen to Closing events
226
- rchan := readChan(s.Context(), c)
227
-
228
- i := 0
229
- for {
230
- select {
231
- case <-s.Closing(): // Swarm closing
232
- return
233
-
234
- case <-c.Closing(): // Conn closing
235
- return
236
-
237
- case data, ok := <-rchan:
238
- if !ok {
239
- log.Infof("%s in channel closed", c)
240
- return // channel closed.
241
- }
242
- i++
243
- log.Debugf("%s received message from %s (%d)", s.local, c.RemotePeer(), i)
244
- s.Incoming <- msg.New(c.RemotePeer(), data)
245
- }
246
- }
247
-}
248
-
249
-// readChan is a temporary fixture to match the old interface. will be removed soon.
250
-func readChan(ctx context.Context, c conn.Conn) <-chan []byte {
251
-
252
- ch := make(chan []byte) // no buffer. sync.
253
-
254
- go func() {
255
- defer close(ch)
256
-
257
- for {
258
- msg, err := c.ReadMsg()
259
- if err != nil {
260
- log.Infof("%s connection failed: %s", c, err)
261
- return
262
- }
263
-
264
- select {
265
- case <-c.Closing():
266
- return
267
- case <-ctx.Done():
268
- return
269
- case ch <- msg:
270
- }
271
- }
272
- }()
273
-
274
- return ch
275
-}
net/swarm/simul_test.go
+2
-4
@@ -12,9 +12,7 @@ import (
12
)
13
14
func TestSimultOpen(t *testing.T) {
15
- if testing.Short() {
16
- t.SkipNow()
17
- }
15
+ // t.Skip("skipping for another test")
16
17
addrs := []string{
18
"/ip4/127.0.0.1/tcp/1244",
@@ -51,7 +49,7 @@ func TestSimultOpen(t *testing.T) {
49
}
50
51
func TestSimultOpenMany(t *testing.T) {
54
- t.Skip("laggy")
52
+ t.Skip("very very slow")
53
54
many := 500
55
addrs := []string{}
net/swarm/swarm.go
+40
-71
@@ -5,17 +5,15 @@ package swarm
5
import (
6
"errors"
7
"fmt"
8
- "sync"
8
9
conn "github.com/jbenet/go-ipfs/net/conn"
11
- msg "github.com/jbenet/go-ipfs/net/message"
10
peer "github.com/jbenet/go-ipfs/peer"
13
- u "github.com/jbenet/go-ipfs/util"
14
- ctxc "github.com/jbenet/go-ipfs/util/ctxcloser"
11
"github.com/jbenet/go-ipfs/util/eventlog"
12
13
+ ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
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
+ router "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-router"
17
)
18
19
var log = eventlog.Logger("swarm")
@@ -46,6 +44,8 @@ func (e *ListenErr) Error() string {
44
// be opened and closed, while still using the same Chan for all
45
// communication. The Chan sends/receives Messages, which note the
46
// destination or source Peer.
47
+//
48
+// Implements router.Node
49
type Swarm struct {
50
51
// local is the peer this swarm represents
@@ -54,43 +54,42 @@ type Swarm struct {
54
// peers is a collection of peers for swarm to use
55
peers peer.Peerstore
56
57
- // Swarm includes a Pipe object.
58
- *msg.Pipe
59
-
60
- // errChan is the channel of errors.
61
- errChan chan error
62
-
63
- // conns are the open connections the swarm is handling.
64
- // these are MultiConns, which multiplex multiple separate underlying Conns.
65
- conns conn.MultiConnMap
66
- connsLock sync.RWMutex
57
+ // rt handles the open connections the swarm is handling.
58
+ rt *swarmRoutingTable
59
60
// listeners for each network address
61
listeners []conn.Listener
62
71
- // ContextCloser
72
- ctxc.ContextCloser
63
+ // ContextGroup
64
+ cg ctxgroup.ContextGroup
65
}
66
67
// NewSwarm constructs a Swarm, with a Chan.
76
-func NewSwarm(ctx context.Context, listenAddrs []ma.Multiaddr, local peer.Peer, ps peer.Peerstore) (*Swarm, error) {
68
+func NewSwarm(ctx context.Context, listenAddrs []ma.Multiaddr,
69
+ local peer.Peer, ps peer.Peerstore, client router.Node) (*Swarm, error) {
70
+
71
s := &Swarm{
78
- Pipe: msg.NewPipe(10),
79
- conns: conn.MultiConnMap{},
80
- local: local,
81
- peers: ps,
82
- errChan: make(chan error, 100),
72
+ local: local,
73
+ peers: ps,
74
+ cg: ctxgroup.WithContext(ctx),
75
+ rt: newRoutingTable(local, client),
76
}
77
85
- // ContextCloser for proper child management.
86
- s.ContextCloser = ctxc.NewContextCloser(ctx, s.close)
87
-
88
- s.Children().Add(1)
89
- go s.fanOut()
78
+ s.cg.SetTeardown(s.close)
79
return s, s.listen(listenAddrs)
80
}
81
93
-// close stops a swarm. It's the underlying function called by ContextCloser
82
+// SetClient assign's the Swarm's client node.
83
+func (s *Swarm) SetClient(n router.Node) {
84
+ s.rt.client = n
85
+}
86
+
87
+// Close stops a swarm. waits till it exits
88
+func (s *Swarm) Close() error {
89
+ return s.cg.Close()
90
+}
91
+
92
+// close stops a swarm. It's the underlying function called by ContextGroup
93
func (s *Swarm) close() error {
94
// close listeners
95
for _, list := range s.listeners {
@@ -139,7 +138,7 @@ func (s *Swarm) Dial(peer peer.Peer) (conn.Conn, error) {
138
// for simplicity, we do this sequentially.
139
// A future commit will do this asynchronously.
140
for _, addr := range peer.Addresses() {
142
- c, err = d.DialAddr(s.Context(), addr, peer)
141
+ c, err = d.DialAddr(s.cg.Context(), addr, peer)
142
if err == nil {
143
break
144
}
@@ -148,7 +147,7 @@ func (s *Swarm) Dial(peer peer.Peer) (conn.Conn, error) {
147
return nil, err
148
}
149
151
- c2, err := s.connSetup(c)
150
+ c2, err := s.connSetup(context.TODO(), c)
151
if err != nil {
152
c.Close()
153
return nil, err
@@ -161,64 +160,34 @@ func (s *Swarm) Dial(peer peer.Peer) (conn.Conn, error) {
160
161
// GetConnection returns the connection in the swarm to given peer.ID
162
func (s *Swarm) GetConnection(pid peer.ID) conn.Conn {
164
- s.connsLock.RLock()
165
- defer s.connsLock.RUnlock()
166
- c, found := s.conns[u.Key(pid)]
167
-
168
- if !found {
163
+ sp := s.rt.getByID(pid)
164
+ if sp == nil {
165
return nil
166
}
171
- return c
167
+ return sp.conn
168
}
169
170
// Connections returns a slice of all connections.
171
func (s *Swarm) Connections() []conn.Conn {
176
- s.connsLock.RLock()
177
-
178
- conns := make([]conn.Conn, 0, len(s.conns))
179
- for _, c := range s.conns {
180
- conns = append(conns, c)
181
- }
182
-
183
- s.connsLock.RUnlock()
184
- return conns
172
+ return s.rt.connList()
173
}
174
175
// CloseConnection removes a given peer from swarm + closes the connection
176
func (s *Swarm) CloseConnection(p peer.Peer) error {
189
- c := s.GetConnection(p.ID())
190
- if c == nil {
191
- return u.ErrNotFound
192
- }
193
-
194
- s.connsLock.Lock()
195
- delete(s.conns, u.Key(p.ID()))
196
- s.connsLock.Unlock()
197
-
198
- return c.Close()
199
-}
200
-
201
-func (s *Swarm) Error(e error) {
202
- s.errChan <- e
203
-}
204
-
205
-// GetErrChan returns the errors chan.
206
-func (s *Swarm) GetErrChan() chan error {
207
- return s.errChan
177
+ return s.closeConn(p)
178
}
179
180
// GetPeerList returns a copy of the set of peers swarm is connected to.
181
func (s *Swarm) GetPeerList() []peer.Peer {
212
- var out []peer.Peer
213
- s.connsLock.RLock()
214
- for _, p := range s.conns {
215
- out = append(out, p.RemotePeer())
216
- }
217
- s.connsLock.RUnlock()
218
- return out
182
+ return s.rt.peerList()
183
}
184
185
// LocalPeer returns the local peer swarm is associated to.
186
func (s *Swarm) LocalPeer() peer.Peer {
187
return s.local
188
}
189
+
190
+// Address returns the address of *this* service.
191
+func (s *Swarm) Address() router.Address {
192
+ return "/ipfs/service/swarm" // for now dont need anything more complicated.
193
+}
net/swarm/swarm_conn.go
new
+128
@@ -0,0 +1,128 @@
1
+package swarm
2
+
3
+import (
4
+ "errors"
5
+ "fmt"
6
+
7
+ conn "github.com/jbenet/go-ipfs/net/conn"
8
+
9
+ ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
10
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
11
+ ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
12
+)
13
+
14
+// Open listeners for each network the swarm should listen on
15
+func (s *Swarm) listen(addrs []ma.Multiaddr) error {
16
+ hasErr := false
17
+ retErr := &ListenErr{
18
+ Errors: make([]error, len(addrs)),
19
+ }
20
+
21
+ // listen on every address
22
+ for i, addr := range addrs {
23
+ err := s.connListen(addr)
24
+ if err != nil {
25
+ hasErr = true
26
+ retErr.Errors[i] = err
27
+ log.Errorf("Failed to listen on: %s - %s", addr, err)
28
+ }
29
+ }
30
+
31
+ if hasErr {
32
+ return retErr
33
+ }
34
+ return nil
35
+}
36
+
37
+// Listen for new connections on the given multiaddr
38
+func (s *Swarm) connListen(maddr ma.Multiaddr) error {
39
+
40
+ resolved, err := resolveUnspecifiedAddresses([]ma.Multiaddr{maddr})
41
+ if err != nil {
42
+ return err
43
+ }
44
+
45
+ list, err := conn.Listen(s.cg.Context(), maddr, s.local, s.peers)
46
+ if err != nil {
47
+ return err
48
+ }
49
+
50
+ // add resolved local addresses to peer
51
+ for _, addr := range resolved {
52
+ s.local.AddAddress(addr)
53
+ }
54
+
55
+ // make sure port can be reused. TOOD this doesn't work...
56
+ // if err := setSocketReuse(list); err != nil {
57
+ // return err
58
+ // }
59
+
60
+ // NOTE: this may require a lock around it later. currently, only run on setup
61
+ s.listeners = append(s.listeners, list)
62
+
63
+ // Accept and handle new connections on this listener until it errors
64
+ // this listener is a child.
65
+ s.cg.AddChildFunc(func(parent ctxgroup.ContextGroup) {
66
+ for {
67
+ select {
68
+ case <-parent.Closing():
69
+ return
70
+
71
+ case conn := <-list.Accept():
72
+ s.handleIncomingConn(parent.Context(), conn)
73
+ }
74
+ }
75
+ })
76
+
77
+ return nil
78
+}
79
+
80
+// Handle getting ID from this peer, handshake, and adding it into the map
81
+func (s *Swarm) handleIncomingConn(ctx context.Context, nconn conn.Conn) {
82
+ // Setup the new connection
83
+ _, err := s.connSetup(ctx, nconn)
84
+ if err != nil {
85
+ log.Errorf("swarm: failed to add incoming connection: %s", err)
86
+ log.Event(ctx, "handleIncomingConn failed", s.LocalPeer(), nconn.RemotePeer())
87
+ nconn.Close()
88
+ }
89
+}
90
+
91
+// connSetup takes a new connection, performs the IPFS handshake (handshake3)
92
+// and then adds it to the appropriate MultiConn.
93
+func (s *Swarm) connSetup(ctx context.Context, c conn.Conn) (conn.Conn, error) {
94
+ if c == nil {
95
+ return nil, errors.New("Tried to start nil connection.")
96
+ }
97
+
98
+ log.Event(ctx, "connSetupBegin", c.LocalPeer(), c.RemotePeer())
99
+
100
+ // add address of connection to Peer. Maybe it should happen in connSecure.
101
+ // NOT adding this address here, because the incoming address in TCP
102
+ // is an EPHEMERAL address, and not the address we want to keep around.
103
+ // addresses should be figured out through the DHT.
104
+ // c.Remote.AddAddress(c.Conn.RemoteMultiaddr())
105
+
106
+ // handshake3
107
+ ctxT, _ := context.WithTimeout(c.Context(), conn.HandshakeTimeout)
108
+ h3result, err := conn.Handshake3(ctxT, c)
109
+ if err != nil {
110
+ c.Close()
111
+ return nil, fmt.Errorf("Handshake3 failed: %s", err)
112
+ }
113
+
114
+ // check for nats. you know, just in case.
115
+ if h3result.LocalObservedAddress != nil {
116
+ s.checkNATWarning(h3result.LocalObservedAddress, c.LocalMultiaddr())
117
+ } else {
118
+ log.Warningf("Received nil observed address from %s", c.RemotePeer())
119
+ }
120
+
121
+ // add to conns
122
+ if err := s.addConn(c); err != nil {
123
+ c.Close()
124
+ return nil, err
125
+ }
126
+ log.Event(ctx, "connSetupSuccess", c.LocalPeer(), c.RemotePeer())
127
+ return c, nil
128
+}
net/swarm/swarm_peer.go
new
+178
@@ -0,0 +1,178 @@
1
+package swarm
2
+
3
+import (
4
+ "fmt"
5
+
6
+ conn "github.com/jbenet/go-ipfs/net/conn"
7
+ netmsg "github.com/jbenet/go-ipfs/net/message"
8
+ peer "github.com/jbenet/go-ipfs/peer"
9
+
10
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
11
+ ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
12
+ router "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-router"
13
+)
14
+
15
+// MaxConcurrentRequestsPerPeer defines the pipelining that we can do per-peer.
16
+// the networking layer makes sure to provide proper backpressure to the remote
17
+// side by only handling a max number of concurrent requests to completion.
18
+const MaxConcurrentRequestsPerPeer = 20
19
+
20
+// swarmPeer represents a connection to the outside world.
21
+// Implements router.Node
22
+type swarmPeer struct {
23
+ swarm *Swarm
24
+ conn *conn.MultiConn
25
+ cg ctxgroup.ContextGroup
26
+}
27
+
28
+// newSwarmPeer constructs a new swarmPeer, and starts is worker.
29
+// this doesn't connect it, or add it to the swarm's routing table.
30
+// Implements router.Node
31
+func newSwarmPeer(s *Swarm, p peer.Peer) (*swarmPeer, error) {
32
+ c, err := conn.NewMultiConn(s.cg.Context(), s.LocalPeer(), p, nil)
33
+ if err != nil {
34
+ return nil, fmt.Errorf("Error creating MultiConn: %s", err)
35
+ }
36
+
37
+ sp := &swarmPeer{
38
+ swarm: s,
39
+ conn: c,
40
+ cg: ctxgroup.WithParent(s.cg), // swarmPeer closes when swarm closes.
41
+ }
42
+
43
+ // kicks off the worker.
44
+ // ctggroup makes sure swarmPeer doesn't close until this func returns
45
+ sp.cg.AddChildFunc(sp.listen)
46
+ return sp, nil
47
+}
48
+
49
+// LocalPeer returns the local peer
50
+func (sp *swarmPeer) LocalPeer() peer.Peer {
51
+ return sp.conn.LocalPeer()
52
+}
53
+
54
+// RemotePeer returns the peer we're connected to.
55
+func (sp *swarmPeer) RemotePeer() peer.Peer {
56
+ return sp.conn.RemotePeer()
57
+}
58
+
59
+// Address is the peer's ID
60
+func (sp *swarmPeer) Address() router.Address {
61
+ return sp.RemotePeer()
62
+}
63
+
64
+// Close closes the swarmPeer service
65
+func (sp *swarmPeer) Close() error {
66
+ return sp.cg.Close()
67
+}
68
+
69
+// list to the multiconn and route packets in.
70
+func (sp *swarmPeer) listen(parent ctxgroup.ContextGroup) {
71
+
72
+ // we listen and pipeline using:
73
+ // - 1x listener (this function, the for loop below)
74
+ // - 1x pipelining semaphore
75
+ // - up to Nx goroutine pipeline workers
76
+ // this approach is chosen over N persistent goroutines because
77
+ // spawning a goroutine every time is cheaper than keeping N
78
+ // additional goroutines all the time, for inactive connections
79
+
80
+ pipelineSema := make(chan struct{}, MaxConcurrentRequestsPerPeer)
81
+ for i := 0; i < MaxConcurrentRequestsPerPeer; i++ {
82
+ pipelineSema <- struct{}{}
83
+ }
84
+
85
+ // the sad part of using io is we still need to consume msgs
86
+ // using a context-less api, which means we need to use an
87
+ // extra goroutine, to make sure we close the connection and
88
+ // unlock our blocked listener.
89
+ // (the Context just does not mix well with io.ReadWriters)
90
+ // - conn.SetDeadline could be explored
91
+ go func() {
92
+ <-parent.Closing()
93
+ sp.conn.Close()
94
+ }()
95
+
96
+ listener := func() {
97
+ for {
98
+ msg, err := sp.conn.ReadMsg()
99
+ // we want this to happen before checking the error, as we may
100
+ // be closing (which is not an error). any last message is dropped.
101
+ select {
102
+ case <-parent.Closing():
103
+ return
104
+ case <-sp.conn.Closing():
105
+ return
106
+ default:
107
+ }
108
+
109
+ if err != nil {
110
+ log.Errorf("error receiving message from multiconn: %s", err)
111
+ continue
112
+ }
113
+
114
+ select {
115
+ case <-parent.Closing():
116
+ return
117
+ case <-pipelineSema: // acquire pipelining resource
118
+ go func(m []byte) {
119
+ sp.handleIncomingMessage(parent.Context(), m)
120
+ pipelineSema <- struct{}{}
121
+ }(msg)
122
+ }
123
+ }
124
+ }
125
+
126
+ // function call so that we can isolate the functionality,
127
+ // and so we can call return in loops above and not confuse flow.
128
+ listener()
129
+}
130
+
131
+func (sp *swarmPeer) handleIncomingMessage(ctx context.Context, msg []byte) {
132
+ // handle incoming message message.
133
+ // we derive a new context for this incoming request.
134
+ ctx, _ = context.WithCancel(ctx)
135
+
136
+ p := netmsg.Packet{
137
+ Src: sp.RemotePeer(),
138
+ Dst: sp.swarm.client().Address(),
139
+ Data: msg,
140
+ Context: ctx,
141
+ }
142
+
143
+ // We also can't yet pass unread io.RW to the clients directly.
144
+ // muxado, SPDY, QUIC, and other stream multiplexors could
145
+ // make this a breeze.
146
+
147
+ // this runs the entire request. it should not return until ALL action
148
+ // is done. this is so that we rate limit and respond to backpressure well.
149
+ // TODO: pipelining (handle up to N concurrent requests).
150
+ // doing pipelining with SPDY or muxado is probably TRTTD.
151
+ if err := sp.HandlePacket(&p, nil); err != nil {
152
+ log.Errorf("error handling incoming request: %v", err)
153
+ }
154
+
155
+ // should be done with the underlying bytes. release (the kraken)!
156
+ // TODO: enable this. there is a bug relating to mpool or something. swarm_tests fail.
157
+ // sp.conn.ReleaseMsg(msg)
158
+}
159
+
160
+func (sp *swarmPeer) HandlePacket(p router.Packet, n router.Node) error {
161
+ switch p.Destination() {
162
+ case sp.swarm.client().Address(): // incoming
163
+ return sp.swarm.HandlePacket(p, sp)
164
+
165
+ case sp.RemotePeer(): // outgoing
166
+ buf, ok := p.Payload().([]byte)
167
+ if !ok {
168
+ return netmsg.ErrInvalidPayload
169
+ }
170
+ if err := sp.conn.WriteMsg(buf); err != nil {
171
+ return fmt.Errorf("swarmPeer error sending: %s", err)
172
+ }
173
+ return nil
174
+
175
+ default: // problem
176
+ return fmt.Errorf("swarmPeer routing error: %v got %v", sp, p)
177
+ }
178
+}
net/swarm/swarm_rt.go
new
+159
@@ -0,0 +1,159 @@
1
+package swarm
2
+
3
+import (
4
+ "sync"
5
+
6
+ conn "github.com/jbenet/go-ipfs/net/conn"
7
+ netmsg "github.com/jbenet/go-ipfs/net/message"
8
+ peer "github.com/jbenet/go-ipfs/peer"
9
+ u "github.com/jbenet/go-ipfs/util"
10
+
11
+ router "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-router"
12
+)
13
+
14
+// swarmRoutingTable collects the peers
15
+type swarmRoutingTable struct {
16
+ local peer.Peer
17
+ client router.Node
18
+ peers map[u.Key]*swarmPeer
19
+ sync.RWMutex
20
+}
21
+
22
+func newRoutingTable(local peer.Peer, client router.Node) *swarmRoutingTable {
23
+ return &swarmRoutingTable{
24
+ local: local,
25
+ client: client,
26
+ peers: map[u.Key]*swarmPeer{},
27
+ }
28
+}
29
+
30
+func (rt *swarmRoutingTable) getOrAdd(s *Swarm, p peer.Peer) (*swarmPeer, error) {
31
+ rt.Lock()
32
+ defer rt.Unlock()
33
+
34
+ sp, ok := rt.peers[p.Key()]
35
+ if ok {
36
+ return sp, nil
37
+ }
38
+
39
+ // newSwarmPeer is what kicks off the reader goroutines.
40
+ sp, err := newSwarmPeer(s, p)
41
+ if err != nil {
42
+ return nil, err
43
+ }
44
+ rt.peers[p.Key()] = sp
45
+ return sp, nil
46
+}
47
+
48
+func (rt *swarmRoutingTable) remove(p peer.Peer) *swarmPeer {
49
+ rt.Lock()
50
+ defer rt.Unlock()
51
+ sp, ok := rt.peers[p.Key()]
52
+ if ok {
53
+ delete(rt.peers, p.Key())
54
+ }
55
+ return sp
56
+}
57
+
58
+func (rt *swarmRoutingTable) getByID(pid peer.ID) *swarmPeer {
59
+ rt.RLock()
60
+ defer rt.RUnlock()
61
+ return rt.peers[u.Key(pid)]
62
+}
63
+
64
+func (rt *swarmRoutingTable) get(p peer.Peer) *swarmPeer {
65
+ rt.RLock()
66
+ defer rt.RUnlock()
67
+ return rt.peers[p.Key()]
68
+}
69
+
70
+func (rt *swarmRoutingTable) connList() []conn.Conn {
71
+ rt.RLock()
72
+ defer rt.RUnlock()
73
+
74
+ var out []conn.Conn
75
+ for _, sp := range rt.peers {
76
+ out = append(out, sp.conn)
77
+ }
78
+ return out
79
+}
80
+
81
+func (rt *swarmRoutingTable) peerList() []peer.Peer {
82
+ rt.RLock()
83
+ defer rt.RUnlock()
84
+
85
+ var out []peer.Peer
86
+ for _, sp := range rt.peers {
87
+ out = append(out, sp.RemotePeer())
88
+ }
89
+ return out
90
+}
91
+
92
+// Route implements routing.Route
93
+func (rt *swarmRoutingTable) Route(p router.Packet) router.Node {
94
+
95
+ // no need to lock :)
96
+ if p.Destination() == rt.client.Address() {
97
+ // log.Debugf("%s swarmRoutingTable route %s to client %s ? ", p.Destination(), rt.client.Address(), p.Payload())
98
+ return rt.client
99
+ }
100
+
101
+ rt.RLock()
102
+ defer rt.RUnlock()
103
+
104
+ for _, sp := range rt.peers {
105
+ if sp.RemotePeer() == p.Destination() {
106
+ // log.Debugf("%s swarmRoutingTable route %s to peer %s ? ", p.Destination(), sp.RemotePeer(), p.Payload())
107
+ return sp
108
+ }
109
+ }
110
+
111
+ return nil // no route
112
+}
113
+
114
+func (s *Swarm) client() router.Node {
115
+ return s.rt.client
116
+}
117
+
118
+func (s *Swarm) addConn(c conn.Conn) error {
119
+ sp, err := s.rt.getOrAdd(s, c.RemotePeer())
120
+ if err != nil {
121
+ return err
122
+ }
123
+
124
+ sp.conn.Add(c)
125
+ return nil
126
+}
127
+
128
+func (s *Swarm) closeConn(p peer.Peer) error {
129
+ sp := s.rt.remove(p)
130
+ if sp == nil {
131
+ return nil
132
+ }
133
+
134
+ return sp.Close()
135
+}
136
+
137
+// HandlePacket routes messages out through connections, or to the client
138
+func (s *Swarm) HandlePacket(p router.Packet, from router.Node) error {
139
+ msg, ok := p.Payload().([]byte)
140
+ if !ok {
141
+ return netmsg.ErrInvalidPayload
142
+ }
143
+
144
+ if len(msg) >= conn.MaxMessageSize {
145
+ log.Criticalf("Attempted to send message bigger than max size. (%d)", len(msg))
146
+ }
147
+
148
+ next := s.rt.Route(p)
149
+ if next == nil {
150
+ // log.Debugf("%s swarm HandlePacket %s -> %s -> %s -> %s: %s",
151
+ // s.local, from.Address(), s.Address(), "????", p.Destination(), p.Payload())
152
+ return router.ErrNoRoute
153
+ }
154
+
155
+ // log.Debugf("%s swarm HandlePacket %s -> %s -> %s -> %s: %s",
156
+ // s.local, from.Address(), s.Address(), next.Address(), p.Destination(), p.Payload())
157
+
158
+ return next.HandlePacket(p, s)
159
+}
net/swarm/swarm_test.go
+96
-31
@@ -7,27 +7,86 @@ import (
7
"time"
8
9
ci "github.com/jbenet/go-ipfs/crypto"
10
- msg "github.com/jbenet/go-ipfs/net/message"
10
+ netmsg "github.com/jbenet/go-ipfs/net/message"
11
peer "github.com/jbenet/go-ipfs/peer"
12
u "github.com/jbenet/go-ipfs/util"
13
testutil "github.com/jbenet/go-ipfs/util/testutil"
14
15
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
16
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
17
+ router "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-router"
18
)
19
19
-func pong(ctx context.Context, swarm *Swarm) {
20
- i := 0
20
+// needed to copy the data. otherwise gets reused... :/
21
+type queueClient struct {
22
+ addr router.Address
23
+ queue chan *netmsg.Packet
24
+}
25
+
26
+func newQueueClient(addr router.Address) *queueClient {
27
+ return &queueClient{addr, make(chan *netmsg.Packet, 20)}
28
+}
29
+
30
+func (qc *queueClient) Address() router.Address {
31
+ return qc.addr
32
+}
33
+
34
+func (qc *queueClient) HandlePacket(p router.Packet, n router.Node) error {
35
+
36
+ pkt1 := p.(*netmsg.Packet)
37
+ pkt2 := netmsg.Packet{}
38
+ pkt2 = *pkt1
39
+ pkt2.Data = make([]byte, len(pkt1.Data))
40
+ copy(pkt2.Data, pkt1.Data)
41
+
42
+ qc.queue <- &pkt2
43
+ return nil
44
+}
45
+
46
+type pongClient struct {
47
+ peer peer.Peer
48
+ count int
49
+ queue chan pongPkt
50
+}
51
+
52
+type pongPkt struct {
53
+ msg netmsg.Packet
54
+ dst router.Node
55
+}
56
+
57
+func newPongClient(ctx context.Context, peer peer.Peer) *pongClient {
58
+ pc := &pongClient{peer: peer, queue: make(chan pongPkt, 10)}
59
+ go pc.echo(ctx)
60
+ return pc
61
+}
62
+
63
+func (pc *pongClient) Address() router.Address {
64
+ return pc.peer.ID().Pretty() + "/pong"
65
+}
66
+
67
+func (pc *pongClient) HandlePacket(p router.Packet, n router.Node) error {
68
+ pkt1 := p.(*netmsg.Packet)
69
+ if !bytes.Equal(pkt1.Data, []byte("ping")) {
70
+ log.Debugf("%s pong dropped pkt: %s (%s -> %s)", pc.Address(), pkt1.Data, pkt1.Src, pkt1.Dst)
71
+ panic("why")
72
+ return nil // drop
73
+ }
74
+
75
+ pc.queue <- pongPkt{pkt1.Response([]byte("pong")), n}
76
+ return nil
77
+}
78
+
79
+func (pc *pongClient) echo(ctx context.Context) {
80
for {
81
select {
82
case <-ctx.Done():
83
return
25
- case m1 := <-swarm.Incoming:
26
- if bytes.Equal(m1.Data(), []byte("ping")) {
27
- m2 := msg.New(m1.Peer(), []byte("pong"))
28
- i++
29
- log.Debugf("%s pong %s (%d)", swarm.local, m1.Peer(), i)
30
- swarm.Outgoing <- m2
84
+
85
+ case pkt := <-pc.queue:
86
+ pc.count++
87
+ log.Debugf("%s pong %s (%d)", pkt.msg.Src, pkt.msg.Dst, pc.count)
88
+ if err := pkt.dst.HandlePacket(&pkt.msg, pc); err != nil {
89
+ log.Errorf("pong error sending: %s", err)
90
}
91
}
92
}
@@ -58,7 +117,8 @@ func makeSwarms(ctx context.Context, t *testing.T, addrs []string) ([]*Swarm, []
117
for _, addr := range addrs {
118
local := setupPeer(t, addr)
119
peerstore := peer.NewPeerstore()
61
- swarm, err := NewSwarm(ctx, local.Addresses(), local, peerstore)
120
+ pong := newPongClient(ctx, local)
121
+ swarm, err := NewSwarm(ctx, local.Addresses(), local, peerstore, pong)
122
if err != nil {
123
t.Fatal(err)
124
}
@@ -91,11 +151,11 @@ func SubtestSwarm(t *testing.T, addrs []string, MsgNum int) {
151
}
152
cp.AddAddress(dst.Addresses()[0])
153
94
- log.Info("SWARM TEST: %s dialing %s", s.local, dst)
154
+ log.Infof("SWARM TEST: %s dialing %s", s.local, dst)
155
if _, err := s.Dial(cp); err != nil {
156
t.Fatal("error swarm dialing to peer", err)
157
}
98
- log.Info("SWARM TEST: %s connected to %s", s.local, dst)
158
+ log.Infof("SWARM TEST: %s connected to %s", s.local, dst)
159
wg.Done()
160
}
161
@@ -109,21 +169,24 @@ func SubtestSwarm(t *testing.T, addrs []string, MsgNum int) {
169
}
170
}
171
wg.Wait()
172
+
173
+ for _, s := range swarms {
174
+ log.Infof("%s swarm routing table: %s", s.local, s.GetPeerList())
175
+ }
176
}
177
178
// ping/pong
179
for _, s1 := range swarms {
116
- ctx, cancel := context.WithCancel(ctx)
117
-
118
- // setup all others to pong
119
- for _, s2 := range swarms {
120
- if s1 == s2 {
121
- continue
122
- }
180
+ log.Debugf("-------------------------------------------------------")
181
+ log.Debugf("%s ping pong round", s1.local)
182
+ log.Debugf("-------------------------------------------------------")
183
124
- go pong(ctx, s2)
125
- }
184
+ // for this test, we'll listen on s1.
185
+ queue := newQueueClient(s1.client().Address())
186
+ pong := s1.client() // set it back at the end.
187
+ s1.SetClient(queue)
188
189
+ ctx, cancel := context.WithCancel(ctx)
190
peers, err := s1.peers.All()
191
if err != nil {
192
t.Fatal(err)
@@ -132,22 +195,26 @@ func SubtestSwarm(t *testing.T, addrs []string, MsgNum int) {
195
for k := 0; k < MsgNum; k++ {
196
for _, p := range *peers {
197
log.Debugf("%s ping %s (%d)", s1.local, p, k)
135
- s1.Outgoing <- msg.New(p, []byte("ping"))
198
+ pkt := netmsg.Packet{Src: s1.local, Dst: p, Data: []byte("ping"), Context: ctx}
199
+ s1.HandlePacket(&pkt, queue)
200
}
201
}
202
203
got := map[u.Key]int{}
204
for k := 0; k < (MsgNum * len(*peers)); k++ {
205
log.Debugf("%s waiting for pong (%d)", s1.local, k)
142
- msg := <-s1.Incoming
143
- if string(msg.Data()) != "pong" {
144
- t.Error("unexpected conn output", msg.Data)
206
+
207
+ msg := <-queue.queue
208
+ if string(msg.Data) != "pong" {
209
+ t.Error("unexpected conn output", string(msg.Data), msg.Data)
210
}
211
147
- n, _ := got[msg.Peer().Key()]
148
- got[msg.Peer().Key()] = n + 1
212
+ p := msg.Src.(peer.Peer)
213
+ n, _ := got[p.Key()]
214
+ got[p.Key()] = n + 1
215
}
216
217
+ log.Debugf("%s got pongs", s1.local)
218
if len(*peers) != len(got) {
219
t.Error("got less messages than sent")
220
}
@@ -159,7 +226,8 @@ func SubtestSwarm(t *testing.T, addrs []string, MsgNum int) {
226
}
227
228
cancel()
162
- <-time.After(50 * time.Microsecond)
229
+ <-time.After(10 * time.Millisecond)
230
+ s1.SetClient(pong)
231
}
232
233
for _, s := range swarms {
@@ -168,9 +236,6 @@ func SubtestSwarm(t *testing.T, addrs []string, MsgNum int) {
236
}
237
238
func TestSwarm(t *testing.T) {
171
- if testing.Short() {
172
- t.SkipNow()
173
- }
239
// t.Skip("skipping for another test")
240
241
addrs := []string{