refactor(mdag, bserv, bs) mocks, etc.
License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>
Brian Tiger Chow committed
Dec 12, 2014 at 20:05 UTC
8e0c8a7a7e427bf0e2d2d33cdd2806f13e455cf6
14 files changed
+104
-92
blockservice/blocks_test.go
+2
-19
@@ -11,10 +11,7 @@ import (
11
blocks "github.com/jbenet/go-ipfs/blocks"
12
blockstore "github.com/jbenet/go-ipfs/blocks/blockstore"
13
blocksutil "github.com/jbenet/go-ipfs/blocks/blocksutil"
14
- bitswap "github.com/jbenet/go-ipfs/exchange/bitswap"
15
- tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
14
offline "github.com/jbenet/go-ipfs/exchange/offline"
17
- "github.com/jbenet/go-ipfs/routing/mock"
15
u "github.com/jbenet/go-ipfs/util"
16
)
17
@@ -63,23 +60,9 @@ func TestBlocks(t *testing.T) {
60
}
61
62
func TestGetBlocksSequential(t *testing.T) {
66
- net := tn.VirtualNetwork()
67
- rs := mock.VirtualRoutingServer()
68
- sg := bitswap.NewSessionGenerator(net, rs)
63
+ var servs = Mocks(t, 4)
64
bg := blocksutil.NewBlockGenerator()
70
-
71
- instances := sg.Instances(4)
65
blks := bg.Blocks(50)
73
- // TODO: verify no duplicates
74
-
75
- var servs []*BlockService
76
- for _, i := range instances {
77
- bserv, err := New(i.Blockstore, i.Exchange)
78
- if err != nil {
79
- t.Fatal(err)
80
- }
81
- servs = append(servs, bserv)
82
- }
66
67
var keys []u.Key
68
for _, blk := range blks {
@@ -89,7 +72,7 @@ func TestGetBlocksSequential(t *testing.T) {
72
73
t.Log("one instance at a time, get blocks concurrently")
74
92
- for i := 1; i < len(instances); i++ {
75
+ for i := 1; i < len(servs); i++ {
76
ctx, _ := context.WithTimeout(context.TODO(), time.Second*5)
77
out := servs[i].GetBlocks(ctx, keys)
78
gotten := make(map[u.Key]*blocks.Block)
blockservice/mock.go
new
+28
@@ -0,0 +1,28 @@
1
+package blockservice
2
+
3
+import (
4
+ "testing"
5
+
6
+ bitswap "github.com/jbenet/go-ipfs/exchange/bitswap"
7
+ tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
8
+ "github.com/jbenet/go-ipfs/routing/mock"
9
+)
10
+
11
+// Mocks returns |n| connected mock Blockservices
12
+func Mocks(t *testing.T, n int) []*BlockService {
13
+ net := tn.VirtualNetwork()
14
+ rs := mock.VirtualRoutingServer()
15
+ sg := bitswap.NewSessionGenerator(net, rs)
16
+
17
+ instances := sg.Instances(n)
18
+
19
+ var servs []*BlockService
20
+ for _, i := range instances {
21
+ bserv, err := New(i.Blockstore(), i.Exchange)
22
+ if err != nil {
23
+ t.Fatal(err)
24
+ }
25
+ servs = append(servs, bserv)
26
+ }
27
+ return servs
28
+}
core/core.go
+3
-3
@@ -28,10 +28,10 @@ import (
28
pin "github.com/jbenet/go-ipfs/pin"
29
routing "github.com/jbenet/go-ipfs/routing"
30
dht "github.com/jbenet/go-ipfs/routing/dht"
31
- u "github.com/jbenet/go-ipfs/util"
31
ctxc "github.com/jbenet/go-ipfs/util/ctxcloser"
32
+ ds2 "github.com/jbenet/go-ipfs/util/datastore2"
33
debugerror "github.com/jbenet/go-ipfs/util/debugerror"
34
- "github.com/jbenet/go-ipfs/util/eventlog"
34
+ eventlog "github.com/jbenet/go-ipfs/util/eventlog"
35
)
36
37
const IpnsValidatorTag = "ipns"
@@ -52,7 +52,7 @@ type IpfsNode struct {
52
Peerstore peer.Peerstore
53
54
// the local datastore
55
- Datastore u.ThreadSafeDatastoreCloser
55
+ Datastore ds2.ThreadSafeDatastoreCloser
56
57
// the network message stream
58
Network inet.Network
core/datastore.go
+5
-4
@@ -10,10 +10,11 @@ import (
10
11
config "github.com/jbenet/go-ipfs/config"
12
u "github.com/jbenet/go-ipfs/util"
13
+ ds2 "github.com/jbenet/go-ipfs/util/datastore2"
14
"github.com/jbenet/go-ipfs/util/debugerror"
15
)
16
16
-func makeDatastore(cfg config.Datastore) (u.ThreadSafeDatastoreCloser, error) {
17
+func makeDatastore(cfg config.Datastore) (ds2.ThreadSafeDatastoreCloser, error) {
18
if len(cfg.Type) == 0 {
19
return nil, debugerror.Errorf("config datastore.type required")
20
}
@@ -23,7 +24,7 @@ func makeDatastore(cfg config.Datastore) (u.ThreadSafeDatastoreCloser, error) {
24
return makeLevelDBDatastore(cfg)
25
26
case "memory":
26
- return u.CloserWrap(syncds.MutexWrap(ds.NewMapDatastore())), nil
27
+ return ds2.CloserWrap(syncds.MutexWrap(ds.NewMapDatastore())), nil
28
29
case "fs":
30
log.Warning("using fs.Datastore at .datastore for testing.")
@@ -32,13 +33,13 @@ func makeDatastore(cfg config.Datastore) (u.ThreadSafeDatastoreCloser, error) {
33
return nil, err
34
}
35
ktd := ktds.Wrap(d, u.B58KeyConverter)
35
- return u.CloserWrap(syncds.MutexWrap(ktd)), nil
36
+ return ds2.CloserWrap(syncds.MutexWrap(ktd)), nil
37
}
38
39
return nil, debugerror.Errorf("Unknown datastore type: %s", cfg.Type)
40
}
41
41
-func makeLevelDBDatastore(cfg config.Datastore) (u.ThreadSafeDatastoreCloser, error) {
42
+func makeLevelDBDatastore(cfg config.Datastore) (ds2.ThreadSafeDatastoreCloser, error) {
43
if len(cfg.Path) == 0 {
44
return nil, debugerror.Errorf("config datastore.path required for leveldb")
45
}
core/mock.go
+2
-2
@@ -12,7 +12,7 @@ import (
12
path "github.com/jbenet/go-ipfs/path"
13
peer "github.com/jbenet/go-ipfs/peer"
14
mdht "github.com/jbenet/go-ipfs/routing/mock"
15
- "github.com/jbenet/go-ipfs/util"
15
+ ds2 "github.com/jbenet/go-ipfs/util/datastore2"
16
)
17
18
// NewMockNode constructs an IpfsNode for use in tests.
@@ -39,7 +39,7 @@ func NewMockNode() (*IpfsNode, error) {
39
40
// Temp Datastore
41
dstore := ds.NewMapDatastore()
42
- nd.Datastore = util.CloserWrap(syncds.MutexWrap(dstore))
42
+ nd.Datastore = ds2.CloserWrap(syncds.MutexWrap(dstore))
43
44
// Routing
45
dht := mdht.NewMockRouter(nd.Identity, nd.Datastore)
exchange/bitswap/bitswap_test.go
+7
-7
@@ -76,7 +76,7 @@ func TestGetBlockFromPeerAfterPeerAnnounces(t *testing.T) {
76
77
hasBlock := g.Next()
78
79
- if err := hasBlock.Blockstore.Put(block); err != nil {
79
+ if err := hasBlock.Blockstore().Put(block); err != nil {
80
t.Fatal(err)
81
}
82
if err := hasBlock.Exchange.HasBlock(context.Background(), block); err != nil {
@@ -135,7 +135,7 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
135
136
first := instances[0]
137
for _, b := range blocks {
138
- first.Blockstore.Put(b)
138
+ first.Blockstore().Put(b)
139
first.Exchange.HasBlock(context.Background(), b)
140
rs.Announce(first.Peer, b.Key())
141
}
@@ -158,7 +158,7 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
158
159
for _, inst := range instances {
160
for _, b := range blocks {
161
- if _, err := inst.Blockstore.Get(b.Key()); err != nil {
161
+ if _, err := inst.Blockstore().Get(b.Key()); err != nil {
162
t.Fatal(err)
163
}
164
}
@@ -166,7 +166,7 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
166
}
167
168
func getOrFail(bitswap Instance, b *blocks.Block, t *testing.T, wg *sync.WaitGroup) {
169
- if _, err := bitswap.Blockstore.Get(b.Key()); err != nil {
169
+ if _, err := bitswap.Blockstore().Get(b.Key()); err != nil {
170
_, err := bitswap.Exchange.GetBlock(context.Background(), b.Key())
171
if err != nil {
172
t.Fatal(err)
@@ -208,7 +208,7 @@ func TestSendToWantingPeer(t *testing.T) {
208
beta := bg.Next()
209
t.Logf("Peer %v announes availability of %v\n", w.Peer, beta.Key())
210
ctx, _ = context.WithTimeout(context.Background(), timeout)
211
- if err := w.Blockstore.Put(beta); err != nil {
211
+ if err := w.Blockstore().Put(beta); err != nil {
212
t.Fatal(err)
213
}
214
w.Exchange.HasBlock(ctx, beta)
@@ -221,7 +221,7 @@ func TestSendToWantingPeer(t *testing.T) {
221
222
t.Logf("%v announces availability of %v\n", o.Peer, alpha.Key())
223
ctx, _ = context.WithTimeout(context.Background(), timeout)
224
- if err := o.Blockstore.Put(alpha); err != nil {
224
+ if err := o.Blockstore().Put(alpha); err != nil {
225
t.Fatal(err)
226
}
227
o.Exchange.HasBlock(ctx, alpha)
@@ -233,7 +233,7 @@ func TestSendToWantingPeer(t *testing.T) {
233
}
234
235
t.Logf("%v should now have %v\n", w.Peer, alpha.Key())
236
- block, err := w.Blockstore.Get(alpha.Key())
236
+ block, err := w.Blockstore().Get(alpha.Key())
237
if err != nil {
238
t.Fatalf("Should not have received an error: %s", err)
239
}
exchange/bitswap/testutils.go
+27
-10
@@ -1,14 +1,18 @@
1
package bitswap
2
3
import (
4
- "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
4
+ "time"
5
+
6
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
8
ds_sync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
7
- "github.com/jbenet/go-ipfs/blocks/blockstore"
8
- "github.com/jbenet/go-ipfs/exchange"
9
+ blockstore "github.com/jbenet/go-ipfs/blocks/blockstore"
10
+ exchange "github.com/jbenet/go-ipfs/exchange"
11
tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
10
- "github.com/jbenet/go-ipfs/peer"
11
- "github.com/jbenet/go-ipfs/routing/mock"
12
+ peer "github.com/jbenet/go-ipfs/peer"
13
+ mock "github.com/jbenet/go-ipfs/routing/mock"
14
+ datastore2 "github.com/jbenet/go-ipfs/util/datastore2"
15
+ delay "github.com/jbenet/go-ipfs/util/delay"
16
)
17
18
func NewSessionGenerator(
@@ -45,7 +49,17 @@ func (g *SessionGenerator) Instances(n int) []Instance {
49
type Instance struct {
50
Peer peer.Peer
51
Exchange exchange.Interface
48
- Blockstore blockstore.Blockstore
52
+ blockstore blockstore.Blockstore
53
+
54
+ blockstoreDelay delay.D
55
+}
56
+
57
+func (i *Instance) Blockstore() blockstore.Blockstore {
58
+ return i.blockstore
59
+}
60
+
61
+func (i *Instance) SetBlockstoreLatency(t time.Duration) time.Duration {
62
+ return i.blockstoreDelay.Set(t)
63
}
64
65
// session creates a test bitswap session.
@@ -58,7 +72,9 @@ func session(net tn.Network, rs mock.RoutingServer, ps peer.Peerstore, id peer.I
72
73
adapter := net.Adapter(p)
74
htc := rs.Client(p)
61
- bstore := blockstore.NewBlockstore(ds_sync.MutexWrap(ds.NewMapDatastore()))
75
+
76
+ bsdelay := delay.Fixed(0)
77
+ bstore := blockstore.NewBlockstore(ds_sync.MutexWrap(datastore2.WithDelay(ds.NewMapDatastore(), bsdelay)))
78
79
const alwaysSendToPeer = true
80
ctx := context.TODO()
@@ -66,8 +82,9 @@ func session(net tn.Network, rs mock.RoutingServer, ps peer.Peerstore, id peer.I
82
bs := New(ctx, p, adapter, htc, bstore, alwaysSendToPeer)
83
84
return Instance{
69
- Peer: p,
70
- Exchange: bs,
71
- Blockstore: bstore,
85
+ Peer: p,
86
+ Exchange: bs,
87
+ blockstore: bstore,
88
+ blockstoreDelay: bsdelay,
89
}
90
}
importer/importer_test.go
+3
-3
@@ -9,10 +9,10 @@ import (
9
"os"
10
"testing"
11
12
- "github.com/jbenet/go-ipfs/importer/chunk"
12
+ chunk "github.com/jbenet/go-ipfs/importer/chunk"
13
+ merkledag "github.com/jbenet/go-ipfs/merkledag"
14
uio "github.com/jbenet/go-ipfs/unixfs/io"
15
u "github.com/jbenet/go-ipfs/util"
15
- testutil "github.com/jbenet/go-ipfs/util/testutil"
16
)
17
18
// NOTE:
@@ -91,7 +91,7 @@ func TestBuilderConsistency(t *testing.T) {
91
buf := new(bytes.Buffer)
92
io.CopyN(buf, u.NewTimeSeededRand(), int64(nbytes))
93
should := dup(buf.Bytes())
94
- dagserv := testutil.GetDAGServ(t)
94
+ dagserv := merkledag.Mock(t)
95
nd, err := BuildDagFromReader(buf, dagserv, nil, chunk.DefaultSplitter)
96
if err != nil {
97
t.Fatal(err)
merkledag/merkledag_test.go
+2
-17
@@ -7,13 +7,10 @@ import (
7
"io/ioutil"
8
"testing"
9
10
- bserv "github.com/jbenet/go-ipfs/blockservice"
11
- bs "github.com/jbenet/go-ipfs/exchange/bitswap"
12
- tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
10
+ blockservice "github.com/jbenet/go-ipfs/blockservice"
11
imp "github.com/jbenet/go-ipfs/importer"
12
chunk "github.com/jbenet/go-ipfs/importer/chunk"
13
. "github.com/jbenet/go-ipfs/merkledag"
16
- "github.com/jbenet/go-ipfs/routing/mock"
14
uio "github.com/jbenet/go-ipfs/unixfs/io"
15
u "github.com/jbenet/go-ipfs/util"
16
)
@@ -79,20 +76,8 @@ func makeTestDag(t *testing.T) *Node {
76
}
77
78
func TestBatchFetch(t *testing.T) {
82
- net := tn.VirtualNetwork()
83
- rs := mock.VirtualRoutingServer()
84
- sg := bs.NewSessionGenerator(net, rs)
85
-
86
- instances := sg.Instances(5)
87
-
88
- var servs []*bserv.BlockService
79
var dagservs []DAGService
90
- for _, i := range instances {
91
- bsi, err := bserv.New(i.Blockstore, i.Exchange)
92
- if err != nil {
93
- t.Fatal(err)
94
- }
95
- servs = append(servs, bsi)
80
+ for _, bsi := range blockservice.Mocks(t, 5) {
81
dagservs = append(dagservs, NewDAGService(bsi))
82
}
83
t.Log("finished setup.")
merkledag/mock.go
new
+20
@@ -0,0 +1,20 @@
1
+package merkledag
2
+
3
+import (
4
+ "testing"
5
+
6
+ ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
7
+ dssync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
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
+)
12
+
13
+func Mock(t testing.TB) DAGService {
14
+ bstore := blockstore.NewBlockstore(dssync.MutexWrap(ds.NewMapDatastore()))
15
+ bserv, err := bsrv.New(bstore, offline.Exchange(bstore))
16
+ if err != nil {
17
+ t.Fatal(err)
18
+ }
19
+ return NewDAGService(bserv)
20
+}
routing/mock/routing.go
+2
-5
@@ -30,10 +30,6 @@ func NewMockRouter(local peer.Peer, dstore ds.Datastore) routing.IpfsRouting {
30
}
31
}
32
33
-func (mr *MockRouter) SetRoutingServer(rs RoutingServer) {
34
- mr.hashTable = rs
35
-}
36
-
33
func (mr *MockRouter) PutValue(ctx context.Context, key u.Key, val []byte) error {
34
log.Debugf("PutValue: %s", key)
35
return mr.datastore.Put(key.DsKey(), val)
@@ -119,7 +115,8 @@ func (rs *hashTable) Announce(p peer.Peer, k u.Key) error {
115
func (rs *hashTable) Providers(k u.Key) []peer.Peer {
116
rs.lock.RLock()
117
defer rs.lock.RUnlock()
122
- ret := make([]peer.Peer, 0)
118
+
119
+ var ret []peer.Peer
120
peerset, ok := rs.providers[k]
121
if !ok {
122
return ret
util/datastore2/datastore_closer.go
renamed
+1
-1
@@ -1,4 +1,4 @@
1
-package util
1
+package datastore2
2
3
import (
4
"io"
util/testutil/gen.go
+1
-19
@@ -2,28 +2,10 @@ package testutil
2
3
import (
4
crand "crypto/rand"
5
- "testing"
6
-
7
- dssync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
8
- "github.com/jbenet/go-ipfs/exchange/offline"
9
- "github.com/jbenet/go-ipfs/peer"
10
-
11
- ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
12
- "github.com/jbenet/go-ipfs/blocks/blockstore"
13
- bsrv "github.com/jbenet/go-ipfs/blockservice"
14
- dag "github.com/jbenet/go-ipfs/merkledag"
5
+ peer "github.com/jbenet/go-ipfs/peer"
6
u "github.com/jbenet/go-ipfs/util"
7
)
8
18
-func GetDAGServ(t testing.TB) dag.DAGService {
19
- bstore := blockstore.NewBlockstore(dssync.MutexWrap(ds.NewMapDatastore()))
20
- bserv, err := bsrv.New(bstore, offline.Exchange(bstore))
21
- if err != nil {
22
- t.Fatal(err)
23
- }
24
- return dag.NewDAGService(bserv)
25
-}
26
-
9
func RandPeer() peer.Peer {
10
id := make([]byte, 16)
11
crand.Read(id)
util/testutil/mock.go
+1
-2
@@ -1,9 +1,8 @@
1
package testutil
2
3
import (
4
- "github.com/jbenet/go-ipfs/peer"
5
-
4
ic "github.com/jbenet/go-ipfs/crypto"
5
+ peer "github.com/jbenet/go-ipfs/peer"
6
)
7
8
func NewPeerWithKeyPair(sk ic.PrivKey, pk ic.PubKey) (peer.Peer, error) {