implement a mock dht for use in testing
Jeromy committed
Sep 22, 2014 at 21:11 UTC
c45cc8c448d9045b9ac7a4d0090fe551dbf6b411
4 files changed
+155
-119
exchange/bitswap/bitswap_test.go
+11
-9
@@ -16,6 +16,7 @@ import (
16
strategy "github.com/jbenet/go-ipfs/exchange/bitswap/strategy"
17
tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
18
peer "github.com/jbenet/go-ipfs/peer"
19
+ mock "github.com/jbenet/go-ipfs/routing/mock"
20
util "github.com/jbenet/go-ipfs/util"
21
testutil "github.com/jbenet/go-ipfs/util/testutil"
22
)
@@ -23,7 +24,7 @@ import (
24
func TestGetBlockTimeout(t *testing.T) {
25
26
net := tn.VirtualNetwork()
26
- rs := tn.VirtualRoutingServer()
27
+ rs := mock.VirtualRoutingServer()
28
g := NewSessionGenerator(net, rs)
29
30
self := g.Next()
@@ -40,7 +41,7 @@ func TestGetBlockTimeout(t *testing.T) {
41
func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
42
43
net := tn.VirtualNetwork()
43
- rs := tn.VirtualRoutingServer()
44
+ rs := mock.VirtualRoutingServer()
45
g := NewSessionGenerator(net, rs)
46
47
block := testutil.NewBlockOrFail(t, "block")
@@ -61,7 +62,7 @@ func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
62
func TestGetBlockFromPeerAfterPeerAnnounces(t *testing.T) {
63
64
net := tn.VirtualNetwork()
64
- rs := tn.VirtualRoutingServer()
65
+ rs := mock.VirtualRoutingServer()
66
block := testutil.NewBlockOrFail(t, "block")
67
g := NewSessionGenerator(net, rs)
68
@@ -90,7 +91,7 @@ func TestGetBlockFromPeerAfterPeerAnnounces(t *testing.T) {
91
92
func TestSwarm(t *testing.T) {
93
net := tn.VirtualNetwork()
93
- rs := tn.VirtualRoutingServer()
94
+ rs := mock.VirtualRoutingServer()
95
sg := NewSessionGenerator(net, rs)
96
bg := NewBlockGenerator(t)
97
@@ -151,7 +152,7 @@ func TestSendToWantingPeer(t *testing.T) {
152
util.Debug = true
153
154
net := tn.VirtualNetwork()
154
- rs := tn.VirtualRoutingServer()
155
+ rs := mock.VirtualRoutingServer()
156
sg := NewSessionGenerator(net, rs)
157
bg := NewBlockGenerator(t)
158
@@ -237,7 +238,7 @@ func (bg *BlockGenerator) Blocks(n int) []*blocks.Block {
238
}
239
240
func NewSessionGenerator(
240
- net tn.Network, rs tn.RoutingServer) SessionGenerator {
241
+ net tn.Network, rs mock.RoutingServer) SessionGenerator {
242
return SessionGenerator{
243
net: net,
244
rs: rs,
@@ -248,7 +249,7 @@ func NewSessionGenerator(
249
type SessionGenerator struct {
250
seq int
251
net tn.Network
251
- rs tn.RoutingServer
252
+ rs mock.RoutingServer
253
}
254
255
func (g *SessionGenerator) Next() instance {
@@ -276,11 +277,12 @@ type instance struct {
277
// NB: It's easy make mistakes by providing the same peer ID to two different
278
// sessions. To safeguard, use the SessionGenerator to generate sessions. It's
279
// just a much better idea.
279
-func session(net tn.Network, rs tn.RoutingServer, id peer.ID) instance {
280
+func session(net tn.Network, rs mock.RoutingServer, id peer.ID) instance {
281
p := &peer.Peer{ID: id}
282
283
adapter := net.Adapter(p)
283
- htc := rs.Client(p)
284
+ htc := mock.NewMockRouter(p, nil)
285
+ htc.SetRoutingServer(rs)
286
287
blockstore := bstore.NewBlockstore(ds.NewMapDatastore())
288
const alwaysSendToPeer = true
exchange/bitswap/testnet/routing.go
-96
@@ -1,97 +1 @@
1
package bitswap
2
-
3
-import (
4
- "math/rand"
5
- "sync"
6
-
7
- context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
- bsnet "github.com/jbenet/go-ipfs/exchange/bitswap/network"
9
- peer "github.com/jbenet/go-ipfs/peer"
10
- u "github.com/jbenet/go-ipfs/util"
11
-)
12
-
13
-type RoutingServer interface {
14
- Announce(*peer.Peer, u.Key) error
15
-
16
- Providers(u.Key) []*peer.Peer
17
-
18
- // Returns a Routing instance configured to query this hash table
19
- Client(*peer.Peer) bsnet.Routing
20
-}
21
-
22
-func VirtualRoutingServer() RoutingServer {
23
- return &hashTable{
24
- providers: make(map[u.Key]peer.Map),
25
- }
26
-}
27
-
28
-type hashTable struct {
29
- lock sync.RWMutex
30
- providers map[u.Key]peer.Map
31
-}
32
-
33
-func (rs *hashTable) Announce(p *peer.Peer, k u.Key) error {
34
- rs.lock.Lock()
35
- defer rs.lock.Unlock()
36
-
37
- _, ok := rs.providers[k]
38
- if !ok {
39
- rs.providers[k] = make(peer.Map)
40
- }
41
- rs.providers[k][p.Key()] = p
42
- return nil
43
-}
44
-
45
-func (rs *hashTable) Providers(k u.Key) []*peer.Peer {
46
- rs.lock.RLock()
47
- defer rs.lock.RUnlock()
48
- ret := make([]*peer.Peer, 0)
49
- peerset, ok := rs.providers[k]
50
- if !ok {
51
- return ret
52
- }
53
- for _, peer := range peerset {
54
- ret = append(ret, peer)
55
- }
56
-
57
- for i := range ret {
58
- j := rand.Intn(i + 1)
59
- ret[i], ret[j] = ret[j], ret[i]
60
- }
61
-
62
- return ret
63
-}
64
-
65
-func (rs *hashTable) Client(p *peer.Peer) bsnet.Routing {
66
- return &routingClient{
67
- peer: p,
68
- hashTable: rs,
69
- }
70
-}
71
-
72
-type routingClient struct {
73
- peer *peer.Peer
74
- hashTable RoutingServer
75
-}
76
-
77
-func (a *routingClient) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-chan *peer.Peer {
78
- out := make(chan *peer.Peer)
79
- go func() {
80
- defer close(out)
81
- for i, p := range a.hashTable.Providers(k) {
82
- if max <= i {
83
- return
84
- }
85
- select {
86
- case out <- p:
87
- case <-ctx.Done():
88
- return
89
- }
90
- }
91
- }()
92
- return out
93
-}
94
-
95
-func (a *routingClient) Provide(_ context.Context, key u.Key) error {
96
- return a.hashTable.Announce(a.peer, key)
97
-}
exchange/bitswap/testnet/routing_test.go
+14
-14
@@ -5,19 +5,15 @@ import (
5
"testing"
6
7
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
-)
9
-import (
8
"github.com/jbenet/go-ipfs/peer"
9
+ mock "github.com/jbenet/go-ipfs/routing/mock"
10
u "github.com/jbenet/go-ipfs/util"
11
)
12
13
func TestKeyNotFound(t *testing.T) {
14
16
- rs := func() RoutingServer {
17
- // TODO fields
18
- return &hashTable{}
19
- }()
20
- empty := rs.Providers(u.Key("not there"))
15
+ vrs := mock.VirtualRoutingServer()
16
+ empty := vrs.Providers(u.Key("not there"))
17
if len(empty) != 0 {
18
t.Fatal("should be empty")
19
}
@@ -29,7 +25,7 @@ func TestSetAndGet(t *testing.T) {
25
ID: pid,
26
}
27
k := u.Key("42")
32
- rs := VirtualRoutingServer()
28
+ rs := mock.VirtualRoutingServer()
29
err := rs.Announce(p, k)
30
if err != nil {
31
t.Fatal(err)
@@ -50,8 +46,9 @@ func TestClientFindProviders(t *testing.T) {
46
peer := &peer.Peer{
47
ID: []byte("42"),
48
}
53
- rs := VirtualRoutingServer()
54
- client := rs.Client(peer)
49
+ rs := mock.VirtualRoutingServer()
50
+ client := mock.NewMockRouter(peer, nil)
51
+ client.SetRoutingServer(rs)
52
k := u.Key("hello")
53
err := client.Provide(context.Background(), k)
54
if err != nil {
@@ -83,7 +80,7 @@ func TestClientFindProviders(t *testing.T) {
80
}
81
82
func TestClientOverMax(t *testing.T) {
86
- rs := VirtualRoutingServer()
83
+ rs := mock.VirtualRoutingServer()
84
k := u.Key("hello")
85
numProvidersForHelloKey := 100
86
for i := 0; i < numProvidersForHelloKey; i++ {
@@ -102,7 +99,8 @@ func TestClientOverMax(t *testing.T) {
99
}
100
101
max := 10
105
- client := rs.Client(&peer.Peer{ID: []byte("TODO")})
102
+ client := mock.NewMockRouter(&peer.Peer{ID: []byte("TODO")}, nil)
103
+ client.SetRoutingServer(rs)
104
providersFromClient := client.FindProvidersAsync(context.Background(), k, max)
105
i := 0
106
for _ = range providersFromClient {
@@ -115,7 +113,7 @@ func TestClientOverMax(t *testing.T) {
113
114
// TODO does dht ensure won't receive self as a provider? probably not.
115
func TestCanceledContext(t *testing.T) {
118
- rs := VirtualRoutingServer()
116
+ rs := mock.VirtualRoutingServer()
117
k := u.Key("hello")
118
119
t.Log("async'ly announce infinite stream of providers for key")
@@ -133,7 +131,9 @@ func TestCanceledContext(t *testing.T) {
131
}
132
}()
133
136
- client := rs.Client(&peer.Peer{ID: []byte("peer id doesn't matter")})
134
+ local := &peer.Peer{ID: []byte("peer id doesn't matter")}
135
+ client := mock.NewMockRouter(local, nil)
136
+ client.SetRoutingServer(rs)
137
138
t.Log("warning: max is finite so this test is non-deterministic")
139
t.Log("context cancellation could simply take lower priority")
routing/mock/routing.go
new
+130
@@ -0,0 +1,130 @@
1
+package mockrouter
2
+
3
+import (
4
+ "errors"
5
+ "math/rand"
6
+ "sync"
7
+
8
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9
+ ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
10
+ peer "github.com/jbenet/go-ipfs/peer"
11
+ routing "github.com/jbenet/go-ipfs/routing"
12
+ u "github.com/jbenet/go-ipfs/util"
13
+)
14
+
15
+var _ routing.IpfsRouting = &MockRouter{}
16
+
17
+type MockRouter struct {
18
+ datastore ds.Datastore
19
+ hashTable RoutingServer
20
+ peer *peer.Peer
21
+}
22
+
23
+func NewMockRouter(local *peer.Peer, dstore ds.Datastore) *MockRouter {
24
+ return &MockRouter{
25
+ datastore: dstore,
26
+ peer: local,
27
+ hashTable: VirtualRoutingServer(),
28
+ }
29
+}
30
+
31
+func (mr *MockRouter) SetRoutingServer(rs RoutingServer) {
32
+ mr.hashTable = rs
33
+}
34
+
35
+func (mr *MockRouter) PutValue(ctx context.Context, key u.Key, val []byte) error {
36
+ return mr.datastore.Put(ds.NewKey(string(key)), val)
37
+}
38
+
39
+func (mr *MockRouter) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
40
+ v, err := mr.datastore.Get(ds.NewKey(string(key)))
41
+ if err != nil {
42
+ return nil, err
43
+ }
44
+
45
+ data, ok := v.([]byte)
46
+ if !ok {
47
+ return nil, errors.New("could not cast value from datastore")
48
+ }
49
+
50
+ return data, nil
51
+}
52
+
53
+func (mr *MockRouter) FindProviders(ctx context.Context, key u.Key) ([]*peer.Peer, error) {
54
+ return nil, nil
55
+}
56
+
57
+func (mr *MockRouter) FindPeer(ctx context.Context, pid peer.ID) (*peer.Peer, error) {
58
+ return nil, nil
59
+}
60
+
61
+func (mr *MockRouter) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-chan *peer.Peer {
62
+ out := make(chan *peer.Peer)
63
+ go func() {
64
+ defer close(out)
65
+ for i, p := range mr.hashTable.Providers(k) {
66
+ if max <= i {
67
+ return
68
+ }
69
+ select {
70
+ case out <- p:
71
+ case <-ctx.Done():
72
+ return
73
+ }
74
+ }
75
+ }()
76
+ return out
77
+}
78
+
79
+func (mr *MockRouter) Provide(_ context.Context, key u.Key) error {
80
+ return mr.hashTable.Announce(mr.peer, key)
81
+}
82
+
83
+type RoutingServer interface {
84
+ Announce(*peer.Peer, u.Key) error
85
+
86
+ Providers(u.Key) []*peer.Peer
87
+}
88
+
89
+func VirtualRoutingServer() RoutingServer {
90
+ return &hashTable{
91
+ providers: make(map[u.Key]peer.Map),
92
+ }
93
+}
94
+
95
+type hashTable struct {
96
+ lock sync.RWMutex
97
+ providers map[u.Key]peer.Map
98
+}
99
+
100
+func (rs *hashTable) Announce(p *peer.Peer, k u.Key) error {
101
+ rs.lock.Lock()
102
+ defer rs.lock.Unlock()
103
+
104
+ _, ok := rs.providers[k]
105
+ if !ok {
106
+ rs.providers[k] = make(peer.Map)
107
+ }
108
+ rs.providers[k][p.Key()] = p
109
+ return nil
110
+}
111
+
112
+func (rs *hashTable) Providers(k u.Key) []*peer.Peer {
113
+ rs.lock.RLock()
114
+ defer rs.lock.RUnlock()
115
+ ret := make([]*peer.Peer, 0)
116
+ peerset, ok := rs.providers[k]
117
+ if !ok {
118
+ return ret
119
+ }
120
+ for _, peer := range peerset {
121
+ ret = append(ret, peer)
122
+ }
123
+
124
+ for i := range ret {
125
+ j := rand.Intn(i + 1)
126
+ ret[i], ret[j] = ret[j], ret[i]
127
+ }
128
+
129
+ return ret
130
+}