@cryptotaxi247 / kubo / commits / abd390b89

core/commands: Made add command show streamed output

Matt Bell committed Dec 17, 2014 at 18:43 UTC abd390b8926a20cfe7e0fb653c68674287b20364
1 file changed +53 -57
core/commands/add.go
+53 -57
@@ -20,10 +20,9 @@ import (
20 // Error indicating the max depth has been exceded.
21 var ErrDepthLimitExceeded = fmt.Errorf("depth limit exceeded")
22
23 -type AddOutput struct {
24 - Objects []*Object
25 - Names []string
26 - Quiet bool
23 +type AddedObject struct {
24 + Name string
25 + Hash string
26 }
27
28 var AddCmd = &cmds.Command{
@@ -45,58 +44,65 @@ remains to be implemented.
44 cmds.BoolOption("quiet", "q", "Write minimal output"),
45 },
46 Run: func(req cmds.Request) (interface{}, error) {
48 - added := &AddOutput{}
47 n, err := req.Context().GetNode()
48 if err != nil {
49 return nil, err
50 }
51
54 - for {
55 - file, err := req.Files().NextFile()
56 - if err != nil && err != io.EOF {
57 - return nil, err
58 - }
59 - if file == nil {
60 - break
61 - }
52 + outChan := make(chan interface{})
53
63 - _, err = addFile(n, file, added)
64 - if err != nil {
65 - return nil, err
66 - }
67 - }
54 + go func() {
55 + defer close(outChan)
56
69 - quiet, _, err := req.Option("quiet").Bool()
70 - if err != nil {
71 - return nil, err
72 - }
57 + for {
58 + file, err := req.Files().NextFile()
59 + if (err != nil && err != io.EOF) || file == nil {
60 + return
61 + }
62
74 - added.Quiet = quiet
63 + _, err = addFile(n, file, outChan)
64 + if err != nil {
65 + return
66 + }
67 + }
68 + }()
69
76 - return added, nil
70 + return outChan, nil
71 },
72 Marshalers: cmds.MarshalerMap{
73 cmds.Text: func(res cmds.Response) (io.Reader, error) {
80 - val, ok := res.Output().(*AddOutput)
74 + outChan, ok := res.Output().(chan interface{})
75 if !ok {
76 return nil, u.ErrCast()
77 }
78
85 - // TODO: use this with an option
86 - // sort.Stable(val)
79 + quiet, _, err := res.Request().Option("quiet").Bool()
80 + if err != nil {
81 + return nil, err
82 + }
83
88 - var buf bytes.Buffer
89 - for i, obj := range val.Objects {
90 - if val.Quiet {
91 - buf.Write([]byte(fmt.Sprintf("%s\n", obj.Hash)))
84 + marshal := func(v interface{}) (io.Reader, error) {
85 + obj, ok := v.(*AddedObject)
86 + if !ok {
87 + return nil, u.ErrCast()
88 + }
89 +
90 + var buf bytes.Buffer
91 + if quiet {
92 + buf.WriteString(fmt.Sprintf("%s\n", obj.Hash))
93 } else {
93 - buf.Write([]byte(fmt.Sprintf("added %s %s\n", obj.Hash, val.Names[i])))
94 + buf.WriteString(fmt.Sprintf("added %s %s\n", obj.Hash, obj.Name))
95 }
96 + return &buf, nil
97 }
96 - return &buf, nil
98 +
99 + return &cmds.ChannelMarshaler{
100 + Channel: outChan,
101 + Marshaler: marshal,
102 + }, nil
103 },
104 },
99 - Type: &AddOutput{},
105 + Type: &AddedObject{},
106 }
107
108 func add(n *core.IpfsNode, readers []io.Reader) ([]*dag.Node, error) {
@@ -132,9 +138,9 @@ func addNode(n *core.IpfsNode, node *dag.Node) error {
138 return nil
139 }
140
135 -func addFile(n *core.IpfsNode, file cmds.File, added *AddOutput) (*dag.Node, error) {
141 +func addFile(n *core.IpfsNode, file cmds.File, out chan interface{}) (*dag.Node, error) {
142 if file.IsDirectory() {
137 - return addDir(n, file, added)
143 + return addDir(n, file, out)
144 }
145
146 dns, err := add(n, []io.Reader{file})
@@ -143,13 +149,13 @@ func addFile(n *core.IpfsNode, file cmds.File, added *AddOutput) (*dag.Node, err
149 }
150
151 log.Infof("adding file: %s", file.FileName())
146 - if err := addDagnode(added, file.FileName(), dns[len(dns)-1]); err != nil {
152 + if err := outputDagnode(out, file.FileName(), dns[len(dns)-1]); err != nil {
153 return nil, err
154 }
155 return dns[len(dns)-1], nil // last dag node is the file.
156 }
157
152 -func addDir(n *core.IpfsNode, dir cmds.File, added *AddOutput) (*dag.Node, error) {
158 +func addDir(n *core.IpfsNode, dir cmds.File, out chan interface{}) (*dag.Node, error) {
159 log.Infof("adding directory: %s", dir.FileName())
160
161 tree := &dag.Node{Data: ft.FolderPBData()}
@@ -163,7 +169,7 @@ func addDir(n *core.IpfsNode, dir cmds.File, added *AddOutput) (*dag.Node, error
169 break
170 }
171
166 - node, err := addFile(n, file, added)
172 + node, err := addFile(n, file, out)
173 if err != nil {
174 return nil, err
175 }
@@ -176,7 +182,7 @@ func addDir(n *core.IpfsNode, dir cmds.File, added *AddOutput) (*dag.Node, error
182 }
183 }
184
179 - err := addDagnode(added, dir.FileName(), tree)
185 + err := outputDagnode(out, dir.FileName(), tree)
186 if err != nil {
187 return nil, err
188 }
@@ -189,27 +195,17 @@ func addDir(n *core.IpfsNode, dir cmds.File, added *AddOutput) (*dag.Node, error
195 return tree, nil
196 }
197
192 -// addDagnode adds dagnode info to an output object
193 -func addDagnode(output *AddOutput, name string, dn *dag.Node) error {
198 +// outputDagnode sends dagnode info over the output channel
199 +func outputDagnode(out chan interface{}, name string, dn *dag.Node) error {
200 o, err := getOutput(dn)
201 if err != nil {
202 return err
203 }
204
199 - output.Objects = append(output.Objects, o)
200 - output.Names = append(output.Names, name)
201 - return nil
202 -}
203 -
204 -// Sort interface implementation to sort add output by name
205 + out <- &AddedObject{
206 + Hash: o.Hash,
207 + Name: name,
208 + }
209
206 -func (a AddOutput) Len() int {
207 - return len(a.Names)
208 -}
209 -func (a AddOutput) Swap(i, j int) {
210 - a.Names[i], a.Names[j] = a.Names[j], a.Names[i]
211 - a.Objects[i], a.Objects[j] = a.Objects[j], a.Objects[i]
212 -}
213 -func (a AddOutput) Less(i, j int) bool {
214 - return a.Names[i] < a.Names[j]
210 + return nil
211 }