@cryptotaxi247 / kubo / commits / 9360f5ca6

Store needed parts of IpfsNode in Adder.

This will make it easier to set up a specialized data pipeline. License: MIT Signed-off-by: Kevin Atkinson <k@kevina.org>

Kevin Atkinson committed Jun 1, 2016 at 16:14 UTC 9360f5ca6fd9393736d5cafb536950b8bc5b4e37
3 files changed +52 -45
core/commands/add.go
+3 -1
@@ -141,11 +141,13 @@ You can now refer to the added file in a gateway, like so:
141 outChan := make(chan interface{}, 8)
142 res.SetOutput((<-chan interface{})(outChan))
143
144 - fileAdder, err := coreunix.NewAdder(req.Context(), n, outChan)
144 + fileAdder, err := coreunix.NewAdder(req.Context(), n.Pinning, n.Blockstore, n.DAG)
145 if err != nil {
146 res.SetError(err, cmds.ErrNormal)
147 return
148 }
149 +
150 + fileAdder.Out = outChan
151 fileAdder.Chunker = chunker
152 fileAdder.Progress = progress
153 fileAdder.Hidden = hidden
core/coreunix/add.go
+47 -43
@@ -67,42 +67,46 @@ type AddedObject struct {
67 Bytes int64 `json:",omitempty"`
68 }
69
70 -func NewAdder(ctx context.Context, n *core.IpfsNode, out chan interface{}) (*Adder, error) {
71 - mr, err := mfs.NewRoot(ctx, n.DAG, newDirNode(), nil)
70 +func NewAdder(ctx context.Context, p pin.Pinner, bs bstore.GCBlockstore, ds dag.DAGService) (*Adder, error) {
71 + mr, err := mfs.NewRoot(ctx, ds, newDirNode(), nil)
72 if err != nil {
73 return nil, err
74 }
75
76 return &Adder{
77 - mr: mr,
78 - ctx: ctx,
79 - node: n,
80 - out: out,
81 - Progress: false,
82 - Hidden: true,
83 - Pin: true,
84 - Trickle: false,
85 - Wrap: false,
86 - Chunker: "",
77 + mr: mr,
78 + ctx: ctx,
79 + pinning: p,
80 + blockstore: bs,
81 + dagService: ds,
82 + Progress: false,
83 + Hidden: true,
84 + Pin: true,
85 + Trickle: false,
86 + Wrap: false,
87 + Chunker: "",
88 }, nil
89 +
90 }
91
92 // Internal structure for holding the switches passed to the `add` call
93 type Adder struct {
92 - ctx context.Context
93 - node *core.IpfsNode
94 - out chan interface{}
95 - Progress bool
96 - Hidden bool
97 - Pin bool
98 - Trickle bool
99 - Silent bool
100 - Wrap bool
101 - Chunker string
102 - root *dag.Node
103 - mr *mfs.Root
104 - unlocker bs.Unlocker
105 - tempRoot key.Key
94 + ctx context.Context
95 + pinning pin.Pinner
96 + blockstore bstore.GCBlockstore
97 + dagService dag.DAGService
98 + Out chan interface{}
99 + Progress bool
100 + Hidden bool
101 + Pin bool
102 + Trickle bool
103 + Silent bool
104 + Wrap bool
105 + Chunker string
106 + root *dag.Node
107 + mr *mfs.Root
108 + unlocker bs.Unlocker
109 + tempRoot key.Key
110 }
111
112 // Perform the actual add & pin locally, outputting results to reader
@@ -114,12 +118,12 @@ func (adder Adder) add(reader io.Reader) (*dag.Node, error) {
118
119 if adder.Trickle {
120 return importer.BuildTrickleDagFromReader(
117 - adder.node.DAG,
121 + adder.dagService,
122 chnk,
123 )
124 }
125 return importer.BuildDagFromReader(
122 - adder.node.DAG,
126 + adder.dagService,
127 chnk,
128 )
129 }
@@ -137,7 +141,7 @@ func (adder *Adder) RootNode() (*dag.Node, error) {
141
142 // if not wrapping, AND one root file, use that hash as root.
143 if !adder.Wrap && len(root.Links) == 1 {
140 - root, err = root.Links[0].GetNode(adder.ctx, adder.node.DAG)
144 + root, err = root.Links[0].GetNode(adder.ctx, adder.dagService)
145 if err != nil {
146 return nil, err
147 }
@@ -156,21 +160,21 @@ func (adder *Adder) PinRoot() error {
160 return nil
161 }
162
159 - rnk, err := adder.node.DAG.Add(root)
163 + rnk, err := adder.dagService.Add(root)
164 if err != nil {
165 return err
166 }
167
168 if adder.tempRoot != "" {
165 - err := adder.node.Pinning.Unpin(adder.ctx, adder.tempRoot, true)
169 + err := adder.pinning.Unpin(adder.ctx, adder.tempRoot, true)
170 if err != nil {
171 return err
172 }
173 adder.tempRoot = rnk
174 }
175
172 - adder.node.Pinning.PinWithMode(rnk, pin.Recursive)
173 - return adder.node.Pinning.Flush()
176 + adder.pinning.PinWithMode(rnk, pin.Recursive)
177 + return adder.pinning.Flush()
178 }
179
180 func (adder *Adder) Finalize() (*dag.Node, error) {
@@ -237,7 +241,7 @@ func (adder *Adder) outputDirs(path string, fs mfs.FSNode) error {
241 }
242 }
243
240 - return outputDagnode(adder.out, path, nd)
244 + return outputDagnode(adder.Out, path, nd)
245 }
246
247 // Add builds a merkledag from the a reader, pinning all objects to the local
@@ -245,7 +249,7 @@ func (adder *Adder) outputDirs(path string, fs mfs.FSNode) error {
249 func Add(n *core.IpfsNode, r io.Reader) (string, error) {
250 defer n.Blockstore.PinLock().Unlock()
251
248 - fileAdder, err := NewAdder(n.Context(), n, nil)
252 + fileAdder, err := NewAdder(n.Context(), n.Pinning, n.Blockstore, n.DAG)
253 if err != nil {
254 return "", err
255 }
@@ -277,7 +281,7 @@ func AddR(n *core.IpfsNode, root string) (key string, err error) {
281 }
282 defer f.Close()
283
280 - fileAdder, err := NewAdder(n.Context(), n, nil)
284 + fileAdder, err := NewAdder(n.Context(), n.Pinning, n.Blockstore, n.DAG)
285 if err != nil {
286 return "", err
287 }
@@ -306,7 +310,7 @@ func AddR(n *core.IpfsNode, root string) (key string, err error) {
310 // the directory, and and error if any.
311 func AddWrapped(n *core.IpfsNode, r io.Reader, filename string) (string, *dag.Node, error) {
312 file := files.NewReaderFile(filename, filename, ioutil.NopCloser(r), nil)
309 - fileAdder, err := NewAdder(n.Context(), n, nil)
313 + fileAdder, err := NewAdder(n.Context(), n.Pinning, n.Blockstore, n.DAG)
314 if err != nil {
315 return "", nil, err
316 }
@@ -355,14 +359,14 @@ func (adder *Adder) addNode(node *dag.Node, path string) error {
359 }
360
361 if !adder.Silent {
358 - return outputDagnode(adder.out, path, node)
362 + return outputDagnode(adder.Out, path, node)
363 }
364 return nil
365 }
366
367 // Add the given file while respecting the adder.
368 func (adder *Adder) AddFile(file files.File) error {
365 - adder.unlocker = adder.node.Blockstore.PinLock()
369 + adder.unlocker = adder.blockstore.PinLock()
370 defer func() {
371 adder.unlocker.Unlock()
372 }()
@@ -388,7 +392,7 @@ func (adder *Adder) addFile(file files.File) error {
392 }
393
394 dagnode := &dag.Node{Data: sdata}
391 - _, err = adder.node.DAG.Add(dagnode)
395 + _, err = adder.dagService.Add(dagnode)
396 if err != nil {
397 return err
398 }
@@ -401,7 +405,7 @@ func (adder *Adder) addFile(file files.File) error {
405 // progress updates to the client (over the output channel)
406 var reader io.Reader = file
407 if adder.Progress {
404 - reader = &progressReader{file: file, out: adder.out}
408 + reader = &progressReader{file: file, out: adder.Out}
409 }
410
411 dagnode, err := adder.add(reader)
@@ -445,14 +449,14 @@ func (adder *Adder) addDir(dir files.File) error {
449 }
450
451 func (adder *Adder) maybePauseForGC() error {
448 - if adder.node.Blockstore.GCRequested() {
452 + if adder.blockstore.GCRequested() {
453 err := adder.PinRoot()
454 if err != nil {
455 return err
456 }
457
458 adder.unlocker.Unlock()
455 - adder.unlocker = adder.node.Blockstore.PinLock()
459 + adder.unlocker = adder.blockstore.PinLock()
460 }
461 return nil
462 }
core/coreunix/add_test.go
+2 -1
@@ -54,10 +54,11 @@ func TestAddGCLive(t *testing.T) {
54
55 errs := make(chan error)
56 out := make(chan interface{})
57 - adder, err := NewAdder(context.Background(), node, out)
57 + adder, err := NewAdder(context.Background(), node.Pinning, node.Blockstore, node.DAG)
58 if err != nil {
59 t.Fatal(err)
60 }
61 + adder.Out = out
62
63 dataa := ioutil.NopCloser(bytes.NewBufferString("testfileA"))
64 rfa := files.NewReaderFile("a", "a", dataa, nil)