Correct pinning for dagmodifier, and a bunch more tests
Jeromy committed
Mar 10, 2015 at 12:57 UTC
50f0b06bdf3600e9af4775ed931ca6417d7daee1
5 files changed
+270
-35
importer/helpers/helpers.go
+6
-1
@@ -7,6 +7,7 @@ import (
7
dag "github.com/jbenet/go-ipfs/merkledag"
8
"github.com/jbenet/go-ipfs/pin"
9
ft "github.com/jbenet/go-ipfs/unixfs"
10
+ u "github.com/jbenet/go-ipfs/util"
11
)
12
13
// BlockSizeLimit specifies the maximum size an imported block can have.
@@ -115,7 +116,11 @@ func (n *UnixfsNode) AddChild(child *UnixfsNode, db *DagBuilderHelper) error {
116
}
117
118
// Removes the child node at the given index
118
-func (n *UnixfsNode) RemoveChild(index int) {
119
+func (n *UnixfsNode) RemoveChild(index int, dbh *DagBuilderHelper) {
120
+ k := u.Key(n.node.Links[index].Hash)
121
+ if dbh.mp != nil {
122
+ dbh.mp.RemovePinWithMode(k, pin.Indirect)
123
+ }
124
n.ufmt.RemoveBlockSize(index)
125
n.node.Links = append(n.node.Links[:index], n.node.Links[index+1:]...)
126
}
importer/trickle/trickledag.go
+1
-1
@@ -134,7 +134,7 @@ func appendFillLastChild(ufsn *h.UnixfsNode, depth int, layerFill int, db *h.Dag
134
}
135
136
// Update changed child in parent node
137
- ufsn.RemoveChild(last)
137
+ ufsn.RemoveChild(last, db)
138
err = ufsn.AddChild(nchild, db)
139
if err != nil {
140
return err
pin/pin.go
+15
@@ -47,6 +47,7 @@ type Pinner interface {
47
// may not be successful
48
type ManualPinner interface {
49
PinWithMode(util.Key, PinMode)
50
+ RemovePinWithMode(util.Key, PinMode)
51
Pinner
52
}
53
@@ -198,6 +199,20 @@ func (p *pinner) IsPinned(key util.Key) bool {
199
p.indirPin.HasKey(key)
200
}
201
202
+func (p *pinner) RemovePinWithMode(key util.Key, mode PinMode) {
203
+ switch mode {
204
+ case Direct:
205
+ p.directPin.RemoveBlock(key)
206
+ case Indirect:
207
+ p.indirPin.Decrement(key)
208
+ case Recursive:
209
+ p.recursePin.RemoveBlock(key)
210
+ default:
211
+ // programmer error, panic OK
212
+ panic("unrecognized pin type")
213
+ }
214
+}
215
+
216
// LoadPinner loads a pinner and its keysets from the given datastore
217
func LoadPinner(d ds.ThreadSafeDatastore, dserv mdag.DAGService) (Pinner, error) {
218
p := new(pinner)
unixfs/mod/dagmodifier.go
+41
-5
@@ -172,13 +172,19 @@ func (dm *DagModifier) Flush() error {
172
// Number of bytes we're going to write
173
buflen := dm.wrBuf.Len()
174
175
+ // Grab key for unpinning after mod operation
176
+ curk, err := dm.curNode.Key()
177
+ if err != nil {
178
+ return err
179
+ }
180
+
181
// overwrite existing dag nodes
176
- k, done, err := dm.modifyDag(dm.curNode, dm.writeStart, dm.wrBuf)
182
+ thisk, done, err := dm.modifyDag(dm.curNode, dm.writeStart, dm.wrBuf)
183
if err != nil {
184
return err
185
}
186
181
- nd, err := dm.dagserv.Get(k)
187
+ nd, err := dm.dagserv.Get(thisk)
188
if err != nil {
189
return err
190
}
@@ -193,7 +199,7 @@ func (dm *DagModifier) Flush() error {
199
return err
200
}
201
196
- _, err := dm.dagserv.Add(nd)
202
+ thisk, err = dm.dagserv.Add(nd)
203
if err != nil {
204
return err
205
}
@@ -201,6 +207,14 @@ func (dm *DagModifier) Flush() error {
207
dm.curNode = nd
208
}
209
210
+ // Finalize correct pinning, and flush pinner
211
+ dm.mp.PinWithMode(thisk, pin.Recursive)
212
+ dm.mp.RemovePinWithMode(curk, pin.Recursive)
213
+ err = dm.mp.Flush()
214
+ if err != nil {
215
+ return err
216
+ }
217
+
218
dm.writeStart += uint64(buflen)
219
220
dm.wrBuf = nil
@@ -237,7 +251,7 @@ func (dm *DagModifier) modifyDag(node *mdag.Node, offset uint64, data io.Reader)
251
252
// Hey look! we're done!
253
var done bool
240
- if n < len(f.Data) {
254
+ if n < len(f.Data[offset:]) {
255
done = true
256
}
257
@@ -249,6 +263,10 @@ func (dm *DagModifier) modifyDag(node *mdag.Node, offset uint64, data io.Reader)
263
for i, bs := range f.GetBlocksizes() {
264
// We found the correct child to write into
265
if cur+bs > offset {
266
+ // Unpin block
267
+ ckey := u.Key(node.Links[i].Hash)
268
+ dm.mp.RemovePinWithMode(ckey, pin.Indirect)
269
+
270
child, err := node.Links[i].GetNode(dm.dagserv)
271
if err != nil {
272
return "", false, err
@@ -258,14 +276,24 @@ func (dm *DagModifier) modifyDag(node *mdag.Node, offset uint64, data io.Reader)
276
return "", false, err
277
}
278
279
+ // pin the new node
280
+ dm.mp.PinWithMode(k, pin.Indirect)
281
+
282
offset += bs
283
node.Links[i].Hash = mh.Multihash(k)
284
285
+ // Recache serialized node
286
+ _, err = node.Encoded(true)
287
+ if err != nil {
288
+ return "", false, err
289
+ }
290
+
291
if sdone {
292
// No more bytes to write!
293
done = true
294
break
295
}
296
+ offset = cur + bs
297
}
298
cur += bs
299
}
@@ -293,7 +321,8 @@ func (dm *DagModifier) Read(b []byte) (int, error) {
321
}
322
323
if dm.read == nil {
296
- dr, err := uio.NewDagReader(dm.ctx, dm.curNode, dm.dagserv)
324
+ ctx, cancel := context.WithCancel(dm.ctx)
325
+ dr, err := uio.NewDagReader(ctx, dm.curNode, dm.dagserv)
326
if err != nil {
327
return 0, err
328
}
@@ -307,6 +336,7 @@ func (dm *DagModifier) Read(b []byte) (int, error) {
336
return 0, ErrSeekFail
337
}
338
339
+ dm.readCancel = cancel
340
dm.read = dr
341
}
342
@@ -451,5 +481,11 @@ func dagTruncate(nd *mdag.Node, size uint64, ds mdag.DAGService) (*mdag.Node, er
481
482
nd.Data = d
483
484
+ // invalidate cache and recompute serialized data
485
+ _, err = nd.Encoded(true)
486
+ if err != nil {
487
+ return nil, err
488
+ }
489
+
490
return nd, nil
491
}
unixfs/mod/dagmodifier_test.go
+207
-28
@@ -4,6 +4,8 @@ import (
4
"fmt"
5
"io"
6
"io/ioutil"
7
+ "math/rand"
8
+ "os"
9
"testing"
10
11
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
@@ -15,6 +17,7 @@ import (
17
h "github.com/jbenet/go-ipfs/importer/helpers"
18
trickle "github.com/jbenet/go-ipfs/importer/trickle"
19
mdag "github.com/jbenet/go-ipfs/merkledag"
20
+ pin "github.com/jbenet/go-ipfs/pin"
21
ft "github.com/jbenet/go-ipfs/unixfs"
22
uio "github.com/jbenet/go-ipfs/unixfs/io"
23
u "github.com/jbenet/go-ipfs/util"
@@ -23,7 +26,7 @@ import (
26
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
27
)
28
26
-func getMockDagServ(t testing.TB) mdag.DAGService {
29
+func getMockDagServ(t testing.TB) (mdag.DAGService, pin.ManualPinner) {
30
dstore := ds.NewMapDatastore()
31
tsds := sync.MutexWrap(dstore)
32
bstore := blockstore.NewBlockstore(tsds)
@@ -31,12 +34,25 @@ func getMockDagServ(t testing.TB) mdag.DAGService {
34
if err != nil {
35
t.Fatal(err)
36
}
34
- return mdag.NewDAGService(bserv)
37
+ dserv := mdag.NewDAGService(bserv)
38
+ return dserv, pin.NewPinner(tsds, dserv).GetManual()
39
}
40
37
-func getNode(t testing.TB, dserv mdag.DAGService, size int64) ([]byte, *mdag.Node) {
41
+func getMockDagServAndBstore(t testing.TB) (mdag.DAGService, blockstore.Blockstore, pin.ManualPinner) {
42
+ dstore := ds.NewMapDatastore()
43
+ tsds := sync.MutexWrap(dstore)
44
+ bstore := blockstore.NewBlockstore(tsds)
45
+ bserv, err := bs.New(bstore, offline.Exchange(bstore))
46
+ if err != nil {
47
+ t.Fatal(err)
48
+ }
49
+ dserv := mdag.NewDAGService(bserv)
50
+ return dserv, bstore, pin.NewPinner(tsds, dserv).GetManual()
51
+}
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)
39
- node, err := imp.BuildTrickleDagFromReader(in, dserv, nil, &chunk.SizeSplitter{500})
55
+ node, err := imp.BuildTrickleDagFromReader(in, dserv, pinner, &chunk.SizeSplitter{500})
56
if err != nil {
57
t.Fatal(err)
58
}
@@ -101,12 +117,12 @@ func testModWrite(t *testing.T, beg, size uint64, orig []byte, dm *DagModifier)
117
}
118
119
func TestDagModifierBasic(t *testing.T) {
104
- dserv := getMockDagServ(t)
105
- b, n := getNode(t, dserv, 50000)
120
+ dserv, pin := getMockDagServ(t)
121
+ b, n := getNode(t, dserv, 50000, pin)
122
ctx, cancel := context.WithCancel(context.Background())
123
defer cancel()
124
109
- dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
125
+ dagmod, err := NewDagModifier(ctx, n, dserv, pin, &chunk.SizeSplitter{Size: 512})
126
if err != nil {
127
t.Fatal(err)
128
}
@@ -155,13 +171,13 @@ func TestDagModifierBasic(t *testing.T) {
171
}
172
173
func TestMultiWrite(t *testing.T) {
158
- dserv := getMockDagServ(t)
159
- _, n := getNode(t, dserv, 0)
174
+ dserv, pins := getMockDagServ(t)
175
+ _, n := getNode(t, dserv, 0, pins)
176
177
ctx, cancel := context.WithCancel(context.Background())
178
defer cancel()
179
164
- dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
180
+ dagmod, err := NewDagModifier(ctx, n, dserv, pins, &chunk.SizeSplitter{Size: 512})
181
if err != nil {
182
t.Fatal(err)
183
}
@@ -208,13 +224,13 @@ func TestMultiWrite(t *testing.T) {
224
}
225
226
func TestMultiWriteAndFlush(t *testing.T) {
211
- dserv := getMockDagServ(t)
212
- _, n := getNode(t, dserv, 0)
227
+ dserv, pins := getMockDagServ(t)
228
+ _, n := getNode(t, dserv, 0, pins)
229
230
ctx, cancel := context.WithCancel(context.Background())
231
defer cancel()
232
217
- dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
233
+ dagmod, err := NewDagModifier(ctx, n, dserv, pins, &chunk.SizeSplitter{Size: 512})
234
if err != nil {
235
t.Fatal(err)
236
}
@@ -256,13 +272,13 @@ func TestMultiWriteAndFlush(t *testing.T) {
272
}
273
274
func TestWriteNewFile(t *testing.T) {
259
- dserv := getMockDagServ(t)
260
- _, n := getNode(t, dserv, 0)
275
+ dserv, pins := getMockDagServ(t)
276
+ _, n := getNode(t, dserv, 0, pins)
277
278
ctx, cancel := context.WithCancel(context.Background())
279
defer cancel()
280
265
- dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
281
+ dagmod, err := NewDagModifier(ctx, n, dserv, pins, &chunk.SizeSplitter{Size: 512})
282
if err != nil {
283
t.Fatal(err)
284
}
@@ -299,13 +315,13 @@ func TestWriteNewFile(t *testing.T) {
315
}
316
317
func TestMultiWriteCoal(t *testing.T) {
302
- dserv := getMockDagServ(t)
303
- _, n := getNode(t, dserv, 0)
318
+ dserv, pins := getMockDagServ(t)
319
+ _, n := getNode(t, dserv, 0, pins)
320
321
ctx, cancel := context.WithCancel(context.Background())
322
defer cancel()
323
308
- dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
324
+ dagmod, err := NewDagModifier(ctx, n, dserv, pins, &chunk.SizeSplitter{Size: 512})
325
if err != nil {
326
t.Fatal(err)
327
}
@@ -345,13 +361,13 @@ func TestMultiWriteCoal(t *testing.T) {
361
}
362
363
func TestLargeWriteChunks(t *testing.T) {
348
- dserv := getMockDagServ(t)
349
- _, n := getNode(t, dserv, 0)
364
+ dserv, pins := getMockDagServ(t)
365
+ _, n := getNode(t, dserv, 0, pins)
366
367
ctx, cancel := context.WithCancel(context.Background())
368
defer cancel()
369
354
- dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
370
+ dagmod, err := NewDagModifier(ctx, n, dserv, pins, &chunk.SizeSplitter{Size: 512})
371
if err != nil {
372
t.Fatal(err)
373
}
@@ -384,12 +400,12 @@ func TestLargeWriteChunks(t *testing.T) {
400
}
401
402
func TestDagTruncate(t *testing.T) {
387
- dserv := getMockDagServ(t)
388
- b, n := getNode(t, dserv, 50000)
403
+ dserv, pins := getMockDagServ(t)
404
+ b, n := getNode(t, dserv, 50000, pins)
405
ctx, cancel := context.WithCancel(context.Background())
406
defer cancel()
407
392
- dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
408
+ dagmod, err := NewDagModifier(ctx, n, dserv, pins, &chunk.SizeSplitter{Size: 512})
409
if err != nil {
410
t.Fatal(err)
411
}
@@ -399,6 +415,11 @@ func TestDagTruncate(t *testing.T) {
415
t.Fatal(err)
416
}
417
418
+ _, err = dagmod.Seek(0, os.SEEK_SET)
419
+ if err != nil {
420
+ t.Fatal(err)
421
+ }
422
+
423
out, err := ioutil.ReadAll(dagmod)
424
if err != nil {
425
t.Fatal(err)
@@ -409,16 +430,174 @@ func TestDagTruncate(t *testing.T) {
430
}
431
}
432
433
+func TestSparseWrite(t *testing.T) {
434
+ dserv, pins := getMockDagServ(t)
435
+ _, n := getNode(t, dserv, 0, pins)
436
+ ctx, cancel := context.WithCancel(context.Background())
437
+ defer cancel()
438
+
439
+ dagmod, err := NewDagModifier(ctx, n, dserv, pins, &chunk.SizeSplitter{Size: 512})
440
+ if err != nil {
441
+ t.Fatal(err)
442
+ }
443
+
444
+ buf := make([]byte, 5000)
445
+ u.NewTimeSeededRand().Read(buf[2500:])
446
+
447
+ wrote, err := dagmod.WriteAt(buf[2500:], 2500)
448
+ if err != nil {
449
+ t.Fatal(err)
450
+ }
451
+
452
+ if wrote != 2500 {
453
+ t.Fatal("incorrect write amount")
454
+ }
455
+
456
+ _, err = dagmod.Seek(0, os.SEEK_SET)
457
+ if err != nil {
458
+ t.Fatal(err)
459
+ }
460
+
461
+ out, err := ioutil.ReadAll(dagmod)
462
+ if err != nil {
463
+ t.Fatal(err)
464
+ }
465
+
466
+ if err = arrComp(out, buf); err != nil {
467
+ t.Fatal(err)
468
+ }
469
+}
470
+
471
+func basicGC(t *testing.T, bs blockstore.Blockstore, pins pin.ManualPinner) {
472
+ ctx, cancel := context.WithCancel(context.Background())
473
+ defer cancel() // in case error occurs during operation
474
+ keychan, err := bs.AllKeysChan(ctx)
475
+ if err != nil {
476
+ t.Fatal(err)
477
+ }
478
+ for k := range keychan { // rely on AllKeysChan to close chan
479
+ if !pins.IsPinned(k) {
480
+ err := bs.DeleteBlock(k)
481
+ if err != nil {
482
+ t.Fatal(err)
483
+ }
484
+ }
485
+ }
486
+}
487
+func TestCorrectPinning(t *testing.T) {
488
+ dserv, bstore, pins := getMockDagServAndBstore(t)
489
+ b, n := getNode(t, dserv, 50000, pins)
490
+ ctx, cancel := context.WithCancel(context.Background())
491
+ defer cancel()
492
+
493
+ dagmod, err := NewDagModifier(ctx, n, dserv, pins, &chunk.SizeSplitter{Size: 512})
494
+ if err != nil {
495
+ t.Fatal(err)
496
+ }
497
+
498
+ buf := make([]byte, 1024)
499
+ for i := 0; i < 100; i++ {
500
+ size, err := dagmod.Size()
501
+ if err != nil {
502
+ t.Fatal(err)
503
+ }
504
+ offset := rand.Intn(int(size))
505
+ u.NewTimeSeededRand().Read(buf)
506
+
507
+ if offset+len(buf) > int(size) {
508
+ b = append(b[:offset], buf...)
509
+ } else {
510
+ copy(b[offset:], buf)
511
+ }
512
+
513
+ n, err := dagmod.WriteAt(buf, int64(offset))
514
+ if err != nil {
515
+ t.Fatal(err)
516
+ }
517
+ if n != len(buf) {
518
+ t.Fatal("wrote incorrect number of bytes")
519
+ }
520
+ }
521
+
522
+ fisize, err := dagmod.Size()
523
+ if err != nil {
524
+ t.Fatal(err)
525
+ }
526
+
527
+ if int(fisize) != len(b) {
528
+ t.Fatal("reported filesize incorrect", fisize, len(b))
529
+ }
530
+
531
+ // Run a GC, then ensure we can still read the file correctly
532
+ basicGC(t, bstore, pins)
533
+
534
+ nd, err := dagmod.GetNode()
535
+ if err != nil {
536
+ t.Fatal(err)
537
+ }
538
+ read, err := uio.NewDagReader(context.Background(), nd, dserv)
539
+ if err != nil {
540
+ t.Fatal(err)
541
+ }
542
+
543
+ out, err := ioutil.ReadAll(read)
544
+ if err != nil {
545
+ t.Fatal(err)
546
+ }
547
+
548
+ if err = arrComp(out, b); err != nil {
549
+ t.Fatal(err)
550
+ }
551
+
552
+ rootk, err := nd.Key()
553
+ if err != nil {
554
+ t.Fatal(err)
555
+ }
556
+
557
+ // Verify only one recursive pin
558
+ recpins := pins.RecursiveKeys()
559
+ if len(recpins) != 1 {
560
+ t.Fatal("Incorrect number of pinned entries")
561
+ }
562
+
563
+ // verify the correct node is pinned
564
+ if recpins[0] != rootk {
565
+ t.Fatal("Incorrect node recursively pinned")
566
+ }
567
+
568
+ indirpins := pins.IndirectKeys()
569
+ children := enumerateChildren(t, nd, dserv)
570
+ if len(indirpins) != len(children) {
571
+ t.Log(len(indirpins), len(children))
572
+ t.Fatal("Incorrect number of indirectly pinned blocks")
573
+ }
574
+
575
+}
576
+
577
+func enumerateChildren(t *testing.T, nd *mdag.Node, ds mdag.DAGService) []u.Key {
578
+ var out []u.Key
579
+ for _, lnk := range nd.Links {
580
+ out = append(out, u.Key(lnk.Hash))
581
+ child, err := lnk.GetNode(ds)
582
+ if err != nil {
583
+ t.Fatal(err)
584
+ }
585
+ children := enumerateChildren(t, child, ds)
586
+ out = append(out, children...)
587
+ }
588
+ return out
589
+}
590
+
591
func BenchmarkDagmodWrite(b *testing.B) {
592
b.StopTimer()
414
- dserv := getMockDagServ(b)
415
- _, n := getNode(b, dserv, 0)
593
+ dserv, pins := getMockDagServ(b)
594
+ _, n := getNode(b, dserv, 0, pins)
595
ctx, cancel := context.WithCancel(context.Background())
596
defer cancel()
597
598
wrsize := 4096
599
421
- dagmod, err := NewDagModifier(ctx, n, dserv, nil, &chunk.SizeSplitter{Size: 512})
600
+ dagmod, err := NewDagModifier(ctx, n, dserv, pins, &chunk.SizeSplitter{Size: 512})
601
if err != nil {
602
b.Fatal(err)
603
}