@cryptotaxi247 / kubo / commits / 1efbc7922

use mfs for adds

License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>

Jeromy committed Dec 4, 2015 at 14:25 UTC 1efbc79223fb581689090d361e1f5fc9bc485192
5 files changed +129 -75
core/commands/add.go
+21 -5
@@ -18,6 +18,7 @@ var ErrDepthLimitExceeded = fmt.Errorf("depth limit exceeded")
18
19 const (
20 quietOptionName = "quiet"
21 + silentOptionName = "silent"
22 progressOptionName = "progress"
23 trickleOptionName = "trickle"
24 wrapOptionName = "wrap-with-directory"
@@ -44,6 +45,7 @@ remains to be implemented.
45 Options: []cmds.Option{
46 cmds.OptionRecursivePath, // a builtin option that allows recursive paths (-r, --recursive)
47 cmds.BoolOption(quietOptionName, "q", "Write minimal output"),
48 + cmds.BoolOption(silentOptionName, "x", "Write no output"),
49 cmds.BoolOption(progressOptionName, "p", "Stream progress data"),
50 cmds.BoolOption(trickleOptionName, "t", "Use trickle-dag format for dag generation"),
51 cmds.BoolOption(onlyHashOptionName, "n", "Only chunk and hash - do not write to disk"),
@@ -59,6 +61,9 @@ remains to be implemented.
61
62 req.SetOption(progressOptionName, true)
63
64 + log.Error("SKIPPING SIZE")
65 + return nil
66 +
67 sizeFile, ok := req.Files().(files.SizeFile)
68 if !ok {
69 // we don't need to error, the progress bar just won't know how big the files are
@@ -100,6 +105,7 @@ remains to be implemented.
105 wrap, _, _ := req.Option(wrapOptionName).Bool()
106 hash, _, _ := req.Option(onlyHashOptionName).Bool()
107 hidden, _, _ := req.Option(hiddenOptionName).Bool()
108 + silent, _, _ := req.Option(silentOptionName).Bool()
109 chunker, _, _ := req.Option(chunkerOptionName).String()
110 dopin, pin_found, _ := req.Option(pinOptionName).Bool()
111
@@ -123,13 +129,18 @@ remains to be implemented.
129 outChan := make(chan interface{}, 8)
130 res.SetOutput((<-chan interface{})(outChan))
131
126 - fileAdder := coreunix.NewAdder(req.Context(), n, outChan)
132 + fileAdder, err := coreunix.NewAdder(req.Context(), n, outChan)
133 + if err != nil {
134 + res.SetError(err, cmds.ErrNormal)
135 + return
136 + }
137 fileAdder.Chunker = chunker
138 fileAdder.Progress = progress
139 fileAdder.Hidden = hidden
140 fileAdder.Trickle = trickle
141 fileAdder.Wrap = wrap
142 fileAdder.Pin = dopin
143 + fileAdder.Silent = silent
144
145 // addAllFiles loops over a convenience slice file to
146 // add each file individually. e.g. 'ipfs add a b c'
@@ -143,7 +154,7 @@ remains to be implemented.
154 return nil // done
155 }
156
146 - if _, err := fileAdder.AddFile(file); err != nil {
157 + if err := fileAdder.AddFile(file); err != nil {
158 return err
159 }
160 }
@@ -159,9 +170,8 @@ remains to be implemented.
170 }
171
172 // copy intermediary nodes from editor to our actual dagservice
162 - _, err := fileAdder.Finalize(n.DAG)
173 + _, err := fileAdder.Finalize()
174 if err != nil {
164 - log.Error("WRITE OUT: ", err)
175 return err
176 }
177
@@ -194,7 +204,13 @@ remains to be implemented.
204 return
205 }
206
197 - showProgressBar := !quiet
207 + progress, _, err := req.Option(progressOptionName).Bool()
208 + if err != nil {
209 + res.SetError(u.ErrCast(), cmds.ErrNormal)
210 + return
211 + }
212 +
213 + showProgressBar := !quiet || progress
214
215 var bar *pb.ProgressBar
216 var terminalWidth int
core/coreunix/add.go
+83 -63
@@ -15,7 +15,7 @@ import (
15 "github.com/ipfs/go-ipfs/exchange/offline"
16 importer "github.com/ipfs/go-ipfs/importer"
17 "github.com/ipfs/go-ipfs/importer/chunk"
18 - dagutils "github.com/ipfs/go-ipfs/merkledag/utils"
18 + mfs "github.com/ipfs/go-ipfs/mfs"
19 "github.com/ipfs/go-ipfs/pin"
20
21 "github.com/ipfs/go-ipfs/commands/files"
@@ -62,12 +62,16 @@ type AddedObject struct {
62 Bytes int64 `json:",omitempty"`
63 }
64
65 -func NewAdder(ctx context.Context, n *core.IpfsNode, out chan interface{}) *Adder {
66 - e := dagutils.NewDagEditor(newDirNode(), nil)
65 +func NewAdder(ctx context.Context, n *core.IpfsNode, out chan interface{}) (*Adder, error) {
66 + mr, err := mfs.NewRoot(ctx, n.DAG, newDirNode(), nil)
67 + if err != nil {
68 + return nil, err
69 + }
70 +
71 return &Adder{
72 + mr: mr,
73 ctx: ctx,
74 node: n,
70 - editor: e,
75 out: out,
76 Progress: false,
77 Hidden: true,
@@ -75,22 +79,23 @@ func NewAdder(ctx context.Context, n *core.IpfsNode, out chan interface{}) *Adde
79 Trickle: false,
80 Wrap: false,
81 Chunker: "",
78 - }
82 + }, nil
83 }
84
85 // Internal structure for holding the switches passed to the `add` call
86 type Adder struct {
87 ctx context.Context
88 node *core.IpfsNode
85 - editor *dagutils.Editor
89 out chan interface{}
90 Progress bool
91 Hidden bool
92 Pin bool
93 Trickle bool
94 + Silent bool
95 Wrap bool
96 Chunker string
97 root *dag.Node
98 + mr *mfs.Root
99 }
100
101 // Perform the actual add & pin locally, outputting results to reader
@@ -113,26 +118,29 @@ func (params Adder) add(reader io.Reader) (*dag.Node, error) {
118 }
119
120 func (params *Adder) RootNode() (*dag.Node, error) {
116 - // for memoizing
117 - if params.root != nil {
118 - return params.root, nil
119 - }
121 + return params.mr.GetValue().GetNode()
122 + /*
123 + // for memoizing
124 + if params.root != nil {
125 + return params.root, nil
126 + }
127
121 - root := params.editor.GetNode()
128 + root := params.editor.GetNode()
129
123 - // if not wrapping, AND one root file, use that hash as root.
124 - if !params.Wrap && len(root.Links) == 1 {
125 - var err error
126 - root, err = root.Links[0].GetNode(params.ctx, params.editor.GetDagService())
130 + // if not wrapping, AND one root file, use that hash as root.
131 + if !params.Wrap && len(root.Links) == 1 {
132 + var err error
133 + root, err = root.Links[0].GetNode(params.ctx, params.editor.GetDagService())
134 + params.root = root
135 + // no need to output, as we've already done so.
136 + return root, err
137 + }
138 +
139 + // otherwise need to output, as we have not.
140 + err := outputDagnode(params.out, "", root)
141 params.root = root
128 - // no need to output, as we've already done so.
142 return root, err
130 - }
131 -
132 - // otherwise need to output, as we have not.
133 - err := outputDagnode(params.out, "", root)
134 - params.root = root
135 - return root, err
143 + */
144 }
145
146 func (params *Adder) PinRoot() error {
@@ -153,8 +161,8 @@ func (params *Adder) PinRoot() error {
161 return params.node.Pinning.Flush()
162 }
163
156 -func (params *Adder) Finalize(DAG dag.DAGService) (*dag.Node, error) {
157 - return params.editor.Finalize(DAG)
164 +func (params *Adder) Finalize() (*dag.Node, error) {
165 + return params.mr.GetValue().GetNode()
166 }
167
168 // Add builds a merkledag from the a reader, pinning all objects to the local
@@ -163,7 +171,10 @@ func Add(n *core.IpfsNode, r io.Reader) (string, error) {
171 unlock := n.Blockstore.PinLock()
172 defer unlock()
173
166 - fileAdder := NewAdder(n.Context(), n, nil)
174 + fileAdder, err := NewAdder(n.Context(), n, nil)
175 + if err != nil {
176 + return "", err
177 + }
178
179 node, err := fileAdder.add(r)
180 if err != nil {
@@ -193,14 +204,22 @@ func AddR(n *core.IpfsNode, root string) (key string, err error) {
204 }
205 defer f.Close()
206
196 - fileAdder := NewAdder(n.Context(), n, nil)
207 + fileAdder, err := NewAdder(n.Context(), n, nil)
208 + if err != nil {
209 + return "", err
210 + }
211
198 - dagnode, err := fileAdder.AddFile(f)
212 + err = fileAdder.AddFile(f)
213 if err != nil {
214 return "", err
215 }
216
203 - k, err := dagnode.Key()
217 + nd, err := fileAdder.Finalize()
218 + if err != nil {
219 + return "", err
220 + }
221 +
222 + k, err := nd.Key()
223 if err != nil {
224 return "", err
225 }
@@ -215,18 +234,29 @@ func AddR(n *core.IpfsNode, root string) (key string, err error) {
234 func AddWrapped(n *core.IpfsNode, r io.Reader, filename string) (string, *dag.Node, error) {
235 file := files.NewReaderFile(filename, filename, ioutil.NopCloser(r), nil)
236 dir := files.NewSliceFile("", "", []files.File{file})
218 - fileAdder := NewAdder(n.Context(), n, nil)
237 + fileAdder, err := NewAdder(n.Context(), n, nil)
238 + if err != nil {
239 + return "", nil, err
240 + }
241
242 unlock := n.Blockstore.PinLock()
243 defer unlock()
222 - dagnode, err := fileAdder.addDir(dir)
244 +
245 + err = fileAdder.addDir(dir)
246 + if err != nil {
247 + return "", nil, err
248 + }
249 +
250 + dagnode, err := fileAdder.Finalize()
251 if err != nil {
252 return "", nil, err
253 }
254 +
255 k, err := dagnode.Key()
256 if err != nil {
257 return "", nil, err
258 }
259 +
260 return gopath.Join(k.String(), filename), dagnode, nil
261 }
262
@@ -241,19 +271,22 @@ func (params *Adder) addNode(node *dag.Node, path string) error {
271 path = key.Pretty()
272 }
273
244 - if err := params.editor.InsertNodeAtPath(params.ctx, path, node, newDirNode); err != nil {
274 + if err := mfs.PutNode(params.mr, path, node); err != nil {
275 return err
276 }
277
248 - return outputDagnode(params.out, path, node)
278 + if !params.Silent {
279 + return outputDagnode(params.out, path, node)
280 + }
281 + return nil
282 }
283
284 // Add the given file while respecting the params.
252 -func (params *Adder) AddFile(file files.File) (*dag.Node, error) {
285 +func (params *Adder) AddFile(file files.File) error {
286 switch {
287 case files.IsHidden(file) && !params.Hidden:
288 log.Debugf("%s is hidden, skipping", file.FileName())
256 - return nil, &hiddenFileError{file.FileName()}
289 + return &hiddenFileError{file.FileName()}
290 case file.IsDirectory():
291 return params.addDir(file)
292 }
@@ -262,17 +295,16 @@ func (params *Adder) AddFile(file files.File) (*dag.Node, error) {
295 if s, ok := file.(*files.Symlink); ok {
296 sdata, err := unixfs.SymlinkData(s.Target)
297 if err != nil {
265 - return nil, err
298 + return err
299 }
300
301 dagnode := &dag.Node{Data: sdata}
302 _, err = params.node.DAG.Add(dagnode)
303 if err != nil {
271 - return nil, err
304 + return err
305 }
306
274 - err = params.addNode(dagnode, s.FileName())
275 - return dagnode, err
307 + return params.addNode(dagnode, s.FileName())
308 }
309
310 // case for regular file
@@ -285,52 +317,40 @@ func (params *Adder) AddFile(file files.File) (*dag.Node, error) {
317
318 dagnode, err := params.add(reader)
319 if err != nil {
288 - return nil, err
320 + return err
321 }
322
323 // patch it into the root
292 - log.Infof("adding file: %s", file.FileName())
293 - err = params.addNode(dagnode, file.FileName())
294 - return dagnode, err
324 + return params.addNode(dagnode, file.FileName())
325 }
326
297 -func (params *Adder) addDir(dir files.File) (*dag.Node, error) {
298 - tree := newDirNode()
327 +func (params *Adder) addDir(dir files.File) error {
328 log.Infof("adding directory: %s", dir.FileName())
329
330 + err := mfs.Mkdir(params.mr, dir.FileName(), true)
331 + if err != nil {
332 + return err
333 + }
334 +
335 for {
336 file, err := dir.NextFile()
337 if err != nil && err != io.EOF {
304 - return nil, err
338 + return err
339 }
340 if file == nil {
341 break
342 }
343
310 - node, err := params.AddFile(file)
344 + err = params.AddFile(file)
345 if _, ok := err.(*hiddenFileError); ok {
346 // hidden file error, skip file
347 continue
348 } else if err != nil {
315 - return nil, err
316 - }
317 -
318 - name := gopath.Base(file.FileName())
319 -
320 - if err := tree.AddNodeLinkClean(name, node); err != nil {
321 - return nil, err
349 + return err
350 }
351 }
352
325 - if err := params.addNode(tree, dir.FileName()); err != nil {
326 - return nil, err
327 - }
328 -
329 - if _, err := params.node.DAG.Add(tree); err != nil {
330 - return nil, err
331 - }
332 -
333 - return tree, nil
353 + return nil
354 }
355
356 // outputDagnode sends dagnode info over the output channel
@@ -379,7 +399,7 @@ func getOutput(dagnode *dag.Node) (*Object, error) {
399 for i, link := range dagnode.Links {
400 output.Links[i] = Link{
401 Name: link.Name,
382 - Hash: link.Hash.B58String(),
402 + //Hash: link.Hash.B58String(),
403 Size: link.Size,
404 }
405 }
exchange/bitswap/workers.go
+1 -1
@@ -89,7 +89,7 @@ func (bs *Bitswap) provideWorker(px process.Process) {
89 defer cancel()
90
91 if err := bs.network.Provide(ctx, k); err != nil {
92 - log.Error(err)
92 + //log.Error(err)
93 }
94 }
95
merkledag/merkledag.go
+9
@@ -3,6 +3,7 @@ package merkledag
3
4 import (
5 "fmt"
6 + "time"
7
8 "github.com/ipfs/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
9 blocks "github.com/ipfs/go-ipfs/blocks"
@@ -48,6 +49,14 @@ func (n *dagService) Add(nd *Node) (key.Key, error) {
49 if n == nil { // FIXME remove this assertion. protect with constructor invariant
50 return "", fmt.Errorf("dagService is nil")
51 }
52 + /*
53 + start := time.Now()
54 + defer func() {
55 + took := time.Now().Sub(start)
56 + log.Error("add took: %s", took)
57 + }()
58 + */
59 + _ = time.Saturday
60
61 d, err := nd.Encoded(false)
62 if err != nil {
mfs/system.go
+15 -6
@@ -71,15 +71,19 @@ func NewRoot(parent context.Context, ds dag.DAGService, node *dag.Node, pf PubFu
71 return nil, err
72 }
73
74 + var repub *Republisher
75 + if pf != nil {
76 + repub = NewRepublisher(parent, pf, time.Millisecond*300, time.Second*3)
77 + repub.setVal(ndk)
78 + go repub.Run()
79 + }
80 +
81 root := &Root{
82 node: node,
76 - repub: NewRepublisher(parent, pf, time.Millisecond*300, time.Second*3),
83 + repub: repub,
84 dserv: ds,
85 }
86
80 - root.repub.setVal(ndk)
81 - go root.repub.Run()
82 -
87 pbn, err := ft.FromBytes(node.Data)
88 if err != nil {
89 log.Error("IPNS pointer was not unixfs node")
@@ -113,12 +117,17 @@ func (kr *Root) closeChild(name string, nd *dag.Node) error {
117 return err
118 }
119
116 - kr.repub.Update(k)
120 + if kr.repub != nil {
121 + kr.repub.Update(k)
122 + }
123 return nil
124 }
125
126 func (kr *Root) Close() error {
121 - return kr.repub.Close()
127 + if kr.repub != nil {
128 + return kr.repub.Close()
129 + }
130 + return nil
131 }
132
133 // Republisher manages when to publish a given entry