@cryptotaxi247 / kubo / commits / 908ffddc1

big squash commit

excerpt of commit messages: - update postrun functions in core/commands - sharness: allow setting -i with TEST_IMMEDIATE=1 - cmds Run func returns error now - gx update cmdkit to 1.1.2 and cmds to 2.0.0-beta1 License: MIT Signed-off-by: keks <keks@cryptoscope.co>

keks committed Apr 13, 2018 at 16:14 UTC 908ffddc1bdce22173099914e36d1698fb3c30b3
38 files changed +491 -649
cmd/ipfs/daemon.go
+21 -36
@@ -21,12 +21,12 @@ import (
21 fsrepo "github.com/ipfs/go-ipfs/repo/fsrepo"
22 migrate "github.com/ipfs/go-ipfs/repo/fsrepo/migrations"
23
24 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
24 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
25 mprome "gx/ipfs/QmUHHsirrDtP6WEHhE8SZeG672CLqDJn6XGzAHnvBHUiA3/go-metrics-prometheus"
26 "gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
27 "gx/ipfs/QmYYv3QFnfQbiwmi1tpkgKF8o4xFnZoBrvpupTiGJwL9nH/client_golang/prometheus"
28 ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
29 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
30 )
31
32 const (
@@ -184,11 +184,11 @@ func defaultMux(path string) corehttp.ServeOption {
184 }
185 }
186
187 -func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) {
187 +func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) error {
188 // Inject metrics before we do anything
189 err := mprome.Inject()
190 if err != nil {
191 - log.Errorf("Injecting prometheus handler for metrics failed with message: %s\n", err.Error())
191 + return fmt.Errorf("Injecting prometheus handler for metrics failed with message %s", err.Error())
192 }
193
194 // let the user know we're going.
@@ -227,8 +227,7 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
227
228 err := initWithDefaults(os.Stdout, cfg, profiles)
229 if err != nil {
230 - re.SetError(err, cmdkit.ErrNormal)
231 - return
230 + return err
231 }
232 }
233 }
@@ -238,8 +237,7 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
237 repo, err := fsrepo.Open(cctx.ConfigRoot)
238 switch err {
239 default:
241 - re.SetError(err, cmdkit.ErrNormal)
242 - return
240 + return err
241 case fsrepo.ErrNeedMigration:
242 domigrate, found := req.Options[migrateKwd].(bool)
243 fmt.Println("Found outdated fs-repo, migrations need to be run.")
@@ -251,8 +249,7 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
249 if !domigrate {
250 fmt.Println("Not running migrations of fs-repo now.")
251 fmt.Println("Please get fs-repo-migrations from https://dist.ipfs.io")
254 - re.SetError(fmt.Errorf("fs-repo requires migration"), cmdkit.ErrNormal)
255 - return
252 + return fmt.Errorf("fs-repo requires migration")
253 }
254
255 err = migrate.RunMigration(fsrepo.RepoVersion)
@@ -261,14 +258,12 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
258 fmt.Printf(" %s\n", err)
259 fmt.Println("If you think this is a bug, please file an issue and include this whole log output.")
260 fmt.Println(" https://github.com/ipfs/fs-repo-migrations")
264 - re.SetError(err, cmdkit.ErrNormal)
265 - return
261 + return err
262 }
263
264 repo, err = fsrepo.Open(cctx.ConfigRoot)
265 if err != nil {
270 - re.SetError(err, cmdkit.ErrNormal)
271 - return
266 + return err
267 }
268 case nil:
269 break
@@ -276,8 +271,7 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
271
272 cfg, err := cctx.GetConfig()
273 if err != nil {
279 - re.SetError(err, cmdkit.ErrNormal)
280 - return
274 + return err
275 }
276
277 offline, _ := req.Options[offlineKwd].(bool)
@@ -303,8 +297,7 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
297 if routingOption == routingOptionDefaultKwd {
298 cfg, err := repo.Config()
299 if err != nil {
306 - re.SetError(err, cmdkit.ErrNormal)
307 - return
300 + return err
301 }
302
303 routingOption = cfg.Routing.Type
@@ -314,8 +307,7 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
307 }
308 switch routingOption {
309 case routingOptionSupernodeKwd:
317 - re.SetError(errors.New("supernode routing was never fully implemented and has been removed"), cmdkit.ErrNormal)
318 - return
310 + return errors.New("supernode routing was never fully implemented and has been removed")
311 case routingOptionDHTClientKwd:
312 ncfg.Routing = core.DHTClientOption
313 case routingOptionDHTKwd:
@@ -323,15 +315,13 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
315 case routingOptionNoneKwd:
316 ncfg.Routing = core.NilRouterOption
317 default:
326 - re.SetError(fmt.Errorf("unrecognized routing option: %s", routingOption), cmdkit.ErrNormal)
327 - return
318 + return fmt.Errorf("unrecognized routing option: %s", routingOption)
319 }
320
321 node, err := core.NewNode(req.Context, ncfg)
322 if err != nil {
323 log.Error("error from node construction: ", err)
333 - re.SetError(err, cmdkit.ErrNormal)
334 - return
324 + return err
325 }
326 node.SetLocal(false)
327
@@ -361,29 +351,24 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
351 // construct api endpoint - every time
352 apiErrc, err := serveHTTPApi(req, cctx)
353 if err != nil {
364 - re.SetError(err, cmdkit.ErrNormal)
365 - return
354 + return err
355 }
356
357 // construct fuse mountpoints - if the user provided the --mount flag
358 mount, _ := req.Options[mountKwd].(bool)
359 if mount && offline {
371 - re.SetError(errors.New("mount is not currently supported in offline mode"),
372 - cmdkit.ErrClient)
373 - return
360 + return cmdkit.Errorf(cmdkit.ErrClient, "mount is not currently supported in offline mode")
361 }
362 if mount {
363 if err := mountFuse(req, cctx); err != nil {
377 - re.SetError(err, cmdkit.ErrNormal)
378 - return
364 + return err
365 }
366 }
367
368 // repo blockstore GC - if --enable-gc flag is present
369 gcErrc, err := maybeRunGC(req, node)
370 if err != nil {
385 - re.SetError(err, cmdkit.ErrNormal)
386 - return
371 + return err
372 }
373
374 // construct http gateway - if it is set in the config
@@ -392,8 +377,7 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
377 var err error
378 gwErrc, err = serveHTTPGateway(req, cctx)
379 if err != nil {
395 - re.SetError(err, cmdkit.ErrNormal)
396 - return
380 + return err
381 }
382 }
383
@@ -405,10 +389,11 @@ func daemonFunc(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment
389 // TODO(cryptix): our fuse currently doesnt follow this pattern for graceful shutdown
390 for err := range merge(apiErrc, gwErrc, gcErrc) {
391 if err != nil {
408 - log.Error(err)
409 - re.SetError(err, cmdkit.ErrNormal)
392 + return err
393 }
394 }
395 +
396 + return nil
397 }
398
399 // serveHTTPApi collects options, creates listener, prints status message and starts serving requests
cmd/ipfs/init.go
+6 -12
@@ -16,9 +16,9 @@ import (
16 namesys "github.com/ipfs/go-ipfs/namesys"
17 fsrepo "github.com/ipfs/go-ipfs/repo/fsrepo"
18
19 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
19 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
20 "gx/ipfs/QmYVqYJTVjetcf1guieEgWpK1PZtHPytP624vKzTF1P3r2/go-ipfs-config"
21 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
22 )
23
24 const (
@@ -72,11 +72,10 @@ environment variable:
72
73 return nil
74 },
75 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
75 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
76 cctx := env.(*oldcmds.Context)
77 if cctx.Online {
78 - res.SetError(errors.New("init must be run offline only"), cmdkit.ErrNormal)
79 - return
78 + return cmdkit.Error{Message: "init must be run offline only"}
79 }
80
81 empty, _ := req.Options["empty-repo"].(bool)
@@ -88,14 +87,12 @@ environment variable:
87 if f != nil {
88 confFile, err := f.NextFile()
89 if err != nil {
91 - res.SetError(err, cmdkit.ErrNormal)
92 - return
90 + return err
91 }
92
93 conf = &config.Config{}
94 if err := json.NewDecoder(confFile).Decode(conf); err != nil {
97 - res.SetError(err, cmdkit.ErrNormal)
98 - return
95 + return err
96 }
97 }
98
@@ -106,10 +103,7 @@ environment variable:
103 profiles = strings.Split(profile, ",")
104 }
105
109 - if err := doInit(os.Stdout, cctx.ConfigRoot, empty, nBitsForKeypair, profiles, conf); err != nil {
110 - res.SetError(err, cmdkit.ErrNormal)
111 - return
112 - }
106 + return doInit(os.Stdout, cctx.ConfigRoot, empty, nBitsForKeypair, profiles, conf)
107 },
108 }
109
cmd/ipfs/ipfs.go
+1 -1
@@ -5,7 +5,7 @@ import (
5
6 commands "github.com/ipfs/go-ipfs/core/commands"
7
8 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
8 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
9 )
10
11 // This is the CLI root, used for executing commands accessible to CLI clients.
cmd/ipfs/main.go
+3 -3
@@ -24,9 +24,6 @@ import (
24 repo "github.com/ipfs/go-ipfs/repo"
25 fsrepo "github.com/ipfs/go-ipfs/repo/fsrepo"
26
27 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
28 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds/cli"
29 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds/http"
27 u "gx/ipfs/QmPdKqUcHGFdeSpvjVoaTRPPstGif9GBZb5Q56RVw9o69A/go-ipfs-util"
28 logging "gx/ipfs/QmRREK2CAZ5Re2Bd9zZFG6FeYDppUWt5cMgsoUEp3ktgSr/go-log"
29 manet "gx/ipfs/QmV6FjemM1K8oXjrvuq3wuVWWoU2TLDPmNnKrxHzY3v6Ai/go-multiaddr-net"
@@ -34,6 +31,9 @@ import (
31 "gx/ipfs/QmYVqYJTVjetcf1guieEgWpK1PZtHPytP624vKzTF1P3r2/go-ipfs-config"
32 ma "gx/ipfs/QmYmsdtJ3HsodkePE3eU3TsCaP2YvPZJ4LoXnNkDE5Tpt7/go-multiaddr"
33 loggables "gx/ipfs/QmZ4zF1mBrt8C2mSCM4ZYE4aAnv78f7GvrzufJC4G5tecK/go-libp2p-loggables"
34 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
35 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds/cli"
36 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds/http"
37 madns "gx/ipfs/QmfXU2MhWoegxHoeMd3A2ytL2P6CY4FfqGWc23LTNWBwZt/go-multiaddr-dns"
38 )
39
commands/legacy/command.go
+5 -7
@@ -3,10 +3,10 @@ package legacy
3 import (
4 "io"
5
6 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
7 -
6 oldcmds "github.com/ipfs/go-ipfs/commands"
7 +
8 logging "gx/ipfs/QmRREK2CAZ5Re2Bd9zZFG6FeYDppUWt5cMgsoUEp3ktgSr/go-log"
9 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
10 )
11
12 var log = logging.Logger("cmds/lgc")
@@ -29,17 +29,15 @@ func NewCommand(oldcmd *oldcmds.Command) *cmds.Command {
29 }
30
31 if oldcmd.Run != nil {
32 - cmd.Run = func(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) {
32 + cmd.Run = func(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) error {
33 oldReq := &requestWrapper{req, OldContext(env)}
34 res := &fakeResponse{req: oldReq, re: re, wait: make(chan struct{})}
35
36 errCh := make(chan error)
37 go res.Send(errCh)
38 +
39 oldcmd.Run(oldReq, res)
39 - err := <-errCh
40 - if err != nil {
41 - log.Error(err)
42 - }
40 + return <-errCh
41 }
42 }
43
commands/legacy/legacy.go
+1 -1
@@ -4,7 +4,7 @@ import (
4 "io"
5 "runtime/debug"
6
7 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
7 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
8
9 oldcmds "github.com/ipfs/go-ipfs/commands"
10 )
commands/legacy/legacy_test.go
+59 -3
@@ -7,8 +7,8 @@ import (
7 "testing"
8
9 oldcmds "github.com/ipfs/go-ipfs/commands"
10 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
10 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
11 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
12 )
13
14 type WriteNopCloser struct {
@@ -80,7 +80,7 @@ func TestNewCommand(t *testing.T) {
80
81 root.Call(req, re, &env)
82
83 - expected := `{"Value":"Test."}
83 + expected := `"Test."
84 `
85
86 if buf.String() != expected {
@@ -114,7 +114,7 @@ func TestNewCommand(t *testing.T) {
114 }
115
116 func TestPipePair(t *testing.T) {
117 - cmd := &cmds.Command{Type: "string"}
117 + cmd := NewCommand(&oldcmds.Command{Type: "string"})
118
119 req, err := cmds.NewRequest(context.TODO(), nil, nil, nil, nil, cmd)
120 if err != nil {
@@ -134,6 +134,11 @@ func TestPipePair(t *testing.T) {
134 t.Fatal(err)
135 }
136
137 + err = re.Close()
138 + if err != nil {
139 + t.Fatal(err)
140 + }
141 +
142 close(wait)
143 }()
144
@@ -149,6 +154,57 @@ func TestPipePair(t *testing.T) {
154 t.Fatalf("expected value %#v but got %#v", expect, v)
155 }
156
157 + _, err = res.Next()
158 + if err != io.EOF {
159 + t.Fatal("expected io.EOF, got:", err)
160 + }
161 +
162 <-wait
163 +}
164 +
165 +func TestChanPair(t *testing.T) {
166 + cmd := NewCommand(&oldcmds.Command{Type: "string"})
167 +
168 + req, err := cmds.NewRequest(context.TODO(), nil, nil, nil, nil, cmd)
169 + if err != nil {
170 + t.Fatal(err)
171 + }
172 +
173 + re, res := cmds.NewChanResponsePair(req)
174
175 + wait := make(chan interface{})
176 +
177 + expect := "abc"
178 + go func() {
179 + err := re.Emit(expect)
180 + if err != nil {
181 + t.Fatal(err)
182 + }
183 +
184 + err = re.Close()
185 + if err != nil {
186 + t.Fatal(err)
187 + }
188 +
189 + close(wait)
190 + }()
191 +
192 + v, err := res.Next()
193 + if err != nil {
194 + t.Fatal(err)
195 + }
196 + str, ok := v.(string)
197 + if !ok {
198 + t.Fatalf("expected type %T but got %T", expect, v)
199 + }
200 + if str != expect {
201 + t.Fatalf("expected value %#v but got %#v", expect, v)
202 + }
203 +
204 + _, err = res.Next()
205 + if err != io.EOF {
206 + t.Fatal("expected io.EOF, got:", err)
207 + }
208 +
209 + <-wait
210 }
commands/legacy/request.go
+1 -1
@@ -7,9 +7,9 @@ import (
7 "os"
8 "reflect"
9
10 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
10 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
11 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit/files"
12 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
13
14 oldcmds "github.com/ipfs/go-ipfs/commands"
15 )
commands/legacy/response.go
+10 -10
@@ -7,8 +7,8 @@ import (
7 "reflect"
8 "sync"
9
10 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
10 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
11 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
12
13 oldcmds "github.com/ipfs/go-ipfs/commands"
14 )
@@ -34,7 +34,10 @@ func (rw *responseWrapper) Output() interface{} {
34 // get first emitted value
35 x, err := rw.Next()
36 if err != nil {
37 - return nil
37 + ch := make(chan interface{})
38 + log.Error(err)
39 + close(ch)
40 + return (<-chan interface{})(ch)
41 }
42 if e, ok := x.(*cmdkit.Error); ok {
43 ch := make(chan interface{})
@@ -120,16 +123,13 @@ func (r *fakeResponse) Send(errCh chan<- error) {
123 defer close(errCh)
124
125 out := r.Output()
123 - if out == nil {
124 - return
125 - }
126
127 - if ch, ok := out.(chan interface{}); ok {
128 - out = (<-chan interface{})(ch)
127 + // don't emit nil or Single{nil}
128 + if out == nil || out == (cmds.Single{Value: nil}) {
129 + return
130 }
131
131 - err := r.re.Emit(out)
132 - errCh <- err
132 + errCh <- r.re.Emit(out)
133 return
134 }
135
@@ -141,7 +141,7 @@ func (r *fakeResponse) Request() oldcmds.Request {
141 // SetError forwards the call to the underlying ResponseEmitter
142 func (r *fakeResponse) SetError(err error, code cmdkit.ErrorType) {
143 defer r.once.Do(func() { close(r.wait) })
144 - r.re.SetError(err, code)
144 + r.re.CloseWithError(cmdkit.Errorf(code, err.Error()))
145 }
146
147 // Error is an empty stub
commands/request.go
+1 -1
@@ -10,10 +10,10 @@ import (
10 coreapi "github.com/ipfs/go-ipfs/core/coreapi"
11 coreiface "github.com/ipfs/go-ipfs/core/coreapi/interface"
12
13 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
13 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
14 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit/files"
15 config "gx/ipfs/QmYVqYJTVjetcf1guieEgWpK1PZtHPytP624vKzTF1P3r2/go-ipfs-config"
16 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
17 )
18
19 type Context struct {
core/commands/add.go
+38 -66
@@ -1,7 +1,6 @@
1 package commands
2
3 import (
4 - "errors"
4 "fmt"
5 "io"
6 "os"
@@ -16,13 +15,13 @@ import (
15 dagtest "gx/ipfs/QmXv5mwmQ74r4aiHcNeQ4GAmfB3aWJuqaE4WyDfDfvkgLM/go-merkledag/test"
16 blockservice "gx/ipfs/Qma2KhbQarYTkmSJAeaMGRAg8HAXAhEWK8ge4SReG7ZSD3/go-blockservice"
17
19 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
18 mh "gx/ipfs/QmPnFwZ2JXKnXgMw8CdBPxn7FWh6LLdjUjxV1fKHuJnkr8/go-multihash"
19 pb "gx/ipfs/QmPtj12fdwuAqj9sBSTNUxBNu8kCGNp8b3o8yUzMm5GHpq/pb"
20 cidutil "gx/ipfs/QmQJSeE3CX4zos9qeaG8EhecEK9zvrTEfTG84J8C5NVRwt/go-cidutil"
21 mfs "gx/ipfs/QmRkrpnhZqDxTxwGCsDbuZMr7uCFZHH6SGfrcjgEQwxF3t/go-mfs"
22 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
23 files "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit/files"
24 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
25 offline "gx/ipfs/QmcRC35JF2pJQneAxa5LdQBQRumWggccWErogSrCkS1h8T/go-ipfs-exchange-offline"
26 bstore "gx/ipfs/QmegPGspn3RpTMQ23Fd3GVVMopo1zsEMurudbFMZ5UXBLH/go-ipfs-blockstore"
27 )
@@ -148,17 +147,15 @@ You can now check what blocks have been created by:
147
148 return nil
149 },
151 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
150 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
151 n, err := cmdenv.GetNode(env)
152 if err != nil {
154 - res.SetError(err, cmdkit.ErrNormal)
155 - return
153 + return err
154 }
155
156 cfg, err := n.Repo.Config()
157 if err != nil {
160 - res.SetError(err, cmdkit.ErrNormal)
161 - return
158 + return err
159 }
160 // check if repo will exceed storage limit if added
161 // TODO: this doesn't handle the case if the hashed file is already in blocks (deduplicated)
@@ -195,19 +192,14 @@ You can now check what blocks have been created by:
192
193 // nocopy -> filestoreEnabled
194 if nocopy && !cfg.Experimental.FilestoreEnabled {
198 - res.SetError(filestore.ErrFilestoreNotEnabled, cmdkit.ErrClient)
199 - return
195 + return cmdkit.Errorf(cmdkit.ErrClient, filestore.ErrFilestoreNotEnabled.Error())
196 }
197
198 // nocopy -> rawblocks
199 if nocopy && !rawblks {
200 // fixed?
201 if rbset {
206 - res.SetError(
207 - fmt.Errorf("nocopy option requires '--raw-leaves' to be enabled as well"),
208 - cmdkit.ErrNormal,
209 - )
210 - return
202 + return fmt.Errorf("nocopy option requires '--raw-leaves' to be enabled as well")
203 }
204 // No, satisfy mandatory constraint.
205 rawblks = true
@@ -216,11 +208,7 @@ You can now check what blocks have been created by:
208 // (hash != "sha2-256") -> CIDv1
209 if hashFunStr != "sha2-256" && cidVer == 0 {
210 if cidVerSet {
219 - res.SetError(
220 - errors.New("CIDv0 only supports sha2-256"),
221 - cmdkit.ErrClient,
222 - )
223 - return
211 + return cmdkit.Errorf(cmdkit.ErrClient, "CIDv0 only supports sha2-256")
212 }
213 cidVer = 1
214 }
@@ -232,14 +220,12 @@ You can now check what blocks have been created by:
220
221 prefix, err := dag.PrefixForCidVersion(cidVer)
222 if err != nil {
235 - res.SetError(err, cmdkit.ErrNormal)
236 - return
223 + return err
224 }
225
226 hashFunCode, ok := mh.Names[strings.ToLower(hashFunStr)]
227 if !ok {
241 - res.SetError(fmt.Errorf("unrecognized hash function: %s", strings.ToLower(hashFunStr)), cmdkit.ErrNormal)
242 - return
228 + return fmt.Errorf("unrecognized hash function: %s", strings.ToLower(hashFunStr))
229 }
230
231 prefix.MhType = hashFunCode
@@ -252,8 +238,7 @@ You can now check what blocks have been created by:
238 NilRepo: true,
239 })
240 if err != nil {
255 - res.SetError(err, cmdkit.ErrNormal)
256 - return
241 + return err
242 }
243 n = nilnode
244 }
@@ -276,8 +261,7 @@ You can now check what blocks have been created by:
261
262 fileAdder, err := coreunix.NewAdder(req.Context, n.Pinning, n.Blockstore, dserv)
263 if err != nil {
279 - res.SetError(err, cmdkit.ErrNormal)
280 - return
264 + return err
265 }
266
267 fileAdder.Out = outChan
@@ -307,8 +291,7 @@ You can now check what blocks have been created by:
291 emptyDirNode.SetCidBuilder(fileAdder.CidBuilder)
292 mr, err := mfs.NewRoot(req.Context, md, emptyDirNode, nil)
293 if err != nil {
310 - res.SetError(err, cmdkit.ErrNormal)
311 - return
294 + return err
295 }
296
297 fileAdder.SetMfsRoot(mr)
@@ -352,24 +335,18 @@ You can now check what blocks have been created by:
335 err = addAllAndPin(req.Files)
336 }()
337
355 - defer res.Close()
356 -
338 err = res.Emit(outChan)
339 if err != nil {
359 - log.Error(err)
360 - return
361 - }
362 - err = <-errCh
363 - if err != nil {
364 - res.SetError(err, cmdkit.ErrNormal)
340 + return err
341 }
342 +
343 + return <-errCh
344 },
345 PostRun: cmds.PostRunMap{
368 - cmds.CLI: func(req *cmds.Request, re cmds.ResponseEmitter) cmds.ResponseEmitter {
369 - reNext, res := cmds.NewChanResponsePair(req)
370 - outChan := make(chan interface{})
371 -
346 + cmds.CLI: func(res cmds.Response, re cmds.ResponseEmitter) error {
347 sizeChan := make(chan int64, 1)
348 + outChan := make(chan interface{})
349 + req := res.Request()
350
351 sizeFile, ok := req.Files.(files.SizeFile)
352 if ok {
@@ -475,38 +452,33 @@ You can now check what blocks have been created by:
452 }
453 }
454
478 - go func() {
479 - // defer order important! First close outChan, then wait for output to finish, then close re
480 - defer re.Close()
481 -
482 - if e := res.Error(); e != nil {
483 - defer close(outChan)
484 - re.SetError(e.Message, e.Code)
485 - return
486 - }
455 + if e := res.Error(); e != nil {
456 + close(outChan)
457 + return e
458 + }
459
488 - wait := make(chan struct{})
489 - go progressBar(wait)
460 + wait := make(chan struct{})
461 + go progressBar(wait)
462
491 - defer func() { <-wait }()
492 - defer close(outChan)
463 + defer func() { <-wait }()
464 + defer close(outChan)
465
494 - for {
495 - v, err := res.Next()
496 - if !cmds.HandleError(err, res, re) {
497 - break
466 + for {
467 + v, err := res.Next()
468 + if err != nil {
469 + if err == io.EOF {
470 + return nil
471 }
472
500 - select {
501 - case outChan <- v:
502 - case <-req.Context.Done():
503 - re.SetError(req.Context.Err(), cmdkit.ErrNormal)
504 - return
505 - }
473 + return err
474 }
507 - }()
475
509 - return reNext
476 + select {
477 + case outChan <- v:
478 + case <-req.Context.Done():
479 + return req.Context.Err()
480 + }
481 + }
482 },
483 },
484 Type: coreunix.AddedObject{},
core/commands/bitswap.go
+7 -11
@@ -13,9 +13,9 @@ import (
13 decision "gx/ipfs/QmUyaGN3WPr3CTLai7DBvMikagK45V4fUi8p8cNRaJQoU1/go-bitswap/decision"
14
15 "gx/ipfs/QmPSBJL4momYnE7DcUyk2DVhD6rH488ZmHBGLbxNdhU44K/go-humanize"
16 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
16 peer "gx/ipfs/QmQsErDt8Qgw1XrsXf2BpEzDgGWtB1YLsTAARBup5b6B9W/go-libp2p-peer"
17 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
18 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
19 )
20
21 var BitswapCmd = &cmds.Command{
@@ -92,31 +92,27 @@ var bitswapStatCmd = &cmds.Command{
92 ShortDescription: ``,
93 },
94 Type: bitswap.Stat{},
95 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
95 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
96 nd, err := cmdenv.GetNode(env)
97 if err != nil {
98 - res.SetError(err, cmdkit.ErrNormal)
99 - return
98 + return err
99 }
100
101 if !nd.OnlineMode() {
103 - res.SetError(ErrNotOnline, cmdkit.ErrClient)
104 - return
102 + return cmdkit.Errorf(cmdkit.ErrClient, ErrNotOnline.Error())
103 }
104
105 bs, ok := nd.Exchange.(*bitswap.Bitswap)
106 if !ok {
109 - res.SetError(e.TypeErr(bs, nd.Exchange), cmdkit.ErrNormal)
110 - return
107 + return e.TypeErr(bs, nd.Exchange)
108 }
109
110 st, err := bs.Stat()
111 if err != nil {
115 - res.SetError(err, cmdkit.ErrNormal)
116 - return
112 + return err
113 }
114
119 - cmds.EmitOnce(res, st)
115 + return cmds.EmitOnce(res, st)
116 },
117 Encoders: cmds.EncoderMap{
118 cmds.Text: cmds.MakeEncoder(func(req *cmds.Request, w io.Writer, v interface{}) error {
core/commands/block.go
+36 -60
@@ -1,6 +1,7 @@
1 package commands
2
3 import (
4 + "errors"
5 "fmt"
6 "io"
7 "os"
@@ -11,9 +12,9 @@ import (
12 coreiface "github.com/ipfs/go-ipfs/core/coreapi/interface"
13 "github.com/ipfs/go-ipfs/core/coreapi/interface/options"
14
14 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
15 mh "gx/ipfs/QmPnFwZ2JXKnXgMw8CdBPxn7FWh6LLdjUjxV1fKHuJnkr8/go-multihash"
16 - "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
16 + cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
17 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
18 )
19
20 type BlockStat struct {
@@ -59,32 +60,26 @@ on raw IPFS blocks. It outputs the following to stdout:
60 Arguments: []cmdkit.Argument{
61 cmdkit.StringArg("key", true, false, "The base58 multihash of an existing block to stat.").EnableStdin(),
62 },
62 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
63 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
64 api, err := cmdenv.GetApi(env)
65 if err != nil {
65 - res.SetError(err, cmdkit.ErrNormal)
66 - return
66 + return err
67 }
68
69 p, err := coreiface.ParsePath(req.Arguments[0])
70 if err != nil {
71 - res.SetError(err, cmdkit.ErrNormal)
72 - return
71 + return err
72 }
73
74 b, err := api.Block().Stat(req.Context, p)
75 if err != nil {
77 - res.SetError(err, cmdkit.ErrNormal)
78 - return
76 + return err
77 }
78
81 - err = cmds.EmitOnce(res, &BlockStat{
79 + return cmds.EmitOnce(res, &BlockStat{
80 Key: b.Path().Cid().String(),
81 Size: b.Size(),
82 })
85 - if err != nil {
86 - log.Error(err)
87 - }
83 },
84 Type: BlockStat{},
85 Encoders: cmds.EncoderMap{
@@ -111,29 +106,23 @@ It outputs to stdout, and <key> is a base58 encoded multihash.
106 Arguments: []cmdkit.Argument{
107 cmdkit.StringArg("key", true, false, "The base58 multihash of an existing block to get.").EnableStdin(),
108 },
114 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
109 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
110 api, err := cmdenv.GetApi(env)
111 if err != nil {
117 - res.SetError(err, cmdkit.ErrNormal)
118 - return
112 + return err
113 }
114
115 p, err := coreiface.ParsePath(req.Arguments[0])
116 if err != nil {
123 - res.SetError(err, cmdkit.ErrNormal)
124 - return
117 + return err
118 }
119
120 r, err := api.Block().Get(req.Context, p)
121 if err != nil {
129 - res.SetError(err, cmdkit.ErrNormal)
130 - return
122 + return err
123 }
124
133 - err = res.Emit(r)
134 - if err != nil {
135 - log.Error(err)
136 - }
125 + return res.Emit(r)
126 },
127 }
128
@@ -157,31 +146,26 @@ than 'sha2-256' or format to anything other than 'v0' will result in CIDv1.
146 cmdkit.StringOption("mhtype", "multihash hash function").WithDefault("sha2-256"),
147 cmdkit.IntOption("mhlen", "multihash hash length").WithDefault(-1),
148 },
160 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
149 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
150 api, err := cmdenv.GetApi(env)
151 if err != nil {
163 - res.SetError(err, cmdkit.ErrNormal)
164 - return
152 + return err
153 }
154
155 file, err := req.Files.NextFile()
156 if err != nil {
169 - res.SetError(err, cmdkit.ErrNormal)
170 - return
157 + return err
158 }
159
160 mhtype, _ := req.Options["mhtype"].(string)
161 mhtval, ok := mh.Names[mhtype]
162 if !ok {
176 - err := fmt.Errorf("unrecognized multihash function: %s", mhtype)
177 - res.SetError(err, cmdkit.ErrNormal)
178 - return
163 + return fmt.Errorf("unrecognized multihash function: %s", mhtype)
164 }
165
166 mhlen, ok := req.Options["mhlen"].(int)
167 if !ok {
183 - res.SetError("missing option \"mhlen\"", cmdkit.ErrNormal)
184 - return
168 + return errors.New("missing option \"mhlen\"")
169 }
170
171 format, formatSet := req.Options["format"].(string)
@@ -195,17 +179,13 @@ than 'sha2-256' or format to anything other than 'v0' will result in CIDv1.
179
180 p, err := api.Block().Put(req.Context, file, options.Block.Hash(mhtval, mhlen), options.Block.Format(format))
181 if err != nil {
198 - res.SetError(err, cmdkit.ErrNormal)
199 - return
182 + return err
183 }
184
202 - err = cmds.EmitOnce(res, &BlockStat{
185 + return cmds.EmitOnce(res, &BlockStat{
186 Key: p.Path().Cid().String(),
187 Size: p.Size(),
188 })
206 - if err != nil {
207 - log.Error(err)
208 - }
189 },
190 Encoders: cmds.EncoderMap{
191 cmds.Text: cmds.MakeEncoder(func(req *cmds.Request, w io.Writer, v interface{}) error {
@@ -235,11 +215,10 @@ It takes a list of base58 encoded multihashes to remove.
215 cmdkit.BoolOption("force", "f", "Ignore nonexistent blocks."),
216 cmdkit.BoolOption("quiet", "q", "Write minimal output."),
217 },
238 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
218 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
219 api, err := cmdenv.GetApi(env)
220 if err != nil {
241 - res.SetError(err, cmdkit.ErrNormal)
242 - return
221 + return err
222 }
223
224 force, _ := req.Options["force"].(bool)
@@ -249,43 +228,40 @@ It takes a list of base58 encoded multihashes to remove.
228 for _, b := range req.Arguments {
229 p, err := coreiface.ParsePath(b)
230 if err != nil {
252 - res.SetError(err, cmdkit.ErrNormal)
253 - return
231 + return err
232 }
233
234 rp, err := api.ResolvePath(req.Context, p)
235 if err != nil {
258 - res.SetError(err, cmdkit.ErrNormal)
259 - return
236 + return err
237 }
238
239 err = api.Block().Rm(req.Context, rp, options.Block.Force(force))
240 if err != nil {
264 - res.Emit(&util.RemovedBlock{
241 + err := res.Emit(&util.RemovedBlock{
242 Hash: rp.Cid().String(),
243 Error: err.Error(),
244 })
245 + if err != nil {
246 + return err
247 + }
248 }
249
250 if !quiet {
271 - res.Emit(&util.RemovedBlock{
251 + err := res.Emit(&util.RemovedBlock{
252 Hash: rp.Cid().String(),
253 })
254 + if err != nil {
255 + return err
256 + }
257 }
258 }
259 +
260 + return nil
261 },
262 PostRun: cmds.PostRunMap{
278 - cmds.CLI: func(req *cmds.Request, re cmds.ResponseEmitter) cmds.ResponseEmitter {
279 - reNext, res := cmds.NewChanResponsePair(req)
280 -
281 - go func() {
282 - defer re.Close()
283 -
284 - err := util.ProcRmOutput(res.Next, os.Stdout, os.Stderr)
285 - cmds.HandleError(err, res, re)
286 - }()
287 -
288 - return reNext
263 + cmds.CLI: func(res cmds.Response, re cmds.ResponseEmitter) error {
264 + return util.ProcRmOutput(res.Next, os.Stdout, os.Stderr)
265 },
266 },
267 Type: util.RemovedBlock{},
core/commands/cat.go
+30 -49
@@ -10,8 +10,8 @@ import (
10 cmdenv "github.com/ipfs/go-ipfs/core/commands/cmdenv"
11 coreunix "github.com/ipfs/go-ipfs/core/coreunix"
12
13 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
13 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
14 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
15 )
16
17 const progressBarMinSize = 1024 * 1024 * 8 // show progress bar for outputs > 8MiB
@@ -29,34 +29,29 @@ var CatCmd = &cmds.Command{
29 cmdkit.IntOption("offset", "o", "Byte offset to begin reading from."),
30 cmdkit.IntOption("length", "l", "Maximum number of bytes to read."),
31 },
32 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
32 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
33 node, err := cmdenv.GetNode(env)
34 if err != nil {
35 - res.SetError(err, cmdkit.ErrNormal)
36 - return
35 + return err
36 }
37
38 if !node.OnlineMode() {
39 if err := node.SetupOfflineRouting(); err != nil {
41 - res.SetError(err, cmdkit.ErrNormal)
42 - return
40 + return err
41 }
42 }
43
44 offset, _ := req.Options["offset"].(int)
45 if offset < 0 {
48 - res.SetError(fmt.Errorf("cannot specify negative offset"), cmdkit.ErrNormal)
49 - return
46 + return fmt.Errorf("cannot specify negative offset")
47 }
48
49 max, found := req.Options["length"].(int)
50 if err != nil {
54 - res.SetError(err, cmdkit.ErrNormal)
55 - return
51 + return err
52 }
53 if max < 0 {
58 - res.SetError(fmt.Errorf("cannot specify negative length"), cmdkit.ErrNormal)
59 - return
54 + return fmt.Errorf("cannot specify negative length")
55 }
56 if !found {
57 max = -1
@@ -64,14 +59,12 @@ var CatCmd = &cmds.Command{
59
60 err = req.ParseBodyArgs()
61 if err != nil {
67 - res.SetError(err, cmdkit.ErrNormal)
68 - return
62 + return err
63 }
64
65 readers, length, err := cat(req.Context, node, req.Arguments, int64(offset), int64(max))
66 if err != nil {
73 - res.SetError(err, cmdkit.ErrNormal)
74 - return
67 + return err
68 }
69
70 /*
@@ -88,48 +81,36 @@ var CatCmd = &cmds.Command{
81 // returned from io.Copy inside Emit, we need to take Emit errors and send
82 // them to the client. Usually we don't do that because it means the connection
83 // is broken or we supplied an illegal argument etc.
91 - err = res.Emit(reader)
92 - if err != nil {
93 - res.SetError(err, cmdkit.ErrNormal)
94 - }
84 + return res.Emit(reader)
85 },
86 PostRun: cmds.PostRunMap{
97 - cmds.CLI: func(req *cmds.Request, re cmds.ResponseEmitter) cmds.ResponseEmitter {
98 - reNext, res := cmds.NewChanResponsePair(req)
87 + cmds.CLI: func(res cmds.Response, re cmds.ResponseEmitter) error {
88 + if res.Length() > 0 && res.Length() < progressBarMinSize {
89 + return cmds.Copy(re, res)
90 + }
91
100 - go func() {
101 - if res.Length() > 0 && res.Length() < progressBarMinSize {
102 - if err := cmds.Copy(re, res); err != nil {
103 - re.SetError(err, cmdkit.ErrNormal)
92 + for {
93 + v, err := res.Next()
94 + if err != nil {
95 + if err == io.EOF {
96 + return nil
97 }
105 -
106 - return
98 + return err
99 }
100
109 - // Copy closes by itself, so we must not do this before
110 - defer re.Close()
111 - for {
112 - v, err := res.Next()
113 - if !cmds.HandleError(err, res, re) {
114 - break
115 - }
101 + switch val := v.(type) {
102 + case io.Reader:
103 + bar, reader := progressBarForReader(os.Stderr, val, int64(res.Length()))
104 + bar.Start()
105
117 - switch val := v.(type) {
118 - case io.Reader:
119 - bar, reader := progressBarForReader(os.Stderr, val, int64(res.Length()))
120 - bar.Start()
121 -
122 - err = re.Emit(reader)
123 - if err != nil {
124 - log.Error(err)
125 - }
126 - default:
127 - log.Warningf("cat postrun: received unexpected type %T", val)
106 + err = re.Emit(reader)
107 + if err != nil {
108 + return err
109 }
110 + default:
111 + log.Warningf("cat postrun: received unexpected type %T", val)
112 }
130 - }()
131 -
132 - return reNext
113 + }
114 },
115 },
116 }
core/commands/cmdenv/env.go
+1 -1
@@ -7,8 +7,8 @@ import (
7 "github.com/ipfs/go-ipfs/core"
8 coreiface "github.com/ipfs/go-ipfs/core/coreapi/interface"
9
10 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
10 config "gx/ipfs/QmYVqYJTVjetcf1guieEgWpK1PZtHPytP624vKzTF1P3r2/go-ipfs-config"
11 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
12 )
13
14 // GetNode extracts the node from the environment.
core/commands/commands.go
+3 -6
@@ -13,8 +13,8 @@ import (
13
14 e "github.com/ipfs/go-ipfs/core/commands/e"
15
16 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
16 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
17 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
18 )
19
20 type commandEncoder struct {
@@ -68,13 +68,10 @@ func CommandsCmd(root *cmds.Command) *cmds.Command {
68 Options: []cmdkit.Option{
69 cmdkit.BoolOption(flagsOptionName, "f", "Show command flags"),
70 },
71 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
71 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
72 rootCmd := cmd2outputCmd("ipfs", root)
73 rootCmd.showOpts, _ = req.Options[flagsOptionName].(bool)
74 - err := cmds.EmitOnce(res, &rootCmd)
75 - if err != nil {
76 - log.Error(err)
77 - }
74 + return cmds.EmitOnce(res, &rootCmd)
75 },
76 Encoders: cmds.EncoderMap{
77 cmds.Text: func(req *cmds.Request) func(io.Writer) cmds.Encoder {
core/commands/commands_test.go
+1 -1
@@ -4,7 +4,7 @@ import (
4 "strings"
5 "testing"
6
7 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
7 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
8 )
9
10 func collectPaths(prefix string, cmd *cmds.Command, out map[string]struct{}) {
core/commands/files.go
+24 -42
@@ -25,11 +25,11 @@ import (
25
26 humanize "gx/ipfs/QmPSBJL4momYnE7DcUyk2DVhD6rH488ZmHBGLbxNdhU44K/go-humanize"
27 cid "gx/ipfs/QmPSQnBKM9g7BaUcZCvswUJVscQ1ipjmwxN5PXCjkp9EQ7/go-cid"
28 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
28 mh "gx/ipfs/QmPnFwZ2JXKnXgMw8CdBPxn7FWh6LLdjUjxV1fKHuJnkr8/go-multihash"
29 logging "gx/ipfs/QmRREK2CAZ5Re2Bd9zZFG6FeYDppUWt5cMgsoUEp3ktgSr/go-log"
30 mfs "gx/ipfs/QmRkrpnhZqDxTxwGCsDbuZMr7uCFZHH6SGfrcjgEQwxF3t/go-mfs"
31 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
32 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
33 offline "gx/ipfs/QmcRC35JF2pJQneAxa5LdQBQRumWggccWErogSrCkS1h8T/go-ipfs-exchange-offline"
34 ipld "gx/ipfs/QmdDXJs4axxefSPgK6Y1QhpJWKuDPnGJiqgq4uncb4rFHL/go-ipld-format"
35 )
@@ -108,23 +108,22 @@ var filesStatCmd = &cmds.Command{
108 cmdkit.BoolOption("size", "Print only size. Implies '--format=<cumulsize>'. Conflicts with other format options."),
109 cmdkit.BoolOption("with-local", "Compute the amount of the dag that is local, and if possible the total size"),
110 },
111 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
111 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
112
113 _, err := statGetFormatOptions(req)
114 if err != nil {
115 - res.SetError(err, cmdkit.ErrClient)
115 + // REVIEW NOTE: We didn't return here before, was that correct?
116 + return cmdkit.Errorf(cmdkit.ErrClient, err.Error())
117 }
118
119 node, err := cmdenv.GetNode(env)
120 if err != nil {
120 - res.SetError(err, cmdkit.ErrNormal)
121 - return
121 + return err
122 }
123
124 path, err := checkPath(req.Arguments[0])
125 if err != nil {
126 - res.SetError(err, cmdkit.ErrNormal)
127 - return
126 + return err
127 }
128
129 withLocal, _ := req.Options["with-local"].(bool)
@@ -142,19 +141,16 @@ var filesStatCmd = &cmds.Command{
141
142 nd, err := getNodeFromPath(req.Context, node, dagserv, path)
143 if err != nil {
145 - res.SetError(err, cmdkit.ErrNormal)
146 - return
144 + return err
145 }
146
147 o, err := statNode(nd)
148 if err != nil {
151 - res.SetError(err, cmdkit.ErrNormal)
152 - return
149 + return err
150 }
151
152 if !withLocal {
156 - cmds.EmitOnce(res, o)
157 - return
153 + return cmds.EmitOnce(res, o)
154 }
155
156 local, sizeLocal, err := walkBlock(req.Context, dagserv, nd)
@@ -163,7 +159,7 @@ var filesStatCmd = &cmds.Command{
159 o.Local = local
160 o.SizeLocal = sizeLocal
161
166 - cmds.EmitOnce(res, o)
162 + return cmds.EmitOnce(res, o)
163 },
164 Encoders: cmds.EncoderMap{
165 cmds.Text: cmds.MakeEncoder(func(req *cmds.Request, w io.Writer, v interface{}) error {
@@ -729,11 +725,10 @@ stat' on the file or any of its ancestors.
725 cidVersionOption,
726 hashOption,
727 },
732 - Run: func(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) {
728 + Run: func(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) error {
729 path, err := checkPath(req.Arguments[0])
730 if err != nil {
735 - re.SetError(err, cmdkit.ErrNormal)
736 - return
731 + return err
732 }
733
734 create, _ := req.Options["create"].(bool)
@@ -744,34 +739,29 @@ stat' on the file or any of its ancestors.
739
740 prefix, err := getPrefixNew(req)
741 if err != nil {
747 - re.SetError(err, cmdkit.ErrNormal)
748 - return
742 + return err
743 }
744
745 nd, err := cmdenv.GetNode(env)
746 if err != nil {
753 - re.SetError(err, cmdkit.ErrNormal)
754 - return
747 + return err
748 }
749
750 offset, _ := req.Options["offset"].(int)
751 if offset < 0 {
759 - re.SetError(fmt.Errorf("cannot have negative write offset"), cmdkit.ErrNormal)
760 - return
752 + return fmt.Errorf("cannot have negative write offset")
753 }
754
755 if mkParents {
756 err := ensureContainingDirectoryExists(nd.FilesRoot, path, prefix)
757 if err != nil {
766 - re.SetError(err, cmdkit.ErrNormal)
767 - return
758 + return err
759 }
760 }
761
762 fi, err := getFileHandle(nd.FilesRoot, path, create, prefix)
763 if err != nil {
773 - re.SetError(err, cmdkit.ErrNormal)
774 - return
764 + return err
765 }
766 if rawLeavesDef {
767 fi.RawLeaves = rawLeaves
@@ -779,41 +769,36 @@ stat' on the file or any of its ancestors.
769
770 wfd, err := fi.Open(mfs.OpenWriteOnly, flush)
771 if err != nil {
782 - re.SetError(err, cmdkit.ErrNormal)
783 - return
772 + return err
773 }
774
775 defer func() {
776 err := wfd.Close()
777 if err != nil {
789 - re.SetError(err, cmdkit.ErrNormal)
778 + re.CloseWithError(cmdkit.Errorf(cmdkit.ErrNormal, err.Error()))
779 }
780 }()
781
782 if trunc {
783 if err := wfd.Truncate(0); err != nil {
795 - re.SetError(err, cmdkit.ErrNormal)
796 - return
784 + return err
785 }
786 }
787
788 count, countfound := req.Options["count"].(int)
789 if countfound && count < 0 {
802 - re.SetError(fmt.Errorf("cannot have negative byte count"), cmdkit.ErrNormal)
803 - return
790 + return fmt.Errorf("cannot have negative byte count")
791 }
792
793 _, err = wfd.Seek(int64(offset), io.SeekStart)
794 if err != nil {
795 flog.Error("seekfail: ", err)
809 - re.SetError(err, cmdkit.ErrNormal)
810 - return
796 + return err
797 }
798
799 input, err := req.Files.NextFile()
800 if err != nil {
815 - re.SetError(err, cmdkit.ErrNormal)
816 - return
801 + return err
802 }
803
804 var r io.Reader = input
@@ -822,10 +807,7 @@ stat' on the file or any of its ancestors.
807 }
808
809 _, err = io.Copy(wfd, r)
825 - if err != nil {
826 - re.SetError(err, cmdkit.ErrNormal)
827 - return
828 - }
810 + return err
811 },
812 }
813
core/commands/filestore.go
+36 -46
@@ -14,8 +14,8 @@ import (
14 "github.com/ipfs/go-ipfs/filestore"
15
16 cid "gx/ipfs/QmPSQnBKM9g7BaUcZCvswUJVscQ1ipjmwxN5PXCjkp9EQ7/go-cid"
17 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
17 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
18 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
19 )
20
21 var FileStoreCmd = &cmds.Command{
@@ -49,11 +49,10 @@ The output is:
49 Options: []cmdkit.Option{
50 cmdkit.BoolOption("file-order", "sort the results based on the path of the backing file"),
51 },
52 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
52 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
53 _, fs, err := getFilestore(env)
54 if err != nil {
55 - res.SetError(err, cmdkit.ErrNormal)
56 - return
55 + return err
56 }
57 args := req.Arguments
58 if len(args) > 0 {
@@ -61,59 +60,50 @@ The output is:
60 return filestore.List(fs, c)
61 })
62
64 - err = res.Emit(out)
65 - if err != nil {
66 - log.Error(err)
67 - }
68 - } else {
69 - fileOrder, _ := req.Options["file-order"].(bool)
70 - next, err := filestore.ListAll(fs, fileOrder)
71 - if err != nil {
72 - res.SetError(err, cmdkit.ErrNormal)
73 - return
74 - }
63 + return res.Emit(out)
64 + }
65
76 - out := listResToChan(req.Context, next)
77 - err = res.Emit(out)
78 - if err != nil {
79 - log.Error(err)
80 - }
66 + fileOrder, _ := req.Options["file-order"].(bool)
67 + next, err := filestore.ListAll(fs, fileOrder)
68 + if err != nil {
69 + return err
70 }
71 +
72 + out := listResToChan(req.Context, next)
73 + return res.Emit(out)
74 },
75 PostRun: cmds.PostRunMap{
84 - cmds.CLI: func(req *cmds.Request, re cmds.ResponseEmitter) cmds.ResponseEmitter {
85 - reNext, res := cmds.NewChanResponsePair(req)
86 -
87 - go func() {
88 - defer re.Close()
89 -
90 - var errors bool
91 - for {
92 - v, err := res.Next()
93 - if !cmds.HandleError(err, res, re) {
76 + cmds.CLI: func(res cmds.Response, re cmds.ResponseEmitter) error {
77 + var errors bool
78 + for {
79 + v, err := res.Next()
80 + if err != nil {
81 + if err == io.EOF {
82 break
83 }
84 + return err
85 + }
86
97 - r, ok := v.(*filestore.ListRes)
98 - if !ok {
99 - log.Error(e.New(e.TypeErr(r, v)))
100 - return
101 - }
102 -
103 - if r.ErrorMsg != "" {
104 - errors = true
105 - fmt.Fprintf(os.Stderr, "%s\n", r.ErrorMsg)
106 - } else {
107 - fmt.Fprintf(os.Stdout, "%s\n", r.FormatLong())
108 - }
87 + r, ok := v.(*filestore.ListRes)
88 + if !ok {
89 + // TODO or just return that error? why didn't we do that before?
90 + log.Error(e.New(e.TypeErr(r, v)))
91 + break
92 }
93
111 - if errors {
112 - re.SetError("errors while displaying some entries", cmdkit.ErrNormal)
94 + if r.ErrorMsg != "" {
95 + errors = true
96 + fmt.Fprintf(os.Stderr, "%s\n", r.ErrorMsg)
97 + } else {
98 + fmt.Fprintf(os.Stdout, "%s\n", r.FormatLong())
99 }
114 - }()
100 + }
101 +
102 + if errors {
103 + return fmt.Errorf("errors while displaying some entries")
104 + }
105
116 - return reNext
106 + return nil
107 },
108 },
109 Type: filestore.ListRes{},
core/commands/get.go
+42 -56
@@ -14,12 +14,12 @@ import (
14 e "github.com/ipfs/go-ipfs/core/commands/e"
15
16 uarchive "gx/ipfs/QmPL8bYtbACcSFFiSr4s2du7Na382NxRADR8hC7D9FkEA2/go-unixfs/archive"
17 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
17 "gx/ipfs/QmPtj12fdwuAqj9sBSTNUxBNu8kCGNp8b3o8yUzMm5GHpq/pb"
18 tar "gx/ipfs/QmQine7gvHncNevKtG9QXxf3nXcwSj6aDDmMm52mHofEEp/tar-utils"
19 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
20 path "gx/ipfs/QmX7uSbkNz76yNwBhuwYwRbhihLnJqM73VTCjS3UMJud9A/go-path"
21 dag "gx/ipfs/QmXv5mwmQ74r4aiHcNeQ4GAmfB3aWJuqaE4WyDfDfvkgLM/go-merkledag"
22 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
23 )
24
25 var ErrInvalidCompressionLevel = errors.New("compression level must be between 1 and 9")
@@ -53,93 +53,79 @@ may also specify the level of compression by specifying '-l=<1-9>'.
53 _, err := getCompressOptions(req)
54 return err
55 },
56 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
56 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
57 cmplvl, err := getCompressOptions(req)
58 if err != nil {
59 - res.SetError(err, cmdkit.ErrNormal)
60 - return
59 + return err
60 }
61
62 node, err := cmdenv.GetNode(env)
63 if err != nil {
65 - res.SetError(err, cmdkit.ErrNormal)
66 - return
64 + return err
65 }
66 p := path.Path(req.Arguments[0])
67 ctx := req.Context
68 dn, err := core.Resolve(ctx, node.Namesys, node.Resolver, p)
69 if err != nil {
72 - res.SetError(err, cmdkit.ErrNormal)
73 - return
70 + return err
71 }
72
73 switch dn := dn.(type) {
74 case *dag.ProtoNode:
75 size, err := dn.Size()
76 if err != nil {
80 - res.SetError(err, cmdkit.ErrNormal)
81 - return
77 + return err
78 }
79
80 res.SetLength(size)
81 case *dag.RawNode:
82 res.SetLength(uint64(len(dn.RawData())))
83 default:
88 - res.SetError(err, cmdkit.ErrNormal)
89 - return
84 + return err
85 }
86
87 archive, _ := req.Options["archive"].(bool)
88 reader, err := uarchive.DagArchive(ctx, dn, p.String(), node.DAG, archive, cmplvl)
89 if err != nil {
95 - res.SetError(err, cmdkit.ErrNormal)
96 - return
90 + return err
91 }
92
99 - res.Emit(reader)
93 + return res.Emit(reader)
94 },
95 PostRun: cmds.PostRunMap{
102 - cmds.CLI: func(req *cmds.Request, re cmds.ResponseEmitter) cmds.ResponseEmitter {
103 - reNext, res := cmds.NewChanResponsePair(req)
104 -
105 - go func() {
106 - defer re.Close()
107 -
108 - v, err := res.Next()
109 - if !cmds.HandleError(err, res, re) {
110 - return
111 - }
112 -
113 - outReader, ok := v.(io.Reader)
114 - if !ok {
115 - log.Error(e.New(e.TypeErr(outReader, v)))
116 - return
117 - }
118 -
119 - outPath := getOutPath(req)
120 -
121 - cmplvl, err := getCompressOptions(req)
122 - if err != nil {
123 - re.SetError(err, cmdkit.ErrNormal)
124 - return
125 - }
126 -
127 - archive, _ := req.Options["archive"].(bool)
128 -
129 - gw := getWriter{
130 - Out: os.Stdout,
131 - Err: os.Stderr,
132 - Archive: archive,
133 - Compression: cmplvl,
134 - Size: int64(res.Length()),
135 - }
136 -
137 - if err := gw.Write(outReader, outPath); err != nil {
138 - re.SetError(err, cmdkit.ErrNormal)
139 - }
140 - }()
141 -
142 - return reNext
96 + cmds.CLI: func(res cmds.Response, re cmds.ResponseEmitter) error {
97 + req := res.Request()
98 +
99 + v, err := res.Next()
100 + if err != nil {
101 + return err
102 + }
103 +
104 + outReader, ok := v.(io.Reader)
105 + if !ok {
106 + // TODO or just return the error here?
107 + log.Error(e.New(e.TypeErr(outReader, v)))
108 + return nil
109 + }
110 +
111 + outPath := getOutPath(req)
112 +
113 + cmplvl, err := getCompressOptions(req)
114 + if err != nil {
115 + return err
116 + }
117 +
118 + archive, _ := req.Options["archive"].(bool)
119 +
120 + gw := getWriter{
121 + Out: os.Stdout,
122 + Err: os.Stderr,
123 + Archive: archive,
124 + Compression: cmplvl,
125 + Size: int64(res.Length()),
126 + }
127 +
128 + return gw.Write(outReader, outPath)
129 },
130 },
131 }
core/commands/get_test.go
+1 -1
@@ -5,8 +5,8 @@ import (
5 "fmt"
6 "testing"
7
8 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
8 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
9 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
10 )
11
12 func TestGetOutputPath(t *testing.T) {
core/commands/helptext_test.go
+1 -1
@@ -4,7 +4,7 @@ import (
4 "strings"
5 "testing"
6
7 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
7 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
8 )
9
10 func checkHelptextRecursive(t *testing.T, name []string, c *cmds.Command) {
core/commands/keystore.go
+19 -29
@@ -9,8 +9,8 @@ import (
9 "github.com/ipfs/go-ipfs/core/commands/e"
10 "github.com/ipfs/go-ipfs/core/coreapi/interface/options"
11
12 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
12 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
13 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
14 )
15
16 var KeyCmd = &cmds.Command{
@@ -66,23 +66,20 @@ var keyGenCmd = &cmds.Command{
66 Arguments: []cmdkit.Argument{
67 cmdkit.StringArg("name", true, false, "name of key to create"),
68 },
69 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
69 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
70 api, err := cmdenv.GetApi(env)
71 if err != nil {
72 - res.SetError(err, cmdkit.ErrNormal)
73 - return
72 + return err
73 }
74
75 typ, f := req.Options["type"].(string)
76 if !f {
78 - res.SetError(fmt.Errorf("please specify a key type with --type"), cmdkit.ErrNormal)
79 - return
77 + return fmt.Errorf("please specify a key type with --type")
78 }
79
80 name := req.Arguments[0]
81 if name == "self" {
84 - res.SetError(fmt.Errorf("cannot create key with name 'self'"), cmdkit.ErrNormal)
85 - return
82 + return fmt.Errorf("cannot create key with name 'self'")
83 }
84
85 opts := []options.KeyGenerateOption{options.Key.Type(typ)}
@@ -95,11 +92,10 @@ var keyGenCmd = &cmds.Command{
92 key, err := api.Key().Generate(req.Context, name, opts...)
93
94 if err != nil {
98 - res.SetError(err, cmdkit.ErrNormal)
99 - return
95 + return err
96 }
97
102 - cmds.EmitOnce(res, &KeyOutput{
98 + return cmds.EmitOnce(res, &KeyOutput{
99 Name: name,
100 Id: key.ID().Pretty(),
101 })
@@ -125,17 +121,15 @@ var keyListCmd = &cmds.Command{
121 Options: []cmdkit.Option{
122 cmdkit.BoolOption("l", "Show extra information about keys."),
123 },
128 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
124 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
125 api, err := cmdenv.GetApi(env)
126 if err != nil {
131 - res.SetError(err, cmdkit.ErrNormal)
132 - return
127 + return err
128 }
129
130 keys, err := api.Key().List(req.Context)
131 if err != nil {
137 - res.SetError(err, cmdkit.ErrNormal)
138 - return
132 + return err
133 }
134
135 list := make([]KeyOutput, 0, len(keys))
@@ -144,7 +138,7 @@ var keyListCmd = &cmds.Command{
138 list = append(list, KeyOutput{Name: key.Name(), Id: key.ID().Pretty()})
139 }
140
147 - cmds.EmitOnce(res, &KeyOutputList{list})
141 + return cmds.EmitOnce(res, &KeyOutputList{list})
142 },
143 Encoders: cmds.EncoderMap{
144 cmds.Text: keyOutputListMarshaler(),
@@ -163,11 +157,10 @@ var keyRenameCmd = &cmds.Command{
157 Options: []cmdkit.Option{
158 cmdkit.BoolOption("force", "f", "Allow to overwrite an existing key."),
159 },
166 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
160 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
161 api, err := cmdenv.GetApi(env)
162 if err != nil {
169 - res.SetError(err, cmdkit.ErrNormal)
170 - return
163 + return err
164 }
165
166 name := req.Arguments[0]
@@ -176,11 +169,10 @@ var keyRenameCmd = &cmds.Command{
169
170 key, overwritten, err := api.Key().Rename(req.Context, name, newName, options.Key.Force(force))
171 if err != nil {
179 - res.SetError(err, cmdkit.ErrNormal)
180 - return
172 + return err
173 }
174
183 - cmds.EmitOnce(res, &KeyRenameOutput{
175 + return cmds.EmitOnce(res, &KeyRenameOutput{
176 Was: name,
177 Now: newName,
178 Id: key.ID().Pretty(),
@@ -215,11 +207,10 @@ var keyRmCmd = &cmds.Command{
207 Options: []cmdkit.Option{
208 cmdkit.BoolOption("l", "Show extra information about keys."),
209 },
218 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
210 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
211 api, err := cmdenv.GetApi(env)
212 if err != nil {
221 - res.SetError(err, cmdkit.ErrNormal)
222 - return
213 + return err
214 }
215
216 names := req.Arguments
@@ -228,14 +219,13 @@ var keyRmCmd = &cmds.Command{
219 for _, name := range names {
220 key, err := api.Key().Remove(req.Context, name)
221 if err != nil {
231 - res.SetError(err, cmdkit.ErrNormal)
232 - return
222 + return err
223 }
224
225 list = append(list, KeyOutput{Name: name, Id: key.ID().Pretty()})
226 }
227
238 - cmds.EmitOnce(res, &KeyOutputList{list})
228 + return cmds.EmitOnce(res, &KeyOutputList{list})
229 },
230 Encoders: cmds.EncoderMap{
231 cmds.Text: keyOutputListMarshaler(),
core/commands/name/ipns.go
+10 -18
@@ -12,11 +12,11 @@ import (
12 namesys "github.com/ipfs/go-ipfs/namesys"
13 nsopts "github.com/ipfs/go-ipfs/namesys/opts"
14
15 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
15 logging "gx/ipfs/QmRREK2CAZ5Re2Bd9zZFG6FeYDppUWt5cMgsoUEp3ktgSr/go-log"
16 offline "gx/ipfs/QmSNe4MWVxZWk6UxxW2z2EKofFo4GdFzud1vfn1iVby3mj/go-ipfs-routing/offline"
17 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
18 path "gx/ipfs/QmX7uSbkNz76yNwBhuwYwRbhihLnJqM73VTCjS3UMJud9A/go-path"
19 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
20 )
21
22 var log = logging.Logger("core/commands/ipns")
@@ -79,18 +79,16 @@ Resolve the value of a dnslink:
79 cmdkit.UintOption(dhtRecordCountOptionName, "dhtrc", "Number of records to request for DHT resolution."),
80 cmdkit.StringOption(dhtTimeoutOptionName, "dhtt", "Max time to collect values during DHT resolution eg \"30s\". Pass 0 for no timeout."),
81 },
82 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
82 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
83 n, err := cmdenv.GetNode(env)
84 if err != nil {
85 - res.SetError(err, cmdkit.ErrNormal)
86 - return
85 + return err
86 }
87
88 if !n.OnlineMode() {
89 err := n.SetupOfflineRouting()
90 if err != nil {
92 - res.SetError(err, cmdkit.ErrNormal)
93 - return
91 + return err
92 }
93 }
94
@@ -101,8 +99,7 @@ Resolve the value of a dnslink:
99 var resolver namesys.Resolver = n.Namesys
100
101 if local && nocache {
104 - res.SetError(errors.New("cannot specify both local and nocache"), cmdkit.ErrNormal)
105 - return
102 + return errors.New("cannot specify both local and nocache")
103 }
104
105 if local {
@@ -117,8 +114,7 @@ Resolve the value of a dnslink:
114 var name string
115 if len(req.Arguments) == 0 {
116 if n.Identity == "" {
120 - res.SetError(errors.New("identity not loaded"), cmdkit.ErrNormal)
121 - return
117 + return errors.New("identity not loaded")
118 }
119 name = n.Identity.Pretty()
120
@@ -140,12 +136,10 @@ Resolve the value of a dnslink:
136 if dhttok {
137 d, err := time.ParseDuration(dhtt)
138 if err != nil {
143 - res.SetError(err, cmdkit.ErrNormal)
144 - return
139 + return err
140 }
141 if d < 0 {
147 - res.SetError(errors.New("DHT timeout value must be >= 0"), cmdkit.ErrNormal)
148 - return
142 + return errors.New("DHT timeout value must be >= 0")
143 }
144 ropts = append(ropts, nsopts.DhtTimeout(d))
145 }
@@ -156,13 +150,11 @@ Resolve the value of a dnslink:
150
151 output, err := resolver.Resolve(req.Context, name, ropts...)
152 if err != nil {
159 - res.SetError(err, cmdkit.ErrNormal)
160 - return
153 + return err
154 }
155
156 // TODO: better errors (in the case of not finding the name, we get "failed to find any peer in table")
164 -
165 - cmds.EmitOnce(res, &ResolvedPath{output})
157 + return cmds.EmitOnce(res, &ResolvedPath{output})
158 },
159 Encoders: cmds.EncoderMap{
160 cmds.Text: cmds.MakeEncoder(func(req *cmds.Request, w io.Writer, v interface{}) error {
core/commands/name/ipnsps.go
+13 -20
@@ -1,7 +1,6 @@
1 package name
2
3 import (
4 - "errors"
4 "fmt"
5 "io"
6 "strings"
@@ -9,9 +8,9 @@ import (
8 "github.com/ipfs/go-ipfs/core/commands/cmdenv"
9 "github.com/ipfs/go-ipfs/core/commands/e"
10
12 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
11 "gx/ipfs/QmQsErDt8Qgw1XrsXf2BpEzDgGWtB1YLsTAARBup5b6B9W/go-libp2p-peer"
12 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
13 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
14 "gx/ipfs/QmdHb9aBELnQKTVhvvA3hsQbRgUAwsWUzBP2vZ6Y5FBYvE/go-libp2p-record"
15 )
16
@@ -48,14 +47,13 @@ var ipnspsStateCmd = &cmds.Command{
47 Helptext: cmdkit.HelpText{
48 Tagline: "Query the state of IPNS pubsub",
49 },
51 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
50 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
51 n, err := cmdenv.GetNode(env)
52 if err != nil {
54 - res.SetError(err, cmdkit.ErrNormal)
55 - return
53 + return err
54 }
55
58 - cmds.EmitOnce(res, &ipnsPubsubState{n.PSRouter != nil})
56 + return cmds.EmitOnce(res, &ipnsPubsubState{n.PSRouter != nil})
57 },
58 Type: ipnsPubsubState{},
59 Encoders: cmds.EncoderMap{
@@ -82,16 +80,14 @@ var ipnspsSubsCmd = &cmds.Command{
80 Helptext: cmdkit.HelpText{
81 Tagline: "Show current name subscriptions",
82 },
85 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
83 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
84 n, err := cmdenv.GetNode(env)
85 if err != nil {
88 - res.SetError(err, cmdkit.ErrNormal)
89 - return
86 + return err
87 }
88
89 if n.PSRouter == nil {
93 - res.SetError(errors.New("IPNS pubsub subsystem is not enabled"), cmdkit.ErrClient)
94 - return
90 + return cmdkit.Errorf(cmdkit.ErrClient, "IPNS pubsub subsystem is not enabled")
91 }
92 var paths []string
93 for _, key := range n.PSRouter.GetSubscriptions() {
@@ -108,7 +104,7 @@ var ipnspsSubsCmd = &cmds.Command{
104 paths = append(paths, "/ipns/"+peer.IDB58Encode(pid))
105 }
106
111 - cmds.EmitOnce(res, &stringList{paths})
107 + return cmds.EmitOnce(res, &stringList{paths})
108 },
109 Type: stringList{},
110 Encoders: cmds.EncoderMap{
@@ -120,28 +116,25 @@ var ipnspsCancelCmd = &cmds.Command{
116 Helptext: cmdkit.HelpText{
117 Tagline: "Cancel a name subscription",
118 },
123 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
119 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
120 n, err := cmdenv.GetNode(env)
121 if err != nil {
126 - res.SetError(err, cmdkit.ErrNormal)
127 - return
122 + return err
123 }
124
125 if n.PSRouter == nil {
131 - res.SetError(errors.New("IPNS pubsub subsystem is not enabled"), cmdkit.ErrClient)
132 - return
126 + return cmdkit.Errorf(cmdkit.ErrClient, "IPNS pubsub subsystem is not enabled")
127 }
128
129 name := req.Arguments[0]
130 name = strings.TrimPrefix(name, "/ipns/")
131 pid, err := peer.IDB58Decode(name)
132 if err != nil {
139 - res.SetError(err, cmdkit.ErrClient)
140 - return
133 + return cmdkit.Errorf(cmdkit.ErrClient, err.Error())
134 }
135
136 ok := n.PSRouter.Cancel("/ipns/" + string(pid))
144 - cmds.EmitOnce(res, &ipnsPubsubCancel{ok})
137 + return cmds.EmitOnce(res, &ipnsPubsubCancel{ok})
138 },
139 Arguments: []cmdkit.Argument{
140 cmdkit.StringArg("name", true, false, "Name to cancel the subscription for."),
core/commands/name/name.go
+1 -1
@@ -1,8 +1,8 @@
1 package name
2
3 import (
4 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
4 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
5 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
6 )
7
8 type IpnsEntry struct {
core/commands/name/publish.go
+14 -23
@@ -12,11 +12,11 @@ import (
12 e "github.com/ipfs/go-ipfs/core/commands/e"
13 keystore "github.com/ipfs/go-ipfs/keystore"
14
15 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
15 crypto "gx/ipfs/QmPvyPwuCgJ7pDmrKDxRtsScJgBaM5h4EpRL2qQJsmXf4n/go-libp2p-crypto"
16 peer "gx/ipfs/QmQsErDt8Qgw1XrsXf2BpEzDgGWtB1YLsTAARBup5b6B9W/go-libp2p-peer"
17 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
18 path "gx/ipfs/QmX7uSbkNz76yNwBhuwYwRbhihLnJqM73VTCjS3UMJud9A/go-path"
19 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
20 )
21
22 var (
@@ -87,36 +87,31 @@ Alternatively, publish an <ipfs-path> using a valid PeerID (as listed by
87 cmdkit.StringOption(ttlOptionName, "Time duration this record should be cached for (caution: experimental)."),
88 cmdkit.StringOption(keyOptionName, "k", "Name of the key to be used or a valid PeerID, as listed by 'ipfs key list -l'. Default: <<default>>.").WithDefault("self"),
89 },
90 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
90 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
91 n, err := cmdenv.GetNode(env)
92 if err != nil {
93 - res.SetError(err, cmdkit.ErrNormal)
94 - return
93 + return err
94 }
95
96 allowOffline, _ := req.Options[allowOfflineOptionName].(bool)
97 if !n.OnlineMode() {
98 if !allowOffline {
100 - res.SetError(errAllowOffline, cmdkit.ErrNormal)
101 - return
99 + return errAllowOffline
100 }
101 err := n.SetupOfflineRouting()
102 if err != nil {
105 - res.SetError(err, cmdkit.ErrNormal)
106 - return
103 + return err
104 }
105 }
106
107 if n.Mounts.Ipns != nil && n.Mounts.Ipns.IsActive() {
111 - res.SetError(errIpnsMount, cmdkit.ErrNormal)
112 - return
108 + return errIpnsMount
109 }
110
111 pstr := req.Arguments[0]
112
113 if n.Identity == "" {
118 - res.SetError(errIdentityLoad, cmdkit.ErrNormal)
119 - return
114 + return errIdentityLoad
115 }
116
117 popts := new(publishOpts)
@@ -126,8 +121,7 @@ Alternatively, publish an <ipfs-path> using a valid PeerID (as listed by
121 validtime, _ := req.Options[lifeTimeOptionName].(string)
122 d, err := time.ParseDuration(validtime)
123 if err != nil {
129 - res.SetError(fmt.Errorf("error parsing lifetime option: %s", err), cmdkit.ErrNormal)
130 - return
124 + return fmt.Errorf("error parsing lifetime option: %s", err)
125 }
126
127 popts.pubValidTime = d
@@ -136,8 +130,7 @@ Alternatively, publish an <ipfs-path> using a valid PeerID (as listed by
130 if ttl, found := req.Options[ttlOptionName].(string); found {
131 d, err := time.ParseDuration(ttl)
132 if err != nil {
139 - res.SetError(err, cmdkit.ErrNormal)
140 - return
133 + return err
134 }
135
136 ctx = context.WithValue(ctx, "ipns-publish-ttl", d)
@@ -146,22 +139,20 @@ Alternatively, publish an <ipfs-path> using a valid PeerID (as listed by
139 kname, _ := req.Options[keyOptionName].(string)
140 k, err := keylookup(n, kname)
141 if err != nil {
149 - res.SetError(err, cmdkit.ErrNormal)
150 - return
142 + return err
143 }
144
145 pth, err := path.ParsePath(pstr)
146 if err != nil {
155 - res.SetError(err, cmdkit.ErrNormal)
156 - return
147 + return err
148 }
149
150 output, err := publish(ctx, n, k, pth, popts)
151 if err != nil {
161 - res.SetError(err, cmdkit.ErrNormal)
162 - return
152 + return err
153 }
164 - cmds.EmitOnce(res, output)
154 +
155 + return cmds.EmitOnce(res, output)
156 },
157 Encoders: cmds.EncoderMap{
158 cmds.Text: cmds.MakeEncoder(func(req *cmds.Request, w io.Writer, v interface{}) error {
core/commands/object/object.go
+2 -2
@@ -17,9 +17,9 @@ import (
17 "github.com/ipfs/go-ipfs/core/coreapi/interface/options"
18
19 cid "gx/ipfs/QmPSQnBKM9g7BaUcZCvswUJVscQ1ipjmwxN5PXCjkp9EQ7/go-cid"
20 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
20 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
21 dag "gx/ipfs/QmXv5mwmQ74r4aiHcNeQ4GAmfB3aWJuqaE4WyDfDfvkgLM/go-merkledag"
22 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
23 ipld "gx/ipfs/QmdDXJs4axxefSPgK6Y1QhpJWKuDPnGJiqgq4uncb4rFHL/go-ipld-format"
24 )
25
@@ -192,7 +192,7 @@ multihash.
192 return buf, nil
193 },
194 },
195 - Type: Object{},
195 + Type: &Object{},
196 }
197
198 var ObjectGetCmd = &oldcmds.Command{
core/commands/object/patch.go
+8 -12
@@ -12,8 +12,8 @@ import (
12 coreiface "github.com/ipfs/go-ipfs/core/coreapi/interface"
13 "github.com/ipfs/go-ipfs/core/coreapi/interface/options"
14
15 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
15 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
16 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
17 )
18
19 var ObjectPatchCmd = &cmds.Command{
@@ -67,34 +67,30 @@ the limit will not be respected by the network.
67 cmdkit.StringArg("root", true, false, "The hash of the node to modify."),
68 cmdkit.FileArg("data", true, false, "Data to append.").EnableStdin(),
69 },
70 - Run: func(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) {
70 + Run: func(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) error {
71 api, err := cmdenv.GetApi(env)
72 if err != nil {
73 - re.SetError(err, cmdkit.ErrNormal)
74 - return
73 + return err
74 }
75
76 root, err := coreiface.ParsePath(req.Arguments[0])
77 if err != nil {
79 - re.SetError(err, cmdkit.ErrNormal)
80 - return
78 + return err
79 }
80
81 data, err := req.Files.NextFile()
82 if err != nil {
85 - re.SetError(err, cmdkit.ErrNormal)
86 - return
83 + return err
84 }
85
86 p, err := api.Object().AppendData(req.Context, root, data)
87 if err != nil {
91 - re.SetError(err, cmdkit.ErrNormal)
92 - return
88 + return err
89 }
90
95 - cmds.EmitOnce(re, &Object{Hash: p.Cid().String()})
91 + return cmds.EmitOnce(re, &Object{Hash: p.Cid().String()})
92 },
97 - Type: Object{},
93 + Type: &Object{},
94 Encoders: cmds.EncoderMap{
95 cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, obj *Object) error {
96 _, err := fmt.Fprintln(w, obj.Hash)
core/commands/pubsub.go
+31 -41
@@ -3,6 +3,7 @@ package commands
3 import (
4 "context"
5 "encoding/binary"
6 + "errors"
7 "fmt"
8 "io"
9 "net/http"
@@ -15,10 +16,10 @@ import (
16 e "github.com/ipfs/go-ipfs/core/commands/e"
17
18 cid "gx/ipfs/QmPSQnBKM9g7BaUcZCvswUJVscQ1ipjmwxN5PXCjkp9EQ7/go-cid"
18 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
19 blocks "gx/ipfs/QmRcHuYzAyswytBuMF78rj3LTChYszomRFXNg4685ZN1WM/go-block-format"
20 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
21 floodsub "gx/ipfs/QmY1L5krVk8dv8d74uESmJTXGpoigVYqBVxXXz1aS8aFSb/go-libp2p-floodsub"
22 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
23 pstore "gx/ipfs/Qmda4cPRvSRyox3SqgJN6DfSZGU5TtHufPTp9uXjFj71X6/go-libp2p-peerstore"
24 )
25
@@ -73,29 +74,25 @@ This command outputs data in the following encodings:
74 Options: []cmdkit.Option{
75 cmdkit.BoolOption("discover", "try to discover other peers subscribed to the same topic"),
76 },
76 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
77 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
78 n, err := cmdenv.GetNode(env)
79 if err != nil {
79 - res.SetError(err, cmdkit.ErrNormal)
80 - return
80 + return err
81 }
82
83 // Must be online!
84 if !n.OnlineMode() {
85 - res.SetError(ErrNotOnline, cmdkit.ErrClient)
86 - return
85 + return cmdkit.Errorf(cmdkit.ErrClient, ErrNotOnline.Error())
86 }
87
88 if n.Floodsub == nil {
90 - res.SetError(fmt.Errorf("experimental pubsub feature not enabled. Run daemon with --enable-pubsub-experiment to use"), cmdkit.ErrNormal)
91 - return
89 + return fmt.Errorf("experimental pubsub feature not enabled. Run daemon with --enable-pubsub-experiment to use")
90 }
91
92 topic := req.Arguments[0]
93 sub, err := n.Floodsub.Subscribe(topic)
94 if err != nil {
97 - res.SetError(err, cmdkit.ErrNormal)
98 - return
95 + return err
96 }
97 defer sub.Cancel()
98
@@ -120,13 +117,15 @@ This command outputs data in the following encodings:
117 for {
118 msg, err := sub.Next(req.Context)
119 if err == io.EOF || err == context.Canceled {
123 - return
120 + return nil
121 } else if err != nil {
125 - res.SetError(err, cmdkit.ErrNormal)
126 - return
122 + return err
123 }
124
129 - res.Emit(msg)
125 + err = res.Emit(msg)
126 + if err != nil {
127 + return err
128 + }
129 }
130 },
131 Encoders: cmds.EncoderMap{
@@ -206,38 +205,35 @@ To use, the daemon must be run with '--enable-pubsub-experiment'.
205 cmdkit.StringArg("topic", true, false, "Topic to publish to."),
206 cmdkit.StringArg("data", true, true, "Payload of message to publish.").EnableStdin(),
207 },
209 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
208 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
209 n, err := cmdenv.GetNode(env)
210 if err != nil {
212 - res.SetError(err, cmdkit.ErrNormal)
213 - return
211 + return err
212 }
213
214 // Must be online!
215 if !n.OnlineMode() {
218 - res.SetError(ErrNotOnline, cmdkit.ErrClient)
219 - return
216 + return cmdkit.Errorf(cmdkit.ErrClient, ErrNotOnline.Error())
217 }
218
219 if n.Floodsub == nil {
223 - res.SetError("experimental pubsub feature not enabled. Run daemon with --enable-pubsub-experiment to use.", cmdkit.ErrNormal)
224 - return
220 + return errors.New("experimental pubsub feature not enabled. Run daemon with --enable-pubsub-experiment to use.")
221 }
222
223 topic := req.Arguments[0]
224
225 err = req.ParseBodyArgs()
226 if err != nil {
231 - res.SetError(err, cmdkit.ErrNormal)
232 - return
227 + return err
228 }
229
230 for _, data := range req.Arguments[1:] {
231 if err := n.Floodsub.Publish(topic, []byte(data)); err != nil {
237 - res.SetError(err, cmdkit.ErrNormal)
238 - return
232 + return err
233 }
234 }
235 +
236 + return nil
237 },
238 }
239
@@ -253,25 +249,22 @@ to be used in a production environment.
249 To use, the daemon must be run with '--enable-pubsub-experiment'.
250 `,
251 },
256 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
252 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
253 n, err := cmdenv.GetNode(env)
254 if err != nil {
259 - res.SetError(err, cmdkit.ErrNormal)
260 - return
255 + return err
256 }
257
258 // Must be online!
259 if !n.OnlineMode() {
265 - res.SetError(ErrNotOnline, cmdkit.ErrClient)
266 - return
260 + return cmdkit.Errorf(cmdkit.ErrClient, ErrNotOnline.Error())
261 }
262
263 if n.Floodsub == nil {
270 - res.SetError("experimental pubsub feature not enabled. Run daemon with --enable-pubsub-experiment to use.", cmdkit.ErrNormal)
271 - return
264 + return errors.New("experimental pubsub feature not enabled. Run daemon with --enable-pubsub-experiment to use.")
265 }
266
274 - cmds.EmitOnce(res, stringList{n.Floodsub.GetTopics()})
267 + return cmds.EmitOnce(res, stringList{n.Floodsub.GetTopics()})
268 },
269 Type: stringList{},
270 Encoders: cmds.EncoderMap{
@@ -310,22 +303,19 @@ To use, the daemon must be run with '--enable-pubsub-experiment'.
303 Arguments: []cmdkit.Argument{
304 cmdkit.StringArg("topic", false, false, "topic to list connected peers of"),
305 },
313 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
306 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
307 n, err := cmdenv.GetNode(env)
308 if err != nil {
316 - res.SetError(err, cmdkit.ErrNormal)
317 - return
309 + return err
310 }
311
312 // Must be online!
313 if !n.OnlineMode() {
322 - res.SetError(ErrNotOnline, cmdkit.ErrClient)
323 - return
314 + return cmdkit.Errorf(cmdkit.ErrClient, ErrNotOnline.Error())
315 }
316
317 if n.Floodsub == nil {
327 - res.SetError(fmt.Errorf("experimental pubsub feature not enabled. Run daemon with --enable-pubsub-experiment to use"), cmdkit.ErrNormal)
328 - return
318 + return errors.New("experimental pubsub feature not enabled. Run daemon with --enable-pubsub-experiment to use")
319 }
320
321 var topic string
@@ -340,7 +330,7 @@ To use, the daemon must be run with '--enable-pubsub-experiment'.
330 list.Strings = append(list.Strings, peer.Pretty())
331 }
332 sort.Strings(list.Strings)
343 - cmds.EmitOnce(res, list)
333 + return cmds.EmitOnce(res, list)
334 },
335 Type: stringList{},
336 Encoders: cmds.EncoderMap{
core/commands/repo.go
+8 -11
@@ -17,9 +17,9 @@ import (
17 fsrepo "github.com/ipfs/go-ipfs/repo/fsrepo"
18
19 cid "gx/ipfs/QmPSQnBKM9g7BaUcZCvswUJVscQ1ipjmwxN5PXCjkp9EQ7/go-cid"
20 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
20 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
21 config "gx/ipfs/QmYVqYJTVjetcf1guieEgWpK1PZtHPytP624vKzTF1P3r2/go-ipfs-config"
22 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
23 bstore "gx/ipfs/QmegPGspn3RpTMQ23Fd3GVVMopo1zsEMurudbFMZ5UXBLH/go-ipfs-blockstore"
24 )
25
@@ -99,7 +99,7 @@ order to reclaim hard disk space.
99 }
100 }
101 if errs {
102 - res.SetError(fmt.Errorf("encountered errors during gc run"), cmdkit.ErrNormal)
102 + outChan <- &GcResult{Error: "encountered errors during gc run"}
103 }
104 } else {
105 err := corerepo.CollectResult(req.Context(), gcOutChan, func(k cid.Cid) {
@@ -165,33 +165,30 @@ Version string The repo version.
165 cmdkit.BoolOption("size-only", "Only report RepoSize and StorageMax."),
166 cmdkit.BoolOption("human", "Output sizes in MiB."),
167 },
168 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
168 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
169 n, err := cmdenv.GetNode(env)
170 if err != nil {
171 - res.SetError(err, cmdkit.ErrNormal)
172 - return
171 + return err
172 }
173
174 sizeOnly, _ := req.Options["size-only"].(bool)
175 if sizeOnly {
176 sizeStat, err := corerepo.RepoSize(req.Context, n)
177 if err != nil {
179 - res.SetError(err, cmdkit.ErrNormal)
180 - return
178 + return err
179 }
180 cmds.EmitOnce(res, &corerepo.Stat{
181 SizeStat: sizeStat,
182 })
185 - return
183 + return nil
184 }
185
186 stat, err := corerepo.RepoStat(req.Context, n)
187 if err != nil {
190 - res.SetError(err, cmdkit.ErrNormal)
191 - return
188 + return err
189 }
190
194 - cmds.EmitOnce(res, &stat)
191 + return cmds.EmitOnce(res, &stat)
192 },
193 Type: &corerepo.Stat{},
194 Encoders: cmds.EncoderMap{
core/commands/root.go
+1 -1
@@ -13,9 +13,9 @@ import (
13 ocmd "github.com/ipfs/go-ipfs/core/commands/object"
14 unixfs "github.com/ipfs/go-ipfs/core/commands/unixfs"
15
16 - "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
16 logging "gx/ipfs/QmRREK2CAZ5Re2Bd9zZFG6FeYDppUWt5cMgsoUEp3ktgSr/go-log"
17 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
18 + "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
19 )
20
21 var log = logging.Logger("core/commands")
core/commands/shutdown.go
+6 -8
@@ -1,32 +1,30 @@
1 package commands
2
3 import (
4 - "fmt"
5 -
4 cmdenv "github.com/ipfs/go-ipfs/core/commands/cmdenv"
5
8 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
6 "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
7 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
8 )
9
10 var daemonShutdownCmd = &cmds.Command{
11 Helptext: cmdkit.HelpText{
12 Tagline: "Shut down the ipfs daemon",
13 },
16 - Run: func(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) {
14 + Run: func(req *cmds.Request, re cmds.ResponseEmitter, env cmds.Environment) error {
15 nd, err := cmdenv.GetNode(env)
16 if err != nil {
19 - re.SetError(err, cmdkit.ErrNormal)
20 - return
17 + return err
18 }
19
20 if nd.LocalMode() {
24 - re.SetError(fmt.Errorf("daemon not running"), cmdkit.ErrClient)
25 - return
21 + return cmdkit.Errorf(cmdkit.ErrClient, "daemon not running")
22 }
23
24 if err := nd.Process().Close(); err != nil {
25 log.Error("error while shutting down ipfs daemon:", err)
26 }
27 +
28 + return nil
29 },
30 }
core/commands/stat.go
+32 -44
@@ -1,7 +1,6 @@
1 package commands
2
3 import (
4 - "errors"
4 "fmt"
5 "io"
6 "os"
@@ -10,10 +9,10 @@ import (
9 cmdenv "github.com/ipfs/go-ipfs/core/commands/cmdenv"
10
11 humanize "gx/ipfs/QmPSBJL4momYnE7DcUyk2DVhD6rH488ZmHBGLbxNdhU44K/go-humanize"
13 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
12 peer "gx/ipfs/QmQsErDt8Qgw1XrsXf2BpEzDgGWtB1YLsTAARBup5b6B9W/go-libp2p-peer"
13 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
14 protocol "gx/ipfs/QmZNkThpqfVXs9GNbexPrfBbXSLNYeKrE7jwFM2oqHbyqN/go-libp2p-protocol"
15 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
16 metrics "gx/ipfs/QmdhwKw53CTV8EJSAsR1bpmMT5kXiWBgeAyv1EXeeDiXqR/go-libp2p-metrics"
17 )
18
@@ -81,37 +80,32 @@ Example:
80 "ns", "us" (or "µs"), "ms", "s", "m", "h".`).WithDefault("1s"),
81 },
82
84 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
83 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
84 nd, err := cmdenv.GetNode(env)
85 if err != nil {
87 - res.SetError(err, cmdkit.ErrNormal)
88 - return
86 + return err
87 }
88
89 // Must be online!
90 if !nd.OnlineMode() {
93 - res.SetError(ErrNotOnline, cmdkit.ErrClient)
94 - return
91 + return cmdkit.Errorf(cmdkit.ErrClient, ErrNotOnline.Error())
92 }
93
94 if nd.Reporter == nil {
98 - res.SetError(fmt.Errorf("bandwidth reporter disabled in config"), cmdkit.ErrNormal)
99 - return
95 + return fmt.Errorf("bandwidth reporter disabled in config")
96 }
97
98 pstr, pfound := req.Options["peer"].(string)
99 tstr, tfound := req.Options["proto"].(string)
100 if pfound && tfound {
105 - res.SetError(errors.New("please only specify peer OR protocol"), cmdkit.ErrClient)
106 - return
101 + return cmdkit.Errorf(cmdkit.ErrClient, "please only specify peer OR protocol")
102 }
103
104 var pid peer.ID
105 if pfound {
106 checkpid, err := peer.IDB58Decode(pstr)
107 if err != nil {
113 - res.SetError(err, cmdkit.ErrNormal)
114 - return
108 + return err
109 }
110 pid = checkpid
111 }
@@ -119,8 +113,7 @@ Example:
113 timeS, _ := req.Options["interval"].(string)
114 interval, err := time.ParseDuration(timeS)
115 if err != nil {
122 - res.SetError(err, cmdkit.ErrNormal)
123 - return
116 + return err
117 }
118
119 doPoll, _ := req.Options["poll"].(bool)
@@ -137,49 +130,44 @@ Example:
130 res.Emit(&totals)
131 }
132 if !doPoll {
140 - return
133 + return nil
134 }
135 select {
136 case <-time.After(interval):
137 case <-req.Context.Done():
145 - return
138 + break
139 }
140 }
148 -
141 },
142 Type: metrics.Stats{},
143 PostRun: cmds.PostRunMap{
152 - cmds.CLI: func(req *cmds.Request, re cmds.ResponseEmitter) cmds.ResponseEmitter {
153 - reNext, res := cmds.NewChanResponsePair(req)
154 -
155 - go func() {
156 - defer re.Close()
157 -
158 - polling, _ := res.Request().Options["poll"].(bool)
159 - if polling {
160 - fmt.Fprintln(os.Stdout, "Total Up Total Down Rate Up Rate Down")
161 - }
162 - for {
163 - v, err := res.Next()
164 - if !cmds.HandleError(err, res, re) {
165 - break
144 + cmds.CLI: func(res cmds.Response, re cmds.ResponseEmitter) error {
145 + polling, _ := res.Request().Options["poll"].(bool)
146 + log.Debug("postrun polling:", polling)
147 + if polling {
148 + fmt.Fprintln(os.Stdout, "Total Up Total Down Rate Up Rate Down")
149 + }
150 + for {
151 + v, err := res.Next()
152 + if err != nil {
153 + if err == io.EOF {
154 + return nil
155 }
156 + return err
157 + }
158
168 - bs := v.(*metrics.Stats)
169 -
170 - if !polling {
171 - printStats(os.Stdout, bs)
172 - return
173 - }
159 + bs := v.(*metrics.Stats)
160
175 - fmt.Fprintf(os.Stdout, "%8s ", humanize.Bytes(uint64(bs.TotalOut)))
176 - fmt.Fprintf(os.Stdout, "%8s ", humanize.Bytes(uint64(bs.TotalIn)))
177 - fmt.Fprintf(os.Stdout, "%8s/s ", humanize.Bytes(uint64(bs.RateOut)))
178 - fmt.Fprintf(os.Stdout, "%8s/s \r", humanize.Bytes(uint64(bs.RateIn)))
161 + if !polling {
162 + printStats(os.Stdout, bs)
163 + return nil
164 }
180 - }()
165
182 - return reNext
166 + fmt.Fprintf(os.Stdout, "%8s ", humanize.Bytes(uint64(bs.TotalOut)))
167 + fmt.Fprintf(os.Stdout, "%8s ", humanize.Bytes(uint64(bs.TotalIn)))
168 + fmt.Fprintf(os.Stdout, "%8s/s ", humanize.Bytes(uint64(bs.RateOut)))
169 + fmt.Fprintf(os.Stdout, "%8s/s \r", humanize.Bytes(uint64(bs.RateIn)))
170 + }
171 },
172 },
173 }
core/commands/urlstore.go
+13 -20
@@ -12,9 +12,9 @@ import (
12 ihelper "gx/ipfs/QmPL8bYtbACcSFFiSr4s2du7Na382NxRADR8hC7D9FkEA2/go-unixfs/importer/helpers"
13 trickle "gx/ipfs/QmPL8bYtbACcSFFiSr4s2du7Na382NxRADR8hC7D9FkEA2/go-unixfs/importer/trickle"
14 cid "gx/ipfs/QmPSQnBKM9g7BaUcZCvswUJVscQ1ipjmwxN5PXCjkp9EQ7/go-cid"
15 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
15 mh "gx/ipfs/QmPnFwZ2JXKnXgMw8CdBPxn7FWh6LLdjUjxV1fKHuJnkr8/go-multihash"
16 cmdkit "gx/ipfs/QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky/go-ipfs-cmdkit"
17 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
18 chunk "gx/ipfs/QmdSeG9s4EQ9TGruJJS9Us38TQDZtMmFGwzTYUDVqNTURm/go-ipfs-chunker"
19 )
20
@@ -50,48 +50,41 @@ time.
50 Arguments: []cmdkit.Argument{
51 cmdkit.StringArg("url", true, false, "URL to add to IPFS"),
52 },
53 - Type: BlockStat{},
53 + Type: &BlockStat{},
54
55 - Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) {
55 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
56 url := req.Arguments[0]
57 n, err := cmdenv.GetNode(env)
58 if err != nil {
59 - res.SetError(err, cmdkit.ErrNormal)
60 - return
59 + return err
60 }
61
62 if !filestore.IsURL(url) {
64 - res.SetError(fmt.Errorf("unsupported url syntax: %s", url), cmdkit.ErrNormal)
65 - return
63 + return fmt.Errorf("unsupported url syntax: %s", url)
64 }
65
66 cfg, err := n.Repo.Config()
67 if err != nil {
70 - res.SetError(err, cmdkit.ErrNormal)
71 - return
68 + return err
69 }
70
71 if !cfg.Experimental.UrlstoreEnabled {
75 - res.SetError(filestore.ErrUrlstoreNotEnabled, cmdkit.ErrNormal)
76 - return
72 + return filestore.ErrUrlstoreNotEnabled
73 }
74
75 useTrickledag, _ := req.Options[trickleOptionName].(bool)
76
77 hreq, err := http.NewRequest("GET", url, nil)
78 if err != nil {
83 - res.SetError(err, cmdkit.ErrNormal)
84 - return
79 + return err
80 }
81
82 hres, err := http.DefaultClient.Do(hreq)
83 if err != nil {
89 - res.SetError(err, cmdkit.ErrNormal)
90 - return
84 + return err
85 }
86 if hres.StatusCode != http.StatusOK {
93 - res.SetError(fmt.Errorf("expected code 200, got: %d", hres.StatusCode), cmdkit.ErrNormal)
94 - return
87 + return fmt.Errorf("expected code 200, got: %d", hres.StatusCode)
88 }
89
90 chk := chunk.NewSizeSplitter(hres.Body, chunk.DefaultBlockSize)
@@ -111,14 +104,14 @@ time.
104 }
105 root, err := layout(dbp.New(chk))
106 if err != nil {
114 - res.SetError(err, cmdkit.ErrNormal)
115 - return
107 + return err
108 }
109
118 - cmds.EmitOnce(res, BlockStat{
110 + err = cmds.EmitOnce(res, &BlockStat{
111 Key: root.Cid().String(),
112 Size: int(hres.ContentLength),
113 })
114 + return err
115 },
116 Encoders: cmds.EncoderMap{
117 cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, bs *BlockStat) error {
core/corehttp/commands.go
+2 -2
@@ -14,10 +14,10 @@ import (
14 "github.com/ipfs/go-ipfs/core"
15 corecommands "github.com/ipfs/go-ipfs/core/commands"
16
17 - cmds "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds"
18 - cmdsHttp "gx/ipfs/QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi/go-ipfs-cmds/http"
17 path "gx/ipfs/QmX7uSbkNz76yNwBhuwYwRbhihLnJqM73VTCjS3UMJud9A/go-path"
18 config "gx/ipfs/QmYVqYJTVjetcf1guieEgWpK1PZtHPytP624vKzTF1P3r2/go-ipfs-config"
19 + cmds "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds"
20 + cmdsHttp "gx/ipfs/QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE/go-ipfs-cmds/http"
21 )
22
23 var (
package.json
+2 -2
@@ -300,9 +300,9 @@
300 "version": "3.0.11"
301 },
302 {
303 - "hash": "QmPTfgFTo9PFr1PvPKyKoeMgBvYPh6cX3aDP7DHKVbnCbi",
303 + "hash": "QmZVPuwGNz2s9THwLS4psrJGam6NSEQMvDTaaZgNfqQBCE",
304 "name": "go-ipfs-cmds",
305 - "version": "1.0.22"
305 + "version": "2.0.0-beta2"
306 },
307 {
308 "hash": "QmSP88ryZkHSRn1fnngAaV2Vcn63WUJzAavnRM9CVdU1Ky",
test/sharness/lib/test-lib.sh
+1
@@ -26,6 +26,7 @@ fi
26 # it's too late to pass in --verbose, and --verbose is harder
27 # to pass through in some cases.
28 test "$TEST_VERBOSE" = 1 && verbose=t
29 +test "$TEST_IMMEDIATE" = 1 && immediate=t
30 # source the common hashes first.
31 . lib/test-lib-hashes.sh
32