@cryptotaxi247 / kubo / commits / 7292c1da2

commands/p2p: use new cmds lib

License: MIT Signed-off-by: Overbool <overbool.xu@gmail.com>

Overbool committed Oct 27, 2018 at 11:39 UTC 7292c1da2335c97e156b819f06e2800196cb9910
2 files changed +89 -144
core/commands/p2p.go
+88 -143
@@ -1,7 +1,6 @@
1 package commands
2
3 import (
4 - "bytes"
4 "context"
5 "errors"
6 "fmt"
@@ -11,15 +10,16 @@ import (
10 "text/tabwriter"
11 "time"
12
14 - cmds "github.com/ipfs/go-ipfs/commands"
13 core "github.com/ipfs/go-ipfs/core"
14 + cmdenv "github.com/ipfs/go-ipfs/core/commands/cmdenv"
15 p2p "github.com/ipfs/go-ipfs/p2p"
16
17 + cmds "gx/ipfs/QmSXUokcP4TJpFfqozT69AVAYRtzXVMUjzQVkYX41R9Svs/go-ipfs-cmds"
18 ma "gx/ipfs/QmT4U94DnD8FRfqr21obWY32HLM5VExccPKMjQHofeYqr9/go-multiaddr"
19 pstore "gx/ipfs/QmTTJcDL3gsnGDALjh2fDGg1onGRUdVgNL2hU2WEZcVrMX/go-libp2p-peerstore"
20 - "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
21 - "gx/ipfs/QmZc5PLgxW61uTPG24TroxHDF6xzgbhZZQf5i53ciQC47Y/go-ipfs-addr"
22 - "gx/ipfs/Qmde5VP1qUkyQXKCfmEUA7bP64V2HAptbJ7phuPp7jXWwg/go-ipfs-cmdkit"
20 + protocol "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
21 + ipfsaddr "gx/ipfs/QmZc5PLgxW61uTPG24TroxHDF6xzgbhZZQf5i53ciQC47Y/go-ipfs-addr"
22 + cmdkit "gx/ipfs/Qmde5VP1qUkyQXKCfmEUA7bP64V2HAptbJ7phuPp7jXWwg/go-ipfs-cmdkit"
23 madns "gx/ipfs/QmeHJXPqCNzSFbVkYM1uQLuM2L5FyJB9zukQ7EeqRP8ZC9/go-multiaddr-dns"
24 )
25
@@ -69,8 +69,7 @@ are refined`,
69 },
70
71 Subcommands: map[string]*cmds.Command{
72 - "stream": p2pStreamCmd,
73 -
72 + "stream": p2pStreamCmd,
73 "forward": p2pForwardCmd,
74 "listen": p2pListenCmd,
75 "close": p2pCloseCmd,
@@ -101,47 +100,35 @@ Example:
100 Options: []cmdkit.Option{
101 cmdkit.BoolOption(allowCustomProtocolOptionName, "Don't require /x/ prefix"),
102 },
104 - Run: func(req cmds.Request, res cmds.Response) {
105 - n, err := p2pGetNode(req)
103 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
104 + n, err := p2pGetNode(env)
105 if err != nil {
107 - res.SetError(err, cmdkit.ErrNormal)
108 - return
106 + return err
107 }
108
111 - protoOpt := req.Arguments()[0]
112 - listenOpt := req.Arguments()[1]
113 - targetOpt := req.Arguments()[2]
109 + protoOpt := req.Arguments[0]
110 + listenOpt := req.Arguments[1]
111 + targetOpt := req.Arguments[2]
112
113 proto := protocol.ID(protoOpt)
114
115 listen, err := ma.NewMultiaddr(listenOpt)
116 if err != nil {
119 - res.SetError(err, cmdkit.ErrNormal)
120 - return
117 + return err
118 }
119
120 targets, err := parseIpfsAddr(targetOpt)
121 if err != nil {
125 - res.SetError(err, cmdkit.ErrNormal)
126 - return
122 + return err
123 }
124
129 - allowCustom, _, err := req.Option(allowCustomProtocolOptionName).Bool()
130 - if err != nil {
131 - res.SetError(err, cmdkit.ErrNormal)
132 - return
133 - }
125 + allowCustom, _ := req.Options[allowCustomProtocolOptionName].(bool)
126
127 if !allowCustom && !strings.HasPrefix(string(proto), P2PProtoPrefix) {
136 - res.SetError(errors.New("protocol name must be within '"+P2PProtoPrefix+"' namespace"), cmdkit.ErrNormal)
137 - return
128 + return errors.New("protocol name must be within '" + P2PProtoPrefix + "' namespace")
129 }
130
140 - if err := forwardLocal(n.Context(), n.P2P, n.Peerstore, proto, listen, targets); err != nil {
141 - res.SetError(err, cmdkit.ErrNormal)
142 - return
143 - }
144 - res.SetOutput(nil)
131 + return forwardLocal(n.Context(), n.P2P, n.Peerstore, proto, listen, targets)
132 },
133 }
134
@@ -197,47 +184,37 @@ Example:
184 Options: []cmdkit.Option{
185 cmdkit.BoolOption(allowCustomProtocolOptionName, "Don't require /x/ prefix"),
186 },
200 - Run: func(req cmds.Request, res cmds.Response) {
201 - n, err := p2pGetNode(req)
187 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
188 + n, err := p2pGetNode(env)
189 if err != nil {
203 - res.SetError(err, cmdkit.ErrNormal)
204 - return
190 + return err
191 }
192
207 - protoOpt := req.Arguments()[0]
208 - targetOpt := req.Arguments()[1]
193 + protoOpt := req.Arguments[0]
194 + targetOpt := req.Arguments[1]
195
196 proto := protocol.ID(protoOpt)
197
198 target, err := ma.NewMultiaddr(targetOpt)
199 if err != nil {
214 - res.SetError(err, cmdkit.ErrNormal)
215 - return
200 + return err
201 }
202
203 // port can't be 0
204 if err := checkPort(target); err != nil {
220 - res.SetError(err, cmdkit.ErrNormal)
221 - return
205 + return err
206 }
207
224 - allowCustom, _, err := req.Option(allowCustomProtocolOptionName).Bool()
208 + allowCustom, _ := req.Options[allowCustomProtocolOptionName].(bool)
209 if err != nil {
226 - res.SetError(err, cmdkit.ErrNormal)
227 - return
210 + return err
211 }
212
213 if !allowCustom && !strings.HasPrefix(string(proto), P2PProtoPrefix) {
231 - res.SetError(errors.New("protocol name must be within '"+P2PProtoPrefix+"' namespace"), cmdkit.ErrNormal)
232 - return
214 + return errors.New("protocol name must be within '" + P2PProtoPrefix + "' namespace")
215 }
216
235 - if err := forwardRemote(n.Context(), n.P2P, proto, target); err != nil {
236 - res.SetError(err, cmdkit.ErrNormal)
237 - return
238 - }
239 -
240 - res.SetOutput(nil)
217 + return forwardRemote(n.Context(), n.P2P, proto, target)
218 },
219 }
220
@@ -305,11 +282,10 @@ var p2pLsCmd = &cmds.Command{
282 Options: []cmdkit.Option{
283 cmdkit.BoolOption(p2pHeadersOptionName, "v", "Print table headers (Protocol, Listen, Target)."),
284 },
308 - Run: func(req cmds.Request, res cmds.Response) {
309 - n, err := p2pGetNode(req)
285 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
286 + n, err := p2pGetNode(env)
287 if err != nil {
311 - res.SetError(err, cmdkit.ErrNormal)
312 - return
288 + return err
289 }
290
291 output := &P2PLsOutput{}
@@ -334,31 +310,24 @@ var p2pLsCmd = &cmds.Command{
310 }
311 n.P2P.ListenersP2P.Unlock()
312
337 - res.SetOutput(output)
313 + return res.Emit(output)
314 },
315 Type: P2PLsOutput{},
340 - Marshalers: cmds.MarshalerMap{
341 - cmds.Text: func(res cmds.Response) (io.Reader, error) {
342 - v, err := unwrapOutput(res.Output())
343 - if err != nil {
344 - return nil, err
345 - }
346 -
347 - headers, _, _ := res.Request().Option(p2pHeadersOptionName).Bool()
348 - list := v.(*P2PLsOutput)
349 - buf := new(bytes.Buffer)
350 - w := tabwriter.NewWriter(buf, 1, 2, 1, ' ', 0)
351 - for _, listener := range list.Listeners {
316 + Encoders: cmds.EncoderMap{
317 + cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *P2PLsOutput) error {
318 + headers, _ := req.Options[p2pHeadersOptionName].(bool)
319 + tw := tabwriter.NewWriter(w, 1, 2, 1, ' ', 0)
320 + for _, listener := range out.Listeners {
321 if headers {
353 - fmt.Fprintln(w, "Protocol\tListen Address\tTarget Address")
322 + fmt.Fprintln(tw, "Protocol\tListen Address\tTarget Address")
323 }
324
356 - fmt.Fprintf(w, "%s\t%s\t%s\n", listener.Protocol, listener.ListenAddress, listener.TargetAddress)
325 + fmt.Fprintf(tw, "%s\t%s\t%s\n", listener.Protocol, listener.ListenAddress, listener.TargetAddress)
326 }
358 - w.Flush()
327 + tw.Flush()
328
360 - return buf, nil
361 - },
329 + return nil
330 + }),
331 },
332 }
333
@@ -379,40 +348,35 @@ var p2pCloseCmd = &cmds.Command{
348 cmdkit.StringOption(p2pListenAddressOptionName, "l", "Match listen address"),
349 cmdkit.StringOption(p2pTargetAddressOptionName, "t", "Match target address"),
350 },
382 - Run: func(req cmds.Request, res cmds.Response) {
383 - n, err := p2pGetNode(req)
351 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
352 + n, err := p2pGetNode(env)
353 if err != nil {
385 - res.SetError(err, cmdkit.ErrNormal)
386 - return
354 + return err
355 }
356
389 - closeAll, _, _ := req.Option(p2pAllOptionName).Bool()
390 - protoOpt, p, _ := req.Option(p2pProtocolOptionName).String()
391 - listenOpt, l, _ := req.Option(p2pListenAddressOptionName).String()
392 - targetOpt, t, _ := req.Option(p2pTargetAddressOptionName).String()
357 + closeAll, _ := req.Options[p2pAllOptionName].(bool)
358 + protoOpt, p := req.Options[p2pProtocolOptionName].(string)
359 + listenOpt, l := req.Options[p2pListenAddressOptionName].(string)
360 + targetOpt, t := req.Options[p2pTargetAddressOptionName].(string)
361
362 proto := protocol.ID(protoOpt)
363
364 listen, err := ma.NewMultiaddr(listenOpt)
365 if err != nil {
398 - res.SetError(err, cmdkit.ErrNormal)
399 - return
366 + return err
367 }
368
369 target, err := ma.NewMultiaddr(targetOpt)
370 if err != nil {
404 - res.SetError(err, cmdkit.ErrNormal)
405 - return
371 + return err
372 }
373
374 if !(closeAll || p || l || t) {
409 - res.SetError(errors.New("no matching options given"), cmdkit.ErrNormal)
410 - return
375 + return errors.New("no matching options given")
376 }
377
378 if closeAll && (p || l || t) {
414 - res.SetError(errors.New("can't combine --all with other matching options"), cmdkit.ErrNormal)
415 - return
379 + return errors.New("can't combine --all with other matching options")
380 }
381
382 match := func(listener p2p.Listener) bool {
@@ -434,22 +398,14 @@ var p2pCloseCmd = &cmds.Command{
398 done := n.P2P.ListenersLocal.Close(match)
399 done += n.P2P.ListenersP2P.Close(match)
400
437 - res.SetOutput(done)
401 + return res.Emit(done)
402 },
403 Type: int(0),
440 - Marshalers: cmds.MarshalerMap{
441 - cmds.Text: func(res cmds.Response) (io.Reader, error) {
442 - v, err := unwrapOutput(res.Output())
443 - if err != nil {
444 - return nil, err
445 - }
446 -
447 - closed := v.(int)
448 - buf := new(bytes.Buffer)
449 - fmt.Fprintf(buf, "Closed %d stream(s)\n", closed)
450 -
451 - return buf, nil
452 - },
404 + Encoders: cmds.EncoderMap{
405 + cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out int) error {
406 + fmt.Fprintf(w, "Closed %d stream(s)\n", out)
407 + return nil
408 + }),
409 },
410 }
411
@@ -477,11 +433,10 @@ var p2pStreamLsCmd = &cmds.Command{
433 Options: []cmdkit.Option{
434 cmdkit.BoolOption(p2pHeadersOptionName, "v", "Print table headers (ID, Protocol, Local, Remote)."),
435 },
480 - Run: func(req cmds.Request, res cmds.Response) {
481 - n, err := p2pGetNode(req)
436 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
437 + n, err := p2pGetNode(env)
438 if err != nil {
483 - res.SetError(err, cmdkit.ErrNormal)
484 - return
439 + return err
440 }
441
442 output := &P2PStreamsOutput{}
@@ -499,31 +454,24 @@ var p2pStreamLsCmd = &cmds.Command{
454 }
455 n.P2P.Streams.Unlock()
456
502 - res.SetOutput(output)
457 + return res.Emit(output)
458 },
459 Type: P2PStreamsOutput{},
505 - Marshalers: cmds.MarshalerMap{
506 - cmds.Text: func(res cmds.Response) (io.Reader, error) {
507 - v, err := unwrapOutput(res.Output())
508 - if err != nil {
509 - return nil, err
510 - }
511 -
512 - headers, _, _ := res.Request().Option(p2pHeadersOptionName).Bool()
513 - list := v.(*P2PStreamsOutput)
514 - buf := new(bytes.Buffer)
515 - w := tabwriter.NewWriter(buf, 1, 2, 1, ' ', 0)
516 - for _, stream := range list.Streams {
460 + Encoders: cmds.EncoderMap{
461 + cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *P2PStreamsOutput) error {
462 + headers, _ := req.Options[p2pHeadersOptionName].(bool)
463 + tw := tabwriter.NewWriter(w, 1, 2, 1, ' ', 0)
464 + for _, stream := range out.Streams {
465 if headers {
518 - fmt.Fprintln(w, "ID\tProtocol\tOrigin\tTarget")
466 + fmt.Fprintln(tw, "ID\tProtocol\tOrigin\tTarget")
467 }
468
521 - fmt.Fprintf(w, "%s\t%s\t%s\t%s\n", stream.HandlerID, stream.Protocol, stream.OriginAddress, stream.TargetAddress)
469 + fmt.Fprintf(tw, "%s\t%s\t%s\t%s\n", stream.HandlerID, stream.Protocol, stream.OriginAddress, stream.TargetAddress)
470 }
523 - w.Flush()
471 + tw.Flush()
472
525 - return buf, nil
526 - },
473 + return nil
474 + }),
475 },
476 }
477
@@ -537,28 +485,23 @@ var p2pStreamCloseCmd = &cmds.Command{
485 Options: []cmdkit.Option{
486 cmdkit.BoolOption(p2pAllOptionName, "a", "Close all streams."),
487 },
540 - Run: func(req cmds.Request, res cmds.Response) {
541 - res.SetOutput(nil)
542 -
543 - n, err := p2pGetNode(req)
488 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
489 + n, err := p2pGetNode(env)
490 if err != nil {
545 - res.SetError(err, cmdkit.ErrNormal)
546 - return
491 + return err
492 }
493
549 - closeAll, _, _ := req.Option(p2pAllOptionName).Bool()
494 + closeAll, _ := req.Options[p2pAllOptionName].(bool)
495 var handlerID uint64
496
497 if !closeAll {
553 - if len(req.Arguments()) == 0 {
554 - res.SetError(errors.New("no id specified"), cmdkit.ErrNormal)
555 - return
498 + if len(req.Arguments) == 0 {
499 + return errors.New("no id specified")
500 }
501
558 - handlerID, err = strconv.ParseUint(req.Arguments()[0], 10, 64)
502 + handlerID, err = strconv.ParseUint(req.Arguments[0], 10, 64)
503 if err != nil {
560 - res.SetError(err, cmdkit.ErrNormal)
561 - return
504 + return err
505 }
506 }
507
@@ -578,16 +521,18 @@ var p2pStreamCloseCmd = &cmds.Command{
521 for _, s := range toClose {
522 n.P2P.Streams.Reset(s)
523 }
524 +
525 + return nil
526 },
527 }
528
584 -func p2pGetNode(req cmds.Request) (*core.IpfsNode, error) {
585 - n, err := req.InvocContext().GetNode()
529 +func p2pGetNode(env cmds.Environment) (*core.IpfsNode, error) {
530 + nd, err := cmdenv.GetNode(env)
531 if err != nil {
532 return nil, err
533 }
534
590 - config, err := n.Repo.Config()
535 + config, err := nd.Repo.Config()
536 if err != nil {
537 return nil, err
538 }
@@ -596,9 +541,9 @@ func p2pGetNode(req cmds.Request) (*core.IpfsNode, error) {
541 return nil, errors.New("libp2p stream mounting not enabled")
542 }
543
599 - if !n.OnlineMode() {
544 + if !nd.OnlineMode() {
545 return nil, ErrNotOnline
546 }
547
603 - return n, nil
548 + return nd, nil
549 }
core/commands/root.go
+1 -1
@@ -137,7 +137,7 @@ var rootSubcommands = map[string]*cmds.Command{
137 "object": ocmd.ObjectCmd,
138 "pin": lgc.NewCommand(PinCmd),
139 "ping": PingCmd,
140 - "p2p": lgc.NewCommand(P2PCmd),
140 + "p2p": P2PCmd,
141 "refs": lgc.NewCommand(RefsCmd),
142 "resolve": ResolveCmd,
143 "swarm": SwarmCmd,