@cryptotaxi247 / kubo / commits / 18ada93ec

rewrite add command to use dagwriter, moved a pinner into the dagwriter for inline pinning

Jeromy committed Oct 30, 2014 at 03:13 UTC 18ada93ec333dfd475139395bbe51e48fa5e61fb
4 files changed +79 -4
core/commands/add.go
+2 -2
@@ -87,7 +87,7 @@ func addDir(n *core.IpfsNode, fpath string, depth int, out io.Writer) (*dag.Node
87 }
88
89 func addFile(n *core.IpfsNode, fpath string, depth int, out io.Writer) (*dag.Node, error) {
90 - root, err := importer.NewDagFromFile(fpath)
90 + root, err := importer.NewDagFromFileWServer(fpath, n.DAG, n.Pinning)
91 if err != nil {
92 return nil, err
93 }
@@ -98,7 +98,7 @@ func addFile(n *core.IpfsNode, fpath string, depth int, out io.Writer) (*dag.Nod
98 log.Info("adding subblock: %s %s", l.Name, l.Hash.B58String())
99 }
100
101 - return root, addNode(n, root, fpath, out)
101 + return root, nil
102 }
103
104 // addNode adds the node to the graph + local storage
importer/importer.go
+41
@@ -1,13 +1,16 @@
1 package importer
2
3 import (
4 + "errors"
5 "fmt"
6 "io"
7 "os"
8
9 "github.com/jbenet/go-ipfs/importer/chunk"
10 dag "github.com/jbenet/go-ipfs/merkledag"
11 + "github.com/jbenet/go-ipfs/pin"
12 ft "github.com/jbenet/go-ipfs/unixfs"
13 + uio "github.com/jbenet/go-ipfs/unixfs/io"
14 "github.com/jbenet/go-ipfs/util"
15 )
16
@@ -72,3 +75,41 @@ func NewDagFromFile(fpath string) (*dag.Node, error) {
75
76 return NewDagFromReader(f)
77 }
78 +
79 +func NewDagFromFileWServer(fpath string, dserv dag.DAGService, p pin.Pinner) (*dag.Node, error) {
80 + stat, err := os.Stat(fpath)
81 + if err != nil {
82 + return nil, err
83 + }
84 +
85 + if stat.IsDir() {
86 + return nil, fmt.Errorf("`%s` is a directory", fpath)
87 + }
88 +
89 + f, err := os.Open(fpath)
90 + if err != nil {
91 + return nil, err
92 + }
93 + defer f.Close()
94 +
95 + return NewDagFromReaderWServer(f, dserv, p)
96 +}
97 +
98 +func NewDagFromReaderWServer(r io.Reader, dserv dag.DAGService, p pin.Pinner) (*dag.Node, error) {
99 + dw := uio.NewDagWriter(dserv, chunk.DefaultSplitter)
100 +
101 + mp, ok := p.(pin.ManualPinner)
102 + if !ok {
103 + return nil, errors.New("Needed to be passed a manual pinner!")
104 + }
105 + dw.Pinner = mp
106 + _, err := io.Copy(dw, r)
107 + if err != nil {
108 + return nil, err
109 + }
110 + err = dw.Close()
111 + if err != nil {
112 + return nil, err
113 + }
114 + return dw.GetNode(), nil
115 +}
pin/pin.go
+26
@@ -20,6 +20,14 @@ var recursePinDatastoreKey = ds.NewKey("/local/pins/recursive/keys")
20 var directPinDatastoreKey = ds.NewKey("/local/pins/direct/keys")
21 var indirectPinDatastoreKey = ds.NewKey("/local/pins/indirect/keys")
22
23 +type PinMode int
24 +
25 +const (
26 + Recursive PinMode = iota
27 + Direct
28 + Indirect
29 +)
30 +
31 type Pinner interface {
32 IsPinned(util.Key) bool
33 Pin(*mdag.Node, bool) error
@@ -27,6 +35,13 @@ type Pinner interface {
35 Flush() error
36 }
37
38 +// ManualPinner is for manually editing the pin structure
39 +// Use with care
40 +type ManualPinner interface {
41 + PinWithMode(util.Key, PinMode)
42 + Pinner
43 +}
44 +
45 type pinner struct {
46 lock sync.RWMutex
47 recursePin set.BlockSet
@@ -228,3 +243,14 @@ func loadSet(d ds.Datastore, k ds.Key, val interface{}) error {
243 }
244 return json.Unmarshal(bf, val)
245 }
246 +
247 +func (p *pinner) PinWithMode(k util.Key, mode PinMode) {
248 + switch mode {
249 + case Recursive:
250 + p.recursePin.AddBlock(k)
251 + case Direct:
252 + p.directPin.AddBlock(k)
253 + case Indirect:
254 + p.indirPin.Increment(k)
255 + }
256 +}
unixfs/io/dagwriter.go
+10 -2
@@ -3,6 +3,7 @@ package io
3 import (
4 "github.com/jbenet/go-ipfs/importer/chunk"
5 dag "github.com/jbenet/go-ipfs/merkledag"
6 + "github.com/jbenet/go-ipfs/pin"
7 ft "github.com/jbenet/go-ipfs/unixfs"
8 "github.com/jbenet/go-ipfs/util"
9 )
@@ -17,6 +18,7 @@ type DagWriter struct {
18 done chan struct{}
19 splitter chunk.BlockSplitter
20 seterr error
21 + Pinner pin.ManualPinner
22 }
23
24 func NewDagWriter(ds dag.DAGService, splitter chunk.BlockSplitter) *DagWriter {
@@ -48,7 +50,10 @@ func (dw *DagWriter) startSplitter() {
50 // Store the block size in the root node
51 mbf.AddBlockSize(uint64(len(blkData)))
52 node := &dag.Node{Data: ft.WrapData(blkData)}
51 - _, err := dw.dagserv.Add(node)
53 + nk, err := dw.dagserv.Add(node)
54 + if dw.Pinner != nil {
55 + dw.Pinner.PinWithMode(nk, pin.Indirect)
56 + }
57 if err != nil {
58 dw.seterr = err
59 log.Critical("Got error adding created node to dagservice: %s", err)
@@ -75,12 +80,15 @@ func (dw *DagWriter) startSplitter() {
80 root.Data = data
81
82 // Add root node to the dagservice
78 - _, err = dw.dagserv.Add(root)
83 + rootk, err := dw.dagserv.Add(root)
84 if err != nil {
85 dw.seterr = err
86 log.Critical("Got error adding created node to dagservice: %s", err)
87 return
88 }
89 + if dw.Pinner != nil {
90 + dw.Pinner.PinWithMode(rootk, pin.Recursive)
91 + }
92 dw.node = root
93 dw.done <- struct{}{}
94 }