@cryptotaxi247 / kubo / commits / 62ce03ed1

cmds/pin: use PostRun in pin add

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

Overbool committed Nov 7, 2018 at 13:33 UTC 62ce03ed14d523199ccbc2730765f6eb03f95ae0
1 file changed +65 -66
core/commands/pin.go
+65 -66
@@ -4,10 +4,12 @@ import (
4 "context"
5 "fmt"
6 "io"
7 + "os"
8 "time"
9
10 core "github.com/ipfs/go-ipfs/core"
11 cmdenv "github.com/ipfs/go-ipfs/core/commands/cmdenv"
12 + e "github.com/ipfs/go-ipfs/core/commands/e"
13 iface "github.com/ipfs/go-ipfs/core/coreapi/interface"
14 options "github.com/ipfs/go-ipfs/core/coreapi/interface/options"
15 corerepo "github.com/ipfs/go-ipfs/core/corerepo"
@@ -19,7 +21,7 @@ import (
21 "gx/ipfs/QmYMQuypUbgsdNHmuCBSUJV6wdQVsBHRivNAp3efHJwZJD/go-verifcid"
22 cmds "gx/ipfs/Qma6uuSyjkecGhMFFLfzyJDPyoDtNJSHJNweDccZhaWkgU/go-ipfs-cmds"
23 dag "gx/ipfs/QmaDBne4KeY3UepeqSVKYpSmQGa3q9zP6x3LfVF2UjF3Hc/go-merkledag"
22 - "gx/ipfs/Qmde5VP1qUkyQXKCfmEUA7bP64V2HAptbJ7phuPp7jXWwg/go-ipfs-cmdkit"
24 + cmdkit "gx/ipfs/Qmde5VP1qUkyQXKCfmEUA7bP64V2HAptbJ7phuPp7jXWwg/go-ipfs-cmdkit"
25 )
26
27 var PinCmd = &cmds.Command{
@@ -65,11 +67,6 @@ var addPinCmd = &cmds.Command{
67 },
68 Type: AddPinOutput{},
69 Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
68 - err := req.ParseBodyArgs()
69 - if err != nil {
70 - return err
71 - }
72 -
70 n, err := cmdenv.GetNode(env)
71 if err != nil {
72 return err
@@ -86,6 +83,10 @@ var addPinCmd = &cmds.Command{
83 recursive, _ := req.Options[pinRecursiveOptionName].(bool)
84 showProgress, _ := req.Options[pinProgressOptionName].(bool)
85
86 + if err := req.ParseBodyArgs(); err != nil {
87 + return err
88 + }
89 +
90 if !showProgress {
91 added, err := corerepo.Pin(n, api, req.Context, req.Arguments, recursive)
92 if err != nil {
@@ -94,8 +95,6 @@ var addPinCmd = &cmds.Command{
95 return cmds.EmitOnce(res, &AddPinOutput{Pins: cidsToStrings(added)})
96 }
97
97 - out := make(chan interface{})
98 -
98 v := new(dag.ProgressTracker)
99 ctx := v.DeriveContext(req.Context)
100
@@ -103,75 +102,76 @@ var addPinCmd = &cmds.Command{
102 pins []cid.Cid
103 err error
104 }
105 +
106 ch := make(chan pinResult, 1)
107 go func() {
108 added, err := corerepo.Pin(n, api, ctx, req.Arguments, recursive)
109 ch <- pinResult{pins: added, err: err}
110 }()
111
112 - errC := make(chan error)
113 - go func() {
114 - var err error
115 - ticker := time.NewTicker(500 * time.Millisecond)
116 - defer ticker.Stop()
117 - defer func() { errC <- err }()
118 - defer close(out)
112 + ticker := time.NewTicker(500 * time.Millisecond)
113 + defer ticker.Stop()
114
120 - for {
121 - select {
122 - case val := <-ch:
123 - if val.err != nil {
124 - err = val.err
125 - return
126 - }
115 + for {
116 + select {
117 + case val := <-ch:
118 + if val.err != nil {
119 + return val.err
120 + }
121
128 - if pv := v.Value(); pv != 0 {
129 - out <- &AddPinOutput{Progress: v.Value()}
122 + if pv := v.Value(); pv != 0 {
123 + if err := res.Emit(&AddPinOutput{Progress: v.Value()}); err != nil {
124 + return err
125 }
131 - out <- &AddPinOutput{Pins: cidsToStrings(val.pins)}
132 - return
133 - case <-ticker.C:
134 - out <- &AddPinOutput{Progress: v.Value()}
135 - case <-ctx.Done():
136 - log.Error(ctx.Err())
137 - err = ctx.Err()
138 - return
126 }
127 + return res.Emit(&AddPinOutput{Pins: cidsToStrings(val.pins)})
128 + case <-ticker.C:
129 + if err := res.Emit(&AddPinOutput{Progress: v.Value()}); err != nil {
130 + return err
131 + }
132 + case <-ctx.Done():
133 + log.Error(ctx.Err())
134 + return ctx.Err()
135 }
141 - }()
142 -
143 - err = res.Emit(out)
144 - if err != nil {
145 - return err
136 }
147 -
148 - return <-errC
137 },
150 - Encoders: cmds.EncoderMap{
151 - cmds.Text: cmds.MakeTypedEncoder(func(req *cmds.Request, w io.Writer, out *AddPinOutput) error {
152 - var added []string
153 -
154 - if out.Pins != nil {
155 - added = out.Pins
156 - } else {
157 - // this can only happen if the progress option is set
158 - return fmt.Errorf("Fetched/Processed %d nodes\r", out.Progress)
159 - }
138 + PostRun: cmds.PostRunMap{
139 + cmds.CLI: func(res cmds.Response, re cmds.ResponseEmitter) error {
140 + for {
141 + v, err := res.Next()
142 + if err != nil {
143 + if err == io.EOF {
144 + return nil
145 + }
146 + return err
147 + }
148
161 - var pintype string
162 - rec, found := req.Options["recursive"].(bool)
163 - if rec || !found {
164 - pintype = "recursively"
165 - } else {
166 - pintype = "directly"
167 - }
149 + out, ok := v.(*AddPinOutput)
150 + if !ok {
151 + return e.TypeErr(out, v)
152 + }
153 + var added []string
154
169 - for _, k := range added {
170 - fmt.Fprintf(w, "pinned %s %s\n", k, pintype)
171 - }
155 + if out.Pins != nil {
156 + added = out.Pins
157 + } else {
158 + // this can only happen if the progress option is set
159 + fmt.Fprintf(os.Stderr, "Fetched/Processed %d nodes\r", out.Progress)
160 + }
161
173 - return nil
174 - }),
162 + var pintype string
163 + rec, found := res.Request().Options["recursive"].(bool)
164 + if rec || !found {
165 + pintype = "recursively"
166 + } else {
167 + pintype = "directly"
168 + }
169 +
170 + for _, k := range added {
171 + fmt.Fprintf(os.Stdout, "pinned %s %s\n", k, pintype)
172 + }
173 + }
174 + },
175 },
176 }
177
@@ -192,11 +192,6 @@ collected if needed. (By default, recursively. Use -r=false for direct pins.)
192 },
193 Type: PinOutput{},
194 Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
195 - err := req.ParseBodyArgs()
196 - if err != nil {
197 - return err
198 - }
199 -
195 n, err := cmdenv.GetNode(env)
196 if err != nil {
197 return err
@@ -210,6 +205,10 @@ collected if needed. (By default, recursively. Use -r=false for direct pins.)
205 // set recursive flag
206 recursive, _ := req.Options[pinRecursiveOptionName].(bool)
207
208 + if err := req.ParseBodyArgs(); err != nil {
209 + return err
210 + }
211 +
212 removed, err := corerepo.Unpin(n, api, req.Context, req.Arguments, recursive)
213 if err != nil {
214 return err