@cryptotaxi247 / kubo / commits / 3ecdec985

refactor(mockrouting) misc

License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>

Brian Tiger Chow committed Dec 12, 2014 at 22:56 UTC 3ecdec985febc4985ef551d80893af87bc7e8ef7
10 files changed +228 -197
blockservice/mock.go
+2 -2
@@ -5,14 +5,14 @@ import (
5
6 bitswap "github.com/jbenet/go-ipfs/exchange/bitswap"
7 tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
8 - mock "github.com/jbenet/go-ipfs/routing/mock"
8 + mockrouting "github.com/jbenet/go-ipfs/routing/mock"
9 delay "github.com/jbenet/go-ipfs/util/delay"
10 )
11
12 // Mocks returns |n| connected mock Blockservices
13 func Mocks(t *testing.T, n int) []*BlockService {
14 net := tn.VirtualNetwork(delay.Fixed(0))
15 - rs := mock.VirtualRoutingServer()
15 + rs := mockrouting.NewServer()
16 sg := bitswap.NewSessionGenerator(net, rs)
17
18 instances := sg.Instances(n)
core/mock.go
+1 -1
@@ -42,7 +42,7 @@ func NewMockNode() (*IpfsNode, error) {
42 nd.Datastore = ds2.CloserWrap(syncds.MutexWrap(dstore))
43
44 // Routing
45 - dht := mdht.NewMockRouter(nd.Identity, nd.Datastore)
45 + dht := mdht.NewServer().ClientWithDatastore(nd.Identity, nd.Datastore)
46 nd.Routing = dht
47
48 // Bitswap
exchange/bitswap/bitswap_test.go
+11 -9
@@ -10,18 +10,20 @@ import (
10 blocks "github.com/jbenet/go-ipfs/blocks"
11 blocksutil "github.com/jbenet/go-ipfs/blocks/blocksutil"
12 tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
13 - mock "github.com/jbenet/go-ipfs/routing/mock"
13 + mockrouting "github.com/jbenet/go-ipfs/routing/mock"
14 delay "github.com/jbenet/go-ipfs/util/delay"
15 testutil "github.com/jbenet/go-ipfs/util/testutil"
16 )
17
18 +// FIXME the tests are really sensitive to the network delay. fix them to work
19 +// well under varying conditions
20 const kNetworkDelay = 0 * time.Millisecond
21
22 func TestClose(t *testing.T) {
23 // TODO
24 t.Skip("TODO Bitswap's Close implementation is a WIP")
25 vnet := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
24 - rout := mock.VirtualRoutingServer()
26 + rout := mockrouting.NewServer()
27 sesgen := NewSessionGenerator(vnet, rout)
28 bgen := blocksutil.NewBlockGenerator()
29
@@ -35,7 +37,7 @@ func TestClose(t *testing.T) {
37 func TestGetBlockTimeout(t *testing.T) {
38
39 net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
38 - rs := mock.VirtualRoutingServer()
40 + rs := mockrouting.NewServer()
41 g := NewSessionGenerator(net, rs)
42
43 self := g.Next()
@@ -52,11 +54,11 @@ func TestGetBlockTimeout(t *testing.T) {
54 func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
55
56 net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
55 - rs := mock.VirtualRoutingServer()
57 + rs := mockrouting.NewServer()
58 g := NewSessionGenerator(net, rs)
59
60 block := blocks.NewBlock([]byte("block"))
59 - rs.Announce(testutil.NewPeerWithIDString("testing"), block.Key()) // but not on network
61 + rs.Client(testutil.NewPeerWithIDString("testing")).Provide(context.Background(), block.Key()) // but not on network
62
63 solo := g.Next()
64
@@ -73,7 +75,7 @@ func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
75 func TestGetBlockFromPeerAfterPeerAnnounces(t *testing.T) {
76
77 net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
76 - rs := mock.VirtualRoutingServer()
78 + rs := mockrouting.NewServer()
79 block := blocks.NewBlock([]byte("block"))
80 g := NewSessionGenerator(net, rs)
81
@@ -125,7 +127,7 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
127 t.SkipNow()
128 }
129 net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
128 - rs := mock.VirtualRoutingServer()
130 + rs := mockrouting.NewServer()
131 sg := NewSessionGenerator(net, rs)
132 bg := blocksutil.NewBlockGenerator()
133
@@ -140,7 +142,7 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
142 for _, b := range blocks {
143 first.Blockstore().Put(b)
144 first.Exchange.HasBlock(context.Background(), b)
143 - rs.Announce(first.Peer, b.Key())
145 + rs.Client(first.Peer).Provide(context.Background(), b.Key())
146 }
147
148 t.Log("Distribute!")
@@ -185,7 +187,7 @@ func TestSendToWantingPeer(t *testing.T) {
187 }
188
189 net := tn.VirtualNetwork(delay.Fixed(kNetworkDelay))
188 - rs := mock.VirtualRoutingServer()
190 + rs := mockrouting.NewServer()
191 sg := NewSessionGenerator(net, rs)
192 bg := blocksutil.NewBlockGenerator()
193
exchange/bitswap/testutils.go
+4 -4
@@ -10,13 +10,13 @@ import (
10 exchange "github.com/jbenet/go-ipfs/exchange"
11 tn "github.com/jbenet/go-ipfs/exchange/bitswap/testnet"
12 peer "github.com/jbenet/go-ipfs/peer"
13 - mock "github.com/jbenet/go-ipfs/routing/mock"
13 + mockrouting "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(
19 - net tn.Network, rs mock.RoutingServer) SessionGenerator {
19 + net tn.Network, rs mockrouting.Server) SessionGenerator {
20 return SessionGenerator{
21 net: net,
22 rs: rs,
@@ -28,7 +28,7 @@ func NewSessionGenerator(
28 type SessionGenerator struct {
29 seq int
30 net tn.Network
31 - rs mock.RoutingServer
31 + rs mockrouting.Server
32 ps peer.Peerstore
33 }
34
@@ -67,7 +67,7 @@ func (i *Instance) SetBlockstoreLatency(t time.Duration) time.Duration {
67 // NB: It's easy make mistakes by providing the same peer ID to two different
68 // sessions. To safeguard, use the SessionGenerator to generate sessions. It's
69 // just a much better idea.
70 -func session(net tn.Network, rs mock.RoutingServer, ps peer.Peerstore, id peer.ID) Instance {
70 +func session(net tn.Network, rs mockrouting.Server, ps peer.Peerstore, id peer.ID) Instance {
71 p := ps.WithID(id)
72
73 adapter := net.Adapter(p)
namesys/resolve_test.go
+2 -4
@@ -3,17 +3,15 @@ package namesys
3 import (
4 "testing"
5
6 - ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
6 ci "github.com/jbenet/go-ipfs/crypto"
8 - mock "github.com/jbenet/go-ipfs/routing/mock"
7 + mockrouting "github.com/jbenet/go-ipfs/routing/mock"
8 u "github.com/jbenet/go-ipfs/util"
9 testutil "github.com/jbenet/go-ipfs/util/testutil"
10 )
11
12 func TestRoutingResolve(t *testing.T) {
13 local := testutil.NewPeerWithIDString("testID")
15 - lds := ds.NewMapDatastore()
16 - d := mock.NewMockRouter(local, lds)
14 + d := mockrouting.NewServer().Client(local)
15
16 resolver := NewRoutingResolver(d)
17 publisher := NewRoutingPublisher(d)
routing/mock/client.go new
+74
@@ -0,0 +1,74 @@
1 +package mockrouting
2 +
3 +import (
4 + "errors"
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 + peer "github.com/jbenet/go-ipfs/peer"
9 + routing "github.com/jbenet/go-ipfs/routing"
10 + u "github.com/jbenet/go-ipfs/util"
11 +)
12 +
13 +var log = u.Logger("mockrouter")
14 +
15 +type client struct {
16 + datastore ds.Datastore
17 + server server
18 + peer peer.Peer
19 +}
20 +
21 +// FIXME(brian): is this method meant to simulate putting a value into the network?
22 +func (c *client) PutValue(ctx context.Context, key u.Key, val []byte) error {
23 + log.Debugf("PutValue: %s", key)
24 + return c.datastore.Put(key.DsKey(), val)
25 +}
26 +
27 +// FIXME(brian): is this method meant to simulate getting a value from the network?
28 +func (c *client) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
29 + log.Debugf("GetValue: %s", key)
30 + v, err := c.datastore.Get(key.DsKey())
31 + if err != nil {
32 + return nil, err
33 + }
34 +
35 + data, ok := v.([]byte)
36 + if !ok {
37 + return nil, errors.New("could not cast value from datastore")
38 + }
39 +
40 + return data, nil
41 +}
42 +
43 +func (c *client) FindProviders(ctx context.Context, key u.Key) ([]peer.Peer, error) {
44 + return c.server.Providers(key), nil
45 +}
46 +
47 +func (c *client) FindPeer(ctx context.Context, pid peer.ID) (peer.Peer, error) {
48 + log.Debugf("FindPeer: %s", pid)
49 + return nil, nil
50 +}
51 +
52 +func (c *client) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-chan peer.Peer {
53 + out := make(chan peer.Peer)
54 + go func() {
55 + defer close(out)
56 + for i, p := range c.server.Providers(k) {
57 + if max <= i {
58 + return
59 + }
60 + select {
61 + case out <- p:
62 + case <-ctx.Done():
63 + return
64 + }
65 + }
66 + }()
67 + return out
68 +}
69 +
70 +func (c *client) Provide(_ context.Context, key u.Key) error {
71 + return c.server.Announce(c.peer, key)
72 +}
73 +
74 +var _ routing.IpfsRouting = &client{}
routing/mock/interface.go new
+40
@@ -0,0 +1,40 @@
1 +// Package mock provides a virtual routing server. To use it, create a virtual
2 +// routing server and use the Client() method to get a routing client
3 +// (IpfsRouting). The server quacks like a DHT but is really a local in-memory
4 +// hash table.
5 +package mockrouting
6 +
7 +import (
8 + context "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/go-datastore"
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 + delay "github.com/jbenet/go-ipfs/util/delay"
14 +)
15 +
16 +// Server provides mockrouting Clients
17 +type Server interface {
18 + Client(p peer.Peer) Client
19 + ClientWithDatastore(peer.Peer, ds.Datastore) Client
20 +}
21 +
22 +// Client implements IpfsRouting
23 +type Client interface {
24 + FindProviders(context.Context, u.Key) ([]peer.Peer, error)
25 +
26 + routing.IpfsRouting
27 +}
28 +
29 +// NewServer returns a mockrouting Server
30 +func NewServer() Server {
31 + return NewServerWithDelay(delay.Fixed(0))
32 +}
33 +
34 +// NewServerWithDelay returns a mockrouting Server with a delay!
35 +func NewServerWithDelay(d delay.D) Server {
36 + return &s{
37 + providers: make(map[u.Key]peer.Map),
38 + delay: d,
39 + }
40 +}
routing/mock/mockrouting_test.go renamed
+18 -36
@@ -1,4 +1,4 @@
1 -package mock
1 +package mockrouting
2
3 import (
4 "bytes"
@@ -12,37 +12,21 @@ import (
12
13 func TestKeyNotFound(t *testing.T) {
14
15 - vrs := VirtualRoutingServer()
16 - empty := vrs.Providers(u.Key("not there"))
17 - if len(empty) != 0 {
18 - t.Fatal("should be empty")
19 - }
20 -}
15 + var peer = testutil.NewPeerWithID(peer.ID([]byte("the peer id")))
16 + var key = u.Key("mock key")
17 + var ctx = context.Background()
18
22 -func TestSetAndGet(t *testing.T) {
23 - pid := peer.ID([]byte("the peer id"))
24 - p := testutil.NewPeerWithID(pid)
25 - k := u.Key("42")
26 - rs := VirtualRoutingServer()
27 - err := rs.Announce(p, k)
28 - if err != nil {
29 - t.Fatal(err)
30 - }
31 - providers := rs.Providers(k)
32 - if len(providers) != 1 {
33 - t.Fatal("should be one")
19 + rs := NewServer()
20 + providers := rs.Client(peer).FindProvidersAsync(ctx, key, 10)
21 + _, ok := <-providers
22 + if ok {
23 + t.Fatal("should be closed")
24 }
35 - for _, elem := range providers {
36 - if bytes.Equal(elem.ID(), pid) {
37 - return
38 - }
39 - }
40 - t.Fatal("ID should have matched")
25 }
26
27 func TestClientFindProviders(t *testing.T) {
28 peer := testutil.NewPeerWithIDString("42")
45 - rs := VirtualRoutingServer()
29 + rs := NewServer()
30 client := rs.Client(peer)
31
32 k := u.Key("hello")
@@ -52,7 +36,10 @@ func TestClientFindProviders(t *testing.T) {
36 }
37 max := 100
38
55 - providersFromHashTable := rs.Providers(k)
39 + providersFromHashTable, err := rs.Client(peer).FindProviders(context.Background(), k)
40 + if err != nil {
41 + t.Fatal(err)
42 + }
43
44 isInHT := false
45 for _, p := range providersFromHashTable {
@@ -76,21 +63,16 @@ func TestClientFindProviders(t *testing.T) {
63 }
64
65 func TestClientOverMax(t *testing.T) {
79 - rs := VirtualRoutingServer()
66 + rs := NewServer()
67 k := u.Key("hello")
68 numProvidersForHelloKey := 100
69 for i := 0; i < numProvidersForHelloKey; i++ {
70 peer := testutil.NewPeerWithIDString(string(i))
84 - err := rs.Announce(peer, k)
71 + err := rs.Client(peer).Provide(context.Background(), k)
72 if err != nil {
73 t.Fatal(err)
74 }
75 }
89 - providersFromHashTable := rs.Providers(k)
90 - if len(providersFromHashTable) != numProvidersForHelloKey {
91 - t.Log(1 == len(providersFromHashTable))
92 - t.Fatal("not all providers were returned")
93 - }
76
77 max := 10
78 peer := testutil.NewPeerWithIDString("TODO")
@@ -108,7 +90,7 @@ func TestClientOverMax(t *testing.T) {
90
91 // TODO does dht ensure won't receive self as a provider? probably not.
92 func TestCanceledContext(t *testing.T) {
111 - rs := VirtualRoutingServer()
93 + rs := NewServer()
94 k := u.Key("hello")
95
96 t.Log("async'ly announce infinite stream of providers for key")
@@ -116,7 +98,7 @@ func TestCanceledContext(t *testing.T) {
98 go func() { // infinite stream
99 for {
100 peer := testutil.NewPeerWithIDString(string(i))
119 - err := rs.Announce(peer, k)
101 + err := rs.Client(peer).Provide(context.Background(), k)
102 if err != nil {
103 t.Fatal(err)
104 }
routing/mock/routing.go deleted
-141
@@ -1,141 +0,0 @@
1 -package mock
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/go-datastore"
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 log = u.Logger("mockrouter")
16 -
17 -var _ routing.IpfsRouting = &MockRouter{}
18 -
19 -type MockRouter struct {
20 - datastore ds.Datastore
21 - hashTable RoutingServer
22 - peer peer.Peer
23 -}
24 -
25 -func NewMockRouter(local peer.Peer, dstore ds.Datastore) routing.IpfsRouting {
26 - return &MockRouter{
27 - datastore: dstore,
28 - peer: local,
29 - hashTable: VirtualRoutingServer(),
30 - }
31 -}
32 -
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)
36 -}
37 -
38 -func (mr *MockRouter) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
39 - log.Debugf("GetValue: %s", key)
40 - v, err := mr.datastore.Get(key.DsKey())
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 - log.Debugf("FindPeer: %s", pid)
59 - return nil, nil
60 -}
61 -
62 -func (mr *MockRouter) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-chan peer.Peer {
63 - out := make(chan peer.Peer)
64 - go func() {
65 - defer close(out)
66 - for i, p := range mr.hashTable.Providers(k) {
67 - if max <= i {
68 - return
69 - }
70 - select {
71 - case out <- p:
72 - case <-ctx.Done():
73 - return
74 - }
75 - }
76 - }()
77 - return out
78 -}
79 -
80 -func (mr *MockRouter) Provide(_ context.Context, key u.Key) error {
81 - return mr.hashTable.Announce(mr.peer, key)
82 -}
83 -
84 -type RoutingServer interface {
85 - Announce(peer.Peer, u.Key) error
86 -
87 - Providers(u.Key) []peer.Peer
88 -
89 - Client(p peer.Peer) routing.IpfsRouting
90 -}
91 -
92 -func VirtualRoutingServer() RoutingServer {
93 - return &hashTable{
94 - providers: make(map[u.Key]peer.Map),
95 - }
96 -}
97 -
98 -type hashTable struct {
99 - lock sync.RWMutex
100 - providers map[u.Key]peer.Map
101 -}
102 -
103 -func (rs *hashTable) Announce(p peer.Peer, k u.Key) error {
104 - rs.lock.Lock()
105 - defer rs.lock.Unlock()
106 -
107 - _, ok := rs.providers[k]
108 - if !ok {
109 - rs.providers[k] = make(peer.Map)
110 - }
111 - rs.providers[k][p.Key()] = p
112 - return nil
113 -}
114 -
115 -func (rs *hashTable) Providers(k u.Key) []peer.Peer {
116 - rs.lock.RLock()
117 - defer rs.lock.RUnlock()
118 -
119 - var ret []peer.Peer
120 - peerset, ok := rs.providers[k]
121 - if !ok {
122 - return ret
123 - }
124 - for _, peer := range peerset {
125 - ret = append(ret, peer)
126 - }
127 -
128 - for i := range ret {
129 - j := rand.Intn(i + 1)
130 - ret[i], ret[j] = ret[j], ret[i]
131 - }
132 -
133 - return ret
134 -}
135 -
136 -func (rs *hashTable) Client(p peer.Peer) routing.IpfsRouting {
137 - return &MockRouter{
138 - peer: p,
139 - hashTable: rs,
140 - }
141 -}
routing/mock/server.go new
+76
@@ -0,0 +1,76 @@
1 +package mockrouting
2 +
3 +import (
4 + "math/rand"
5 + "sync"
6 +
7 + ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
8 + peer "github.com/jbenet/go-ipfs/peer"
9 + u "github.com/jbenet/go-ipfs/util"
10 + delay "github.com/jbenet/go-ipfs/util/delay"
11 +)
12 +
13 +// server is the mockrouting.Client's private interface to the routing server
14 +type server interface {
15 + Announce(peer.Peer, u.Key) error
16 + Providers(u.Key) []peer.Peer
17 +
18 + Server
19 +}
20 +
21 +// s is an implementation of the private server interface
22 +type s struct {
23 + delay delay.D
24 +
25 + lock sync.RWMutex
26 + providers map[u.Key]peer.Map
27 +}
28 +
29 +func (rs *s) Announce(p peer.Peer, k u.Key) error {
30 + rs.delay.Wait() // before locking
31 +
32 + rs.lock.Lock()
33 + defer rs.lock.Unlock()
34 +
35 + _, ok := rs.providers[k]
36 + if !ok {
37 + rs.providers[k] = make(peer.Map)
38 + }
39 + rs.providers[k][p.Key()] = p
40 + return nil
41 +}
42 +
43 +func (rs *s) Providers(k u.Key) []peer.Peer {
44 + rs.delay.Wait() // before locking
45 +
46 + rs.lock.RLock()
47 + defer rs.lock.RUnlock()
48 +
49 + var ret []peer.Peer
50 + peerset, ok := rs.providers[k]
51 + if !ok {
52 + return ret
53 + }
54 + for _, peer := range peerset {
55 + ret = append(ret, peer)
56 + }
57 +
58 + for i := range ret {
59 + j := rand.Intn(i + 1)
60 + ret[i], ret[j] = ret[j], ret[i]
61 + }
62 +
63 + return ret
64 +}
65 +
66 +func (rs *s) Client(p peer.Peer) Client {
67 + return rs.ClientWithDatastore(p, ds.NewMapDatastore())
68 +}
69 +
70 +func (rs *s) ClientWithDatastore(p peer.Peer, datastore ds.Datastore) Client {
71 + return &client{
72 + peer: p,
73 + datastore: ds.NewMapDatastore(),
74 + server: rs,
75 + }
76 +}