refactor importer package with trickle and balanced dag generation
Jeromy committed
Feb 4, 2015 at 03:53 UTC
bc79ae17a1987c19e9cd258c4c12d4ba9c3470c6
11 files changed
+1313
-863
importer/balanced/balanced_test.go
new
+443
@@ -0,0 +1,443 @@
1
+package balanced
2
+
3
+import (
4
+ "bytes"
5
+ "crypto/rand"
6
+ "fmt"
7
+ "io"
8
+ "io/ioutil"
9
+ mrand "math/rand"
10
+ "os"
11
+ "testing"
12
+
13
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
14
+ chunk "github.com/jbenet/go-ipfs/importer/chunk"
15
+ h "github.com/jbenet/go-ipfs/importer/helpers"
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
+ uio "github.com/jbenet/go-ipfs/unixfs/io"
20
+ u "github.com/jbenet/go-ipfs/util"
21
+)
22
+
23
+func buildTestDag(r io.Reader, ds merkledag.DAGService, spl chunk.BlockSplitter) (*merkledag.Node, error) {
24
+ // Start the splitter
25
+ blkch := spl.Split(r)
26
+
27
+ dbp := h.DagBuilderParams{
28
+ Dagserv: ds,
29
+ Maxlinks: h.DefaultLinksPerBlock,
30
+ }
31
+
32
+ return BalancedLayout(dbp.New(blkch))
33
+}
34
+
35
+//Test where calls to read are smaller than the chunk size
36
+func TestSizeBasedSplit(t *testing.T) {
37
+ if testing.Short() {
38
+ t.SkipNow()
39
+ }
40
+ bs := &chunk.SizeSplitter{Size: 512}
41
+ testFileConsistency(t, bs, 32*512)
42
+ bs = &chunk.SizeSplitter{Size: 4096}
43
+ testFileConsistency(t, bs, 32*4096)
44
+
45
+ // Uneven offset
46
+ testFileConsistency(t, bs, 31*4095)
47
+}
48
+
49
+func dup(b []byte) []byte {
50
+ o := make([]byte, len(b))
51
+ copy(o, b)
52
+ return o
53
+}
54
+
55
+func testFileConsistency(t *testing.T, bs chunk.BlockSplitter, nbytes int) {
56
+ should := make([]byte, nbytes)
57
+ u.NewTimeSeededRand().Read(should)
58
+
59
+ read := bytes.NewReader(should)
60
+ ds := mdtest.Mock(t)
61
+ nd, err := buildTestDag(read, ds, bs)
62
+ if err != nil {
63
+ t.Fatal(err)
64
+ }
65
+
66
+ r, err := uio.NewDagReader(context.Background(), nd, ds)
67
+ if err != nil {
68
+ t.Fatal(err)
69
+ }
70
+
71
+ out, err := ioutil.ReadAll(r)
72
+ if err != nil {
73
+ t.Fatal(err)
74
+ }
75
+
76
+ err = arrComp(out, should)
77
+ if err != nil {
78
+ t.Fatal(err)
79
+ }
80
+}
81
+
82
+func TestBuilderConsistency(t *testing.T) {
83
+ nbytes := 100000
84
+ buf := new(bytes.Buffer)
85
+ io.CopyN(buf, u.NewTimeSeededRand(), int64(nbytes))
86
+ should := dup(buf.Bytes())
87
+ dagserv := mdtest.Mock(t)
88
+ nd, err := buildTestDag(buf, dagserv, chunk.DefaultSplitter)
89
+ if err != nil {
90
+ t.Fatal(err)
91
+ }
92
+ r, err := uio.NewDagReader(context.Background(), nd, dagserv)
93
+ if err != nil {
94
+ t.Fatal(err)
95
+ }
96
+
97
+ out, err := ioutil.ReadAll(r)
98
+ if err != nil {
99
+ t.Fatal(err)
100
+ }
101
+
102
+ err = arrComp(out, should)
103
+ if err != nil {
104
+ t.Fatal(err)
105
+ }
106
+}
107
+
108
+func arrComp(a, b []byte) error {
109
+ if len(a) != len(b) {
110
+ return fmt.Errorf("Arrays differ in length. %d != %d", len(a), len(b))
111
+ }
112
+ for i, v := range a {
113
+ if v != b[i] {
114
+ return fmt.Errorf("Arrays differ at index: %d", i)
115
+ }
116
+ }
117
+ return nil
118
+}
119
+
120
+func TestMaybeRabinConsistency(t *testing.T) {
121
+ if testing.Short() {
122
+ t.SkipNow()
123
+ }
124
+ testFileConsistency(t, chunk.NewMaybeRabin(4096), 256*4096)
125
+}
126
+
127
+func TestRabinBlockSize(t *testing.T) {
128
+ if testing.Short() {
129
+ t.SkipNow()
130
+ }
131
+ buf := new(bytes.Buffer)
132
+ nbytes := 1024 * 1024
133
+ io.CopyN(buf, rand.Reader, int64(nbytes))
134
+ rab := chunk.NewMaybeRabin(4096)
135
+ blkch := rab.Split(buf)
136
+
137
+ var blocks [][]byte
138
+ for b := range blkch {
139
+ blocks = append(blocks, b)
140
+ }
141
+
142
+ fmt.Printf("Avg block size: %d\n", nbytes/len(blocks))
143
+
144
+}
145
+
146
+type dagservAndPinner struct {
147
+ ds merkledag.DAGService
148
+ mp pin.ManualPinner
149
+}
150
+
151
+func TestIndirectBlocks(t *testing.T) {
152
+ splitter := &chunk.SizeSplitter{512}
153
+ nbytes := 1024 * 1024
154
+ buf := make([]byte, nbytes)
155
+ u.NewTimeSeededRand().Read(buf)
156
+
157
+ read := bytes.NewReader(buf)
158
+
159
+ ds := mdtest.Mock(t)
160
+ dag, err := buildTestDag(read, ds, splitter)
161
+ if err != nil {
162
+ t.Fatal(err)
163
+ }
164
+
165
+ reader, err := uio.NewDagReader(context.Background(), dag, ds)
166
+ if err != nil {
167
+ t.Fatal(err)
168
+ }
169
+
170
+ out, err := ioutil.ReadAll(reader)
171
+ if err != nil {
172
+ t.Fatal(err)
173
+ }
174
+
175
+ if !bytes.Equal(out, buf) {
176
+ t.Fatal("Not equal!")
177
+ }
178
+}
179
+
180
+func TestSeekingBasic(t *testing.T) {
181
+ nbytes := int64(10 * 1024)
182
+ should := make([]byte, nbytes)
183
+ u.NewTimeSeededRand().Read(should)
184
+
185
+ read := bytes.NewReader(should)
186
+ ds := mdtest.Mock(t)
187
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
188
+ if err != nil {
189
+ t.Fatal(err)
190
+ }
191
+
192
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
193
+ if err != nil {
194
+ t.Fatal(err)
195
+ }
196
+
197
+ start := int64(4000)
198
+ n, err := rs.Seek(start, os.SEEK_SET)
199
+ if err != nil {
200
+ t.Fatal(err)
201
+ }
202
+ if n != start {
203
+ t.Fatal("Failed to seek to correct offset")
204
+ }
205
+
206
+ out, err := ioutil.ReadAll(rs)
207
+ if err != nil {
208
+ t.Fatal(err)
209
+ }
210
+
211
+ err = arrComp(out, should[start:])
212
+ if err != nil {
213
+ t.Fatal(err)
214
+ }
215
+}
216
+
217
+func TestSeekToBegin(t *testing.T) {
218
+ nbytes := int64(10 * 1024)
219
+ should := make([]byte, nbytes)
220
+ u.NewTimeSeededRand().Read(should)
221
+
222
+ read := bytes.NewReader(should)
223
+ ds := mdtest.Mock(t)
224
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
225
+ if err != nil {
226
+ t.Fatal(err)
227
+ }
228
+
229
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
230
+ if err != nil {
231
+ t.Fatal(err)
232
+ }
233
+
234
+ n, err := io.CopyN(ioutil.Discard, rs, 1024*4)
235
+ if err != nil {
236
+ t.Fatal(err)
237
+ }
238
+ if n != 4096 {
239
+ t.Fatal("Copy didnt copy enough bytes")
240
+ }
241
+
242
+ seeked, err := rs.Seek(0, os.SEEK_SET)
243
+ if err != nil {
244
+ t.Fatal(err)
245
+ }
246
+ if seeked != 0 {
247
+ t.Fatal("Failed to seek to beginning")
248
+ }
249
+
250
+ out, err := ioutil.ReadAll(rs)
251
+ if err != nil {
252
+ t.Fatal(err)
253
+ }
254
+
255
+ err = arrComp(out, should)
256
+ if err != nil {
257
+ t.Fatal(err)
258
+ }
259
+}
260
+
261
+func TestSeekToAlmostBegin(t *testing.T) {
262
+ nbytes := int64(10 * 1024)
263
+ should := make([]byte, nbytes)
264
+ u.NewTimeSeededRand().Read(should)
265
+
266
+ read := bytes.NewReader(should)
267
+ ds := mdtest.Mock(t)
268
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
269
+ if err != nil {
270
+ t.Fatal(err)
271
+ }
272
+
273
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
274
+ if err != nil {
275
+ t.Fatal(err)
276
+ }
277
+
278
+ n, err := io.CopyN(ioutil.Discard, rs, 1024*4)
279
+ if err != nil {
280
+ t.Fatal(err)
281
+ }
282
+ if n != 4096 {
283
+ t.Fatal("Copy didnt copy enough bytes")
284
+ }
285
+
286
+ seeked, err := rs.Seek(1, os.SEEK_SET)
287
+ if err != nil {
288
+ t.Fatal(err)
289
+ }
290
+ if seeked != 1 {
291
+ t.Fatal("Failed to seek to almost beginning")
292
+ }
293
+
294
+ out, err := ioutil.ReadAll(rs)
295
+ if err != nil {
296
+ t.Fatal(err)
297
+ }
298
+
299
+ err = arrComp(out, should[1:])
300
+ if err != nil {
301
+ t.Fatal(err)
302
+ }
303
+}
304
+
305
+func TestSeekEnd(t *testing.T) {
306
+ nbytes := int64(50 * 1024)
307
+ should := make([]byte, nbytes)
308
+ u.NewTimeSeededRand().Read(should)
309
+
310
+ read := bytes.NewReader(should)
311
+ ds := mdtest.Mock(t)
312
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
313
+ if err != nil {
314
+ t.Fatal(err)
315
+ }
316
+
317
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
318
+ if err != nil {
319
+ t.Fatal(err)
320
+ }
321
+
322
+ seeked, err := rs.Seek(0, os.SEEK_END)
323
+ if err != nil {
324
+ t.Fatal(err)
325
+ }
326
+ if seeked != nbytes {
327
+ t.Fatal("Failed to seek to end")
328
+ }
329
+}
330
+
331
+func TestSeekEndSingleBlockFile(t *testing.T) {
332
+ nbytes := int64(100)
333
+ should := make([]byte, nbytes)
334
+ u.NewTimeSeededRand().Read(should)
335
+
336
+ read := bytes.NewReader(should)
337
+ ds := mdtest.Mock(t)
338
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{5000})
339
+ if err != nil {
340
+ t.Fatal(err)
341
+ }
342
+
343
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
344
+ if err != nil {
345
+ t.Fatal(err)
346
+ }
347
+
348
+ seeked, err := rs.Seek(0, os.SEEK_END)
349
+ if err != nil {
350
+ t.Fatal(err)
351
+ }
352
+ if seeked != nbytes {
353
+ t.Fatal("Failed to seek to end")
354
+ }
355
+}
356
+
357
+func TestSeekingStress(t *testing.T) {
358
+ nbytes := int64(1024 * 1024)
359
+ should := make([]byte, nbytes)
360
+ u.NewTimeSeededRand().Read(should)
361
+
362
+ read := bytes.NewReader(should)
363
+ ds := mdtest.Mock(t)
364
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{1000})
365
+ if err != nil {
366
+ t.Fatal(err)
367
+ }
368
+
369
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
370
+ if err != nil {
371
+ t.Fatal(err)
372
+ }
373
+
374
+ testbuf := make([]byte, nbytes)
375
+ for i := 0; i < 50; i++ {
376
+ offset := mrand.Intn(int(nbytes))
377
+ l := int(nbytes) - offset
378
+ n, err := rs.Seek(int64(offset), os.SEEK_SET)
379
+ if err != nil {
380
+ t.Fatal(err)
381
+ }
382
+ if n != int64(offset) {
383
+ t.Fatal("Seek failed to move to correct position")
384
+ }
385
+
386
+ nread, err := rs.Read(testbuf[:l])
387
+ if err != nil {
388
+ t.Fatal(err)
389
+ }
390
+ if nread != l {
391
+ t.Fatal("Failed to read enough bytes")
392
+ }
393
+
394
+ err = arrComp(testbuf[:l], should[offset:offset+l])
395
+ if err != nil {
396
+ t.Fatal(err)
397
+ }
398
+ }
399
+
400
+}
401
+
402
+func TestSeekingConsistency(t *testing.T) {
403
+ nbytes := int64(128 * 1024)
404
+ should := make([]byte, nbytes)
405
+ u.NewTimeSeededRand().Read(should)
406
+
407
+ read := bytes.NewReader(should)
408
+ ds := mdtest.Mock(t)
409
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
410
+ if err != nil {
411
+ t.Fatal(err)
412
+ }
413
+
414
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
415
+ if err != nil {
416
+ t.Fatal(err)
417
+ }
418
+
419
+ out := make([]byte, nbytes)
420
+
421
+ for coff := nbytes - 4096; coff >= 0; coff -= 4096 {
422
+ t.Log(coff)
423
+ n, err := rs.Seek(coff, os.SEEK_SET)
424
+ if err != nil {
425
+ t.Fatal(err)
426
+ }
427
+ if n != coff {
428
+ t.Fatal("wasnt able to seek to the right position")
429
+ }
430
+ nread, err := rs.Read(out[coff : coff+4096])
431
+ if err != nil {
432
+ t.Fatal(err)
433
+ }
434
+ if nread != 4096 {
435
+ t.Fatal("didnt read the correct number of bytes")
436
+ }
437
+ }
438
+
439
+ err = arrComp(out, should)
440
+ if err != nil {
441
+ t.Fatal(err)
442
+ }
443
+}
importer/balanced/builder.go
new
+66
@@ -0,0 +1,66 @@
1
+package balanced
2
+
3
+import (
4
+ "errors"
5
+
6
+ h "github.com/jbenet/go-ipfs/importer/helpers"
7
+ dag "github.com/jbenet/go-ipfs/merkledag"
8
+)
9
+
10
+func BalancedLayout(db *h.DagBuilderHelper) (*dag.Node, error) {
11
+ var root *h.UnixfsNode
12
+ for level := 0; !db.Done(); level++ {
13
+
14
+ nroot := h.NewUnixfsNode()
15
+
16
+ // add our old root as a child of the new root.
17
+ if root != nil { // nil if it's the first node.
18
+ if err := nroot.AddChild(root, db); err != nil {
19
+ return nil, err
20
+ }
21
+ }
22
+
23
+ // fill it up.
24
+ if err := fillNodeRec(db, nroot, level); err != nil {
25
+ return nil, err
26
+ }
27
+
28
+ root = nroot
29
+ }
30
+ if root == nil {
31
+ root = h.NewUnixfsNode()
32
+ }
33
+
34
+ return db.Add(root)
35
+}
36
+
37
+// fillNodeRec will fill the given node with data from the dagBuilders input
38
+// source down to an indirection depth as specified by 'depth'
39
+// it returns the total dataSize of the node, and a potential error
40
+//
41
+// warning: **children** pinned indirectly, but input node IS NOT pinned.
42
+func fillNodeRec(db *h.DagBuilderHelper, node *h.UnixfsNode, depth int) error {
43
+ if depth < 0 {
44
+ return errors.New("attempt to fillNode at depth < 0")
45
+ }
46
+
47
+ // Base case
48
+ if depth <= 0 { // catch accidental -1's in case error above is removed.
49
+ return db.FillNodeWithData(node)
50
+ }
51
+
52
+ // while we have room AND we're not done
53
+ for node.NumChildren() < db.Maxlinks() && !db.Done() {
54
+ child := h.NewUnixfsNode()
55
+
56
+ if err := fillNodeRec(db, child, depth-1); err != nil {
57
+ return err
58
+ }
59
+
60
+ if err := node.AddChild(child, db); err != nil {
61
+ return err
62
+ }
63
+ }
64
+
65
+ return nil
66
+}
importer/calc_test.go
deleted
-50
@@ -1,50 +0,0 @@
1
-package importer
2
-
3
-import (
4
- "math"
5
- "testing"
6
-
7
- humanize "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/dustin/go-humanize"
8
-)
9
-
10
-func TestCalculateSizes(t *testing.T) {
11
-
12
- // d := ((lbs/271) ^ layer) * dbs
13
-
14
- increments := func(a, b int) []int {
15
- ints := []int{}
16
- for ; a <= b; a *= 2 {
17
- ints = append(ints, a)
18
- }
19
- return ints
20
- }
21
-
22
- layers := 7
23
- roughLinkSize := roughLinkSize // from importer pkg
24
- dataBlockSizes := increments(1<<12, 1<<18)
25
- linkBlockSizes := increments(1<<12, 1<<14)
26
-
27
- t.Logf("rough link size: %d", roughLinkSize)
28
- t.Logf("data block sizes: %v", dataBlockSizes)
29
- t.Logf("link block sizes: %v", linkBlockSizes)
30
- for _, dbs := range dataBlockSizes {
31
- t.Logf("")
32
- t.Logf("with data block size: %d", dbs)
33
- for _, lbs := range linkBlockSizes {
34
- t.Logf("")
35
- t.Logf("\twith data block size: %d", dbs)
36
- t.Logf("\twith link block size: %d", lbs)
37
-
38
- lpb := lbs / roughLinkSize
39
- t.Logf("\tlinks per block: %d", lpb)
40
-
41
- for l := 1; l < layers; l++ {
42
- total := int(math.Pow(float64(lpb), float64(l))) * dbs
43
- htotal := humanize.Bytes(uint64(total))
44
- t.Logf("\t\t\tlayer %d: %s\t%d", l, htotal, total)
45
- }
46
-
47
- }
48
- }
49
-
50
-}
importer/helpers/dagbuilder.go
new
+141
@@ -0,0 +1,141 @@
1
+package helpers
2
+
3
+import (
4
+ dag "github.com/jbenet/go-ipfs/merkledag"
5
+ "github.com/jbenet/go-ipfs/pin"
6
+)
7
+
8
+// DagBuilderHelper wraps together a bunch of objects needed to
9
+// efficiently create unixfs dag trees
10
+type DagBuilderHelper struct {
11
+ dserv dag.DAGService
12
+ mp pin.ManualPinner
13
+ in <-chan []byte
14
+ nextData []byte // the next item to return.
15
+ maxlinks int
16
+}
17
+
18
+type DagBuilderParams struct {
19
+ // Maximum number of links per intermediate node
20
+ Maxlinks int
21
+
22
+ // DAGService to write blocks to (required)
23
+ Dagserv dag.DAGService
24
+
25
+ // Pinner to use for pinning files (optionally nil)
26
+ Pinner pin.ManualPinner
27
+}
28
+
29
+// Generate a new DagBuilderHelper from the given params, using 'in' as a
30
+// data source
31
+func (dbp *DagBuilderParams) New(in <-chan []byte) *DagBuilderHelper {
32
+ return &DagBuilderHelper{
33
+ dserv: dbp.Dagserv,
34
+ mp: dbp.Pinner,
35
+ in: in,
36
+ maxlinks: dbp.Maxlinks,
37
+ }
38
+}
39
+
40
+// prepareNext consumes the next item from the channel and puts it
41
+// in the nextData field. it is idempotent-- if nextData is full
42
+// it will do nothing.
43
+//
44
+// i realized that building the dag becomes _a lot_ easier if we can
45
+// "peek" the "are done yet?" (i.e. not consume it from the channel)
46
+func (db *DagBuilderHelper) prepareNext() {
47
+ if db.in == nil {
48
+ // if our input is nil, there is "nothing to do". we're done.
49
+ // as if there was no data at all. (a sort of zero-value)
50
+ return
51
+ }
52
+
53
+ // if we already have data waiting to be consumed, we're ready.
54
+ if db.nextData != nil {
55
+ return
56
+ }
57
+
58
+ // if it's closed, nextData will be correctly set to nil, signaling
59
+ // that we're done consuming from the channel.
60
+ db.nextData = <-db.in
61
+}
62
+
63
+// Done returns whether or not we're done consuming the incoming data.
64
+func (db *DagBuilderHelper) Done() bool {
65
+ // ensure we have an accurate perspective on data
66
+ // as `done` this may be called before `next`.
67
+ db.prepareNext() // idempotent
68
+ return db.nextData == nil
69
+}
70
+
71
+// Next returns the next chunk of data to be inserted into the dag
72
+// if it returns nil, that signifies that the stream is at an end, and
73
+// that the current building operation should finish
74
+func (db *DagBuilderHelper) Next() []byte {
75
+ db.prepareNext() // idempotent
76
+ d := db.nextData
77
+ db.nextData = nil // signal we've consumed it
78
+ return d
79
+}
80
+
81
+// FillNodeLayer will add datanodes as children to the give node until
82
+// at most db.indirSize ndoes are added
83
+//
84
+// warning: **children** pinned indirectly, but input node IS NOT pinned.
85
+func (db *DagBuilderHelper) FillNodeLayer(node *UnixfsNode) error {
86
+
87
+ // while we have room AND we're not done
88
+ for node.NumChildren() < db.maxlinks && !db.Done() {
89
+ child := NewUnixfsNode()
90
+
91
+ if err := db.FillNodeWithData(child); err != nil {
92
+ return err
93
+ }
94
+
95
+ if err := node.AddChild(child, db); err != nil {
96
+ return err
97
+ }
98
+ }
99
+
100
+ return nil
101
+}
102
+
103
+func (db *DagBuilderHelper) FillNodeWithData(node *UnixfsNode) error {
104
+ data := db.Next()
105
+ if data == nil { // we're done!
106
+ return nil
107
+ }
108
+
109
+ if len(data) > BlockSizeLimit {
110
+ return ErrSizeLimitExceeded
111
+ }
112
+
113
+ node.setData(data)
114
+ return nil
115
+}
116
+
117
+func (db *DagBuilderHelper) Add(node *UnixfsNode) (*dag.Node, error) {
118
+ dn, err := node.GetDagNode()
119
+ if err != nil {
120
+ return nil, err
121
+ }
122
+
123
+ key, err := db.dserv.Add(dn)
124
+ if err != nil {
125
+ return nil, err
126
+ }
127
+
128
+ if db.mp != nil {
129
+ db.mp.PinWithMode(key, pin.Recursive)
130
+ err := db.mp.Flush()
131
+ if err != nil {
132
+ return nil, err
133
+ }
134
+ }
135
+
136
+ return dn, nil
137
+}
138
+
139
+func (db *DagBuilderHelper) Maxlinks() int {
140
+ return db.maxlinks
141
+}
importer/helpers/helpers.go
new
+99
@@ -0,0 +1,99 @@
1
+package helpers
2
+
3
+import (
4
+ "fmt"
5
+
6
+ chunk "github.com/jbenet/go-ipfs/importer/chunk"
7
+ dag "github.com/jbenet/go-ipfs/merkledag"
8
+ "github.com/jbenet/go-ipfs/pin"
9
+ ft "github.com/jbenet/go-ipfs/unixfs"
10
+)
11
+
12
+// BlockSizeLimit specifies the maximum size an imported block can have.
13
+var BlockSizeLimit = 1048576 // 1 MB
14
+
15
+// rough estimates on expected sizes
16
+var roughDataBlockSize = chunk.DefaultBlockSize
17
+var roughLinkBlockSize = 1 << 13 // 8KB
18
+var roughLinkSize = 258 + 8 + 5 // sha256 multihash + size + no name + protobuf framing
19
+
20
+// DefaultLinksPerBlock governs how the importer decides how many links there
21
+// will be per block. This calculation is based on expected distributions of:
22
+// * the expected distribution of block sizes
23
+// * the expected distribution of link sizes
24
+// * desired access speed
25
+// For now, we use:
26
+//
27
+// var roughLinkBlockSize = 1 << 13 // 8KB
28
+// var roughLinkSize = 288 // sha256 + framing + name
29
+// var DefaultLinksPerBlock = (roughLinkBlockSize / roughLinkSize)
30
+//
31
+// See calc_test.go
32
+var DefaultLinksPerBlock = (roughLinkBlockSize / roughLinkSize)
33
+
34
+// ErrSizeLimitExceeded signals that a block is larger than BlockSizeLimit.
35
+var ErrSizeLimitExceeded = fmt.Errorf("object size limit exceeded")
36
+
37
+// UnixfsNode is a struct created to aid in the generation
38
+// of unixfs DAG trees
39
+type UnixfsNode struct {
40
+ node *dag.Node
41
+ ufmt *ft.MultiBlock
42
+}
43
+
44
+func NewUnixfsNode() *UnixfsNode {
45
+ return &UnixfsNode{
46
+ node: new(dag.Node),
47
+ ufmt: new(ft.MultiBlock),
48
+ }
49
+}
50
+
51
+func (n *UnixfsNode) NumChildren() int {
52
+ return n.ufmt.NumChildren()
53
+}
54
+
55
+// addChild will add the given UnixfsNode as a child of the receiver.
56
+// the passed in DagBuilderHelper is used to store the child node an
57
+// pin it locally so it doesnt get lost
58
+func (n *UnixfsNode) AddChild(child *UnixfsNode, db *DagBuilderHelper) error {
59
+ n.ufmt.AddBlockSize(child.ufmt.FileSize())
60
+
61
+ childnode, err := child.GetDagNode()
62
+ if err != nil {
63
+ return err
64
+ }
65
+
66
+ // Add a link to this node without storing a reference to the memory
67
+ // This way, we avoid nodes building up and consuming all of our RAM
68
+ err = n.node.AddNodeLinkClean("", childnode)
69
+ if err != nil {
70
+ return err
71
+ }
72
+
73
+ childkey, err := db.dserv.Add(childnode)
74
+ if err != nil {
75
+ return err
76
+ }
77
+
78
+ // Pin the child node indirectly
79
+ if db.mp != nil {
80
+ db.mp.PinWithMode(childkey, pin.Indirect)
81
+ }
82
+
83
+ return nil
84
+}
85
+
86
+func (n *UnixfsNode) setData(data []byte) {
87
+ n.ufmt.Data = data
88
+}
89
+
90
+// getDagNode fills out the proper formatting for the unixfs node
91
+// inside of a DAG node and returns the dag node
92
+func (n *UnixfsNode) GetDagNode() (*dag.Node, error) {
93
+ data, err := n.ufmt.GetBytes()
94
+ if err != nil {
95
+ return nil, err
96
+ }
97
+ n.node.Data = data
98
+ return n.node, nil
99
+}
importer/importer.go
+16
-265
@@ -3,70 +3,21 @@
3
package importer
4
5
import (
6
- "errors"
6
"fmt"
7
"io"
8
"os"
9
10
+ bal "github.com/jbenet/go-ipfs/importer/balanced"
11
"github.com/jbenet/go-ipfs/importer/chunk"
12
+ h "github.com/jbenet/go-ipfs/importer/helpers"
13
+ trickle "github.com/jbenet/go-ipfs/importer/trickle"
14
dag "github.com/jbenet/go-ipfs/merkledag"
15
"github.com/jbenet/go-ipfs/pin"
14
- ft "github.com/jbenet/go-ipfs/unixfs"
16
"github.com/jbenet/go-ipfs/util"
17
)
18
19
var log = util.Logger("importer")
20
20
-// BlockSizeLimit specifies the maximum size an imported block can have.
21
-var BlockSizeLimit = 1048576 // 1 MB
22
-
23
-// rough estimates on expected sizes
24
-var roughDataBlockSize = chunk.DefaultBlockSize
25
-var roughLinkBlockSize = 1 << 13 // 8KB
26
-var roughLinkSize = 258 + 8 + 5 // sha256 multihash + size + no name + protobuf framing
27
-
28
-// DefaultLinksPerBlock governs how the importer decides how many links there
29
-// will be per block. This calculation is based on expected distributions of:
30
-// * the expected distribution of block sizes
31
-// * the expected distribution of link sizes
32
-// * desired access speed
33
-// For now, we use:
34
-//
35
-// var roughLinkBlockSize = 1 << 13 // 8KB
36
-// var roughLinkSize = 288 // sha256 + framing + name
37
-// var DefaultLinksPerBlock = (roughLinkBlockSize / roughLinkSize)
38
-//
39
-// See calc_test.go
40
-var DefaultLinksPerBlock = (roughLinkBlockSize / roughLinkSize)
41
-
42
-// ErrSizeLimitExceeded signals that a block is larger than BlockSizeLimit.
43
-var ErrSizeLimitExceeded = fmt.Errorf("object size limit exceeded")
44
-
45
-// IndirectBlocksCopyData governs whether indirect blocks should copy over
46
-// data from their first child, and how much. If this is 0, indirect blocks
47
-// have no data, only links. If this is larger, Indirect blocks will copy
48
-// as much as (maybe less than) this many bytes.
49
-//
50
-// This number should be <= (BlockSizeLimit - (DefaultLinksPerBlock * LinkSize))
51
-// Note that it is not known here what the LinkSize is, because the hash function
52
-// could vary wildly in size. Exercise caution when setting this option. For
53
-// safety, it will be clipped to (BlockSizeLimit - (DefaultLinksPerBlock * 256))
54
-var IndirectBlockDataSize = 0
55
-
56
-// this check is here to ensure the conditions on IndirectBlockDataSize hold.
57
-// returns int because it will be used as an input to `make()` later on. if
58
-// `int` will flip over to negative, better know here.
59
-func defaultIndirectBlockDataSize() int {
60
- max := BlockSizeLimit - (DefaultLinksPerBlock * 256)
61
- if IndirectBlockDataSize < max {
62
- max = IndirectBlockDataSize
63
- }
64
- if max < 0 {
65
- return 0
66
- }
67
- return max
68
-}
69
-
21
// Builds a DAG from the given file, writing created blocks to disk as they are
22
// created
23
func BuildDagFromFile(fpath string, ds dag.DAGService, mp pin.ManualPinner) (*dag.Node, error) {
@@ -88,228 +39,28 @@ func BuildDagFromFile(fpath string, ds dag.DAGService, mp pin.ManualPinner) (*da
39
return BuildDagFromReader(f, ds, mp, chunk.DefaultSplitter)
40
}
41
91
-// unixfsNode is a struct created to aid in the generation
92
-// of unixfs DAG trees
93
-type unixfsNode struct {
94
- node *dag.Node
95
- ufmt *ft.MultiBlock
96
-}
97
-
98
-func newUnixfsNode() *unixfsNode {
99
- return &unixfsNode{
100
- node: new(dag.Node),
101
- ufmt: new(ft.MultiBlock),
102
- }
103
-}
104
-
105
-func (n *unixfsNode) numChildren() int {
106
- return n.ufmt.NumChildren()
107
-}
108
-
109
-// addChild will add the given unixfsNode as a child of the receiver.
110
-// the passed in dagBuilderHelper is used to store the child node an
111
-// pin it locally so it doesnt get lost
112
-func (n *unixfsNode) addChild(child *unixfsNode, db *dagBuilderHelper) error {
113
- n.ufmt.AddBlockSize(child.ufmt.FileSize())
114
-
115
- childnode, err := child.getDagNode()
116
- if err != nil {
117
- return err
118
- }
119
-
120
- // Add a link to this node without storing a reference to the memory
121
- // This way, we avoid nodes building up and consuming all of our RAM
122
- err = n.node.AddNodeLinkClean("", childnode)
123
- if err != nil {
124
- return err
125
- }
126
-
127
- childkey, err := db.dserv.Add(childnode)
128
- if err != nil {
129
- return err
130
- }
131
-
132
- // Pin the child node indirectly
133
- if db.mp != nil {
134
- db.mp.PinWithMode(childkey, pin.Indirect)
135
- }
136
-
137
- return nil
138
-}
139
-
140
-func (n *unixfsNode) setData(data []byte) {
141
- n.ufmt.Data = data
142
-}
143
-
144
-// getDagNode fills out the proper formatting for the unixfs node
145
-// inside of a DAG node and returns the dag node
146
-func (n *unixfsNode) getDagNode() (*dag.Node, error) {
147
- data, err := n.ufmt.GetBytes()
148
- if err != nil {
149
- return nil, err
150
- }
151
- n.node.Data = data
152
- return n.node, nil
153
-}
154
-
42
func BuildDagFromReader(r io.Reader, ds dag.DAGService, mp pin.ManualPinner, spl chunk.BlockSplitter) (*dag.Node, error) {
43
// Start the splitter
44
blkch := spl.Split(r)
45
159
- // Create our builder helper
160
- db := &dagBuilderHelper{
161
- dserv: ds,
162
- mp: mp,
163
- in: blkch,
164
- maxlinks: DefaultLinksPerBlock,
165
- indrSize: defaultIndirectBlockDataSize(),
166
- }
167
-
168
- var root *unixfsNode
169
- for level := 0; !db.done(); level++ {
170
-
171
- nroot := newUnixfsNode()
172
-
173
- // add our old root as a child of the new root.
174
- if root != nil { // nil if it's the first node.
175
- if err := nroot.addChild(root, db); err != nil {
176
- return nil, err
177
- }
178
- }
179
-
180
- // fill it up.
181
- if err := db.fillNodeRec(nroot, level); err != nil {
182
- return nil, err
183
- }
184
-
185
- root = nroot
186
- }
187
- if root == nil {
188
- root = newUnixfsNode()
189
- }
190
-
191
- rootnode, err := root.getDagNode()
192
- if err != nil {
193
- return nil, err
194
- }
195
-
196
- rootkey, err := ds.Add(rootnode)
197
- if err != nil {
198
- return nil, err
199
- }
200
-
201
- if mp != nil {
202
- mp.PinWithMode(rootkey, pin.Recursive)
203
- err := mp.Flush()
204
- if err != nil {
205
- return nil, err
206
- }
207
- }
208
-
209
- return root.getDagNode()
210
-}
211
-
212
-// dagBuilderHelper wraps together a bunch of objects needed to
213
-// efficiently create unixfs dag trees
214
-type dagBuilderHelper struct {
215
- dserv dag.DAGService
216
- mp pin.ManualPinner
217
- in <-chan []byte
218
- nextData []byte // the next item to return.
219
- maxlinks int
220
- indrSize int // see IndirectBlockData
221
-}
222
-
223
-// prepareNext consumes the next item from the channel and puts it
224
-// in the nextData field. it is idempotent-- if nextData is full
225
-// it will do nothing.
226
-//
227
-// i realized that building the dag becomes _a lot_ easier if we can
228
-// "peek" the "are done yet?" (i.e. not consume it from the channel)
229
-func (db *dagBuilderHelper) prepareNext() {
230
- if db.in == nil {
231
- // if our input is nil, there is "nothing to do". we're done.
232
- // as if there was no data at all. (a sort of zero-value)
233
- return
234
- }
235
-
236
- // if we already have data waiting to be consumed, we're ready.
237
- if db.nextData != nil {
238
- return
46
+ dbp := h.DagBuilderParams{
47
+ Dagserv: ds,
48
+ Maxlinks: h.DefaultLinksPerBlock,
49
+ Pinner: mp,
50
}
51
241
- // if it's closed, nextData will be correctly set to nil, signaling
242
- // that we're done consuming from the channel.
243
- db.nextData = <-db.in
52
+ return bal.BalancedLayout(dbp.New(blkch))
53
}
54
246
-// done returns whether or not we're done consuming the incoming data.
247
-func (db *dagBuilderHelper) done() bool {
248
- // ensure we have an accurate perspective on data
249
- // as `done` this may be called before `next`.
250
- db.prepareNext() // idempotent
251
- return db.nextData == nil
252
-}
253
-
254
-// next returns the next chunk of data to be inserted into the dag
255
-// if it returns nil, that signifies that the stream is at an end, and
256
-// that the current building operation should finish
257
-func (db *dagBuilderHelper) next() []byte {
258
- db.prepareNext() // idempotent
259
- d := db.nextData
260
- db.nextData = nil // signal we've consumed it
261
- return d
262
-}
263
-
264
-// fillNodeRec will fill the given node with data from the dagBuilders input
265
-// source down to an indirection depth as specified by 'depth'
266
-// it returns the total dataSize of the node, and a potential error
267
-//
268
-// warning: **children** pinned indirectly, but input node IS NOT pinned.
269
-func (db *dagBuilderHelper) fillNodeRec(node *unixfsNode, depth int) error {
270
- if depth < 0 {
271
- return errors.New("attempt to fillNode at depth < 0")
272
- }
273
-
274
- // Base case
275
- if depth <= 0 { // catch accidental -1's in case error above is removed.
276
- return db.fillNodeWithData(node)
277
- }
278
-
279
- // while we have room AND we're not done
280
- for node.numChildren() < db.maxlinks && !db.done() {
281
- child := newUnixfsNode()
282
-
283
- if err := db.fillNodeRec(child, depth-1); err != nil {
284
- return err
285
- }
286
-
287
- if err := node.addChild(child, db); err != nil {
288
- return err
289
- }
290
- }
291
-
292
- return nil
293
-}
294
-
295
-func (db *dagBuilderHelper) fillNodeWithData(node *unixfsNode) error {
296
- data := db.next()
297
- if data == nil { // we're done!
298
- return nil
299
- }
55
+func BuildTrickleDagFromReader(r io.Reader, ds dag.DAGService, mp pin.ManualPinner, spl chunk.BlockSplitter) (*dag.Node, error) {
56
+ // Start the splitter
57
+ blkch := spl.Split(r)
58
301
- if len(data) > BlockSizeLimit {
302
- return ErrSizeLimitExceeded
59
+ dbp := h.DagBuilderParams{
60
+ Dagserv: ds,
61
+ Maxlinks: h.DefaultLinksPerBlock,
62
+ Pinner: mp,
63
}
64
305
- node.setData(data)
306
- return nil
307
-}
308
-
309
-// why is intmin not in math?
310
-func min(a, b int) int {
311
- if a > b {
312
- return a
313
- }
314
- return b
65
+ return trickle.TrickleLayout(dbp.New(blkch))
66
}
importer/importer_test.go
+43
-454
@@ -2,511 +2,100 @@ package importer
2
3
import (
4
"bytes"
5
- "crypto/rand"
6
- "fmt"
5
"io"
6
"io/ioutil"
9
- mrand "math/rand"
10
- "os"
7
"testing"
8
13
- "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
14
- ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
15
- dssync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
16
- bstore "github.com/jbenet/go-ipfs/blocks/blockstore"
17
- bserv "github.com/jbenet/go-ipfs/blockservice"
18
- offline "github.com/jbenet/go-ipfs/exchange/offline"
9
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
10
chunk "github.com/jbenet/go-ipfs/importer/chunk"
20
- merkledag "github.com/jbenet/go-ipfs/merkledag"
21
- pin "github.com/jbenet/go-ipfs/pin"
11
+ dag "github.com/jbenet/go-ipfs/merkledag"
12
+ mdtest "github.com/jbenet/go-ipfs/merkledag/test"
13
uio "github.com/jbenet/go-ipfs/unixfs/io"
14
u "github.com/jbenet/go-ipfs/util"
15
)
16
26
-//Test where calls to read are smaller than the chunk size
27
-func TestSizeBasedSplit(t *testing.T) {
28
- if testing.Short() {
29
- t.SkipNow()
30
- }
31
- bs := &chunk.SizeSplitter{Size: 512}
32
- testFileConsistency(t, bs, 32*512)
33
- bs = &chunk.SizeSplitter{Size: 4096}
34
- testFileConsistency(t, bs, 32*4096)
35
-
36
- // Uneven offset
37
- testFileConsistency(t, bs, 31*4095)
38
-}
39
-
40
-func dup(b []byte) []byte {
41
- o := make([]byte, len(b))
42
- copy(o, b)
43
- return o
44
-}
45
-
46
-func testFileConsistency(t *testing.T, bs chunk.BlockSplitter, nbytes int) {
47
- should := make([]byte, nbytes)
48
- u.NewTimeSeededRand().Read(should)
49
-
50
- read := bytes.NewReader(should)
51
- dnp := getDagservAndPinner(t)
52
- nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, bs)
53
- if err != nil {
54
- t.Fatal(err)
55
- }
56
-
57
- r, err := uio.NewDagReader(context.Background(), nd, dnp.ds)
58
- if err != nil {
59
- t.Fatal(err)
60
- }
61
-
62
- out, err := ioutil.ReadAll(r)
63
- if err != nil {
64
- t.Fatal(err)
65
- }
66
-
67
- err = arrComp(out, should)
68
- if err != nil {
69
- t.Fatal(err)
70
- }
71
-}
72
-
73
-func TestBuilderConsistency(t *testing.T) {
74
- nbytes := 100000
75
- buf := new(bytes.Buffer)
76
- io.CopyN(buf, u.NewTimeSeededRand(), int64(nbytes))
77
- should := dup(buf.Bytes())
78
- dagserv := merkledag.Mock(t)
79
- nd, err := BuildDagFromReader(buf, dagserv, nil, chunk.DefaultSplitter)
80
- if err != nil {
81
- t.Fatal(err)
82
- }
83
- r, err := uio.NewDagReader(context.Background(), nd, dagserv)
17
+func getBalancedDag(t testing.TB, size int64) (*dag.Node, dag.DAGService) {
18
+ ds := mdtest.Mock(t)
19
+ r := io.LimitReader(u.NewTimeSeededRand(), size)
20
+ nd, err := BuildDagFromReader(r, ds, nil, chunk.DefaultSplitter)
21
if err != nil {
22
t.Fatal(err)
23
}
87
-
88
- out, err := ioutil.ReadAll(r)
89
- if err != nil {
90
- t.Fatal(err)
91
- }
92
-
93
- err = arrComp(out, should)
94
- if err != nil {
95
- t.Fatal(err)
96
- }
97
-}
98
-
99
-func TestTrickleBuilderConsistency(t *testing.T) {
100
- nbytes := 100000
101
- buf := new(bytes.Buffer)
102
- io.CopyN(buf, u.NewTimeSeededRand(), int64(nbytes))
103
- should := dup(buf.Bytes())
104
- dagserv := merkledag.Mock(t)
105
- nd, err := BuildTrickleDagFromReader(buf, dagserv, nil, chunk.DefaultSplitter)
106
- if err != nil {
107
- t.Fatal(err)
108
- }
109
- r, err := uio.NewDagReader(context.Background(), nd, dagserv)
110
- if err != nil {
111
- t.Fatal(err)
112
- }
113
-
114
- out, err := ioutil.ReadAll(r)
115
- if err != nil {
116
- t.Fatal(err)
117
- }
118
-
119
- err = arrComp(out, should)
120
- if err != nil {
121
- t.Fatal(err)
122
- }
123
-}
124
-
125
-func arrComp(a, b []byte) error {
126
- if len(a) != len(b) {
127
- return fmt.Errorf("Arrays differ in length. %d != %d", len(a), len(b))
128
- }
129
- for i, v := range a {
130
- if v != b[i] {
131
- return fmt.Errorf("Arrays differ at index: %d", i)
132
- }
133
- }
134
- return nil
135
-}
136
-
137
-func TestMaybeRabinConsistency(t *testing.T) {
138
- if testing.Short() {
139
- t.SkipNow()
140
- }
141
- testFileConsistency(t, chunk.NewMaybeRabin(4096), 256*4096)
142
-}
143
-
144
-func TestRabinBlockSize(t *testing.T) {
145
- if testing.Short() {
146
- t.SkipNow()
147
- }
148
- buf := new(bytes.Buffer)
149
- nbytes := 1024 * 1024
150
- io.CopyN(buf, rand.Reader, int64(nbytes))
151
- rab := chunk.NewMaybeRabin(4096)
152
- blkch := rab.Split(buf)
153
-
154
- var blocks [][]byte
155
- for b := range blkch {
156
- blocks = append(blocks, b)
157
- }
158
-
159
- fmt.Printf("Avg block size: %d\n", nbytes/len(blocks))
160
-
24
+ return nd, ds
25
}
26
163
-type dagservAndPinner struct {
164
- ds merkledag.DAGService
165
- mp pin.ManualPinner
166
-}
167
-
168
-func getDagservAndPinner(t *testing.T) dagservAndPinner {
169
- db := dssync.MutexWrap(ds.NewMapDatastore())
170
- bs := bstore.NewBlockstore(db)
171
- blockserv, err := bserv.New(bs, offline.Exchange(bs))
27
+func getTrickleDag(t testing.TB, size int64) (*dag.Node, dag.DAGService) {
28
+ ds := mdtest.Mock(t)
29
+ r := io.LimitReader(u.NewTimeSeededRand(), size)
30
+ nd, err := BuildTrickleDagFromReader(r, ds, nil, chunk.DefaultSplitter)
31
if err != nil {
32
t.Fatal(err)
33
}
175
- dserv := merkledag.NewDAGService(blockserv)
176
- mpin := pin.NewPinner(db, dserv).GetManual()
177
- return dagservAndPinner{
178
- ds: dserv,
179
- mp: mpin,
180
- }
34
+ return nd, ds
35
}
36
183
-func TestIndirectBlocks(t *testing.T) {
184
- splitter := &chunk.SizeSplitter{512}
185
- nbytes := 1024 * 1024
186
- buf := make([]byte, nbytes)
37
+func TestBalancedDag(t *testing.T) {
38
+ ds := mdtest.Mock(t)
39
+ buf := make([]byte, 10000)
40
u.NewTimeSeededRand().Read(buf)
41
+ r := bytes.NewReader(buf)
42
189
- read := bytes.NewReader(buf)
190
-
191
- dnp := getDagservAndPinner(t)
192
- dag, err := BuildDagFromReader(read, dnp.ds, dnp.mp, splitter)
43
+ nd, err := BuildDagFromReader(r, ds, nil, chunk.DefaultSplitter)
44
if err != nil {
45
t.Fatal(err)
46
}
47
197
- reader, err := uio.NewDagReader(context.Background(), dag, dnp.ds)
48
+ dr, err := uio.NewDagReader(context.TODO(), nd, ds)
49
if err != nil {
50
t.Fatal(err)
51
}
52
202
- out, err := ioutil.ReadAll(reader)
53
+ out, err := ioutil.ReadAll(dr)
54
if err != nil {
55
t.Fatal(err)
56
}
57
58
if !bytes.Equal(out, buf) {
208
- t.Fatal("Not equal!")
209
- }
210
-}
211
-
212
-func TestSeekingBasic(t *testing.T) {
213
- nbytes := int64(10 * 1024)
214
- should := make([]byte, nbytes)
215
- u.NewTimeSeededRand().Read(should)
216
-
217
- read := bytes.NewReader(should)
218
- dnp := getDagservAndPinner(t)
219
- nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{500})
220
- if err != nil {
221
- t.Fatal(err)
222
- }
223
-
224
- rs, err := uio.NewDagReader(context.Background(), nd, dnp.ds)
225
- if err != nil {
226
- t.Fatal(err)
227
- }
228
-
229
- start := int64(4000)
230
- n, err := rs.Seek(start, os.SEEK_SET)
231
- if err != nil {
232
- t.Fatal(err)
233
- }
234
- if n != start {
235
- t.Fatal("Failed to seek to correct offset")
236
- }
237
-
238
- out, err := ioutil.ReadAll(rs)
239
- if err != nil {
240
- t.Fatal(err)
241
- }
242
-
243
- err = arrComp(out, should[start:])
244
- if err != nil {
245
- t.Fatal(err)
59
+ t.Fatal("bad read")
60
}
61
}
62
249
-func TestTrickleSeekingBasic(t *testing.T) {
250
- nbytes := int64(10 * 1024)
251
- should := make([]byte, nbytes)
252
- u.NewTimeSeededRand().Read(should)
63
+func BenchmarkBalancedRead(b *testing.B) {
64
+ b.StopTimer()
65
+ nd, ds := getBalancedDag(b, int64(b.N))
66
254
- read := bytes.NewReader(should)
255
- dnp := getDagservAndPinner(t)
256
- nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{500})
67
+ read, err := uio.NewDagReader(context.TODO(), nd, ds)
68
if err != nil {
258
- t.Fatal(err)
69
+ b.Fatal(err)
70
}
71
261
- rs, err := uio.NewDagReader(context.Background(), nd, dnp.ds)
72
+ b.StartTimer()
73
+ b.SetBytes(int64(b.N))
74
+ n, err := io.Copy(ioutil.Discard, read)
75
if err != nil {
263
- t.Fatal(err)
76
+ b.Fatal(err)
77
}
265
-
266
- start := int64(4000)
267
- n, err := rs.Seek(start, os.SEEK_SET)
268
- if err != nil {
269
- t.Fatal(err)
270
- }
271
- if n != start {
272
- t.Fatal("Failed to seek to correct offset")
273
- }
274
-
275
- out, err := ioutil.ReadAll(rs)
276
- if err != nil {
277
- t.Fatal(err)
278
- }
279
-
280
- err = arrComp(out, should[start:])
281
- if err != nil {
282
- t.Fatal(err)
78
+ if n != int64(b.N) {
79
+ b.Fatal("Failed to read correct amount")
80
}
81
}
82
286
-func TestSeekToBegin(t *testing.T) {
287
- nbytes := int64(10 * 1024)
288
- should := make([]byte, nbytes)
289
- u.NewTimeSeededRand().Read(should)
83
+func BenchmarkTrickleRead(b *testing.B) {
84
+ b.StopTimer()
85
+ nd, ds := getTrickleDag(b, int64(b.N))
86
291
- read := bytes.NewReader(should)
292
- dnp := getDagservAndPinner(t)
293
- nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{500})
87
+ read, err := uio.NewDagReader(context.TODO(), nd, ds)
88
if err != nil {
295
- t.Fatal(err)
89
+ b.Fatal(err)
90
}
91
298
- rs, err := uio.NewDagReader(context.Background(), nd, dnp.ds)
92
+ b.StartTimer()
93
+ b.SetBytes(int64(b.N))
94
+ n, err := io.Copy(new(bytes.Buffer), read)
95
if err != nil {
300
- t.Fatal(err)
96
+ b.Fatal(err)
97
}
302
-
303
- n, err := io.CopyN(ioutil.Discard, rs, 1024*4)
304
- if err != nil {
305
- t.Fatal(err)
306
- }
307
- if n != 4096 {
308
- t.Fatal("Copy didnt copy enough bytes")
309
- }
310
-
311
- seeked, err := rs.Seek(0, os.SEEK_SET)
312
- if err != nil {
313
- t.Fatal(err)
314
- }
315
- if seeked != 0 {
316
- t.Fatal("Failed to seek to beginning")
317
- }
318
-
319
- out, err := ioutil.ReadAll(rs)
320
- if err != nil {
321
- t.Fatal(err)
322
- }
323
-
324
- err = arrComp(out, should)
325
- if err != nil {
326
- t.Fatal(err)
327
- }
328
-}
329
-
330
-func TestSeekToAlmostBegin(t *testing.T) {
331
- nbytes := int64(10 * 1024)
332
- should := make([]byte, nbytes)
333
- u.NewTimeSeededRand().Read(should)
334
-
335
- read := bytes.NewReader(should)
336
- dnp := getDagservAndPinner(t)
337
- nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{500})
338
- if err != nil {
339
- t.Fatal(err)
340
- }
341
-
342
- rs, err := uio.NewDagReader(context.Background(), nd, dnp.ds)
343
- if err != nil {
344
- t.Fatal(err)
345
- }
346
-
347
- n, err := io.CopyN(ioutil.Discard, rs, 1024*4)
348
- if err != nil {
349
- t.Fatal(err)
350
- }
351
- if n != 4096 {
352
- t.Fatal("Copy didnt copy enough bytes")
353
- }
354
-
355
- seeked, err := rs.Seek(1, os.SEEK_SET)
356
- if err != nil {
357
- t.Fatal(err)
358
- }
359
- if seeked != 1 {
360
- t.Fatal("Failed to seek to almost beginning")
361
- }
362
-
363
- out, err := ioutil.ReadAll(rs)
364
- if err != nil {
365
- t.Fatal(err)
366
- }
367
-
368
- err = arrComp(out, should[1:])
369
- if err != nil {
370
- t.Fatal(err)
371
- }
372
-}
373
-
374
-func TestSeekEnd(t *testing.T) {
375
- nbytes := int64(50 * 1024)
376
- should := make([]byte, nbytes)
377
- u.NewTimeSeededRand().Read(should)
378
-
379
- read := bytes.NewReader(should)
380
- dnp := getDagservAndPinner(t)
381
- nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{500})
382
- if err != nil {
383
- t.Fatal(err)
384
- }
385
-
386
- rs, err := uio.NewDagReader(context.Background(), nd, dnp.ds)
387
- if err != nil {
388
- t.Fatal(err)
389
- }
390
-
391
- seeked, err := rs.Seek(0, os.SEEK_END)
392
- if err != nil {
393
- t.Fatal(err)
394
- }
395
- if seeked != nbytes {
396
- t.Fatal("Failed to seek to end")
397
- }
398
-}
399
-
400
-func TestSeekEndSingleBlockFile(t *testing.T) {
401
- nbytes := int64(100)
402
- should := make([]byte, nbytes)
403
- u.NewTimeSeededRand().Read(should)
404
-
405
- read := bytes.NewReader(should)
406
- dnp := getDagservAndPinner(t)
407
- nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{5000})
408
- if err != nil {
409
- t.Fatal(err)
410
- }
411
-
412
- rs, err := uio.NewDagReader(context.Background(), nd, dnp.ds)
413
- if err != nil {
414
- t.Fatal(err)
415
- }
416
-
417
- seeked, err := rs.Seek(0, os.SEEK_END)
418
- if err != nil {
419
- t.Fatal(err)
420
- }
421
- if seeked != nbytes {
422
- t.Fatal("Failed to seek to end")
423
- }
424
-}
425
-
426
-func TestSeekingStress(t *testing.T) {
427
- nbytes := int64(1024 * 1024)
428
- should := make([]byte, nbytes)
429
- u.NewTimeSeededRand().Read(should)
430
-
431
- read := bytes.NewReader(should)
432
- dnp := getDagservAndPinner(t)
433
- nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{1000})
434
- if err != nil {
435
- t.Fatal(err)
436
- }
437
-
438
- rs, err := uio.NewDagReader(context.Background(), nd, dnp.ds)
439
- if err != nil {
440
- t.Fatal(err)
441
- }
442
-
443
- testbuf := make([]byte, nbytes)
444
- for i := 0; i < 50; i++ {
445
- offset := mrand.Intn(int(nbytes))
446
- l := int(nbytes) - offset
447
- n, err := rs.Seek(int64(offset), os.SEEK_SET)
448
- if err != nil {
449
- t.Fatal(err)
450
- }
451
- if n != int64(offset) {
452
- t.Fatal("Seek failed to move to correct position")
453
- }
454
-
455
- nread, err := rs.Read(testbuf[:l])
456
- if err != nil {
457
- t.Fatal(err)
458
- }
459
- if nread != l {
460
- t.Fatal("Failed to read enough bytes")
461
- }
462
-
463
- err = arrComp(testbuf[:l], should[offset:offset+l])
464
- if err != nil {
465
- t.Fatal(err)
466
- }
467
- }
468
-
469
-}
470
-
471
-func TestSeekingConsistency(t *testing.T) {
472
- nbytes := int64(128 * 1024)
473
- should := make([]byte, nbytes)
474
- u.NewTimeSeededRand().Read(should)
475
-
476
- read := bytes.NewReader(should)
477
- dnp := getDagservAndPinner(t)
478
- nd, err := BuildDagFromReader(read, dnp.ds, dnp.mp, &chunk.SizeSplitter{500})
479
- if err != nil {
480
- t.Fatal(err)
481
- }
482
-
483
- rs, err := uio.NewDagReader(context.Background(), nd, dnp.ds)
484
- if err != nil {
485
- t.Fatal(err)
486
- }
487
-
488
- out := make([]byte, nbytes)
489
-
490
- for coff := nbytes - 4096; coff >= 0; coff -= 4096 {
491
- t.Log(coff)
492
- n, err := rs.Seek(coff, os.SEEK_SET)
493
- if err != nil {
494
- t.Fatal(err)
495
- }
496
- if n != coff {
497
- t.Fatal("wasnt able to seek to the right position")
498
- }
499
- nread, err := rs.Read(out[coff : coff+4096])
500
- if err != nil {
501
- t.Fatal(err)
502
- }
503
- if nread != 4096 {
504
- t.Fatal("didnt read the correct number of bytes")
505
- }
506
- }
507
-
508
- err = arrComp(out, should)
509
- if err != nil {
510
- t.Fatal(err)
98
+ if n != int64(b.N) {
99
+ b.Fatal("Failed to read correct amount")
100
}
101
}
importer/trickle/trickle_test.go
new
+443
@@ -0,0 +1,443 @@
1
+package trickle
2
+
3
+import (
4
+ "bytes"
5
+ "crypto/rand"
6
+ "fmt"
7
+ "io"
8
+ "io/ioutil"
9
+ mrand "math/rand"
10
+ "os"
11
+ "testing"
12
+
13
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
14
+ chunk "github.com/jbenet/go-ipfs/importer/chunk"
15
+ h "github.com/jbenet/go-ipfs/importer/helpers"
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
+ uio "github.com/jbenet/go-ipfs/unixfs/io"
20
+ u "github.com/jbenet/go-ipfs/util"
21
+)
22
+
23
+func buildTestDag(r io.Reader, ds merkledag.DAGService, spl chunk.BlockSplitter) (*merkledag.Node, error) {
24
+ // Start the splitter
25
+ blkch := spl.Split(r)
26
+
27
+ dbp := h.DagBuilderParams{
28
+ Dagserv: ds,
29
+ Maxlinks: h.DefaultLinksPerBlock,
30
+ }
31
+
32
+ return TrickleLayout(dbp.New(blkch))
33
+}
34
+
35
+//Test where calls to read are smaller than the chunk size
36
+func TestSizeBasedSplit(t *testing.T) {
37
+ if testing.Short() {
38
+ t.SkipNow()
39
+ }
40
+ bs := &chunk.SizeSplitter{Size: 512}
41
+ testFileConsistency(t, bs, 32*512)
42
+ bs = &chunk.SizeSplitter{Size: 4096}
43
+ testFileConsistency(t, bs, 32*4096)
44
+
45
+ // Uneven offset
46
+ testFileConsistency(t, bs, 31*4095)
47
+}
48
+
49
+func dup(b []byte) []byte {
50
+ o := make([]byte, len(b))
51
+ copy(o, b)
52
+ return o
53
+}
54
+
55
+func testFileConsistency(t *testing.T, bs chunk.BlockSplitter, nbytes int) {
56
+ should := make([]byte, nbytes)
57
+ u.NewTimeSeededRand().Read(should)
58
+
59
+ read := bytes.NewReader(should)
60
+ ds := mdtest.Mock(t)
61
+ nd, err := buildTestDag(read, ds, bs)
62
+ if err != nil {
63
+ t.Fatal(err)
64
+ }
65
+
66
+ r, err := uio.NewDagReader(context.Background(), nd, ds)
67
+ if err != nil {
68
+ t.Fatal(err)
69
+ }
70
+
71
+ out, err := ioutil.ReadAll(r)
72
+ if err != nil {
73
+ t.Fatal(err)
74
+ }
75
+
76
+ err = arrComp(out, should)
77
+ if err != nil {
78
+ t.Fatal(err)
79
+ }
80
+}
81
+
82
+func TestBuilderConsistency(t *testing.T) {
83
+ nbytes := 100000
84
+ buf := new(bytes.Buffer)
85
+ io.CopyN(buf, u.NewTimeSeededRand(), int64(nbytes))
86
+ should := dup(buf.Bytes())
87
+ dagserv := mdtest.Mock(t)
88
+ nd, err := buildTestDag(buf, dagserv, chunk.DefaultSplitter)
89
+ if err != nil {
90
+ t.Fatal(err)
91
+ }
92
+ r, err := uio.NewDagReader(context.Background(), nd, dagserv)
93
+ if err != nil {
94
+ t.Fatal(err)
95
+ }
96
+
97
+ out, err := ioutil.ReadAll(r)
98
+ if err != nil {
99
+ t.Fatal(err)
100
+ }
101
+
102
+ err = arrComp(out, should)
103
+ if err != nil {
104
+ t.Fatal(err)
105
+ }
106
+}
107
+
108
+func arrComp(a, b []byte) error {
109
+ if len(a) != len(b) {
110
+ return fmt.Errorf("Arrays differ in length. %d != %d", len(a), len(b))
111
+ }
112
+ for i, v := range a {
113
+ if v != b[i] {
114
+ return fmt.Errorf("Arrays differ at index: %d", i)
115
+ }
116
+ }
117
+ return nil
118
+}
119
+
120
+func TestMaybeRabinConsistency(t *testing.T) {
121
+ if testing.Short() {
122
+ t.SkipNow()
123
+ }
124
+ testFileConsistency(t, chunk.NewMaybeRabin(4096), 256*4096)
125
+}
126
+
127
+func TestRabinBlockSize(t *testing.T) {
128
+ if testing.Short() {
129
+ t.SkipNow()
130
+ }
131
+ buf := new(bytes.Buffer)
132
+ nbytes := 1024 * 1024
133
+ io.CopyN(buf, rand.Reader, int64(nbytes))
134
+ rab := chunk.NewMaybeRabin(4096)
135
+ blkch := rab.Split(buf)
136
+
137
+ var blocks [][]byte
138
+ for b := range blkch {
139
+ blocks = append(blocks, b)
140
+ }
141
+
142
+ fmt.Printf("Avg block size: %d\n", nbytes/len(blocks))
143
+
144
+}
145
+
146
+type dagservAndPinner struct {
147
+ ds merkledag.DAGService
148
+ mp pin.ManualPinner
149
+}
150
+
151
+func TestIndirectBlocks(t *testing.T) {
152
+ splitter := &chunk.SizeSplitter{512}
153
+ nbytes := 1024 * 1024
154
+ buf := make([]byte, nbytes)
155
+ u.NewTimeSeededRand().Read(buf)
156
+
157
+ read := bytes.NewReader(buf)
158
+
159
+ ds := mdtest.Mock(t)
160
+ dag, err := buildTestDag(read, ds, splitter)
161
+ if err != nil {
162
+ t.Fatal(err)
163
+ }
164
+
165
+ reader, err := uio.NewDagReader(context.Background(), dag, ds)
166
+ if err != nil {
167
+ t.Fatal(err)
168
+ }
169
+
170
+ out, err := ioutil.ReadAll(reader)
171
+ if err != nil {
172
+ t.Fatal(err)
173
+ }
174
+
175
+ if !bytes.Equal(out, buf) {
176
+ t.Fatal("Not equal!")
177
+ }
178
+}
179
+
180
+func TestSeekingBasic(t *testing.T) {
181
+ nbytes := int64(10 * 1024)
182
+ should := make([]byte, nbytes)
183
+ u.NewTimeSeededRand().Read(should)
184
+
185
+ read := bytes.NewReader(should)
186
+ ds := mdtest.Mock(t)
187
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
188
+ if err != nil {
189
+ t.Fatal(err)
190
+ }
191
+
192
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
193
+ if err != nil {
194
+ t.Fatal(err)
195
+ }
196
+
197
+ start := int64(4000)
198
+ n, err := rs.Seek(start, os.SEEK_SET)
199
+ if err != nil {
200
+ t.Fatal(err)
201
+ }
202
+ if n != start {
203
+ t.Fatal("Failed to seek to correct offset")
204
+ }
205
+
206
+ out, err := ioutil.ReadAll(rs)
207
+ if err != nil {
208
+ t.Fatal(err)
209
+ }
210
+
211
+ err = arrComp(out, should[start:])
212
+ if err != nil {
213
+ t.Fatal(err)
214
+ }
215
+}
216
+
217
+func TestSeekToBegin(t *testing.T) {
218
+ nbytes := int64(10 * 1024)
219
+ should := make([]byte, nbytes)
220
+ u.NewTimeSeededRand().Read(should)
221
+
222
+ read := bytes.NewReader(should)
223
+ ds := mdtest.Mock(t)
224
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
225
+ if err != nil {
226
+ t.Fatal(err)
227
+ }
228
+
229
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
230
+ if err != nil {
231
+ t.Fatal(err)
232
+ }
233
+
234
+ n, err := io.CopyN(ioutil.Discard, rs, 1024*4)
235
+ if err != nil {
236
+ t.Fatal(err)
237
+ }
238
+ if n != 4096 {
239
+ t.Fatal("Copy didnt copy enough bytes")
240
+ }
241
+
242
+ seeked, err := rs.Seek(0, os.SEEK_SET)
243
+ if err != nil {
244
+ t.Fatal(err)
245
+ }
246
+ if seeked != 0 {
247
+ t.Fatal("Failed to seek to beginning")
248
+ }
249
+
250
+ out, err := ioutil.ReadAll(rs)
251
+ if err != nil {
252
+ t.Fatal(err)
253
+ }
254
+
255
+ err = arrComp(out, should)
256
+ if err != nil {
257
+ t.Fatal(err)
258
+ }
259
+}
260
+
261
+func TestSeekToAlmostBegin(t *testing.T) {
262
+ nbytes := int64(10 * 1024)
263
+ should := make([]byte, nbytes)
264
+ u.NewTimeSeededRand().Read(should)
265
+
266
+ read := bytes.NewReader(should)
267
+ ds := mdtest.Mock(t)
268
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
269
+ if err != nil {
270
+ t.Fatal(err)
271
+ }
272
+
273
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
274
+ if err != nil {
275
+ t.Fatal(err)
276
+ }
277
+
278
+ n, err := io.CopyN(ioutil.Discard, rs, 1024*4)
279
+ if err != nil {
280
+ t.Fatal(err)
281
+ }
282
+ if n != 4096 {
283
+ t.Fatal("Copy didnt copy enough bytes")
284
+ }
285
+
286
+ seeked, err := rs.Seek(1, os.SEEK_SET)
287
+ if err != nil {
288
+ t.Fatal(err)
289
+ }
290
+ if seeked != 1 {
291
+ t.Fatal("Failed to seek to almost beginning")
292
+ }
293
+
294
+ out, err := ioutil.ReadAll(rs)
295
+ if err != nil {
296
+ t.Fatal(err)
297
+ }
298
+
299
+ err = arrComp(out, should[1:])
300
+ if err != nil {
301
+ t.Fatal(err)
302
+ }
303
+}
304
+
305
+func TestSeekEnd(t *testing.T) {
306
+ nbytes := int64(50 * 1024)
307
+ should := make([]byte, nbytes)
308
+ u.NewTimeSeededRand().Read(should)
309
+
310
+ read := bytes.NewReader(should)
311
+ ds := mdtest.Mock(t)
312
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
313
+ if err != nil {
314
+ t.Fatal(err)
315
+ }
316
+
317
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
318
+ if err != nil {
319
+ t.Fatal(err)
320
+ }
321
+
322
+ seeked, err := rs.Seek(0, os.SEEK_END)
323
+ if err != nil {
324
+ t.Fatal(err)
325
+ }
326
+ if seeked != nbytes {
327
+ t.Fatal("Failed to seek to end")
328
+ }
329
+}
330
+
331
+func TestSeekEndSingleBlockFile(t *testing.T) {
332
+ nbytes := int64(100)
333
+ should := make([]byte, nbytes)
334
+ u.NewTimeSeededRand().Read(should)
335
+
336
+ read := bytes.NewReader(should)
337
+ ds := mdtest.Mock(t)
338
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{5000})
339
+ if err != nil {
340
+ t.Fatal(err)
341
+ }
342
+
343
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
344
+ if err != nil {
345
+ t.Fatal(err)
346
+ }
347
+
348
+ seeked, err := rs.Seek(0, os.SEEK_END)
349
+ if err != nil {
350
+ t.Fatal(err)
351
+ }
352
+ if seeked != nbytes {
353
+ t.Fatal("Failed to seek to end")
354
+ }
355
+}
356
+
357
+func TestSeekingStress(t *testing.T) {
358
+ nbytes := int64(1024 * 1024)
359
+ should := make([]byte, nbytes)
360
+ u.NewTimeSeededRand().Read(should)
361
+
362
+ read := bytes.NewReader(should)
363
+ ds := mdtest.Mock(t)
364
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{1000})
365
+ if err != nil {
366
+ t.Fatal(err)
367
+ }
368
+
369
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
370
+ if err != nil {
371
+ t.Fatal(err)
372
+ }
373
+
374
+ testbuf := make([]byte, nbytes)
375
+ for i := 0; i < 50; i++ {
376
+ offset := mrand.Intn(int(nbytes))
377
+ l := int(nbytes) - offset
378
+ n, err := rs.Seek(int64(offset), os.SEEK_SET)
379
+ if err != nil {
380
+ t.Fatal(err)
381
+ }
382
+ if n != int64(offset) {
383
+ t.Fatal("Seek failed to move to correct position")
384
+ }
385
+
386
+ nread, err := rs.Read(testbuf[:l])
387
+ if err != nil {
388
+ t.Fatal(err)
389
+ }
390
+ if nread != l {
391
+ t.Fatal("Failed to read enough bytes")
392
+ }
393
+
394
+ err = arrComp(testbuf[:l], should[offset:offset+l])
395
+ if err != nil {
396
+ t.Fatal(err)
397
+ }
398
+ }
399
+
400
+}
401
+
402
+func TestSeekingConsistency(t *testing.T) {
403
+ nbytes := int64(128 * 1024)
404
+ should := make([]byte, nbytes)
405
+ u.NewTimeSeededRand().Read(should)
406
+
407
+ read := bytes.NewReader(should)
408
+ ds := mdtest.Mock(t)
409
+ nd, err := buildTestDag(read, ds, &chunk.SizeSplitter{500})
410
+ if err != nil {
411
+ t.Fatal(err)
412
+ }
413
+
414
+ rs, err := uio.NewDagReader(context.Background(), nd, ds)
415
+ if err != nil {
416
+ t.Fatal(err)
417
+ }
418
+
419
+ out := make([]byte, nbytes)
420
+
421
+ for coff := nbytes - 4096; coff >= 0; coff -= 4096 {
422
+ t.Log(coff)
423
+ n, err := rs.Seek(coff, os.SEEK_SET)
424
+ if err != nil {
425
+ t.Fatal(err)
426
+ }
427
+ if n != coff {
428
+ t.Fatal("wasnt able to seek to the right position")
429
+ }
430
+ nread, err := rs.Read(out[coff : coff+4096])
431
+ if err != nil {
432
+ t.Fatal(err)
433
+ }
434
+ if nread != 4096 {
435
+ t.Fatal("didnt read the correct number of bytes")
436
+ }
437
+ }
438
+
439
+ err = arrComp(out, should)
440
+ if err != nil {
441
+ t.Fatal(err)
442
+ }
443
+}
importer/trickle/trickledag.go
new
+58
@@ -0,0 +1,58 @@
1
+package trickle
2
+
3
+import (
4
+ h "github.com/jbenet/go-ipfs/importer/helpers"
5
+ dag "github.com/jbenet/go-ipfs/merkledag"
6
+)
7
+
8
+// layerRepeat specifies how many times to append a child tree of a
9
+// given depth. Higher values increase the width of a given node, which
10
+// improves seek speeds.
11
+const layerRepeat = 4
12
+
13
+func TrickleLayout(db *h.DagBuilderHelper) (*dag.Node, error) {
14
+ root := h.NewUnixfsNode()
15
+ err := db.FillNodeLayer(root)
16
+ if err != nil {
17
+ return nil, err
18
+ }
19
+ for level := 1; !db.Done(); level++ {
20
+ for i := 0; i < layerRepeat && !db.Done(); i++ {
21
+ next := h.NewUnixfsNode()
22
+ err := fillTrickleRec(db, next, level)
23
+ if err != nil {
24
+ return nil, err
25
+ }
26
+ err = root.AddChild(next, db)
27
+ if err != nil {
28
+ return nil, err
29
+ }
30
+ }
31
+ }
32
+
33
+ return db.Add(root)
34
+}
35
+
36
+func fillTrickleRec(db *h.DagBuilderHelper, node *h.UnixfsNode, depth int) error {
37
+ // Always do this, even in the base case
38
+ err := db.FillNodeLayer(node)
39
+ if err != nil {
40
+ return err
41
+ }
42
+
43
+ for i := 1; i < depth && !db.Done(); i++ {
44
+ for j := 0; j < layerRepeat; j++ {
45
+ next := h.NewUnixfsNode()
46
+ err := fillTrickleRec(db, next, i)
47
+ if err != nil {
48
+ return err
49
+ }
50
+
51
+ err = node.AddChild(next, db)
52
+ if err != nil {
53
+ return err
54
+ }
55
+ }
56
+ }
57
+ return nil
58
+}
importer/trickledag.go
deleted
-91
@@ -1,91 +0,0 @@
1
-package importer
2
-
3
-import (
4
- "io"
5
-
6
- "github.com/jbenet/go-ipfs/importer/chunk"
7
- dag "github.com/jbenet/go-ipfs/merkledag"
8
- "github.com/jbenet/go-ipfs/pin"
9
-)
10
-
11
-// layerRepeat specifies how many times to append a child tree of a
12
-// given depth. Higher values increase the width of a given node, which
13
-// improves seek speeds.
14
-const layerRepeat = 4
15
-
16
-func BuildTrickleDagFromReader(r io.Reader, ds dag.DAGService, mp pin.ManualPinner, spl chunk.BlockSplitter) (*dag.Node, error) {
17
- // Start the splitter
18
- blkch := spl.Split(r)
19
-
20
- // Create our builder helper
21
- db := &dagBuilderHelper{
22
- dserv: ds,
23
- mp: mp,
24
- in: blkch,
25
- maxlinks: DefaultLinksPerBlock,
26
- indrSize: defaultIndirectBlockDataSize(),
27
- }
28
-
29
- root := newUnixfsNode()
30
- err := db.fillNodeRec(root, 1)
31
- if err != nil {
32
- return nil, err
33
- }
34
- for level := 1; !db.done(); level++ {
35
- for i := 0; i < layerRepeat && !db.done(); i++ {
36
- next := newUnixfsNode()
37
- err := db.fillTrickleRec(next, level)
38
- if err != nil {
39
- return nil, err
40
- }
41
- err = root.addChild(next, db)
42
- if err != nil {
43
- return nil, err
44
- }
45
- }
46
- }
47
-
48
- rootnode, err := root.getDagNode()
49
- if err != nil {
50
- return nil, err
51
- }
52
-
53
- rootkey, err := ds.Add(rootnode)
54
- if err != nil {
55
- return nil, err
56
- }
57
-
58
- if mp != nil {
59
- mp.PinWithMode(rootkey, pin.Recursive)
60
- err := mp.Flush()
61
- if err != nil {
62
- return nil, err
63
- }
64
- }
65
-
66
- return root.getDagNode()
67
-}
68
-
69
-func (db *dagBuilderHelper) fillTrickleRec(node *unixfsNode, depth int) error {
70
- // Always do this, even in the base case
71
- err := db.fillNodeRec(node, 1)
72
- if err != nil {
73
- return err
74
- }
75
-
76
- for i := 1; i < depth && !db.done(); i++ {
77
- for j := 0; j < layerRepeat; j++ {
78
- next := newUnixfsNode()
79
- err := db.fillTrickleRec(next, i)
80
- if err != nil {
81
- return err
82
- }
83
-
84
- err = node.addChild(next, db)
85
- if err != nil {
86
- return err
87
- }
88
- }
89
- }
90
- return nil
91
-}
merkledag/test/utils.go
renamed
+4
-3
@@ -1,4 +1,4 @@
1
-package merkledag
1
+package mdutils
2
3
import (
4
"testing"
@@ -8,13 +8,14 @@ import (
8
"github.com/jbenet/go-ipfs/blocks/blockstore"
9
bsrv "github.com/jbenet/go-ipfs/blockservice"
10
"github.com/jbenet/go-ipfs/exchange/offline"
11
+ dag "github.com/jbenet/go-ipfs/merkledag"
12
)
13
13
-func Mock(t testing.TB) DAGService {
14
+func Mock(t testing.TB) dag.DAGService {
15
bstore := blockstore.NewBlockstore(dssync.MutexWrap(ds.NewMapDatastore()))
16
bserv, err := bsrv.New(bstore, offline.Exchange(bstore))
17
if err != nil {
18
t.Fatal(err)
19
}
19
- return NewDAGService(bserv)
20
+ return dag.NewDAGService(bserv)
21
}