@cryptotaxi247 / kubo / commits / dd928a2b1

change pinning to happen in a callback

Jeromy committed May 19, 2015 at 08:52 UTC dd928a2b1d9bfde65653d3053aef53249fd8a164
12 files changed +94 -46
core/commands/add.go
+15 -4
@@ -7,7 +7,6 @@ import (
7 "strings"
8
9 "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/cheggaaa/pb"
10 - context "github.com/ipfs/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
10
11 cmds "github.com/ipfs/go-ipfs/commands"
12 files "github.com/ipfs/go-ipfs/commands/files"
@@ -16,6 +15,7 @@ import (
15 importer "github.com/ipfs/go-ipfs/importer"
16 "github.com/ipfs/go-ipfs/importer/chunk"
17 dag "github.com/ipfs/go-ipfs/merkledag"
18 + pin "github.com/ipfs/go-ipfs/pin"
19 ft "github.com/ipfs/go-ipfs/unixfs"
20 u "github.com/ipfs/go-ipfs/util"
21 )
@@ -113,12 +113,16 @@ remains to be implemented.
113 return
114 }
115
116 - err = n.Pinning.Pin(context.Background(), rootnd, true)
116 + rnk, err := rootnd.Key()
117 if err != nil {
118 res.SetError(err, cmds.ErrNormal)
119 return
120 }
121
122 + mp := n.Pinning.GetManual()
123 + mp.RemovePinWithMode(rnk, pin.Indirect)
124 + mp.PinWithMode(rnk, pin.Recursive)
125 +
126 err = n.Pinning.Flush()
127 if err != nil {
128 res.SetError(err, cmds.ErrNormal)
@@ -214,7 +218,12 @@ remains to be implemented.
218 }
219
220 func add(n *core.IpfsNode, reader io.Reader) (*dag.Node, error) {
217 - node, err := importer.BuildDagFromReader(reader, n.DAG, nil, chunk.DefaultSplitter)
221 + node, err := importer.BuildDagFromReader(
222 + reader,
223 + n.DAG,
224 + chunk.DefaultSplitter,
225 + importer.PinIndirectCB(n.Pinning.GetManual()),
226 + )
227 if err != nil {
228 return nil, err
229 }
@@ -290,11 +299,13 @@ func addDir(n *core.IpfsNode, dir files.File, out chan interface{}, progress boo
299 return nil, err
300 }
301
293 - _, err = n.DAG.Add(tree)
302 + k, err := n.DAG.Add(tree)
303 if err != nil {
304 return nil, err
305 }
306
307 + n.Pinning.GetManual().PinWithMode(k, pin.Indirect)
308 +
309 return tree, nil
310 }
311
core/corehttp/gateway_handler.go
+1 -1
@@ -72,7 +72,7 @@ func (i *gatewayHandler) newDagFromReader(r io.Reader) (*dag.Node, error) {
72 // TODO(cryptix): change and remove this helper once PR1136 is merged
73 // return ufs.AddFromReader(i.node, r.Body)
74 return importer.BuildDagFromReader(
75 - r, i.node.DAG, i.node.Pinning.GetManual(), chunk.DefaultSplitter)
75 + r, i.node.DAG, chunk.DefaultSplitter, importer.BasicPinnerCB(i.node.Pinning.GetManual()))
76 }
77
78 // TODO(btc): break this apart into separate handlers using a more expressive muxer
core/coreunix/add.go
+22 -15
@@ -1,7 +1,6 @@
1 package coreunix
2
3 import (
4 - "errors"
4 "io"
5 "io/ioutil"
6 "os"
@@ -18,6 +17,7 @@ import (
17 "github.com/ipfs/go-ipfs/pin"
18 "github.com/ipfs/go-ipfs/thirdparty/eventlog"
19 unixfs "github.com/ipfs/go-ipfs/unixfs"
20 + u "github.com/ipfs/go-ipfs/util"
21 )
22
23 var log = eventlog.Logger("coreunix")
@@ -29,15 +29,12 @@ func Add(n *core.IpfsNode, r io.Reader) (string, error) {
29 dagNode, err := importer.BuildDagFromReader(
30 r,
31 n.DAG,
32 - n.Pinning.GetManual(), // Fix this interface
32 chunk.DefaultSplitter,
33 + importer.BasicPinnerCB(n.Pinning.GetManual()),
34 )
35 if err != nil {
36 return "", err
37 }
38 - if err := n.Pinning.Flush(); err != nil {
39 - return "", err
40 - }
38 k, err := dagNode.Key()
39 if err != nil {
40 return "", err
@@ -53,18 +50,28 @@ func AddR(n *core.IpfsNode, root string) (key string, err error) {
50 return "", err
51 }
52 defer f.Close()
53 +
54 ff, err := files.NewSerialFile(root, f)
55 if err != nil {
56 return "", err
57 }
58 +
59 dagnode, err := addFile(n, ff)
60 if err != nil {
61 return "", err
62 }
63 +
64 k, err := dagnode.Key()
65 if err != nil {
66 return "", err
67 }
68 +
69 + n.Pinning.GetManual().RemovePinWithMode(k, pin.Indirect)
70 + err = n.Pinning.Flush()
71 + if err != nil {
72 + return "", err
73 + }
74 +
75 return k.String(), nil
76 }
77
@@ -87,17 +94,17 @@ func AddWrapped(n *core.IpfsNode, r io.Reader, filename string) (string, *merkle
94 }
95
96 func add(n *core.IpfsNode, reader io.Reader) (*merkledag.Node, error) {
90 - mp, ok := n.Pinning.(pin.ManualPinner)
91 - if !ok {
92 - return nil, errors.New("invalid pinner type! expected manual pinner")
93 - }
94 -
95 - node, err := importer.BuildDagFromReader(reader, n.DAG, mp, chunk.DefaultSplitter)
96 - if err != nil {
97 - return nil, err
98 - }
97 + mp := n.Pinning.GetManual()
98
100 - err = n.Pinning.Flush()
99 + node, err := importer.BuildDagFromReader(
100 + reader,
101 + n.DAG,
102 + chunk.DefaultSplitter,
103 + func(k u.Key, root bool) error {
104 + mp.PinWithMode(k, pin.Indirect)
105 + return nil
106 + },
107 + )
108 if err != nil {
109 return nil, err
110 }
core/coreunix/metadata_test.go
+1 -1
@@ -37,7 +37,7 @@ func TestMetadata(t *testing.T) {
37 data := make([]byte, 1000)
38 u.NewTimeSeededRand().Read(data)
39 r := bytes.NewReader(data)
40 - nd, err := importer.BuildDagFromReader(r, ds, nil, chunk.DefaultSplitter)
40 + nd, err := importer.BuildDagFromReader(r, ds, chunk.DefaultSplitter, nil)
41 if err != nil {
42 t.Fatal(err)
43 }
fuse/readonly/ipfs_test.go
+1 -1
@@ -34,7 +34,7 @@ func randObj(t *testing.T, nd *core.IpfsNode, size int64) (*dag.Node, []byte) {
34 buf := make([]byte, size)
35 u.NewTimeSeededRand().Read(buf)
36 read := bytes.NewReader(buf)
37 - obj, err := importer.BuildTrickleDagFromReader(read, nd.DAG, nil, chunk.DefaultSplitter)
37 + obj, err := importer.BuildTrickleDagFromReader(read, nd.DAG, chunk.DefaultSplitter, nil)
38 if err != nil {
39 t.Fatal(err)
40 }
importer/helpers/dagbuilder.go
+18 -9
@@ -3,8 +3,13 @@ package helpers
3 import (
4 dag "github.com/ipfs/go-ipfs/merkledag"
5 "github.com/ipfs/go-ipfs/pin"
6 + u "github.com/ipfs/go-ipfs/util"
7 )
8
9 +type BlockCB func(u.Key, bool) error
10 +
11 +var nilFunc BlockCB = func(_ u.Key, _ bool) error { return nil }
12 +
13 // DagBuilderHelper wraps together a bunch of objects needed to
14 // efficiently create unixfs dag trees
15 type DagBuilderHelper struct {
@@ -13,6 +18,7 @@ type DagBuilderHelper struct {
18 in <-chan []byte
19 nextData []byte // the next item to return.
20 maxlinks int
21 + bcb BlockCB
22 }
23
24 type DagBuilderParams struct {
@@ -22,18 +28,23 @@ type DagBuilderParams struct {
28 // DAGService to write blocks to (required)
29 Dagserv dag.DAGService
30
25 - // Pinner to use for pinning files (optionally nil)
26 - Pinner pin.ManualPinner
31 + // Callback for each block added
32 + BlockCB BlockCB
33 }
34
35 // Generate a new DagBuilderHelper from the given params, using 'in' as a
36 // data source
37 func (dbp *DagBuilderParams) New(in <-chan []byte) *DagBuilderHelper {
38 + bcb := dbp.BlockCB
39 + if bcb == nil {
40 + bcb = nilFunc
41 + }
42 +
43 return &DagBuilderHelper{
44 dserv: dbp.Dagserv,
34 - mp: dbp.Pinner,
45 in: in,
46 maxlinks: dbp.Maxlinks,
47 + bcb: bcb,
48 }
49 }
50
@@ -130,12 +141,10 @@ func (db *DagBuilderHelper) Add(node *UnixfsNode) (*dag.Node, error) {
141 return nil, err
142 }
143
133 - if db.mp != nil {
134 - db.mp.PinWithMode(key, pin.Recursive)
135 - err := db.mp.Flush()
136 - if err != nil {
137 - return nil, err
138 - }
144 + // block callback
145 + err = db.bcb(key, true)
146 + if err != nil {
147 + return nil, err
148 }
149
150 return dn, nil
importer/helpers/helpers.go
+3 -2
@@ -113,8 +113,9 @@ func (n *UnixfsNode) AddChild(child *UnixfsNode, db *DagBuilderHelper) error {
113 }
114
115 // Pin the child node indirectly
116 - if db.mp != nil {
117 - db.mp.PinWithMode(childkey, pin.Indirect)
116 + err = db.bcb(childkey, false)
117 + if err != nil {
118 + return err
119 }
120
121 return nil
importer/importer.go
+26 -7
@@ -13,10 +13,10 @@ import (
13 trickle "github.com/ipfs/go-ipfs/importer/trickle"
14 dag "github.com/ipfs/go-ipfs/merkledag"
15 "github.com/ipfs/go-ipfs/pin"
16 - "github.com/ipfs/go-ipfs/util"
16 + u "github.com/ipfs/go-ipfs/util"
17 )
18
19 -var log = util.Logger("importer")
19 +var log = u.Logger("importer")
20
21 // Builds a DAG from the given file, writing created blocks to disk as they are
22 // created
@@ -36,31 +36,50 @@ func BuildDagFromFile(fpath string, ds dag.DAGService, mp pin.ManualPinner) (*da
36 }
37 defer f.Close()
38
39 - return BuildDagFromReader(f, ds, mp, chunk.DefaultSplitter)
39 + return BuildDagFromReader(f, ds, chunk.DefaultSplitter, BasicPinnerCB(mp))
40 }
41
42 -func BuildDagFromReader(r io.Reader, ds dag.DAGService, mp pin.ManualPinner, spl chunk.BlockSplitter) (*dag.Node, error) {
42 +func BuildDagFromReader(r io.Reader, ds dag.DAGService, spl chunk.BlockSplitter, bcb h.BlockCB) (*dag.Node, error) {
43 // Start the splitter
44 blkch := spl.Split(r)
45
46 dbp := h.DagBuilderParams{
47 Dagserv: ds,
48 Maxlinks: h.DefaultLinksPerBlock,
49 - Pinner: mp,
49 + BlockCB: bcb,
50 }
51
52 return bal.BalancedLayout(dbp.New(blkch))
53 }
54
55 -func BuildTrickleDagFromReader(r io.Reader, ds dag.DAGService, mp pin.ManualPinner, spl chunk.BlockSplitter) (*dag.Node, error) {
55 +func BuildTrickleDagFromReader(r io.Reader, ds dag.DAGService, spl chunk.BlockSplitter, bcb h.BlockCB) (*dag.Node, error) {
56 // Start the splitter
57 blkch := spl.Split(r)
58
59 dbp := h.DagBuilderParams{
60 Dagserv: ds,
61 Maxlinks: h.DefaultLinksPerBlock,
62 - Pinner: mp,
62 + BlockCB: bcb,
63 }
64
65 return trickle.TrickleLayout(dbp.New(blkch))
66 }
67 +
68 +func BasicPinnerCB(p pin.ManualPinner) h.BlockCB {
69 + return func(k u.Key, root bool) error {
70 + if root {
71 + p.PinWithMode(k, pin.Recursive)
72 + return p.Flush()
73 + } else {
74 + p.PinWithMode(k, pin.Indirect)
75 + return nil
76 + }
77 + }
78 +}
79 +
80 +func PinIndirectCB(p pin.ManualPinner) h.BlockCB {
81 + return func(k u.Key, root bool) error {
82 + p.PinWithMode(k, pin.Indirect)
83 + return nil
84 + }
85 +}
importer/importer_test.go
+3 -3
@@ -17,7 +17,7 @@ import (
17 func getBalancedDag(t testing.TB, size int64, blksize int) (*dag.Node, dag.DAGService) {
18 ds := mdtest.Mock(t)
19 r := io.LimitReader(u.NewTimeSeededRand(), size)
20 - nd, err := BuildDagFromReader(r, ds, nil, &chunk.SizeSplitter{blksize})
20 + nd, err := BuildDagFromReader(r, ds, &chunk.SizeSplitter{blksize}, nil)
21 if err != nil {
22 t.Fatal(err)
23 }
@@ -27,7 +27,7 @@ func getBalancedDag(t testing.TB, size int64, blksize int) (*dag.Node, dag.DAGSe
27 func getTrickleDag(t testing.TB, size int64, blksize int) (*dag.Node, dag.DAGService) {
28 ds := mdtest.Mock(t)
29 r := io.LimitReader(u.NewTimeSeededRand(), size)
30 - nd, err := BuildTrickleDagFromReader(r, ds, nil, &chunk.SizeSplitter{blksize})
30 + nd, err := BuildTrickleDagFromReader(r, ds, &chunk.SizeSplitter{blksize}, nil)
31 if err != nil {
32 t.Fatal(err)
33 }
@@ -40,7 +40,7 @@ func TestBalancedDag(t *testing.T) {
40 u.NewTimeSeededRand().Read(buf)
41 r := bytes.NewReader(buf)
42
43 - nd, err := BuildDagFromReader(r, ds, nil, chunk.DefaultSplitter)
43 + nd, err := BuildDagFromReader(r, ds, chunk.DefaultSplitter, nil)
44 if err != nil {
45 t.Fatal(err)
46 }
merkledag/merkledag_test.go
+1 -1
@@ -156,7 +156,7 @@ func runBatchFetchTest(t *testing.T, read io.Reader) {
156
157 spl := &chunk.SizeSplitter{512}
158
159 - root, err := imp.BuildDagFromReader(read, dagservs[0], nil, spl)
159 + root, err := imp.BuildDagFromReader(read, dagservs[0], spl, nil)
160 if err != nil {
161 t.Fatal(err)
162 }
unixfs/mod/dagmodifier.go
+2 -1
@@ -11,6 +11,7 @@ import (
11 mh "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
12 context "github.com/ipfs/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
13
14 + imp "github.com/ipfs/go-ipfs/importer"
15 chunk "github.com/ipfs/go-ipfs/importer/chunk"
16 help "github.com/ipfs/go-ipfs/importer/helpers"
17 trickle "github.com/ipfs/go-ipfs/importer/trickle"
@@ -308,7 +309,7 @@ func (dm *DagModifier) appendData(node *mdag.Node, blks <-chan []byte) (*mdag.No
309 dbp := &help.DagBuilderParams{
310 Dagserv: dm.dagserv,
311 Maxlinks: help.DefaultLinksPerBlock,
311 - Pinner: dm.mp,
312 + BlockCB: imp.BasicPinnerCB(dm.mp),
313 }
314
315 return trickle.TrickleAppend(node, dbp.New(blks))
unixfs/mod/dagmodifier_test.go
+1 -1
@@ -52,7 +52,7 @@ func getMockDagServAndBstore(t testing.TB) (mdag.DAGService, blockstore.Blocksto
52
53 func getNode(t testing.TB, dserv mdag.DAGService, size int64, pinner pin.ManualPinner) ([]byte, *mdag.Node) {
54 in := io.LimitReader(u.NewTimeSeededRand(), size)
55 - node, err := imp.BuildTrickleDagFromReader(in, dserv, pinner, &chunk.SizeSplitter{500})
55 + node, err := imp.BuildTrickleDagFromReader(in, dserv, &chunk.SizeSplitter{500}, imp.BasicPinnerCB(pinner))
56 if err != nil {
57 t.Fatal(err)
58 }