@cryptotaxi247 / kubo / commits / 68b85c992

broke out dial + listen

Juan Batiz-Benet committed Oct 19, 2014 at 03:33 UTC 68b85c992b39ee32fe2329194516bece559867af
5 files changed +213 -179
net/conn/conn.go
+9 -179
@@ -101,6 +101,10 @@ func (c *singleConn) ID() string {
101 return ID(c)
102 }
103
104 +func (c *singleConn) String() string {
105 + return String(c, "singleConn")
106 +}
107 +
108 // LocalMultiaddr is the Multiaddr on this side
109 func (c *singleConn) LocalMultiaddr() ma.Multiaddr {
110 return c.maconn.LocalMultiaddr()
@@ -131,7 +135,7 @@ func (c *singleConn) Out() chan<- []byte {
135 return c.msgio.outgoing.MsgChan
136 }
137
134 -// ID returns the
138 +// ID returns the ID of a given Conn.
139 func ID(c Conn) string {
140 l := fmt.Sprintf("%s/%s", c.LocalMultiaddr(), c.LocalPeer().ID)
141 r := fmt.Sprintf("%s/%s", c.RemoteMultiaddr(), c.RemotePeer().ID)
@@ -141,182 +145,8 @@ func ID(c Conn) string {
145 return u.Key(ch).Pretty()
146 }
147
144 -// Dialer is an object that can open connections. We could have a "convenience"
145 -// Dial function as before, but it would have many arguments, as dialing is
146 -// no longer simple (need a peerstore, a local peer, a context, a network, etc)
147 -type Dialer struct {
148 -
149 - // LocalPeer is the identity of the local Peer.
150 - LocalPeer *peer.Peer
151 -
152 - // Peerstore is the set of peers we know about locally. The Dialer needs it
153 - // because when an incoming connection is identified, we should reuse the
154 - // same peer objects (otherwise things get inconsistent).
155 - Peerstore peer.Peerstore
156 -}
157 -
158 -// Dial connects to a particular peer, over a given network
159 -// Example: d.Dial(ctx, "udp", peer)
160 -func (d *Dialer) Dial(ctx context.Context, network string, remote *peer.Peer) (Conn, error) {
161 - laddr := d.LocalPeer.NetAddress(network)
162 - if laddr == nil {
163 - return nil, fmt.Errorf("No local address for network %s", network)
164 - }
165 -
166 - raddr := remote.NetAddress(network)
167 - if raddr == nil {
168 - return nil, fmt.Errorf("No remote address for network %s", network)
169 - }
170 -
171 - // TODO: try to get reusing addr/ports to work.
172 - // madialer := manet.Dialer{LocalAddr: laddr}
173 - madialer := manet.Dialer{}
174 -
175 - log.Info("%s dialing %s %s", d.LocalPeer, remote, raddr)
176 - maconn, err := madialer.Dial(raddr)
177 - if err != nil {
178 - return nil, err
179 - }
180 -
181 - if err := d.Peerstore.Put(remote); err != nil {
182 - log.Error("Error putting peer into peerstore: %s", remote)
183 - }
184 -
185 - c, err := newSingleConn(ctx, d.LocalPeer, remote, maconn)
186 - if err != nil {
187 - return nil, err
188 - }
189 -
190 - return newSecureConn(ctx, c, d.Peerstore)
191 -}
192 -
193 -// listener is an object that can accept connections. It implements Listener
194 -type listener struct {
195 - manet.Listener
196 -
197 - // chansize is the size of the internal channels for concurrency
198 - chansize int
199 -
200 - // channel of incoming conections
201 - conns chan Conn
202 -
203 - // Local multiaddr to listen on
204 - maddr ma.Multiaddr
205 -
206 - // LocalPeer is the identity of the local Peer.
207 - local *peer.Peer
208 -
209 - // Peerstore is the set of peers we know about locally
210 - peers peer.Peerstore
211 -
212 - // embedded ContextCloser
213 - ContextCloser
214 -}
215 -
216 -// disambiguate
217 -func (l *listener) Close() error {
218 - return l.ContextCloser.Close()
219 -}
220 -
221 -// close called by ContextCloser.Close
222 -func (l *listener) close() error {
223 - log.Info("listener closing: %s %s", l.local, l.maddr)
224 - return l.Listener.Close()
225 -}
226 -
227 -func (l *listener) listen() {
228 - l.Children().Add(1)
229 - defer l.Children().Done()
230 -
231 - // handle at most chansize concurrent handshakes
232 - sem := make(chan struct{}, l.chansize)
233 -
234 - // handle is a goroutine work function that handles the handshake.
235 - // it's here only so that accepting new connections can happen quickly.
236 - handle := func(maconn manet.Conn) {
237 - defer func() { <-sem }() // release
238 -
239 - c, err := newSingleConn(l.Context(), l.local, nil, maconn)
240 - if err != nil {
241 - log.Error("Error accepting connection: %v", err)
242 - return
243 - }
244 -
245 - sc, err := newSecureConn(l.Context(), c, l.peers)
246 - if err != nil {
247 - log.Error("Error securing connection: %v", err)
248 - return
249 - }
250 -
251 - l.conns <- sc
252 - }
253 -
254 - for {
255 - maconn, err := l.Listener.Accept()
256 - if err != nil {
257 -
258 - // if closing, we should exit.
259 - select {
260 - case <-l.Closing():
261 - return // done.
262 - default:
263 - }
264 -
265 - log.Error("Failed to accept connection: %v", err)
266 - continue
267 - }
268 -
269 - sem <- struct{}{} // acquire
270 - go handle(maconn)
271 - }
272 -}
273 -
274 -// Accept waits for and returns the next connection to the listener.
275 -// Note that unfortunately this
276 -func (l *listener) Accept() <-chan Conn {
277 - return l.conns
278 -}
279 -
280 -// Multiaddr is the identity of the local Peer.
281 -func (l *listener) Multiaddr() ma.Multiaddr {
282 - return l.maddr
283 -}
284 -
285 -// LocalPeer is the identity of the local Peer.
286 -func (l *listener) LocalPeer() *peer.Peer {
287 - return l.local
288 -}
289 -
290 -// Peerstore is the set of peers we know about locally. The Listener needs it
291 -// because when an incoming connection is identified, we should reuse the
292 -// same peer objects (otherwise things get inconsistent).
293 -func (l *listener) Peerstore() peer.Peerstore {
294 - return l.peers
295 -}
296 -
297 -// Listen listens on the particular multiaddr, with given peer and peerstore.
298 -func Listen(ctx context.Context, addr ma.Multiaddr, local *peer.Peer, peers peer.Peerstore) (Listener, error) {
299 -
300 - ml, err := manet.Listen(addr)
301 - if err != nil {
302 - return nil, err
303 - }
304 -
305 - // todo make this a variable
306 - chansize := 10
307 -
308 - l := &listener{
309 - Listener: ml,
310 - maddr: addr,
311 - peers: peers,
312 - local: local,
313 - conns: make(chan Conn, chansize),
314 - chansize: chansize,
315 - }
316 -
317 - l.ContextCloser = NewContextCloser(ctx, l.close)
318 -
319 - go l.listen()
320 -
321 - return l, nil
148 +// String returns the user-friendly String representation of a conn
149 +func String(c Conn, typ string) string {
150 + return fmt.Sprintf("%s (%s) <-- %s --> (%s) %s",
151 + c.LocalPeer(), c.LocalMultiaddr(), typ, c.RemoteMultiaddr(), c.RemotePeer())
152 }
net/conn/dial.go new
+46
@@ -0,0 +1,46 @@
1 +package conn
2 +
3 +import (
4 + "fmt"
5 +
6 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7 +
8 + manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr/net"
9 +
10 + peer "github.com/jbenet/go-ipfs/peer"
11 +)
12 +
13 +// Dial connects to a particular peer, over a given network
14 +// Example: d.Dial(ctx, "udp", peer)
15 +func (d *Dialer) Dial(ctx context.Context, network string, remote *peer.Peer) (Conn, error) {
16 + laddr := d.LocalPeer.NetAddress(network)
17 + if laddr == nil {
18 + return nil, fmt.Errorf("No local address for network %s", network)
19 + }
20 +
21 + raddr := remote.NetAddress(network)
22 + if raddr == nil {
23 + return nil, fmt.Errorf("No remote address for network %s", network)
24 + }
25 +
26 + // TODO: try to get reusing addr/ports to work.
27 + // madialer := manet.Dialer{LocalAddr: laddr}
28 + madialer := manet.Dialer{}
29 +
30 + log.Info("%s dialing %s %s", d.LocalPeer, remote, raddr)
31 + maconn, err := madialer.Dial(raddr)
32 + if err != nil {
33 + return nil, err
34 + }
35 +
36 + if err := d.Peerstore.Put(remote); err != nil {
37 + log.Error("Error putting peer into peerstore: %s", remote)
38 + }
39 +
40 + c, err := newSingleConn(ctx, d.LocalPeer, remote, maconn)
41 + if err != nil {
42 + return nil, err
43 + }
44 +
45 + return newSecureConn(ctx, c, d.Peerstore)
46 +}
net/conn/interface.go
+14
@@ -40,6 +40,20 @@ type Conn interface {
40 // Close() error -- already in ContextCloser
41 }
42
43 +// Dialer is an object that can open connections. We could have a "convenience"
44 +// Dial function as before, but it would have many arguments, as dialing is
45 +// no longer simple (need a peerstore, a local peer, a context, a network, etc)
46 +type Dialer struct {
47 +
48 + // LocalPeer is the identity of the local Peer.
49 + LocalPeer *peer.Peer
50 +
51 + // Peerstore is the set of peers we know about locally. The Dialer needs it
52 + // because when an incoming connection is identified, we should reuse the
53 + // same peer objects (otherwise things get inconsistent).
54 + Peerstore peer.Peerstore
55 +}
56 +
57 // Listener is an object that can accept connections. It matches net.Listener
58 type Listener interface {
59
net/conn/listen.go new
+140
@@ -0,0 +1,140 @@
1 +package conn
2 +
3 +import (
4 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5 + ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
6 + manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr/net"
7 +
8 + peer "github.com/jbenet/go-ipfs/peer"
9 +)
10 +
11 +// listener is an object that can accept connections. It implements Listener
12 +type listener struct {
13 + manet.Listener
14 +
15 + // chansize is the size of the internal channels for concurrency
16 + chansize int
17 +
18 + // channel of incoming conections
19 + conns chan Conn
20 +
21 + // Local multiaddr to listen on
22 + maddr ma.Multiaddr
23 +
24 + // LocalPeer is the identity of the local Peer.
25 + local *peer.Peer
26 +
27 + // Peerstore is the set of peers we know about locally
28 + peers peer.Peerstore
29 +
30 + // embedded ContextCloser
31 + ContextCloser
32 +}
33 +
34 +// disambiguate
35 +func (l *listener) Close() error {
36 + return l.ContextCloser.Close()
37 +}
38 +
39 +// close called by ContextCloser.Close
40 +func (l *listener) close() error {
41 + log.Info("listener closing: %s %s", l.local, l.maddr)
42 + return l.Listener.Close()
43 +}
44 +
45 +func (l *listener) listen() {
46 + l.Children().Add(1)
47 + defer l.Children().Done()
48 +
49 + // handle at most chansize concurrent handshakes
50 + sem := make(chan struct{}, l.chansize)
51 +
52 + // handle is a goroutine work function that handles the handshake.
53 + // it's here only so that accepting new connections can happen quickly.
54 + handle := func(maconn manet.Conn) {
55 + defer func() { <-sem }() // release
56 +
57 + c, err := newSingleConn(l.Context(), l.local, nil, maconn)
58 + if err != nil {
59 + log.Error("Error accepting connection: %v", err)
60 + return
61 + }
62 +
63 + sc, err := newSecureConn(l.Context(), c, l.peers)
64 + if err != nil {
65 + log.Error("Error securing connection: %v", err)
66 + return
67 + }
68 +
69 + l.conns <- sc
70 + }
71 +
72 + for {
73 + maconn, err := l.Listener.Accept()
74 + if err != nil {
75 +
76 + // if closing, we should exit.
77 + select {
78 + case <-l.Closing():
79 + return // done.
80 + default:
81 + }
82 +
83 + log.Error("Failed to accept connection: %v", err)
84 + continue
85 + }
86 +
87 + sem <- struct{}{} // acquire
88 + go handle(maconn)
89 + }
90 +}
91 +
92 +// Accept waits for and returns the next connection to the listener.
93 +// Note that unfortunately this
94 +func (l *listener) Accept() <-chan Conn {
95 + return l.conns
96 +}
97 +
98 +// Multiaddr is the identity of the local Peer.
99 +func (l *listener) Multiaddr() ma.Multiaddr {
100 + return l.maddr
101 +}
102 +
103 +// LocalPeer is the identity of the local Peer.
104 +func (l *listener) LocalPeer() *peer.Peer {
105 + return l.local
106 +}
107 +
108 +// Peerstore is the set of peers we know about locally. The Listener needs it
109 +// because when an incoming connection is identified, we should reuse the
110 +// same peer objects (otherwise things get inconsistent).
111 +func (l *listener) Peerstore() peer.Peerstore {
112 + return l.peers
113 +}
114 +
115 +// Listen listens on the particular multiaddr, with given peer and peerstore.
116 +func Listen(ctx context.Context, addr ma.Multiaddr, local *peer.Peer, peers peer.Peerstore) (Listener, error) {
117 +
118 + ml, err := manet.Listen(addr)
119 + if err != nil {
120 + return nil, err
121 + }
122 +
123 + // todo make this a variable
124 + chansize := 10
125 +
126 + l := &listener{
127 + Listener: ml,
128 + maddr: addr,
129 + peers: peers,
130 + local: local,
131 + conns: make(chan Conn, chansize),
132 + chansize: chansize,
133 + }
134 +
135 + l.ContextCloser = NewContextCloser(ctx, l.close)
136 +
137 + go l.listen()
138 +
139 + return l, nil
140 +}
net/conn/secure_conn.go
+4
@@ -98,6 +98,10 @@ func (c *secureConn) ID() string {
98 return ID(c)
99 }
100
101 +func (c *secureConn) String() string {
102 + return String(c, "secureConn")
103 +}
104 +
105 // LocalMultiaddr is the Multiaddr on this side
106 func (c *secureConn) LocalMultiaddr() ma.Multiaddr {
107 return c.insecure.LocalMultiaddr()