implement append to trickledag
Jeromy committed
Mar 6, 2015 at 15:17 UTC
298305e4579341f0f86c37c43cb3492a4a1e3c74
2 files changed
+390
-2
importer/trickle/trickle_test.go
+180
-1
@@ -16,6 +16,7 @@ import (
16
merkledag "github.com/jbenet/go-ipfs/merkledag"
17
mdtest "github.com/jbenet/go-ipfs/merkledag/test"
18
pin "github.com/jbenet/go-ipfs/pin"
19
+ ft "github.com/jbenet/go-ipfs/unixfs"
20
uio "github.com/jbenet/go-ipfs/unixfs/io"
21
u "github.com/jbenet/go-ipfs/util"
22
)
@@ -29,7 +30,12 @@ func buildTestDag(r io.Reader, ds merkledag.DAGService, spl chunk.BlockSplitter)
30
Maxlinks: h.DefaultLinksPerBlock,
31
}
32
32
- return TrickleLayout(dbp.New(blkch))
33
+ nd, err := TrickleLayout(dbp.New(blkch))
34
+ if err != nil {
35
+ return nil, err
36
+ }
37
+
38
+ return nd, VerifyTrickleDagStructure(nd, ds, dbp.Maxlinks, layerRepeat)
39
}
40
41
//Test where calls to read are smaller than the chunk size
@@ -441,3 +447,176 @@ func TestSeekingConsistency(t *testing.T) {
447
t.Fatal(err)
448
}
449
}
450
+
451
+func TestAppend(t *testing.T) {
452
+ nbytes := int64(128 * 1024)
453
+ should := make([]byte, nbytes)
454
+ u.NewTimeSeededRand().Read(should)
455
+
456
+ // Reader for half the bytes
457
+ read := bytes.NewReader(should[:nbytes/2])
458
+ ds := mdtest.Mock(t)
459
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
460
+ if err != nil {
461
+ t.Fatal(err)
462
+ }
463
+
464
+ dbp := &h.DagBuilderParams{
465
+ Dagserv: ds,
466
+ Maxlinks: h.DefaultLinksPerBlock,
467
+ }
468
+
469
+ spl := &chunk.SizeSplitter{500}
470
+ blks := spl.Split(bytes.NewReader(should[nbytes/2:]))
471
+
472
+ nnode, err := TrickleAppend(nd, dbp.New(blks))
473
+ if err != nil {
474
+ t.Fatal(err)
475
+ }
476
+
477
+ err = VerifyTrickleDagStructure(nnode, ds, dbp.Maxlinks, layerRepeat)
478
+ if err != nil {
479
+ t.Fatal(err)
480
+ }
481
+
482
+ fread, err := uio.NewDagReader(context.TODO(), nnode, ds)
483
+ if err != nil {
484
+ t.Fatal(err)
485
+ }
486
+
487
+ out, err := ioutil.ReadAll(fread)
488
+ if err != nil {
489
+ t.Fatal(err)
490
+ }
491
+
492
+ err = arrComp(out, should)
493
+ if err != nil {
494
+ t.Fatal(err)
495
+ }
496
+}
497
+
498
+// This test appends one byte at a time to an empty file
499
+func TestMultipleAppends(t *testing.T) {
500
+ ds := mdtest.Mock(t)
501
+
502
+ // TODO: fix small size appends and make this number bigger
503
+ nbytes := int64(1000)
504
+ should := make([]byte, nbytes)
505
+ u.NewTimeSeededRand().Read(should)
506
+
507
+ read := bytes.NewReader(nil)
508
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
509
+ if err != nil {
510
+ t.Fatal(err)
511
+ }
512
+
513
+ dbp := &h.DagBuilderParams{
514
+ Dagserv: ds,
515
+ Maxlinks: 4,
516
+ }
517
+
518
+ spl := &chunk.SizeSplitter{500}
519
+
520
+ for i := 0; i < len(should); i++ {
521
+ blks := spl.Split(bytes.NewReader(should[i : i+1]))
522
+
523
+ nnode, err := TrickleAppend(nd, dbp.New(blks))
524
+ if err != nil {
525
+ t.Fatal(err)
526
+ }
527
+
528
+ err = VerifyTrickleDagStructure(nnode, ds, dbp.Maxlinks, layerRepeat)
529
+ if err != nil {
530
+ t.Fatal(err)
531
+ }
532
+
533
+ fread, err := uio.NewDagReader(context.TODO(), nnode, ds)
534
+ if err != nil {
535
+ t.Fatal(err)
536
+ }
537
+
538
+ out, err := ioutil.ReadAll(fread)
539
+ if err != nil {
540
+ t.Fatal(err)
541
+ }
542
+
543
+ err = arrComp(out, should[:i+1])
544
+ if err != nil {
545
+ t.Fatal(err)
546
+ }
547
+ }
548
+}
549
+
550
+func TestAppendSingleBytesToEmpty(t *testing.T) {
551
+ ds := mdtest.Mock(t)
552
+
553
+ data := []byte("AB")
554
+
555
+ nd := new(merkledag.Node)
556
+ nd.Data = ft.FilePBData(nil, 0)
557
+
558
+ dbp := &h.DagBuilderParams{
559
+ Dagserv: ds,
560
+ Maxlinks: 4,
561
+ }
562
+
563
+ spl := &chunk.SizeSplitter{500}
564
+
565
+ blks := spl.Split(bytes.NewReader(data[:1]))
566
+
567
+ nnode, err := TrickleAppend(nd, dbp.New(blks))
568
+ if err != nil {
569
+ t.Fatal(err)
570
+ }
571
+
572
+ blks = spl.Split(bytes.NewReader(data[1:]))
573
+
574
+ nnode, err = TrickleAppend(nnode, dbp.New(blks))
575
+ if err != nil {
576
+ t.Fatal(err)
577
+ }
578
+
579
+ fread, err := uio.NewDagReader(context.TODO(), nnode, ds)
580
+ if err != nil {
581
+ t.Fatal(err)
582
+ }
583
+
584
+ out, err := ioutil.ReadAll(fread)
585
+ if err != nil {
586
+ t.Fatal(err)
587
+ }
588
+
589
+ fmt.Println(out, data)
590
+ err = arrComp(out, data)
591
+ if err != nil {
592
+ t.Fatal(err)
593
+ }
594
+}
595
+
596
+func printDag(nd *merkledag.Node, ds merkledag.DAGService, indent int) {
597
+ pbd, err := ft.FromBytes(nd.Data)
598
+ if err != nil {
599
+ panic(err)
600
+ }
601
+
602
+ for i := 0; i < indent; i++ {
603
+ fmt.Print(" ")
604
+ }
605
+ fmt.Printf("{size = %d, type = %s, nc = %d", pbd.GetFilesize(), pbd.GetType().String(), len(pbd.GetBlocksizes()))
606
+ if len(nd.Links) > 0 {
607
+ fmt.Println()
608
+ }
609
+ for _, lnk := range nd.Links {
610
+ child, err := lnk.GetNode(ds)
611
+ if err != nil {
612
+ panic(err)
613
+ }
614
+ printDag(child, ds, indent+1)
615
+ }
616
+ if len(nd.Links) > 0 {
617
+ for i := 0; i < indent; i++ {
618
+ fmt.Print(" ")
619
+ }
620
+ }
621
+ fmt.Println("}")
622
+}
importer/trickle/trickledag.go
+210
-1
@@ -1,8 +1,10 @@
1
package trickle
2
3
import (
4
+ "errors"
5
h "github.com/jbenet/go-ipfs/importer/helpers"
6
dag "github.com/jbenet/go-ipfs/merkledag"
7
+ ft "github.com/jbenet/go-ipfs/unixfs"
8
)
9
10
// layerRepeat specifies how many times to append a child tree of a
@@ -41,7 +43,7 @@ func fillTrickleRec(db *h.DagBuilderHelper, node *h.UnixfsNode, depth int) error
43
}
44
45
for i := 1; i < depth && !db.Done(); i++ {
44
- for j := 0; j < layerRepeat; j++ {
46
+ for j := 0; j < layerRepeat && !db.Done(); j++ {
47
next := h.NewUnixfsNode()
48
err := fillTrickleRec(db, next, i)
49
if err != nil {
@@ -56,3 +58,210 @@ func fillTrickleRec(db *h.DagBuilderHelper, node *h.UnixfsNode, depth int) error
58
}
59
return nil
60
}
61
+
62
+// TrickleAppend appends the data in `db` to the dag, using the Trickledag format
63
+func TrickleAppend(base *dag.Node, db *h.DagBuilderHelper) (*dag.Node, error) {
64
+ // Convert to unixfs node for working with easily
65
+ ufsn, err := h.NewUnixfsNodeFromDag(base)
66
+ if err != nil {
67
+ return nil, err
68
+ }
69
+
70
+ // Get depth of this 'tree'
71
+ n, j := trickleDepthInfo(ufsn, db.Maxlinks())
72
+ if n == 0 {
73
+ // If direct blocks not filled...
74
+ err := db.FillNodeLayer(ufsn)
75
+ if err != nil {
76
+ return nil, err
77
+ }
78
+
79
+ if db.Done() {
80
+ return ufsn.GetDagNode()
81
+ }
82
+
83
+ // If continuing, our depth has increased by one
84
+ n++
85
+ }
86
+
87
+ err = appendFillLastChild(ufsn, n-1, j, db)
88
+ if err != nil {
89
+ return nil, err
90
+ }
91
+
92
+ // after appendFillLastChild, our depth is now increased by one
93
+ if !db.Done() {
94
+ n++
95
+ }
96
+
97
+ // Now, continue filling out tree like normal
98
+ for i := n; !db.Done(); i++ {
99
+ for j := 0; j < layerRepeat && !db.Done(); j++ {
100
+ next := h.NewUnixfsNode()
101
+ err := fillTrickleRec(db, next, i)
102
+ if err != nil {
103
+ return nil, err
104
+ }
105
+
106
+ err = ufsn.AddChild(next, db)
107
+ if err != nil {
108
+ return nil, err
109
+ }
110
+ }
111
+ }
112
+
113
+ return ufsn.GetDagNode()
114
+}
115
+
116
+func appendFillLastChild(ufsn *h.UnixfsNode, depth int, layerFill int, db *h.DagBuilderHelper) error {
117
+ if ufsn.NumChildren() > db.Maxlinks() {
118
+ // Recursive step, grab last child
119
+ last := ufsn.NumChildren() - 1
120
+ lastChild, err := ufsn.GetChild(last, db.GetDagServ())
121
+ if err != nil {
122
+ return err
123
+ }
124
+
125
+ // Fill out last child (may not be full tree)
126
+ nchild, err := trickleAppendRec(lastChild, db, depth-1)
127
+ if err != nil {
128
+ return err
129
+ }
130
+
131
+ // Update changed child in parent node
132
+ ufsn.RemoveChild(last)
133
+ err = ufsn.AddChild(nchild, db)
134
+ if err != nil {
135
+ return err
136
+ }
137
+
138
+ // Partially filled depth layer
139
+ if layerFill != 0 {
140
+ for ; layerFill < layerRepeat && !db.Done(); layerFill++ {
141
+ next := h.NewUnixfsNode()
142
+ err := fillTrickleRec(db, next, depth)
143
+ if err != nil {
144
+ return err
145
+ }
146
+
147
+ err = ufsn.AddChild(next, db)
148
+ if err != nil {
149
+ return err
150
+ }
151
+ }
152
+ }
153
+ }
154
+
155
+ return nil
156
+}
157
+
158
+func trickleAppendRec(ufsn *h.UnixfsNode, db *h.DagBuilderHelper, depth int) (*h.UnixfsNode, error) {
159
+ if depth == 0 || db.Done() {
160
+ return ufsn, nil
161
+ }
162
+
163
+ // Get depth of this 'tree'
164
+ n, j := trickleDepthInfo(ufsn, db.Maxlinks())
165
+ if n == 0 {
166
+ // If direct blocks not filled...
167
+ err := db.FillNodeLayer(ufsn)
168
+ if err != nil {
169
+ return nil, err
170
+ }
171
+ n++
172
+ }
173
+
174
+ // If at correct depth, no need to continue
175
+ if n == depth {
176
+ return ufsn, nil
177
+ }
178
+
179
+ err := appendFillLastChild(ufsn, n, j, db)
180
+ if err != nil {
181
+ return nil, err
182
+ }
183
+
184
+ // after appendFillLastChild, our depth is now increased by one
185
+ if !db.Done() {
186
+ n++
187
+ }
188
+
189
+ // Now, continue filling out tree like normal
190
+ for i := n; i < depth && !db.Done(); i++ {
191
+ for j := 0; j < layerRepeat && !db.Done(); j++ {
192
+ next := h.NewUnixfsNode()
193
+ err := fillTrickleRec(db, next, i)
194
+ if err != nil {
195
+ return nil, err
196
+ }
197
+
198
+ err = ufsn.AddChild(next, db)
199
+ if err != nil {
200
+ return nil, err
201
+ }
202
+ }
203
+ }
204
+
205
+ return ufsn, nil
206
+}
207
+
208
+func trickleDepthInfo(node *h.UnixfsNode, maxlinks int) (int, int) {
209
+ n := node.NumChildren()
210
+ if n < maxlinks {
211
+ return 0, 0
212
+ }
213
+
214
+ return ((n - maxlinks) / layerRepeat) + 1, (n - maxlinks) % layerRepeat
215
+}
216
+
217
+// VerifyTrickleDagStructure checks that the given dag matches exactly the trickle dag datastructure
218
+// layout
219
+func VerifyTrickleDagStructure(nd *dag.Node, ds dag.DAGService, direct int, layerRepeat int) error {
220
+ return verifyTDagRec(nd, -1, direct, layerRepeat, ds)
221
+}
222
+
223
+// Recursive call for verifying the structure of a trickledag
224
+func verifyTDagRec(nd *dag.Node, depth, direct, layerRepeat int, ds dag.DAGService) error {
225
+ if depth == 0 {
226
+ // zero depth dag is raw data block
227
+ if len(nd.Links) > 0 {
228
+ return errors.New("expected direct block")
229
+ }
230
+
231
+ pbn, err := ft.FromBytes(nd.Data)
232
+ if err != nil {
233
+ return err
234
+ }
235
+
236
+ if pbn.GetType() != ft.TRaw {
237
+ return errors.New("Expected raw block")
238
+ }
239
+ return nil
240
+ }
241
+
242
+ for i := 0; i < len(nd.Links); i++ {
243
+ child, err := nd.Links[i].GetNode(ds)
244
+ if err != nil {
245
+ return nil
246
+ }
247
+
248
+ if i < direct {
249
+ // Direct blocks
250
+ err := verifyTDagRec(child, 0, direct, layerRepeat, ds)
251
+ if err != nil {
252
+ return err
253
+ }
254
+ } else {
255
+ // Recursive trickle dags
256
+ rdepth := ((i - direct) / layerRepeat) + 1
257
+ if rdepth >= depth && depth > 0 {
258
+ return errors.New("Child dag was too deep!")
259
+ }
260
+ err := verifyTDagRec(child, rdepth, direct, layerRepeat, ds)
261
+ if err != nil {
262
+ return err
263
+ }
264
+ }
265
+ }
266
+ return nil
267
+}