@cryptotaxi247 / kubo / commits / 2462bf913

commands/filestore: use res.Emit directly

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

Overbool committed Oct 28, 2018 at 11:30 UTC 2462bf913af34ccf05655ad1d691450f65922615
1 file changed +52 -76
core/commands/filestore.go
+52 -76
@@ -1,7 +1,6 @@
1 package commands
2
3 import (
4 - "context"
4 "fmt"
5 "io"
6
@@ -56,11 +55,7 @@ The output is:
55 }
56 args := req.Arguments
57 if len(args) > 0 {
59 - out := perKeyActionToChan(req.Context, args, func(c cid.Cid) *filestore.ListRes {
60 - return filestore.List(fs, c)
61 - })
62 -
63 - return res.Emit(out)
58 + return listByArgs(res, fs, args)
59 }
60
61 fileOrder, _ := req.Options[fileOrderOptionName].(bool)
@@ -69,8 +64,17 @@ The output is:
64 return err
65 }
66
72 - out := listResToChan(req.Context, next)
73 - return res.Emit(out)
67 + for {
68 + r := next()
69 + if r == nil {
70 + break
71 + }
72 + if err := res.Emit(r); err != nil {
73 + return err
74 + }
75 + }
76 +
77 + return nil
78 },
79 PostRun: cmds.PostRunMap{
80 cmds.CLI: streamResult(func(v interface{}, out io.Writer) nonFatalError {
@@ -122,10 +126,7 @@ For ERROR entries the error will also be printed to stderr.
126 }
127 args := req.Arguments
128 if len(args) > 0 {
125 - out := perKeyActionToChan(req.Context, args, func(c cid.Cid) *filestore.ListRes {
126 - return filestore.Verify(fs, c)
127 - })
128 - return res.Emit(out)
129 + return listByArgs(res, fs, args)
130 }
131
132 fileOrder, _ := req.Options[fileOrderOptionName].(bool)
@@ -133,8 +134,18 @@ For ERROR entries the error will also be printed to stderr.
134 if err != nil {
135 return err
136 }
136 - out := listResToChan(req.Context, next)
137 - return res.Emit(out)
137 +
138 + for {
139 + r := next()
140 + if r == nil {
141 + break
142 + }
143 + if err := res.Emit(r); err != nil {
144 + return err
145 + }
146 + }
147 +
148 + return nil
149 },
150 Encoders: cmds.EncoderMap{
151 cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *filestore.ListRes) error {
@@ -162,29 +173,19 @@ var dupsFileStore = &cmds.Command{
173 return err
174 }
175
165 - out := make(chan interface{}, 128)
166 -
167 - go func() {
168 - defer close(out)
169 - for cid := range ch {
170 - have, err := fs.MainBlockstore().Has(cid)
171 - if err != nil {
172 - select {
173 - case out <- &RefWrapper{Err: err.Error()}:
174 - case <-req.Context.Done():
175 - }
176 - return
177 - }
178 - if have {
179 - select {
180 - case out <- &RefWrapper{Ref: cid.String()}:
181 - case <-req.Context.Done():
182 - return
183 - }
176 + for cid := range ch {
177 + have, err := fs.MainBlockstore().Has(cid)
178 + if err != nil {
179 + return res.Emit(&RefWrapper{Err: err.Error()})
180 + }
181 + if have {
182 + if err := res.Emit(&RefWrapper{Ref: cid.String()}); err != nil {
183 + return err
184 }
185 }
186 - }()
187 - return res.Emit(out)
186 + }
187 +
188 + return nil
189 },
190 Encoders: cmds.EncoderMap{
191 cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *RefWrapper) error {
@@ -212,49 +213,24 @@ func getFilestore(env cmds.Environment) (*core.IpfsNode, *filestore.Filestore, e
213 return n, fs, err
214 }
215
215 -func listResToChan(ctx context.Context, next func() *filestore.ListRes) <-chan interface{} {
216 - out := make(chan interface{}, 128)
217 - go func() {
218 - defer close(out)
219 - for {
220 - r := next()
221 - if r == nil {
222 - return
216 +func listByArgs(res cmds.ResponseEmitter, fs *filestore.Filestore, args []string) error {
217 + for _, arg := range args {
218 + c, err := cid.Decode(arg)
219 + if err != nil {
220 + ret := &filestore.ListRes{
221 + Status: filestore.StatusOtherError,
222 + ErrorMsg: fmt.Sprintf("%s: %v", arg, err),
223 }
224 - select {
225 - case out <- r:
226 - case <-ctx.Done():
227 - return
224 + if err := res.Emit(ret); err != nil {
225 + return err
226 }
227 + continue
228 }
230 - }()
231 - return out
232 -}
233 -
234 -func perKeyActionToChan(ctx context.Context, args []string, action func(cid.Cid) *filestore.ListRes) <-chan interface{} {
235 - out := make(chan interface{}, 128)
236 - go func() {
237 - defer close(out)
238 - for _, arg := range args {
239 - c, err := cid.Decode(arg)
240 - if err != nil {
241 - select {
242 - case out <- &filestore.ListRes{
243 - Status: filestore.StatusOtherError,
244 - ErrorMsg: fmt.Sprintf("%s: %v", arg, err),
245 - }:
246 - case <-ctx.Done():
247 - }
248 -
249 - continue
250 - }
251 - r := action(c)
252 - select {
253 - case out <- r:
254 - case <-ctx.Done():
255 - return
256 - }
229 + r := filestore.Verify(fs, c)
230 + if err := res.Emit(r); err != nil {
231 + return err
232 }
258 - }()
259 - return out
233 + }
234 +
235 + return nil
236 }