@cryptotaxi247 / kubo / commits / 043aa5f30

PTP API: Address review comments

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

Łukasz Magiera committed Jun 7, 2017 at 19:21 UTC 043aa5f308148a5c79226d7f7d03246e56a34d9e
5 files changed +152 -134
core/commands/ptp.go
+102 -99
@@ -51,15 +51,40 @@ Note: this command is experimental and subject to change as usecases and APIs ar
51 },
52
53 Subcommands: map[string]*cmds.Command{
54 - "ls": ptpLsCmd,
55 - "streams": ptpStreamsCmd,
56 - "dial": ptpDialCmd,
57 - "listen": ptpListenCmd,
58 - "close": ptpCloseCmd,
54 + "listener": ptpListenerCmd,
55 + "stream": ptpStreamCmd,
56 },
57 }
58
62 -var ptpLsCmd = &cmds.Command{
59 +// ptpListenerCmd is the 'ipfs ptp listener' command
60 +var ptpListenerCmd = &cmds.Command{
61 + Helptext: cmds.HelpText{
62 + Tagline: "P2P listener management.",
63 + ShortDescription: "Create and manage listener p2p endpoints",
64 + },
65 +
66 + Subcommands: map[string]*cmds.Command{
67 + "ls": ptpListenerLsCmd,
68 + "open": ptpListenerListenCmd,
69 + "close": ptpListenerCloseCmd,
70 + },
71 +}
72 +
73 +// ptpStreamCmd is the 'ipfs ptp stream' command
74 +var ptpStreamCmd = &cmds.Command{
75 + Helptext: cmds.HelpText{
76 + Tagline: "P2P stream management.",
77 + ShortDescription: "Create and manage p2p streams",
78 + },
79 +
80 + Subcommands: map[string]*cmds.Command{
81 + "ls": ptpStreamLsCmd,
82 + "dial": ptpStreamDialCmd,
83 + "close": ptpStreamCloseCmd,
84 + },
85 +}
86 +
87 +var ptpListenerLsCmd = &cmds.Command{
88 Helptext: cmds.HelpText{
89 Tagline: "List active p2p listeners.",
90 },
@@ -67,23 +92,13 @@ var ptpLsCmd = &cmds.Command{
92 cmds.BoolOption("headers", "v", "Print table headers (HandlerID, Protocol, Local, Remote).").Default(false),
93 },
94 Run: func(req cmds.Request, res cmds.Response) {
70 - n, err := req.InvocContext().GetNode()
71 - if err != nil {
72 - res.SetError(err, cmds.ErrNormal)
73 - return
74 - }
95
76 - err = checkEnabled(n)
96 + n, err := getNode(req)
97 if err != nil {
98 res.SetError(err, cmds.ErrNormal)
99 return
100 }
101
82 - if !n.OnlineMode() {
83 - res.SetError(errNotOnline, cmds.ErrClient)
84 - return
85 - }
86 -
102 output := &PTPLsOutput{}
103
104 for _, listener := range n.PTP.Listeners.Listeners {
@@ -116,7 +131,7 @@ var ptpLsCmd = &cmds.Command{
131 },
132 }
133
119 -var ptpStreamsCmd = &cmds.Command{
134 +var ptpStreamLsCmd = &cmds.Command{
135 Helptext: cmds.HelpText{
136 Tagline: "List active p2p streams.",
137 },
@@ -124,23 +139,12 @@ var ptpStreamsCmd = &cmds.Command{
139 cmds.BoolOption("headers", "v", "Print table headers (HagndlerID, Protocol, Local, Remote).").Default(false),
140 },
141 Run: func(req cmds.Request, res cmds.Response) {
127 - n, err := req.InvocContext().GetNode()
142 + n, err := getNode(req)
143 if err != nil {
144 res.SetError(err, cmds.ErrNormal)
145 return
146 }
147
133 - err = checkEnabled(n)
134 - if err != nil {
135 - res.SetError(err, cmds.ErrNormal)
136 - return
137 - }
138 -
139 - if !n.OnlineMode() {
140 - res.SetError(errNotOnline, cmds.ErrClient)
141 - return
142 - }
143 -
148 output := &PTPStreamsOutput{}
149
150 for _, s := range n.PTP.Streams.Streams {
@@ -180,7 +184,7 @@ var ptpStreamsCmd = &cmds.Command{
184 },
185 }
186
183 -var ptpListenCmd = &cmds.Command{
187 +var ptpListenerListenCmd = &cmds.Command{
188 Helptext: cmds.HelpText{
189 Tagline: "Forward p2p connections to a network multiaddr.",
190 ShortDescription: `
@@ -194,23 +198,12 @@ Note that the connections originate from the ipfs daemon process.
198 cmds.StringArg("Address", true, false, "Request handling application address."),
199 },
200 Run: func(req cmds.Request, res cmds.Response) {
197 - n, err := req.InvocContext().GetNode()
198 - if err != nil {
199 - res.SetError(err, cmds.ErrNormal)
200 - return
201 - }
202 -
203 - err = checkEnabled(n)
201 + n, err := getNode(req)
202 if err != nil {
203 res.SetError(err, cmds.ErrNormal)
204 return
205 }
206
209 - if !n.OnlineMode() {
210 - res.SetError(errNotOnline, cmds.ErrClient)
211 - return
212 - }
213 -
207 proto := "/ptp/" + req.Arguments()[0]
208 if n.PTP.CheckProtoExists(proto) {
209 res.SetError(errors.New("protocol handler already registered"), cmds.ErrNormal)
@@ -237,7 +230,7 @@ Note that the connections originate from the ipfs daemon process.
230 },
231 }
232
240 -var ptpDialCmd = &cmds.Command{
233 +var ptpStreamDialCmd = &cmds.Command{
234 Helptext: cmds.HelpText{
235 Tagline: "Dial to a p2p listener.",
236
@@ -255,23 +248,12 @@ transparently connect to a p2p service.
248 cmds.StringArg("BindAddress", false, false, "Address to listen for connection/s (default: /ip4/127.0.0.1/tcp/0)."),
249 },
250 Run: func(req cmds.Request, res cmds.Response) {
258 - n, err := req.InvocContext().GetNode()
251 + n, err := getNode(req)
252 if err != nil {
253 res.SetError(err, cmds.ErrNormal)
254 return
255 }
256
264 - err = checkEnabled(n)
265 - if err != nil {
266 - res.SetError(err, cmds.ErrNormal)
267 - return
268 - }
269 -
270 - if !n.OnlineMode() {
271 - res.SetError(errNotOnline, cmds.ErrClient)
272 - return
273 - }
274 -
257 addr, peer, err := ParsePeerParam(req.Arguments()[0])
258 if err != nil {
259 res.SetError(err, cmds.ErrNormal)
@@ -304,89 +286,110 @@ transparently connect to a p2p service.
286 },
287 }
288
307 -var ptpCloseCmd = &cmds.Command{
289 +var ptpListenerCloseCmd = &cmds.Command{
290 Helptext: cmds.HelpText{
309 - Tagline: "Closes an active p2p stream or listener.",
291 + Tagline: "Close active p2p listener.",
292 },
293 Arguments: []cmds.Argument{
312 - cmds.StringArg("Identifier", false, false, "Stream HandlerID or p2p listener protocol"),
294 + cmds.StringArg("Protocol", false, false, "P2P listener protocol"),
295 },
296 Options: []cmds.Option{
315 - cmds.BoolOption("all", "a", "Close all streams and listeners.").Default(false),
297 + cmds.BoolOption("all", "a", "Close all listeners.").Default(false),
298 },
299 Run: func(req cmds.Request, res cmds.Response) {
318 - n, err := req.InvocContext().GetNode()
300 + n, err := getNode(req)
301 if err != nil {
302 res.SetError(err, cmds.ErrNormal)
303 return
304 }
305
324 - err = checkEnabled(n)
325 - if err != nil {
326 - res.SetError(err, cmds.ErrNormal)
327 - return
306 + closeAll, _, _ := req.Option("all").Bool()
307 + var proto string
308 +
309 + if !closeAll {
310 + if len(req.Arguments()) == 0 {
311 + res.SetError(errors.New("no protocol name specified"), cmds.ErrNormal)
312 + return
313 + }
314 +
315 + proto = "/ptp/" + req.Arguments()[0]
316 }
317
330 - if !n.OnlineMode() {
331 - res.SetError(errNotOnline, cmds.ErrClient)
318 + for _, listener := range n.PTP.Listeners.Listeners {
319 + if !closeAll && listener.Protocol != proto {
320 + continue
321 + }
322 + listener.Close()
323 + if !closeAll {
324 + break
325 + }
326 + }
327 + },
328 +}
329 +
330 +var ptpStreamCloseCmd = &cmds.Command{
331 + Helptext: cmds.HelpText{
332 + Tagline: "Close active p2p stream.",
333 + },
334 + Arguments: []cmds.Argument{
335 + cmds.StringArg("HandlerID", false, false, "Stream HandlerID"),
336 + },
337 + Options: []cmds.Option{
338 + cmds.BoolOption("all", "a", "Close all streams.").Default(false),
339 + },
340 + Run: func(req cmds.Request, res cmds.Response) {
341 + n, err := getNode(req)
342 + if err != nil {
343 + res.SetError(err, cmds.ErrNormal)
344 return
345 }
346
347 closeAll, _, _ := req.Option("all").Bool()
336 -
337 - var proto string
348 var handlerID uint64
349
340 - useHandlerID := false
341 -
350 if !closeAll {
351 if len(req.Arguments()) == 0 {
344 - res.SetError(errors.New("no handlerID nor listener protocol specified"), cmds.ErrNormal)
352 + res.SetError(errors.New("no HandlerID specified"), cmds.ErrNormal)
353 return
354 }
355
356 handlerID, err = strconv.ParseUint(req.Arguments()[0], 10, 64)
357 if err != nil {
350 - proto = "/ptp/" + req.Arguments()[0]
351 - } else {
352 - useHandlerID = true
358 + res.SetError(err, cmds.ErrNormal)
359 + return
360 }
361 }
362
356 - if closeAll || useHandlerID {
357 - for _, stream := range n.PTP.Streams.Streams {
358 - if !closeAll && handlerID != stream.HandlerID {
359 - continue
360 - }
361 - stream.Close()
362 - if !closeAll {
363 - break
364 - }
363 + for _, stream := range n.PTP.Streams.Streams {
364 + if !closeAll && handlerID != stream.HandlerID {
365 + continue
366 }
366 - }
367 -
368 - if closeAll || !useHandlerID {
369 - for _, listener := range n.PTP.Listeners.Listeners {
370 - if !closeAll && listener.Protocol != proto {
371 - continue
372 - }
373 - listener.Close()
374 - if !closeAll {
375 - break
376 - }
367 + stream.Close()
368 + if !closeAll {
369 + break
370 }
371 }
372 },
373 }
374
382 -func checkEnabled(n *core.IpfsNode) error {
375 +func getNode(req cmds.Request) (*core.IpfsNode, error) {
376 + n, err := req.InvocContext().GetNode()
377 + if err != nil {
378 + return nil, err
379 + }
380 +
381 config, err := n.Repo.Config()
382 if err != nil {
385 - return err
383 + return nil, err
384 }
385
386 if !config.Experimental.Libp2pStreamMounting {
389 - return errors.New("libp2p stream mounting not enabled")
387 + return nil, errors.New("libp2p stream mounting not enabled")
388 }
391 - return nil
389 +
390 + if !n.OnlineMode() {
391 + return nil, errNotOnline
392 + }
393 +
394 + return n, nil
395 }
ptp/ptp.go
+1
@@ -43,6 +43,7 @@ func (ptp *PTP) newStreamTo(ctx2 context.Context, p peer.ID, protocol string) (n
43 return ptp.peerHost.NewStream(ctx2, p, pro.ID(protocol))
44 }
45
46 +// Dial creates new P2P stream to a remote listener
47 func (ptp *PTP) Dial(ctx context.Context, addr ma.Multiaddr, peer peer.ID, proto string, bindAddr ma.Multiaddr) (*ListenerInfo, error) {
48 lnet, _, err := manet.DialArgs(bindAddr)
49 if err != nil {
ptp/registry.go
+4 -4
@@ -83,10 +83,10 @@ type StreamInfo struct {
83 }
84
85 // Close closes stream endpoints and deregisters it
86 -func (c *StreamInfo) Close() error {
87 - c.Local.Close()
88 - c.Remote.Close()
89 - c.Registry.Deregister(c.HandlerID)
86 +func (s *StreamInfo) Close() error {
87 + s.Local.Close()
88 + s.Remote.Close()
89 + s.Registry.Deregister(s.HandlerID)
90 return nil
91 }
92
test/dependencies/ma-pipe-unidir/main.go
+5 -2
@@ -38,6 +38,11 @@ func app() int {
38 mode := args[0]
39 addr := args[1]
40
41 + if mode != "send" && mode != "recv" {
42 + fmt.Print(USAGE)
43 + return 1
44 + }
45 +
46 if len(opts.PidFile) > 0 {
47 data := []byte(strconv.Itoa(os.Getpid()))
48 err := ioutil.WriteFile(opts.PidFile, data, 0644)
@@ -80,8 +85,6 @@ func app() int {
85 case "send":
86 io.Copy(conn, os.Stdin)
87 default:
83 - //TODO: a bit late
84 - fmt.Print(USAGE)
88 return 1
89 }
90 return 0
test/sharness/t0180-ptp.sh
+40 -29
@@ -27,7 +27,7 @@ test_expect_success "test ports are closed" '
27 '
28
29 test_must_fail 'fail without config option being enabled' '
30 - ipfsi 0 ptp ls
30 + ipfsi 0 ptp stream ls
31 '
32
33 test_expect_success "enable filestore config setting" '
@@ -36,14 +36,14 @@ test_expect_success "enable filestore config setting" '
36 '
37
38 test_expect_success 'start ptp listener' '
39 - ipfsi 0 ptp listen ptp-test /ip4/127.0.0.1/tcp/10101 2>&1 > listener-stdouterr.log
39 + ipfsi 0 ptp listener open ptp-test /ip4/127.0.0.1/tcp/10101 2>&1 > listener-stdouterr.log
40 '
41
42 test_expect_success 'Test server to client communications' '
43 ma-pipe-unidir --listen send /ip4/127.0.0.1/tcp/10101 < test0.bin &
44 SERVER_PID=$!
45
46 - ipfsi 1 ptp dial $PEERID_0 ptp-test /ip4/127.0.0.1/tcp/10102 2>&1 > dialer-stdouterr.log &&
46 + ipfsi 1 ptp stream dial $PEERID_0 ptp-test /ip4/127.0.0.1/tcp/10102 2>&1 > dialer-stdouterr.log &&
47 ma-pipe-unidir recv /ip4/127.0.0.1/tcp/10102 > client.out &&
48 wait $SERVER_PID
49 '
@@ -52,7 +52,7 @@ test_expect_success 'Test client to server communications' '
52 ma-pipe-unidir --listen recv /ip4/127.0.0.1/tcp/10101 > server.out &
53 SERVER_PID=$!
54
55 - ipfsi 1 ptp dial $PEERID_0 ptp-test /ip4/127.0.0.1/tcp/10102 2>&1 > dialer-stdouterr.log &&
55 + ipfsi 1 ptp stream dial $PEERID_0 ptp-test /ip4/127.0.0.1/tcp/10102 2>&1 > dialer-stdouterr.log &&
56 ma-pipe-unidir send /ip4/127.0.0.1/tcp/10102 < test1.bin
57 wait $SERVER_PID
58 '
@@ -65,79 +65,90 @@ test_expect_success 'client to server output looks good' '
65 test_cmp server.out test1.bin
66 '
67
68 -test_expect_success "'ipfs ptp ls' succeeds" '
68 +test_expect_success "'ipfs listener ptp ls' succeeds" '
69 echo "/ip4/127.0.0.1/tcp/10101 /ptp/ptp-test" > expected &&
70 - ipfsi 0 ptp ls > actual
70 + ipfsi 0 ptp listener ls > actual
71 '
72
73 -test_expect_success "'ipfs ptp ls' output looks good" '
73 +test_expect_success "'ipfs ptp listener ls' output looks good" '
74 test_cmp expected actual
75 '
76
77 test_expect_success "Cannot re-register app handler" '
78 - (! ipfsi 0 ptp listen ptp-test /ip4/127.0.0.1/tcp/10101)
78 + (! ipfsi 0 ptp listener open ptp-test /ip4/127.0.0.1/tcp/10101)
79 '
80
81 -test_expect_success "'ipfs ptp streams' output is empty" '
82 - ipfsi 0 ptp streams > actual &&
81 +test_expect_success "'ipfs ptp stream ls' output is empty" '
82 + ipfsi 0 ptp stream ls > actual &&
83 test_must_be_empty actual
84 '
85
86 test_expect_success "Setup: Idle stream" '
87 ma-pipe-unidir --listen --pidFile=listener.pid recv /ip4/127.0.0.1/tcp/10101 &
88
89 - ipfsi 1 ptp dial $PEERID_0 ptp-test /ip4/127.0.0.1/tcp/10102 2>&1 > dialer-stdouterr.log &&
89 + ipfsi 1 ptp stream dial $PEERID_0 ptp-test /ip4/127.0.0.1/tcp/10102 2>&1 > dialer-stdouterr.log &&
90 ma-pipe-unidir --pidFile=client.pid recv /ip4/127.0.0.1/tcp/10102 &
91
92 go-sleep 500ms &&
93 kill -0 $(cat listener.pid) && kill -0 $(cat client.pid)
94 '
95
96 -test_expect_success "'ipfs ptp streams' succeeds" '
96 +test_expect_success "'ipfs ptp stream ls' succeeds" '
97 echo "2 /ptp/ptp-test /ip4/127.0.0.1/tcp/10101 $PEERID_1" > expected
98 - ipfsi 0 ptp streams > actual
98 + ipfsi 0 ptp stream ls > actual
99 '
100
101 -test_expect_success "'ipfs ptp streams' output looks good" '
101 +test_expect_success "'ipfs ptp stream ls' output looks good" '
102 test_cmp expected actual
103 '
104
105 -test_expect_success "'ipfs ptp close' closes stream" '
106 - ipfsi 0 ptp close 2 &&
107 - ipfsi 0 ptp streams > actual &&
105 +test_expect_success "'ipfs ptp stream close' closes stream" '
106 + ipfsi 0 ptp stream close 2 &&
107 + ipfsi 0 ptp stream ls > actual &&
108 [ ! -f listener.pid ] && [ ! -f client.pid ] &&
109 test_must_be_empty actual
110 '
111
112 -test_expect_success "'ipfs ptp close' closes app handler" '
113 - ipfsi 0 ptp close ptp-test &&
114 - ipfsi 0 ptp ls > actual &&
112 +test_expect_success "'ipfs ptp listener close' closes app handler" '
113 + ipfsi 0 ptp listener close ptp-test &&
114 + ipfsi 0 ptp listener ls > actual &&
115 test_must_be_empty actual
116 '
117
118 test_expect_success "Setup: Idle stream(2)" '
119 ma-pipe-unidir --listen --pidFile=listener.pid recv /ip4/127.0.0.1/tcp/10101 &
120
121 - ipfsi 0 ptp listen ptp-test2 /ip4/127.0.0.1/tcp/10101 2>&1 > listener-stdouterr.log &&
122 - ipfsi 1 ptp dial $PEERID_0 ptp-test2 /ip4/127.0.0.1/tcp/10102 2>&1 > dialer-stdouterr.log &&
121 + ipfsi 0 ptp listener open ptp-test2 /ip4/127.0.0.1/tcp/10101 2>&1 > listener-stdouterr.log &&
122 + ipfsi 1 ptp stream dial $PEERID_0 ptp-test2 /ip4/127.0.0.1/tcp/10102 2>&1 > dialer-stdouterr.log &&
123 ma-pipe-unidir --pidFile=client.pid recv /ip4/127.0.0.1/tcp/10102 &
124
125 go-sleep 500ms &&
126 kill -0 $(cat listener.pid) && kill -0 $(cat client.pid)
127 '
128
129 -test_expect_success "'ipfs ptp streams' succeeds(2)" '
129 +test_expect_success "'ipfs ptp stream ls' succeeds(2)" '
130 echo "3 /ptp/ptp-test2 /ip4/127.0.0.1/tcp/10101 $PEERID_1" > expected
131 - ipfsi 0 ptp streams > actual
131 + ipfsi 0 ptp stream ls > actual
132 test_cmp expected actual
133 '
134
135 -test_expect_success "'ipfs ptp close -a' closes streams and app handlers" '
136 - ipfsi 0 ptp close -a &&
137 - ipfsi 0 ptp streams > actual &&
135 +test_expect_success "'ipfs ptp listener close -a' closes app handlers" '
136 + ipfsi 0 ptp listener close -a &&
137 + ipfsi 0 ptp listener ls > actual &&
138 + test_must_be_empty actual
139 +'
140 +
141 +test_expect_success "'ipfs ptp stream close -a' closes streams" '
142 + ipfsi 0 ptp stream close -a &&
143 + ipfsi 0 ptp stream ls > actual &&
144 [ ! -f listener.pid ] && [ ! -f client.pid ] &&
139 - test_must_be_empty actual &&
140 - ipfsi 0 ptp ls > actual &&
145 + test_must_be_empty actual
146 +'
147 +
148 +test_expect_success "'ipfs ptp listener close' closes app numeric handlers" '
149 + ipfsi 0 ptp listener open 1234 /ip4/127.0.0.1/tcp/10101 &&
150 + ipfsi 0 ptp listener close 1234 &&
151 + ipfsi 0 ptp listener ls > actual &&
152 test_must_be_empty actual
153 '
154