@cryptotaxi247 / kubo / commits / a3c84e20e

p2p: refactor review

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

Łukasz Magiera committed Jun 20, 2018 at 15:17 UTC a3c84e20ef187492581bcc6cadf8517289f365dc
9 files changed +259 -172
core/commands/p2p.go
+80 -58
@@ -14,9 +14,11 @@ import (
14 core "github.com/ipfs/go-ipfs/core"
15 p2p "github.com/ipfs/go-ipfs/p2p"
16
17 - pstore "gx/ipfs/QmXauCuJzmzapetmC6W4TuDJLL1yFFrVzSHoWv8YdbmnxH/go-libp2p-peerstore"
17 ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
18 + "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
19 + pstore "gx/ipfs/QmZR2XWVVBCtbgBWnQhWk2xcQfaR3W8faQPriAiaaj7rsr/go-libp2p-peerstore"
20 "gx/ipfs/QmdE4gMduCKCGAcczM2F5ioYDfdeKuPix138wrES1YSr7f/go-ipfs-cmdkit"
21 + "gx/ipfs/Qme4QgoVPyQqxVc4G1c2L2wc9TDa6o294rtspGMnBNRujm/go-ipfs-addr"
22 )
23
24 // P2PProtoPrefix is the default required prefix for protocol names
@@ -98,9 +100,23 @@ Example:
100 return
101 }
102
101 - proto := req.Arguments()[0]
102 - listen := req.Arguments()[1]
103 - target := req.Arguments()[2]
103 + protoOpt := req.Arguments()[0]
104 + listenOpt := req.Arguments()[1]
105 + targetOpt := req.Arguments()[2]
106 +
107 + proto := protocol.ID(protoOpt)
108 +
109 + listen, err := ma.NewMultiaddr(listenOpt)
110 + if err != nil {
111 + res.SetError(err, cmdkit.ErrNormal)
112 + return
113 + }
114 +
115 + target, err := ipfsaddr.ParseString(targetOpt)
116 + if err != nil {
117 + res.SetError(err, cmdkit.ErrNormal)
118 + return
119 + }
120
121 allowCustom, _, err := req.Option("allow-custom-protocol").Bool()
122 if err != nil {
@@ -108,7 +124,7 @@ Example:
124 return
125 }
126
111 - if !allowCustom && !strings.HasPrefix(proto, P2PProtoPrefix) {
127 + if !allowCustom && !strings.HasPrefix(string(proto), P2PProtoPrefix) {
128 res.SetError(errors.New("protocol name must be within '"+P2PProtoPrefix+"' namespace"), cmdkit.ErrNormal)
129 return
130 }
@@ -149,8 +165,16 @@ Example:
165 return
166 }
167
152 - proto := req.Arguments()[0]
153 - target := req.Arguments()[1]
168 + protoOpt := req.Arguments()[0]
169 + targetOpt := req.Arguments()[1]
170 +
171 + proto := protocol.ID(protoOpt)
172 +
173 + target, err := ma.NewMultiaddr(targetOpt)
174 + if err != nil {
175 + res.SetError(err, cmdkit.ErrNormal)
176 + return
177 + }
178
179 allowCustom, _, err := req.Option("allow-custom-protocol").Bool()
180 if err != nil {
@@ -158,7 +182,7 @@ Example:
182 return
183 }
184
161 - if !allowCustom && !strings.HasPrefix(proto, P2PProtoPrefix) {
185 + if !allowCustom && !strings.HasPrefix(string(proto), P2PProtoPrefix) {
186 res.SetError(errors.New("protocol name must be within '"+P2PProtoPrefix+"' namespace"), cmdkit.ErrNormal)
187 return
188 }
@@ -173,39 +197,20 @@ Example:
197 }
198
199 // forwardRemote forwards libp2p service connections to a manet address
176 -func forwardRemote(ctx context.Context, p *p2p.P2P, proto string, target string) error {
177 - if strings.HasPrefix(target, "/ipfs") {
178 - return errors.New("cannot forward libp2p service connections to another libp2p service")
179 - }
180 -
181 - addr, err := ma.NewMultiaddr(target)
182 - if err != nil {
183 - return err
184 - }
185 -
200 +func forwardRemote(ctx context.Context, p *p2p.P2P, proto protocol.ID, target ma.Multiaddr) error {
201 // TODO: return some info
187 - _, err = p.ForwardRemote(ctx, proto, addr)
202 + _, err := p.ForwardRemote(ctx, proto, target)
203 return err
204 }
205
206 // forwardLocal forwards local connections to a libp2p service
192 -func forwardLocal(ctx context.Context, p *p2p.P2P, ps pstore.Peerstore, proto string, listen string, target string) error {
193 - bindAddr, err := ma.NewMultiaddr(listen)
194 - if err != nil {
195 - return err
196 - }
197 -
198 - addr, peer, err := ParsePeerParam(target)
199 - if err != nil {
200 - return err
201 - }
202 -
207 +func forwardLocal(ctx context.Context, p *p2p.P2P, ps pstore.Peerstore, proto protocol.ID, bindAddr ma.Multiaddr, addr ipfsaddr.IPFSAddr) error {
208 if addr != nil {
204 - ps.AddAddr(peer, addr, pstore.TempAddrTTL)
209 + ps.AddAddr(addr.ID(), addr.Multiaddr(), pstore.TempAddrTTL)
210 }
211
212 // TODO: return some info
208 - _, err = p.ForwardLocal(ctx, peer, proto, bindAddr)
213 + _, err := p.ForwardLocal(ctx, addr.ID(), proto, bindAddr)
214 return err
215 }
216
@@ -227,9 +232,9 @@ var p2pLsCmd = &cmds.Command{
232
233 for _, listener := range n.P2P.Listeners.Listeners {
234 output.Listeners = append(output.Listeners, P2PListenerInfoOutput{
230 - Protocol: listener.Protocol(),
231 - ListenAddress: listener.ListenAddress(),
232 - TargetAddress: listener.TargetAddress(),
235 + Protocol: string(listener.Protocol()),
236 + ListenAddress: listener.ListenAddress().String(),
237 + TargetAddress: listener.TargetAddress().String(),
238 })
239 }
240
@@ -272,8 +277,6 @@ var p2pCloseCmd = &cmds.Command{
277 cmdkit.StringOption("target-address", "t", "Match target address"),
278 },
279 Run: func(req cmds.Request, res cmds.Response) {
275 - res.SetOutput(nil)
276 -
280 n, err := p2pGetNode(req)
281 if err != nil {
282 res.SetError(err, cmdkit.ErrNormal)
@@ -281,12 +284,26 @@ var p2pCloseCmd = &cmds.Command{
284 }
285
286 closeAll, _, _ := req.Option("all").Bool()
284 - proto, p, _ := req.Option("protocol").String()
285 - listen, l, _ := req.Option("listen-address").String()
286 - target, t, _ := req.Option("target-address").String()
287 + protoOpt, p, _ := req.Option("protocol").String()
288 + listenOpt, l, _ := req.Option("listen-address").String()
289 + targetOpt, t, _ := req.Option("target-address").String()
290 +
291 + proto := protocol.ID(protoOpt)
292 +
293 + listen, err := ma.NewMultiaddr(listenOpt)
294 + if err != nil {
295 + res.SetError(err, cmdkit.ErrNormal)
296 + return
297 + }
298 +
299 + target, err := ma.NewMultiaddr(targetOpt)
300 + if err != nil {
301 + res.SetError(err, cmdkit.ErrNormal)
302 + return
303 + }
304
305 if !(closeAll || p || l || t) {
289 - res.SetError(errors.New("no connection matching options given"), cmdkit.ErrNormal)
306 + res.SetError(errors.New("no matching options given"), cmdkit.ErrNormal)
307 return
308 }
309
@@ -296,31 +313,36 @@ var p2pCloseCmd = &cmds.Command{
313 }
314
315 match := func(listener p2p.Listener) bool {
299 - out := true
300 -
301 - if p {
302 - out = out && (proto == listener.Protocol())
316 + if closeAll {
317 + return true
318 }
304 - if l {
305 - out = out && (listen == listener.ListenAddress())
319 + if p && proto != listener.Protocol() {
320 + return false
321 }
307 - if t {
308 - out = out && (target == listener.TargetAddress())
322 + if l && !listen.Equal(listener.ListenAddress()) {
323 + return false
324 }
310 -
311 - out = out || closeAll
312 - return out
325 + if t && !target.Equal(listener.TargetAddress()) {
326 + return false
327 + }
328 + return true
329 }
330
315 - var closed int
316 - for _, listener := range n.P2P.Listeners.Listeners {
317 - if !match(listener) {
331 + todo := make([]p2p.Listener, 0)
332 + n.P2P.Listeners.Lock()
333 + for _, l := range n.P2P.Listeners.Listeners {
334 + if !match(l) {
335 continue
336 }
320 - listener.Close()
321 - closed++
337 + todo = append(todo, l)
338 }
323 - res.SetOutput(closed)
339 + n.P2P.Listeners.Unlock()
340 +
341 + for _, l := range todo {
342 + l.Close()
343 + }
344 +
345 + res.SetOutput(len(todo))
346 },
347 Type: int(0),
348 Marshalers: cmds.MarshalerMap{
docs/experimental-features.md
+1 -1
@@ -357,7 +357,7 @@ with `ssh [user]@127.0.0.1 -p 2222`.
357 ### Road to being a real feature
358 - [ ] Needs more people to use and report on how well it works / fits use cases
359 - [ ] More documentation
360 -- [ ] Support other protocols (e.g, unix domain sockets)
360 +- [ ] Support other protocols (e.g, unix domain sockets, websockets, etc.)
361
362 ---
363
p2p/listener.go
+33 -22
@@ -3,13 +3,18 @@ package p2p
3 import (
4 "errors"
5 "sync"
6 +
7 + ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
8 + "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
9 )
10
11 // Listener listens for connections and proxies them to a target
12 type Listener interface {
10 - Protocol() string
11 - ListenAddress() string
12 - TargetAddress() string
13 + Protocol() protocol.ID
14 + ListenAddress() ma.Multiaddr
15 + TargetAddress() ma.Multiaddr
16 +
17 + start() error
18
19 // Close closes the listener. Does not affect child streams
20 Close() error
@@ -23,43 +28,49 @@ type listenerKey struct {
28
29 // ListenerRegistry is a collection of local application proto listeners.
30 type ListenerRegistry struct {
31 + sync.Mutex
32 +
33 Listeners map[listenerKey]Listener
27 - lk sync.Mutex
34 }
35
30 -func (r *ListenerRegistry) lock(l Listener) error {
31 - r.lk.Lock()
36 +// Register registers listenerInfo into this registry and starts it
37 +func (r *ListenerRegistry) Register(l Listener) error {
38 + r.Lock()
39
40 if _, ok := r.Listeners[getListenerKey(l)]; ok {
34 - r.lk.Unlock()
41 + r.Unlock()
42 return errors.New("listener already registered")
43 }
37 - return nil
38 -}
44
40 -func (r *ListenerRegistry) unlock() {
41 - r.lk.Unlock()
42 -}
45 + r.Listeners[getListenerKey(l)] = l
46
44 -// Register registers listenerInfo in this registry
45 -func (r *ListenerRegistry) Register(l Listener) {
46 - defer r.lk.Unlock()
47 + r.Unlock()
48
48 - r.Listeners[getListenerKey(l)] = l
49 + if err := l.start(); err != nil {
50 + r.Lock()
51 + defer r.Lock()
52 +
53 + delete(r.Listeners, getListenerKey(l))
54 + return err
55 + }
56 +
57 + return nil
58 }
59
60 // Deregister removes p2p listener from this registry
52 -func (r *ListenerRegistry) Deregister(k listenerKey) {
53 - r.lk.Lock()
54 - defer r.lk.Unlock()
61 +func (r *ListenerRegistry) Deregister(k listenerKey) bool {
62 + r.Lock()
63 + defer r.Unlock()
64
65 + _, ok := r.Listeners[k]
66 delete(r.Listeners, k)
67 + return ok
68 }
69
70 func getListenerKey(l Listener) listenerKey {
71 return listenerKey{
61 - proto: l.Protocol(),
62 - listen: l.ListenAddress(),
63 - target: l.TargetAddress(),
72 + proto: string(l.Protocol()),
73 + listen: l.ListenAddress().String(),
74 + target: l.TargetAddress().String(),
75 }
76 }
p2p/local.go
+63 -54
@@ -2,14 +2,15 @@ package p2p
2
3 import (
4 "context"
5 + "errors"
6 "time"
7
7 - manet "gx/ipfs/QmNqRnejxJxjRroz7buhrjfU8i3yNBLa81hFtmf2pXEffN/go-multiaddr-net"
8 - ma "gx/ipfs/QmUxSEGbv2nmYNnfXi7839wwQqTN3kwQeUxe8dTjZWZs7J/go-multiaddr"
9 - peer "gx/ipfs/QmVf8hTAsLLFtn4WPCRNdnaF2Eag2qTBS6uR8AiHPZARXy/go-libp2p-peer"
10 - net "gx/ipfs/QmXdgNhVEgjLxjUoMs5ViQL7pboAt3Y7V7eGHRiE4qrmTE/go-libp2p-net"
11 - protocol "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
12 - pstore "gx/ipfs/QmZhsmorLpD9kmQ4ynbAu4vbKv2goMUnXazwGA4gnWHDjB/go-libp2p-peerstore"
8 + "gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
9 + ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
10 + "gx/ipfs/QmdVrMn1LhB4ybb8hMVaMLXnA8XRSewMnK6YqXKXoTcRvN/go-libp2p-peer"
11 + tec "gx/ipfs/QmWHgLqrghM9zw77nF6gdvT9ExQ2RB9pLxkd8sDHZf1rWb/go-temp-err-catcher"
12 + "gx/ipfs/QmPjvxTpVH8qJyQDnxnsxF9kv9jezKD1kozz1hs3fCGsNh/go-libp2p-net"
13 + "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
14 )
15
16 // localListener manet streams and proxies them to libp2p services
@@ -27,98 +28,106 @@ type localListener struct {
28 }
29
30 // ForwardLocal creates new P2P stream to a remote listener
30 -func (p2p *P2P) ForwardLocal(ctx context.Context, peer peer.ID, proto string, bindAddr ma.Multiaddr) (Listener, error) {
31 +func (p2p *P2P) ForwardLocal(ctx context.Context, peer peer.ID, proto protocol.ID, bindAddr ma.Multiaddr) (Listener, error) {
32 listener := &localListener{
33 ctx: ctx,
34
35 p2p: p2p,
36 id: p2p.identity,
37
37 - proto: protocol.ID(proto),
38 + proto: proto,
39 laddr: bindAddr,
40 peer: peer,
41 }
42
42 - if err := p2p.Listeners.lock(listener); err != nil {
43 + if err := p2p.Listeners.Register(listener); err != nil {
44 return nil, err
45 }
46
46 - maListener, err := manet.Listen(bindAddr)
47 - if err != nil {
48 - p2p.Listeners.unlock()
49 - return nil, err
50 - }
51 -
52 - listener.listener = maListener
53 -
54 - p2p.Listeners.Register(listener)
47 go listener.acceptConns()
48
49 return listener, nil
50 }
51
60 -func (l *localListener) dial() (net.Stream, error) {
61 - ctx, cancel := context.WithTimeout(l.ctx, time.Second*30) //TODO: configurable?
52 +func (l *localListener) dial(ctx context.Context) (net.Stream, error) {
53 + cctx, cancel := context.WithTimeout(ctx, time.Second*30) //TODO: configurable?
54 defer cancel()
55
64 - err := l.p2p.peerHost.Connect(ctx, pstore.PeerInfo{ID: l.peer})
65 - if err != nil {
66 - return nil, err
67 - }
68 -
69 - return l.p2p.peerHost.NewStream(l.ctx, l.peer, l.proto)
56 + return l.p2p.peerHost.NewStream(cctx, l.peer, l.proto)
57 }
58
59 func (l *localListener) acceptConns() {
60 for {
61 local, err := l.listener.Accept()
62 if err != nil {
63 + if tec.ErrIsTemporary(err) {
64 + continue
65 + }
66 return
67 }
68
79 - remote, err := l.dial()
80 - if err != nil {
81 - local.Close()
82 - return
83 - }
69 + go l.setupStream(local)
70 + }
71 +}
72
85 - tgt, err := ma.NewMultiaddr(l.TargetAddress())
86 - if err != nil {
87 - local.Close()
88 - return
89 - }
73 +func (l *localListener) setupStream(local manet.Conn) {
74 + remote, err := l.dial(l.ctx)
75 + if err != nil {
76 + local.Close()
77 + log.Warningf("failed to dial to remote %s/%s", l.peer.Pretty(), l.proto)
78 + return
79 + }
80
91 - stream := &Stream{
92 - Protocol: l.proto,
81 + stream := &Stream{
82 + Protocol: l.proto,
83
94 - OriginAddr: local.RemoteMultiaddr(),
95 - TargetAddr: tgt,
84 + OriginAddr: local.RemoteMultiaddr(),
85 + TargetAddr: l.TargetAddress(),
86
97 - Local: local,
98 - Remote: remote,
87 + Local: local,
88 + Remote: remote,
89
100 - Registry: l.p2p.Streams,
101 - }
90 + Registry: l.p2p.Streams,
91 + }
92 +
93 + l.p2p.Streams.Register(stream)
94 + stream.startStreaming()
95 +}
96
103 - l.p2p.Streams.Register(stream)
104 - stream.startStreaming()
97 +func (l *localListener) start() error {
98 + maListener, err := manet.Listen(l.laddr)
99 + if err != nil {
100 + return err
101 }
102 +
103 + l.listener = maListener
104 + return nil
105 }
106
107 func (l *localListener) Close() error {
109 - l.listener.Close()
110 - l.p2p.Listeners.Deregister(getListenerKey(l))
108 + if l.listener == nil {
109 + return errors.New("uninitialized")
110 + }
111 +
112 + if l.p2p.Listeners.Deregister(getListenerKey(l)) {
113 + l.listener.Close()
114 + l.listener = nil
115 + }
116 return nil
117 }
118
114 -func (l *localListener) Protocol() string {
115 - return string(l.proto)
119 +func (l *localListener) Protocol() protocol.ID {
120 + return l.proto
121 }
122
118 -func (l *localListener) ListenAddress() string {
119 - return l.laddr.String()
123 +func (l *localListener) ListenAddress() ma.Multiaddr {
124 + return l.laddr
125 }
126
122 -func (l *localListener) TargetAddress() string {
123 - return "/ipfs/" + l.peer.Pretty()
127 +func (l *localListener) TargetAddress() ma.Multiaddr {
128 + addr, err := ma.NewMultiaddr(maPrefix + l.peer.Pretty())
129 + if err != nil {
130 + panic(err)
131 + }
132 + return addr
133 }
p2p/p2p.go
+3 -4
@@ -1,13 +1,14 @@
1 package p2p
2
3 import (
4 - "sync"
5 -
4 + logging "gx/ipfs/QmcVVHfdyv15GVPk7NrxdWjh2hLVccXnoD8j2tyQShiXJb/go-log"
5 pstore "gx/ipfs/QmZR2XWVVBCtbgBWnQhWk2xcQfaR3W8faQPriAiaaj7rsr/go-libp2p-peerstore"
6 p2phost "gx/ipfs/Qmb8T6YBBsjYsVGfrihQLfCJveczZnneSBqBKkYEBWDjge/go-libp2p-host"
7 peer "gx/ipfs/QmdVrMn1LhB4ybb8hMVaMLXnA8XRSewMnK6YqXKXoTcRvN/go-libp2p-peer"
8 )
9
10 +var log = logging.Logger("p2p-mount")
11 +
12 // P2P structure holds information on currently running streams/listeners
13 type P2P struct {
14 Listeners *ListenerRegistry
@@ -27,11 +28,9 @@ func NewP2P(identity peer.ID, peerHost p2phost.Host, peerstore pstore.Peerstore)
28
29 Listeners: &ListenerRegistry{
30 Listeners: map[listenerKey]Listener{},
30 - lk: sync.Mutex{},
31 },
32 Streams: &StreamRegistry{
33 Streams: map[uint64]*Stream{},
34 - lk: sync.Mutex{},
34 },
35 }
36 }
p2p/remote.go
+43 -25
@@ -2,13 +2,16 @@ package p2p
2
3 import (
4 "context"
5 + "errors"
6
6 - manet "gx/ipfs/QmNqRnejxJxjRroz7buhrjfU8i3yNBLa81hFtmf2pXEffN/go-multiaddr-net"
7 - ma "gx/ipfs/QmUxSEGbv2nmYNnfXi7839wwQqTN3kwQeUxe8dTjZWZs7J/go-multiaddr"
8 - net "gx/ipfs/QmXdgNhVEgjLxjUoMs5ViQL7pboAt3Y7V7eGHRiE4qrmTE/go-libp2p-net"
7 + manet "gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
8 + ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
9 + net "gx/ipfs/QmPjvxTpVH8qJyQDnxnsxF9kv9jezKD1kozz1hs3fCGsNh/go-libp2p-net"
10 protocol "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
11 )
12
13 +var maPrefix = "/" + ma.ProtocolWithCode(ma.P_IPFS).Name + "/"
14 +
15 // remoteListener accepts libp2p streams and proxies them to a manet host
16 type remoteListener struct {
17 p2p *P2P
@@ -18,70 +21,85 @@ type remoteListener struct {
21
22 // Address to proxy the incoming connections to
23 addr ma.Multiaddr
24 +
25 + initialized bool
26 }
27
28 // ForwardRemote creates new p2p listener
24 -func (p2p *P2P) ForwardRemote(ctx context.Context, proto string, addr ma.Multiaddr) (Listener, error) {
29 +func (p2p *P2P) ForwardRemote(ctx context.Context, proto protocol.ID, addr ma.Multiaddr) (Listener, error) {
30 listener := &remoteListener{
31 p2p: p2p,
32
28 - proto: protocol.ID(proto),
33 + proto: proto,
34 addr: addr,
35 }
36
32 - if err := p2p.Listeners.lock(listener); err != nil {
37 + if err := p2p.Listeners.Register(listener); err != nil {
38 return nil, err
39 }
40
36 - p2p.peerHost.SetStreamHandler(listener.proto, func(remote net.Stream) {
37 - local, err := manet.Dial(addr)
41 + return listener, nil
42 +}
43 +
44 +func (l *remoteListener) start() error {
45 + // TODO: handle errors when https://github.com/libp2p/go-libp2p-host/issues/16 will be done
46 + l.p2p.peerHost.SetStreamHandler(l.proto, func(remote net.Stream) {
47 + local, err := manet.Dial(l.addr)
48 if err != nil {
49 remote.Reset()
50 return
51 }
52
43 - //TODO: review: is there a better way to do this?
44 - peerMa, err := ma.NewMultiaddr("/ipfs/" + remote.Conn().RemotePeer().Pretty())
53 + peerMa, err := ma.NewMultiaddr(maPrefix + remote.Conn().RemotePeer().Pretty())
54 if err != nil {
55 remote.Reset()
56 return
57 }
58
59 stream := &Stream{
51 - Protocol: listener.proto,
60 + Protocol: l.proto,
61
62 OriginAddr: peerMa,
54 - TargetAddr: addr,
63 + TargetAddr: l.addr,
64
65 Local: local,
66 Remote: remote,
67
59 - Registry: p2p.Streams,
68 + Registry: l.p2p.Streams,
69 }
70
62 - p2p.Streams.Register(stream)
71 + l.p2p.Streams.Register(stream)
72 stream.startStreaming()
73 })
74
66 - p2p.Listeners.Register(listener)
67 -
68 - return listener, nil
75 + l.initialized = true
76 + return nil
77 }
78
71 -func (l *remoteListener) Protocol() string {
72 - return string(l.proto)
79 +func (l *remoteListener) Protocol() protocol.ID {
80 + return l.proto
81 }
82
75 -func (l *remoteListener) ListenAddress() string {
76 - return "/ipfs"
83 +func (l *remoteListener) ListenAddress() ma.Multiaddr {
84 + addr, err := ma.NewMultiaddr(maPrefix + l.p2p.identity.Pretty())
85 + if err != nil {
86 + panic(err)
87 + }
88 + return addr
89 }
90
79 -func (l *remoteListener) TargetAddress() string {
80 - return l.addr.String()
91 +func (l *remoteListener) TargetAddress() ma.Multiaddr {
92 + return l.addr
93 }
94
95 func (l *remoteListener) Close() error {
84 - l.p2p.peerHost.RemoveStreamHandler(protocol.ID(l.proto))
85 - l.p2p.Listeners.Deregister(getListenerKey(l))
96 + if !l.initialized {
97 + return errors.New("uninitialized")
98 + }
99 +
100 + if l.p2p.Listeners.Deregister(getListenerKey(l)) {
101 + l.p2p.peerHost.RemoveStreamHandler(l.proto)
102 + l.initialized = false
103 + }
104 return nil
105 }
p2p/stream.go
+9 -5
@@ -4,9 +4,9 @@ import (
4 "io"
5 "sync"
6
7 - manet "gx/ipfs/QmNqRnejxJxjRroz7buhrjfU8i3yNBLa81hFtmf2pXEffN/go-multiaddr-net"
8 - ma "gx/ipfs/QmUxSEGbv2nmYNnfXi7839wwQqTN3kwQeUxe8dTjZWZs7J/go-multiaddr"
9 - net "gx/ipfs/QmXdgNhVEgjLxjUoMs5ViQL7pboAt3Y7V7eGHRiE4qrmTE/go-libp2p-net"
7 + manet "gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
8 + ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
9 + net "gx/ipfs/QmPjvxTpVH8qJyQDnxnsxF9kv9jezKD1kozz1hs3fCGsNh/go-libp2p-net"
10 "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
11 )
12
@@ -43,8 +43,12 @@ func (s *Stream) Reset() error {
43
44 func (s *Stream) startStreaming() {
45 go func() {
46 - io.Copy(s.Local, s.Remote)
47 - s.Reset()
46 + _, err := io.Copy(s.Local, s.Remote)
47 + if err != nil {
48 + s.Reset()
49 + } else {
50 + s.Close()
51 + }
52 }()
53
54 go func() {
package.json
+6
@@ -486,6 +486,12 @@
486 "name": "go-ipns",
487 "version": "0.1.8"
488 },
489 + {
490 + "author": "whyrusleeping",
491 + "hash": "QmWHgLqrghM9zw77nF6gdvT9ExQ2RB9pLxkd8sDHZf1rWb",
492 + "name": "go-temp-err-catcher",
493 + "version": "0.0.0"
494 + },
495 {
496 "author": "why",
497 "hash": "QmVDDgboX5nPUE4pBcK2xC1b9XbStA4t2KrUWBRMr9AiFd",
test/sharness/t0180-p2p.sh
+21 -3
@@ -127,7 +127,7 @@ check_test_ports
127 # Listing streams
128
129 test_expect_success "'ipfs p2p ls' succeeds" '
130 - echo "/x/p2p-test /ipfs /ip4/127.0.0.1/tcp/10101" > expected &&
130 + echo "/x/p2p-test /ipfs/$PEERID_0 /ip4/127.0.0.1/tcp/10101" > expected &&
131 ipfsi 0 p2p ls > actual
132 '
133
@@ -144,10 +144,12 @@ test_expect_success "'ipfs p2p stream ls' output is empty" '
144 test_must_be_empty actual
145 '
146
147 +check_test_ports
148 +
149 test_expect_success "Setup: Idle stream" '
150 ma-pipe-unidir --listen --pidFile=listener.pid recv /ip4/127.0.0.1/tcp/10101 &
151
150 - ipfsi 1 p2p forward /x/p2p-test /ip4/127.0.0.1/tcp/10102 /ipfs/$PEERID_0 2>&1 > dialer-stdouterr.log &&
152 + ipfsi 1 p2p forward /x/p2p-test /ip4/127.0.0.1/tcp/10102 /ipfs/$PEERID_0 &&
153 ma-pipe-unidir --pidFile=client.pid recv /ip4/127.0.0.1/tcp/10102 &
154
155 test_wait_for_file 30 100ms listener.pid &&
@@ -231,8 +233,22 @@ test_expect_success "'ipfs p2p close' closes app numeric handlers" '
233 test_must_be_empty actual
234 '
235
236 +test_expect_success "'ipfs p2p close' closes by listen addr" '
237 + ipfsi 0 p2p listen /x/p2p-test /ip4/127.0.0.1/tcp/10101 &&
238 + ipfsi 0 p2p close -l /ipfs/$PEERID_0 &&
239 + ipfsi 0 p2p ls > actual &&
240 + test_must_be_empty actual
241 +'
242 +
243 +test_expect_success "'ipfs p2p close' closes by target addr" '
244 + ipfsi 0 p2p listen /x/p2p-test /ip4/127.0.0.1/tcp/10101 &&
245 + ipfsi 0 p2p close -t /ip4/127.0.0.1/tcp/10101 &&
246 + ipfsi 0 p2p ls > actual &&
247 + test_must_be_empty actual
248 +'
249 +
250 test_expect_success "non /x/ scoped protocols are not allowed" '
235 - test_must_fail ipfsi 0 p2p forward /its/not/a/x/path /ipfs /ip4/127.0.0.1/tcp/10101 2> actual &&
251 + test_must_fail ipfsi 0 p2p listen /its/not/a/x/path /ip4/127.0.0.1/tcp/10101 2> actual &&
252 echo "Error: protocol name must be within '"'"'/x/'"'"' namespace" > expected
253 test_cmp expected actual
254 '
@@ -246,6 +262,8 @@ test_expect_success 'start p2p listener on custom proto' '
262
263 test_expect_success 'C->S Close local listener' '
264 ipfsi 0 p2p close -p /p2p-test
265 + ipfsi 0 p2p ls > actual &&
266 + test_must_be_empty actual
267 '
268
269 test_expect_success 'stop iptb' '