Add test for running gc during an add
License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>
Jeromy committed
Dec 5, 2015 at 23:42 UTC
5dd32d6491e672d1f5ee174036647133b0fed571
2 files changed
+131
-1
core/coreunix/add.go
+1
-1
@@ -155,7 +155,7 @@ func (params *Adder) PinRoot() error {
155
return nil
156
}
157
158
- rnk, err := root.Key()
158
+ rnk, err := params.node.DAG.Add(root)
159
if err != nil {
160
return err
161
}
core/coreunix/add_test.go
+130
@@ -1,10 +1,18 @@
1
package coreunix
2
3
import (
4
+ "bytes"
5
+ "io"
6
+ "io/ioutil"
7
"testing"
8
+ "time"
9
10
"github.com/ipfs/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
11
+ "github.com/ipfs/go-ipfs/blocks/key"
12
+ "github.com/ipfs/go-ipfs/commands/files"
13
"github.com/ipfs/go-ipfs/core"
14
+ dag "github.com/ipfs/go-ipfs/merkledag"
15
+ "github.com/ipfs/go-ipfs/pin/gc"
16
"github.com/ipfs/go-ipfs/repo"
17
"github.com/ipfs/go-ipfs/repo/config"
18
"github.com/ipfs/go-ipfs/util/testutil"
@@ -29,3 +37,125 @@ func TestAddRecursive(t *testing.T) {
37
t.Fatal("keys do not match: ", k)
38
}
39
}
40
+
41
+func TestAddGCLive(t *testing.T) {
42
+ r := &repo.Mock{
43
+ C: config.Config{
44
+ Identity: config.Identity{
45
+ PeerID: "Qmfoo", // required by offline node
46
+ },
47
+ },
48
+ D: testutil.ThreadSafeCloserMapDatastore(),
49
+ }
50
+ node, err := core.NewNode(context.Background(), &core.BuildCfg{Repo: r})
51
+ if err != nil {
52
+ t.Fatal(err)
53
+ }
54
+
55
+ errs := make(chan error)
56
+ out := make(chan interface{})
57
+ adder, err := NewAdder(context.Background(), node, out)
58
+ if err != nil {
59
+ t.Fatal(err)
60
+ }
61
+
62
+ dataa := ioutil.NopCloser(bytes.NewBufferString("testfileA"))
63
+ rfa := files.NewReaderFile("a", "a", dataa, nil)
64
+
65
+ // make two files with pipes so we can 'pause' the add for timing of the test
66
+ piper, pipew := io.Pipe()
67
+ hangfile := files.NewReaderFile("b", "b", piper, nil)
68
+
69
+ datad := ioutil.NopCloser(bytes.NewBufferString("testfileD"))
70
+ rfd := files.NewReaderFile("d", "d", datad, nil)
71
+
72
+ slf := files.NewSliceFile("files", "files", []files.File{rfa, hangfile, rfd})
73
+
74
+ addDone := make(chan struct{})
75
+ go func() {
76
+ defer close(addDone)
77
+ defer close(out)
78
+ err := adder.AddFile(slf)
79
+
80
+ if err != nil {
81
+ t.Fatal(err)
82
+ }
83
+
84
+ }()
85
+
86
+ addedHashes := make(map[string]struct{})
87
+ select {
88
+ case o := <-out:
89
+ addedHashes[o.(*AddedObject).Hash] = struct{}{}
90
+ case <-addDone:
91
+ t.Fatal("add shouldnt complete yet")
92
+ }
93
+
94
+ var gcout <-chan key.Key
95
+ gcstarted := make(chan struct{})
96
+ go func() {
97
+ defer close(gcstarted)
98
+ gcchan, err := gc.GC(context.Background(), node.Blockstore, node.Pinning)
99
+ if err != nil {
100
+ log.Error("GC ERROR:", err)
101
+ errs <- err
102
+ return
103
+ }
104
+
105
+ gcout = gcchan
106
+ }()
107
+
108
+ // gc shouldnt start until we let the add finish its current file.
109
+ pipew.Write([]byte("some data for file b"))
110
+
111
+ select {
112
+ case <-gcstarted:
113
+ t.Fatal("gc shouldnt have started yet")
114
+ case err := <-errs:
115
+ t.Fatal(err)
116
+ default:
117
+ }
118
+
119
+ time.Sleep(time.Millisecond * 100) // make sure gc gets to requesting lock
120
+
121
+ // finish write and unblock gc
122
+ pipew.Close()
123
+
124
+ // receive next object from adder
125
+ select {
126
+ case o := <-out:
127
+ addedHashes[o.(*AddedObject).Hash] = struct{}{}
128
+ case err := <-errs:
129
+ t.Fatal(err)
130
+ }
131
+
132
+ select {
133
+ case <-gcstarted:
134
+ case err := <-errs:
135
+ t.Fatal(err)
136
+ }
137
+
138
+ for k := range gcout {
139
+ if _, ok := addedHashes[k.B58String()]; ok {
140
+ t.Fatal("gc'ed a hash we just added")
141
+ }
142
+ }
143
+
144
+ var last key.Key
145
+ for a := range out {
146
+ // wait for it to finish
147
+ last = key.B58KeyDecode(a.(*AddedObject).Hash)
148
+ }
149
+
150
+ ctx, cancel := context.WithTimeout(context.Background(), time.Second*5)
151
+ defer cancel()
152
+ root, err := node.DAG.Get(ctx, last)
153
+ if err != nil {
154
+ t.Fatal(err)
155
+ }
156
+
157
+ err = dag.EnumerateChildren(ctx, node.DAG, root, key.NewKeySet())
158
+ if err != nil {
159
+ t.Fatal(err)
160
+ }
161
+}