@cryptotaxi247 / kubo / commits / ef6e9cf28

feat(commands): --stream option for ls

Convert LS Command to use current cmds lib Update LS Command to support streaming Rebase fixes License: MIT Signed-off-by: hannahhoward <hannah@hannahhoward.net>

hannahhoward committed Oct 18, 2018 at 10:49 UTC ef6e9cf2839b25bd13f2bcf2aff9530ee4e0161b
2 files changed +184 -108
core/commands/ls.go
+182 -105
@@ -1,12 +1,11 @@
1 package commands
2
3 import (
4 - "bytes"
4 "fmt"
5 "io"
6 "text/tabwriter"
7
9 - cmds "github.com/ipfs/go-ipfs/commands"
8 + cmdenv "github.com/ipfs/go-ipfs/core/commands/cmdenv"
9 e "github.com/ipfs/go-ipfs/core/commands/e"
10 iface "github.com/ipfs/go-ipfs/core/coreapi/interface"
11
@@ -16,29 +15,40 @@ import (
15 unixfspb "gx/ipfs/QmUnHNqhSB1JgzVCxL1Kz3yb4bdyB4q1Z9AD5AUBVmt3fZ/go-unixfs/pb"
16 blockservice "gx/ipfs/QmVDTbzzTwnuBwNbJdhW3u7LoBQp46bezm9yp4z1RoEepM/go-blockservice"
17 offline "gx/ipfs/QmYZwey1thDTynSrvd6qQkX24UpTka6TFhQ2v569UpoqxD/go-ipfs-exchange-offline"
18 + cmds "gx/ipfs/Qma6uuSyjkecGhMFFLfzyJDPyoDtNJSHJNweDccZhaWkgU/go-ipfs-cmds"
19 merkledag "gx/ipfs/QmcGt25mrjuB2kKW2zhPbXVZNHc4yoTDQ65NA8m6auP2f1/go-merkledag"
20 ipld "gx/ipfs/QmcKKBwfz6FyQdHR2jsXrrF6XeSBXYL86anmWNewpFpoF5/go-ipld-format"
21 "gx/ipfs/Qmde5VP1qUkyQXKCfmEUA7bP64V2HAptbJ7phuPp7jXWwg/go-ipfs-cmdkit"
22 )
23
24 +// LsLink contains printable data for a single ipld link in ls output
25 type LsLink struct {
26 Name, Hash string
27 Size uint64
28 Type unixfspb.Data_DataType
29 }
30
31 +// LsObject is an element of LsOutput
32 +// It can represent a whole directory, a directory header, one or more links,
33 +// Or a the end of a directory
34 type LsObject struct {
31 - Hash string
32 - Links []LsLink
35 + Hash string
36 + Links []LsLink
37 + HasHeader bool
38 + HasLinks bool
39 + HasFooter bool
40 }
41
42 +// LsOutput is a set of printable data for directories
43 type LsOutput struct {
36 - Objects []LsObject
44 + MultipleFolders bool
45 + Objects []LsObject
46 }
47
48 const (
49 lsHeadersOptionNameTime = "headers"
50 lsResolveTypeOptionName = "resolve-type"
51 + lsStreamOptionName = "stream"
52 )
53
54 var LsCmd = &cmds.Command{
@@ -60,32 +70,20 @@ The JSON output contains type information.
70 Options: []cmdkit.Option{
71 cmdkit.BoolOption(lsHeadersOptionNameTime, "v", "Print table headers (Hash, Size, Name)."),
72 cmdkit.BoolOption(lsResolveTypeOptionName, "Resolve linked objects to find out their types.").WithDefault(true),
73 + cmdkit.BoolOption(lsStreamOptionName, "s", "Stream directory entries as they are found."),
74 },
64 - Run: func(req cmds.Request, res cmds.Response) {
65 - nd, err := req.InvocContext().GetNode()
75 + Run: func(req *cmds.Request, res cmds.ResponseEmitter, env cmds.Environment) error {
76 + nd, err := cmdenv.GetNode(env)
77 if err != nil {
67 - res.SetError(err, cmdkit.ErrNormal)
68 - return
78 + return err
79 }
80
71 - api, err := req.InvocContext().GetApi()
81 + api, err := cmdenv.GetApi(env)
82 if err != nil {
73 - res.SetError(err, cmdkit.ErrNormal)
74 - return
75 - }
76 -
77 - // get options early -> exit early in case of error
78 - if _, _, err := req.Option(lsHeadersOptionNameTime).Bool(); err != nil {
79 - res.SetError(err, cmdkit.ErrNormal)
80 - return
81 - }
82 -
83 - resolve, _, err := req.Option(lsResolveTypeOptionName).Bool()
84 - if err != nil {
85 - res.SetError(err, cmdkit.ErrNormal)
86 - return
83 + return err
84 }
85
86 + resolve, _ := req.Options[lsResolveTypeOptionName].(bool)
87 dserv := nd.DAG
88 if !resolve {
89 offlineexch := offline.Exchange(nd.Blockstore)
@@ -93,125 +91,204 @@ The JSON output contains type information.
91 dserv = merkledag.NewDAGService(bserv)
92 }
93
96 - paths := req.Arguments()
94 + err = req.ParseBodyArgs()
95 + if err != nil {
96 + return err
97 + }
98 +
99 + paths := req.Arguments
100
101 var dagnodes []ipld.Node
102 for _, fpath := range paths {
103 p, err := iface.ParsePath(fpath)
104 if err != nil {
102 - res.SetError(err, cmdkit.ErrNormal)
103 - return
105 + return err
106 }
107
106 - dagnode, err := api.ResolveNode(req.Context(), p)
108 + dagnode, err := api.ResolveNode(req.Context, p)
109 if err != nil {
108 - res.SetError(err, cmdkit.ErrNormal)
109 - return
110 + return err
111 }
112 dagnodes = append(dagnodes, dagnode)
113 }
113 -
114 - output := make([]LsObject, len(req.Arguments()))
115 - ng := merkledag.NewSession(req.Context(), nd.DAG)
114 + ng := merkledag.NewSession(req.Context, nd.DAG)
115 ro := merkledag.NewReadOnlyDagService(ng)
116
117 + stream, _ := req.Options[lsStreamOptionName].(bool)
118 + multipleFolders := len(req.Arguments) > 1
119 + if !stream {
120 + output := make([]LsObject, len(req.Arguments))
121 +
122 + for i, dagnode := range dagnodes {
123 + dir, err := uio.NewDirectoryFromNode(ro, dagnode)
124 + if err != nil && err != uio.ErrNotADir {
125 + return fmt.Errorf("the data in %s (at %q) is not a UnixFS directory: %s", dagnode.Cid(), paths[i], err)
126 + }
127 +
128 + var links []*ipld.Link
129 + if dir == nil {
130 + links = dagnode.Links()
131 + } else {
132 + links, err = dir.Links(req.Context)
133 + if err != nil {
134 + return err
135 + }
136 + }
137 + outputLinks := make([]LsLink, len(links))
138 + for j, link := range links {
139 + lsLink, err := makeLsLink(req, dserv, resolve, link)
140 + if err != nil {
141 + return err
142 + }
143 + outputLinks[j] = *lsLink
144 + }
145 + output[i] = newFullDirectoryLsObject(paths[i], outputLinks)
146 + }
147 +
148 + return cmds.EmitOnce(res, &LsOutput{multipleFolders, output})
149 + }
150 +
151 for i, dagnode := range dagnodes {
152 dir, err := uio.NewDirectoryFromNode(ro, dagnode)
153 if err != nil && err != uio.ErrNotADir {
121 - res.SetError(fmt.Errorf("the data in %s (at %q) is not a UnixFS directory: %s", dagnode.Cid(), paths[i], err), cmdkit.ErrNormal)
122 - return
154 + return fmt.Errorf("the data in %s (at %q) is not a UnixFS directory: %s", dagnode.Cid(), paths[i], err)
155 }
156
125 - var links []*ipld.Link
157 + var linkResults <-chan unixfs.LinkResult
158 if dir == nil {
127 - links = dagnode.Links()
159 + linkResults = makeDagNodeLinkResults(req, dagnode)
160 } else {
129 - links, err = dir.Links(req.Context())
130 - if err != nil {
131 - res.SetError(err, cmdkit.ErrNormal)
132 - return
133 - }
161 + linkResults = dir.EnumLinksAsync(req.Context)
162 }
163
136 - output[i] = LsObject{
137 - Hash: paths[i],
138 - Links: make([]LsLink, len(links)),
139 - }
140 -
141 - for j, link := range links {
142 - t := unixfspb.Data_DataType(-1)
143 -
144 - switch link.Cid.Type() {
145 - case cid.Raw:
146 - // No need to check with raw leaves
147 - t = unixfs.TFile
148 - case cid.DagProtobuf:
149 - linkNode, err := link.GetNode(req.Context(), dserv)
150 - if err == ipld.ErrNotFound && !resolve {
151 - // not an error
152 - linkNode = nil
153 - } else if err != nil {
154 - res.SetError(err, cmdkit.ErrNormal)
155 - return
156 - }
164 + output := make([]LsObject, 1)
165 + outputLinks := make([]LsLink, 1)
166
158 - if pn, ok := linkNode.(*merkledag.ProtoNode); ok {
159 - d, err := unixfs.FSNodeFromBytes(pn.Data())
160 - if err != nil {
161 - res.SetError(err, cmdkit.ErrNormal)
162 - return
163 - }
164 - t = d.Type()
165 - }
167 + output[0] = newDirectoryHeaderLsObject(paths[i])
168 + if err = res.Emit(&LsOutput{multipleFolders, output}); err != nil {
169 + return nil
170 + }
171 + for linkResult := range linkResults {
172 + if linkResult.Err != nil {
173 + return linkResult.Err
174 + }
175 + link := linkResult.Link
176 + lsLink, err := makeLsLink(req, dserv, resolve, link)
177 + if err != nil {
178 + return err
179 }
167 - output[i].Links[j] = LsLink{
168 - Name: link.Name,
169 - Hash: link.Cid.String(),
170 - Size: link.Size,
171 - Type: t,
180 + outputLinks[0] = *lsLink
181 + output[0] = newDirectoryLinksLsObject(outputLinks)
182 + if err = res.Emit(&LsOutput{multipleFolders, output}); err != nil {
183 + return err
184 }
185 }
186 + output[0] = newDirectoryFooterLsObject()
187 + if err = res.Emit(&LsOutput{multipleFolders, output}); err != nil {
188 + return err
189 + }
190 }
175 -
176 - res.SetOutput(&LsOutput{output})
191 + return nil
192 },
178 - Marshalers: cmds.MarshalerMap{
179 - cmds.Text: func(res cmds.Response) (io.Reader, error) {
180 -
181 - v, err := unwrapOutput(res.Output())
182 - if err != nil {
183 - return nil, err
184 - }
185 -
186 - headers, _, _ := res.Request().Option(lsHeadersOptionNameTime).Bool()
193 + Encoders: cmds.EncoderMap{
194 + cmds.Text: cmds.MakeEncoder(func(req *cmds.Request, w io.Writer, v interface{}) error {
195 + headers, _ := req.Options[lsHeadersOptionNameTime].(bool)
196 output, ok := v.(*LsOutput)
197 if !ok {
189 - return nil, e.TypeErr(output, v)
198 + return e.TypeErr(output, v)
199 }
200
192 - buf := new(bytes.Buffer)
193 - w := tabwriter.NewWriter(buf, 1, 2, 1, ' ', 0)
201 + tw := tabwriter.NewWriter(w, 1, 2, 1, ' ', 0)
202 for _, object := range output.Objects {
195 - if len(output.Objects) > 1 {
196 - fmt.Fprintf(w, "%s:\n", object.Hash)
197 - }
198 - if headers {
199 - fmt.Fprintln(w, "Hash\tSize\tName")
203 + if object.HasHeader {
204 + if output.MultipleFolders {
205 + fmt.Fprintf(tw, "%s:\n", object.Hash)
206 + }
207 + if headers {
208 + fmt.Fprintln(tw, "Hash\tSize\tName")
209 + }
210 }
201 - for _, link := range object.Links {
202 - if link.Type == unixfs.TDirectory {
203 - link.Name += "/"
211 + if object.HasLinks {
212 + for _, link := range object.Links {
213 + if link.Type == unixfs.TDirectory {
214 + link.Name += "/"
215 + }
216 +
217 + fmt.Fprintf(tw, "%s\t%v\t%s\n", link.Hash, link.Size, link.Name)
218 }
205 - fmt.Fprintf(w, "%s\t%v\t%s\n", link.Hash, link.Size, link.Name)
219 }
207 - if len(output.Objects) > 1 {
208 - fmt.Fprintln(w)
220 + if object.HasFooter {
221 + if output.MultipleFolders {
222 + fmt.Fprintln(tw)
223 + }
224 }
225 }
211 - w.Flush()
212 -
213 - return buf, nil
214 - },
226 + tw.Flush()
227 + return nil
228 + }),
229 },
230 Type: LsOutput{},
231 }
232 +
233 +func makeDagNodeLinkResults(req *cmds.Request, dagnode ipld.Node) <-chan unixfs.LinkResult {
234 + linkResults := make(chan unixfs.LinkResult)
235 + go func() {
236 + defer close(linkResults)
237 + for _, l := range dagnode.Links() {
238 + select {
239 + case linkResults <- unixfs.LinkResult{
240 + Link: l,
241 + Err: nil,
242 + }:
243 + case <-req.Context.Done():
244 + return
245 + }
246 + }
247 + }()
248 + return linkResults
249 +}
250 +
251 +func newFullDirectoryLsObject(hash string, links []LsLink) LsObject {
252 + return LsObject{hash, links, true, true, true}
253 +}
254 +func newDirectoryHeaderLsObject(hash string) LsObject {
255 + return LsObject{hash, nil, true, false, false}
256 +}
257 +func newDirectoryLinksLsObject(links []LsLink) LsObject {
258 + return LsObject{"", links, false, true, false}
259 +}
260 +func newDirectoryFooterLsObject() LsObject {
261 + return LsObject{"", nil, false, false, true}
262 +}
263 +
264 +func makeLsLink(req *cmds.Request, dserv ipld.DAGService, resolve bool, link *ipld.Link) (*LsLink, error) {
265 + t := unixfspb.Data_DataType(-1)
266 +
267 + switch link.Cid.Type() {
268 + case cid.Raw:
269 + // No need to check with raw leaves
270 + t = unixfs.TFile
271 + case cid.DagProtobuf:
272 + linkNode, err := link.GetNode(req.Context, dserv)
273 + if err == ipld.ErrNotFound && !resolve {
274 + // not an error
275 + linkNode = nil
276 + } else if err != nil {
277 + return nil, err
278 + }
279 +
280 + if pn, ok := linkNode.(*merkledag.ProtoNode); ok {
281 + d, err := unixfs.FSNodeFromBytes(pn.Data())
282 + if err != nil {
283 + return nil, err
284 + }
285 + t = d.Type()
286 + }
287 + }
288 + return &LsLink{
289 + Name: link.Name,
290 + Hash: link.Cid.String(),
291 + Size: link.Size,
292 + Type: t,
293 + }, nil
294 +}
core/commands/root.go
+2 -3
@@ -3,7 +3,6 @@ package commands
3 import (
4 "errors"
5
6 - lgc "github.com/ipfs/go-ipfs/commands/legacy"
6 dag "github.com/ipfs/go-ipfs/core/commands/dag"
7 name "github.com/ipfs/go-ipfs/core/commands/name"
8 ocmd "github.com/ipfs/go-ipfs/core/commands/object"
@@ -127,7 +126,7 @@ var rootSubcommands = map[string]*cmds.Command{
126 "id": IDCmd,
127 "key": KeyCmd,
128 "log": LogCmd,
130 - "ls": lgc.NewCommand(LsCmd),
129 + "ls": LsCmd,
130 "mount": MountCmd,
131 "name": name.NameCmd,
132 "object": ocmd.ObjectCmd,
@@ -165,7 +164,7 @@ var rootROSubcommands = map[string]*cmds.Command{
164 },
165 "get": GetCmd,
166 "dns": DNSCmd,
168 - "ls": lgc.NewCommand(LsCmd),
167 + "ls": LsCmd,
168 "name": {
169 Subcommands: map[string]*cmds.Command{
170 "resolve": name.IpnsCmd,