@cryptotaxi247 / kubo / commits / df6540e01

p2p: rework stream/listener registration

License: MIT Signed-off-by: Łukasz Magiera <magik6k@gmail.com>

Łukasz Magiera committed May 25, 2018 at 15:54 UTC df6540e014841501ebce0c43735404ffdb0ddfc4
6 files changed +109 -195
core/commands/p2p.go
+66 -162
@@ -30,10 +30,8 @@ type P2PListenerInfoOutput struct {
30 type P2PStreamInfoOutput struct {
31 HandlerID string
32 Protocol string
33 - LocalPeer string
34 - LocalAddress string
35 - RemotePeer string
36 - RemoteAddress string
33 + OriginAddress string
34 + TargetAddress string
35 }
36
37 // P2PLsOutput is output type of ls command
@@ -61,6 +59,7 @@ are refined`,
59 "stream": p2pStreamCmd,
60
61 "forward": p2pForwardCmd,
62 + "close": p2pCloseCmd,
63 "ls": p2pLsCmd,
64 },
65 }
@@ -209,6 +208,65 @@ var p2pLsCmd = &cmds.Command{
208 },
209 }
210
211 +var p2pCloseCmd = &cmds.Command{
212 + Helptext: cmdkit.HelpText{
213 + Tagline: "Stop listening for new connections to forward.",
214 + },
215 + Options: []cmdkit.Option{
216 + cmdkit.BoolOption("all", "a", "Close all listeners."),
217 + cmdkit.StringOption("protocol", "p", "Match protocol name"),
218 + cmdkit.StringOption("listen-address", "l", "Match listen address"),
219 + cmdkit.StringOption("target-address", "t", "Match target address"),
220 + },
221 + Run: func(req cmds.Request, res cmds.Response) {
222 + res.SetOutput(nil)
223 +
224 + n, err := getNode(req)
225 + if err != nil {
226 + res.SetError(err, cmdkit.ErrNormal)
227 + return
228 + }
229 +
230 + closeAll, _, _ := req.Option("all").Bool()
231 + proto, p, _ := req.Option("protocol").String()
232 + listen, l, _ := req.Option("listen-address").String()
233 + target, t, _ := req.Option("target-address").String()
234 +
235 + if !(closeAll || p || l || t) {
236 + res.SetError(errors.New("no connection matching options given"), cmdkit.ErrNormal)
237 + return
238 + }
239 +
240 + if closeAll && (p || l || t) {
241 + res.SetError(errors.New("can't combine --all with other matching options"), cmdkit.ErrNormal)
242 + return
243 + }
244 +
245 + match := func(listener p2p.Listener) bool {
246 + out := true
247 + if p {
248 + out = out && (proto == listener.Protocol())
249 + }
250 + if l {
251 + out = out && (listen == listener.ListenAddress())
252 + }
253 + if t {
254 + out = out && (target == listener.TargetAddress())
255 + }
256 +
257 + out = out || closeAll
258 + return out
259 + }
260 +
261 + for _, listener := range n.P2P.Listeners.Listeners {
262 + if !match(listener) {
263 + continue
264 + }
265 + listener.Close()
266 + }
267 + },
268 +}
269 +
270 ///////
271 // Listener
272 //
@@ -222,7 +280,6 @@ var p2pStreamCmd = &cmds.Command{
280
281 Subcommands: map[string]*cmds.Command{
282 "ls": p2pStreamLsCmd,
225 - "dial": p2pStreamDialCmd,
283 "close": p2pStreamCloseCmd,
284 },
285 }
@@ -249,11 +306,8 @@ var p2pStreamLsCmd = &cmds.Command{
306
307 Protocol: s.Protocol,
308
252 - LocalPeer: s.LocalPeer.Pretty(),
253 - LocalAddress: s.LocalAddr.String(),
254 -
255 - RemotePeer: s.RemotePeer.Pretty(),
256 - RemoteAddress: s.RemoteAddr.String(),
309 + OriginAddress: s.OriginAddr.String(),
310 + TargetAddress: s.TargetAddr.String(),
311 })
312 }
313
@@ -273,10 +327,10 @@ var p2pStreamLsCmd = &cmds.Command{
327 w := tabwriter.NewWriter(buf, 1, 2, 1, ' ', 0)
328 for _, stream := range list.Streams {
329 if headers {
276 - fmt.Fprintln(w, "Id\tProtocol\tLocal\tRemote")
330 + fmt.Fprintln(w, "Id\tProtocol\tOrigin\tTarget")
331 }
332
279 - fmt.Fprintf(w, "%s\t%s\t%s\t%s\n", stream.HandlerID, stream.Protocol, stream.LocalAddress, stream.RemotePeer)
333 + fmt.Fprintf(w, "%s\t%s\t%s\t%s\n", stream.HandlerID, stream.Protocol, stream.OriginAddress, stream.TargetAddress)
334 }
335 w.Flush()
336
@@ -285,156 +339,6 @@ var p2pStreamLsCmd = &cmds.Command{
339 },
340 }
341
288 -var p2pListenerListenCmd = &cmds.Command{
289 - Helptext: cmdkit.HelpText{
290 - Tagline: "Forward p2p connections to a network multiaddr.",
291 - ShortDescription: `
292 -Register a p2p connection handler and forward the connections to a specified
293 -address.
294 -
295 -Note that the connections originate from the ipfs daemon process.
296 - `,
297 - },
298 - Arguments: []cmdkit.Argument{
299 - cmdkit.StringArg("Protocol", true, false, "Protocol identifier."),
300 - cmdkit.StringArg("Address", true, false, "Request handling application address."),
301 - },
302 - Run: func(req cmds.Request, res cmds.Response) {
303 - n, err := getNode(req)
304 - if err != nil {
305 - res.SetError(err, cmdkit.ErrNormal)
306 - return
307 - }
308 -
309 - proto := "/p2p/" + req.Arguments()[0]
310 - if n.P2P.CheckProtoExists(proto) {
311 - res.SetError(errors.New("protocol handler already registered"), cmdkit.ErrNormal)
312 - return
313 - }
314 -
315 - addr, err := ma.NewMultiaddr(req.Arguments()[1])
316 - if err != nil {
317 - res.SetError(err, cmdkit.ErrNormal)
318 - return
319 - }
320 -
321 - _, err = n.P2P.ForwardRemote(n.Context(), proto, addr)
322 - if err != nil {
323 - res.SetError(err, cmdkit.ErrNormal)
324 - return
325 - }
326 -
327 - // Successful response.
328 - res.SetOutput(&P2PListenerInfoOutput{
329 - Protocol: proto,
330 - TargetAddress: addr.String(),
331 - })
332 - },
333 -}
334 -
335 -var p2pStreamDialCmd = &cmds.Command{
336 - Helptext: cmdkit.HelpText{
337 - Tagline: "Dial to a p2p listener.",
338 -
339 - ShortDescription: `
340 -Establish a new connection to a peer service.
341 -
342 -When a connection is made to a peer service the ipfs daemon will setup one
343 -time TCP listener and return it's bind port, this way a dialing application
344 -can transparently connect to a p2p service.
345 - `,
346 - },
347 - Arguments: []cmdkit.Argument{
348 - cmdkit.StringArg("Peer", true, false, "Remote peer to connect to"),
349 - cmdkit.StringArg("Protocol", true, false, "Protocol identifier."),
350 - cmdkit.StringArg("BindAddress", false, false, "Address to listen for connection/s (default: /ip4/127.0.0.1/tcp/0)."),
351 - },
352 - Run: func(req cmds.Request, res cmds.Response) {
353 - n, err := getNode(req)
354 - if err != nil {
355 - res.SetError(err, cmdkit.ErrNormal)
356 - return
357 - }
358 -
359 - addr, peer, err := ParsePeerParam(req.Arguments()[0])
360 - if err != nil {
361 - res.SetError(err, cmdkit.ErrNormal)
362 - return
363 - }
364 -
365 - if addr != nil {
366 - n.Peerstore.AddAddr(peer, addr, pstore.TempAddrTTL)
367 - }
368 -
369 - proto := "/p2p/" + req.Arguments()[1]
370 -
371 - bindAddr, _ := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/0")
372 - if len(req.Arguments()) == 3 {
373 - bindAddr, err = ma.NewMultiaddr(req.Arguments()[2])
374 - if err != nil {
375 - res.SetError(err, cmdkit.ErrNormal)
376 - return
377 - }
378 - }
379 -
380 - listenerInfo, err := n.P2P.ForwardLocal(n.Context(), peer, proto, bindAddr)
381 - if err != nil {
382 - res.SetError(err, cmdkit.ErrNormal)
383 - return
384 - }
385 -
386 - output := P2PListenerInfoOutput{
387 - Protocol: listenerInfo.Protocol(),
388 - ListenAddress: listenerInfo.ListenAddress(),
389 - }
390 -
391 - res.SetOutput(&output)
392 - },
393 -}
394 -
395 -var p2pListenerCloseCmd = &cmds.Command{
396 - Helptext: cmdkit.HelpText{
397 - Tagline: "Close active p2p listener.",
398 - },
399 - Arguments: []cmdkit.Argument{
400 - cmdkit.StringArg("Protocol", false, false, "P2P listener protocol"),
401 - },
402 - Options: []cmdkit.Option{
403 - cmdkit.BoolOption("all", "a", "Close all listeners."),
404 - },
405 - Run: func(req cmds.Request, res cmds.Response) {
406 - res.SetOutput(nil)
407 -
408 - n, err := getNode(req)
409 - if err != nil {
410 - res.SetError(err, cmdkit.ErrNormal)
411 - return
412 - }
413 -
414 - closeAll, _, _ := req.Option("all").Bool()
415 - var proto string
416 -
417 - if !closeAll {
418 - if len(req.Arguments()) == 0 {
419 - res.SetError(errors.New("no protocol name specified"), cmdkit.ErrNormal)
420 - return
421 - }
422 -
423 - proto = "/p2p/" + req.Arguments()[0]
424 - }
425 -
426 - for _, listener := range n.P2P.Listeners.Listeners {
427 - if !closeAll && listener.Protocol() != proto {
428 - continue
429 - }
430 - listener.Close()
431 - if !closeAll {
432 - break
433 - }
434 - }
435 - },
436 -}
437 -
342 var p2pStreamCloseCmd = &cmds.Command{
343 Helptext: cmdkit.HelpText{
344 Tagline: "Close active p2p stream.",
p2p/listener.go
+20 -6
@@ -9,17 +9,31 @@ type Listener interface {
9 Close() error
10 }
11
12 +type listenerKey struct {
13 + proto string
14 + listen string
15 + target string
16 +}
17 +
18 // ListenerRegistry is a collection of local application proto listeners.
19 type ListenerRegistry struct {
14 - Listeners map[string]Listener
20 + Listeners map[listenerKey]Listener
21 }
22
17 -// Register registers listenerInfo2 in this registry
18 -func (c *ListenerRegistry) Register(listenerInfo Listener) {
19 - c.Listeners[listenerInfo.Protocol()] = listenerInfo
23 +// Register registers listenerInfo in this registry
24 +func (c *ListenerRegistry) Register(l Listener) {
25 + c.Listeners[getListenerKey(l)] = l
26 }
27
28 // Deregister removes p2p listener from this registry
23 -func (c *ListenerRegistry) Deregister(proto string) {
24 - delete(c.Listeners, proto)
29 +func (c *ListenerRegistry) Deregister(k listenerKey) {
30 + delete(c.Listeners, k)
31 +}
32 +
33 +func getListenerKey(l Listener) listenerKey {
34 + return listenerKey{
35 + proto: l.Protocol(),
36 + listen: l.ListenAddress(),
37 + target: l.TargetAddress(),
38 + }
39 }
p2p/local.go
+12 -9
@@ -75,29 +75,32 @@ func (l *localListener) acceptConns() {
75 return
76 }
77
78 - stream := Stream{
79 - Protocol: l.proto,
78 + tgt, err := ma.NewMultiaddr(l.TargetAddress())
79 + if err != nil {
80 + local.Close()
81 + return
82 + }
83
81 - LocalPeer: l.id,
82 - LocalAddr: l.listener.Multiaddr(),
84 + stream := &Stream{
85 + Protocol: l.proto,
86
84 - RemotePeer: remote.Conn().RemotePeer(),
85 - RemoteAddr: remote.Conn().RemoteMultiaddr(),
87 + OriginAddr: local.RemoteMultiaddr(),
88 + TargetAddr: tgt,
89
90 Local: local,
91 Remote: remote,
92
90 - Registry: &l.p2p.Streams,
93 + Registry: l.p2p.Streams,
94 }
95
93 - l.p2p.Streams.Register(&stream)
96 + l.p2p.Streams.Register(stream)
97 stream.startStreaming()
98 }
99 }
100
101 func (l *localListener) Close() error {
102 l.listener.Close()
100 - l.p2p.Listeners.Deregister(l.proto)
103 + l.p2p.Listeners.Deregister(getListenerKey(l))
104 return nil
105 }
106
p2p/p2p.go
+5 -5
@@ -8,8 +8,8 @@ import (
8
9 // P2P structure holds information on currently running streams/listeners
10 type P2P struct {
11 - Listeners ListenerRegistry
12 - Streams StreamRegistry
11 + Listeners *ListenerRegistry
12 + Streams *StreamRegistry
13
14 identity peer.ID
15 peerHost p2phost.Host
@@ -23,10 +23,10 @@ func NewP2P(identity peer.ID, peerHost p2phost.Host, peerstore pstore.Peerstore)
23 peerHost: peerHost,
24 peerstore: peerstore,
25
26 - Listeners: ListenerRegistry{
27 - Listeners: map[string]Listener{},
26 + Listeners: &ListenerRegistry{
27 + Listeners: map[listenerKey]Listener{},
28 },
29 - Streams: StreamRegistry{
29 + Streams: &StreamRegistry{
30 Streams: map[uint64]*Stream{},
31 },
32 }
p2p/remote.go
+4 -7
@@ -41,16 +41,13 @@ func (p2p *P2P) ForwardRemote(ctx context.Context, proto string, addr ma.Multiad
41 stream := Stream{
42 Protocol: proto,
43
44 - LocalPeer: p2p.identity,
45 - LocalAddr: addr,
46 -
47 - RemotePeer: remote.Conn().RemotePeer(),
48 - RemoteAddr: remote.Conn().RemoteMultiaddr(),
44 + OriginAddr: remote.Conn().RemoteMultiaddr(),
45 + TargetAddr: addr,
46
47 Local: local,
48 Remote: remote,
49
53 - Registry: &p2p.Streams,
50 + Registry: p2p.Streams,
51 }
52
53 p2p.Streams.Register(&stream)
@@ -74,6 +71,6 @@ func (l *remoteListener) TargetAddress() string {
71
72 func (l *remoteListener) Close() error {
73 l.p2p.peerHost.RemoveStreamHandler(protocol.ID(l.proto))
77 - l.p2p.Listeners.Deregister(l.proto)
74 + l.p2p.Listeners.Deregister(getListenerKey(l))
75 return nil
76 }
p2p/stream.go
+2 -6
@@ -6,7 +6,6 @@ import (
6 ma "gx/ipfs/QmWWQ2Txc2c6tqjsBpzg5Ar652cHPGNsQQp2SejkNmkUMb/go-multiaddr"
7 net "gx/ipfs/QmYj8wdn5sZEHX2XMDWGBvcXJNdzVbaVpHmXvhHBVZepen/go-libp2p-net"
8 manet "gx/ipfs/QmcGXGdw9BWDysPJQHxJinjGHha3eEg4vzFETre4woNwcX/go-multiaddr-net"
9 - peer "gx/ipfs/QmcJukH2sAFjY3HdBKq35WDzWoL3UUu2gt9wdfqZTUyM74/go-libp2p-peer"
9 )
10
11 // Stream holds information on active incoming and outgoing p2p streams.
@@ -15,11 +14,8 @@ type Stream struct {
14
15 Protocol string
16
18 - LocalPeer peer.ID
19 - LocalAddr ma.Multiaddr
20 -
21 - RemotePeer peer.ID
22 - RemoteAddr ma.Multiaddr
17 + OriginAddr ma.Multiaddr
18 + TargetAddr ma.Multiaddr
19
20 Local manet.Conn
21 Remote net.Stream