test(bitswap:testnet)
misc: * test network client getting more than max * test for find providers * rename factory method * local network * misc test improvements * test bitswap get block timeout * test provider exists but cannot connect to peer * test sending a message async over local network
Brian Tiger Chow committed
Sep 19, 2014 at 08:11 UTC
c80c8aa9773d4d3279540606b861ad0a5f6f6f5f
6 files changed
+653
exchange/bitswap/bitswap_test.go
new
+81
@@ -0,0 +1,81 @@
1
+package bitswap
2
+
3
+import (
4
+ "testing"
5
+ "time"
6
+
7
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
+
9
+ ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
10
+ bstore "github.com/jbenet/go-ipfs/blockstore"
11
+ exchange "github.com/jbenet/go-ipfs/exchange"
12
+ notifications "github.com/jbenet/go-ipfs/exchange/bitswap/notifications"
13
+ strategy "github.com/jbenet/go-ipfs/exchange/bitswap/strategy"
14
+ peer "github.com/jbenet/go-ipfs/peer"
15
+ testutil "github.com/jbenet/go-ipfs/util/testutil"
16
+)
17
+
18
+func TestGetBlockTimeout(t *testing.T) {
19
+
20
+ net := LocalNetwork()
21
+ rs := newRoutingServer()
22
+ ipfs := session(net, rs, []byte("peer id"))
23
+ ctx, _ := context.WithTimeout(context.Background(), time.Nanosecond)
24
+ block := testutil.NewBlockOrFail(t, "block")
25
+
26
+ _, err := ipfs.exchange.Block(ctx, block.Key())
27
+ if err != context.DeadlineExceeded {
28
+ t.Fatal("Expected DeadlineExceeded error")
29
+ }
30
+}
31
+
32
+func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
33
+
34
+ net := LocalNetwork()
35
+ rs := newRoutingServer()
36
+ ipfs := session(net, rs, []byte("peer id"))
37
+ // ctx := context.Background()
38
+ ctx, _ := context.WithTimeout(context.Background(), time.Nanosecond)
39
+ block := testutil.NewBlockOrFail(t, "block")
40
+
41
+ rs.Announce(&peer.Peer{}, block.Key()) // but not on network
42
+
43
+ _, err := ipfs.exchange.Block(ctx, block.Key())
44
+ if err != context.DeadlineExceeded {
45
+ t.Fatal("Expected DeadlineExceeded error")
46
+ }
47
+}
48
+
49
+type ipfs struct {
50
+ peer *peer.Peer
51
+ exchange exchange.Interface
52
+ blockstore bstore.Blockstore
53
+}
54
+
55
+func session(net Network, rs RoutingServer, id peer.ID) ipfs {
56
+ p := &peer.Peer{}
57
+
58
+ adapter := net.Adapter(p)
59
+ htc := rs.Client(p)
60
+
61
+ blockstore := bstore.NewBlockstore(ds.NewMapDatastore())
62
+ bs := &bitswap{
63
+ blockstore: blockstore,
64
+ notifications: notifications.New(),
65
+ strategy: strategy.New(),
66
+ routing: htc,
67
+ sender: adapter,
68
+ }
69
+ adapter.SetDelegate(bs)
70
+ return ipfs{
71
+ peer: p,
72
+ exchange: bs,
73
+ blockstore: blockstore,
74
+ }
75
+}
76
+
77
+func TestSendToWantingPeer(t *testing.T) {
78
+ t.Log("Peer |w| tells me it wants file, but I don't have it")
79
+ t.Log("Then another peer |o| sends it to me")
80
+ t.Log("After receiving the file from |o|, I send it to the wanting peer |w|")
81
+}
exchange/bitswap/hash_table.go
new
+96
@@ -0,0 +1,96 @@
1
+package bitswap
2
+
3
+import (
4
+ "errors"
5
+ "sync"
6
+
7
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
+ peer "github.com/jbenet/go-ipfs/peer"
9
+ u "github.com/jbenet/go-ipfs/util"
10
+)
11
+
12
+type RoutingServer interface {
13
+ // TODO
14
+ Announce(*peer.Peer, u.Key) error
15
+
16
+ // TODO
17
+ Providers(u.Key) []*peer.Peer
18
+
19
+ // TODO
20
+ // Returns a Routing instance configured to query this hash table
21
+ Client(*peer.Peer) Routing
22
+}
23
+
24
+func newRoutingServer() RoutingServer {
25
+ return &hashTable{
26
+ m: make(map[u.Key]map[*peer.Peer]bool),
27
+ }
28
+}
29
+
30
+type hashTable struct {
31
+ lock sync.RWMutex
32
+ m map[u.Key]map[*peer.Peer]bool
33
+}
34
+
35
+var TODO = errors.New("TODO")
36
+
37
+func (rs *hashTable) Announce(p *peer.Peer, k u.Key) error {
38
+ rs.lock.Lock()
39
+ defer rs.lock.Unlock()
40
+
41
+ _, ok := rs.m[k]
42
+ if !ok {
43
+ rs.m[k] = make(map[*peer.Peer]bool)
44
+ }
45
+ rs.m[k][p] = true
46
+ return nil
47
+}
48
+
49
+func (rs *hashTable) Providers(k u.Key) []*peer.Peer {
50
+ rs.lock.RLock()
51
+ defer rs.lock.RUnlock()
52
+ ret := make([]*peer.Peer, 0)
53
+ peerset, ok := rs.m[k]
54
+ if !ok {
55
+ return ret
56
+ }
57
+ for peer, _ := range peerset {
58
+ ret = append(ret, peer)
59
+ }
60
+ return ret
61
+}
62
+
63
+// TODO
64
+func (rs *hashTable) Client(p *peer.Peer) Routing {
65
+ return &routingClient{
66
+ peer: p,
67
+ hashTable: rs,
68
+ }
69
+}
70
+
71
+type routingClient struct {
72
+ peer *peer.Peer
73
+ hashTable RoutingServer
74
+}
75
+
76
+func (a *routingClient) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-chan *peer.Peer {
77
+ out := make(chan *peer.Peer)
78
+ go func() {
79
+ defer close(out)
80
+ for i, p := range a.hashTable.Providers(k) {
81
+ if max <= i {
82
+ return
83
+ }
84
+ select {
85
+ case out <- p:
86
+ case <-ctx.Done():
87
+ return
88
+ }
89
+ }
90
+ }()
91
+ return out
92
+}
93
+
94
+func (a *routingClient) Provide(key u.Key) error {
95
+ return a.hashTable.Announce(a.peer, key)
96
+}
exchange/bitswap/hash_table_test.go
new
+157
@@ -0,0 +1,157 @@
1
+package bitswap
2
+
3
+import (
4
+ "bytes"
5
+ "testing"
6
+
7
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
+)
9
+import (
10
+ "github.com/jbenet/go-ipfs/peer"
11
+ u "github.com/jbenet/go-ipfs/util"
12
+)
13
+
14
+func TestKeyNotFound(t *testing.T) {
15
+
16
+ rs := func() RoutingServer {
17
+ // TODO fields
18
+ return &hashTable{}
19
+ }()
20
+ empty := rs.Providers(u.Key("not there"))
21
+ if len(empty) != 0 {
22
+ t.Fatal("should be empty")
23
+ }
24
+}
25
+
26
+func TestSetAndGet(t *testing.T) {
27
+ pid := peer.ID([]byte("the peer id"))
28
+ p := &peer.Peer{
29
+ ID: pid,
30
+ }
31
+ k := u.Key("42")
32
+ rs := newRoutingServer()
33
+ err := rs.Announce(p, k)
34
+ if err != nil {
35
+ t.Fatal(err)
36
+ }
37
+ providers := rs.Providers(k)
38
+ if len(providers) != 1 {
39
+ t.Fatal("should be one")
40
+ }
41
+ for _, elem := range providers {
42
+ if bytes.Equal(elem.ID, pid) {
43
+ return
44
+ }
45
+ }
46
+ t.Fatal("ID should have matched")
47
+}
48
+
49
+func TestClientFindProviders(t *testing.T) {
50
+ peer := &peer.Peer{
51
+ ID: []byte("42"),
52
+ }
53
+ rs := newRoutingServer()
54
+ client := rs.Client(peer)
55
+ k := u.Key("hello")
56
+ err := client.Provide(k)
57
+ if err != nil {
58
+ t.Fatal(err)
59
+ }
60
+ max := 100
61
+
62
+ providersFromHashTable := rs.Providers(k)
63
+
64
+ isInHT := false
65
+ for _, p := range providersFromHashTable {
66
+ if bytes.Equal(p.ID, peer.ID) {
67
+ isInHT = true
68
+ }
69
+ }
70
+ if !isInHT {
71
+ t.Fatal("Despite client providing key, peer wasn't in hash table as a provider")
72
+ }
73
+ providersFromClient := client.FindProvidersAsync(context.Background(), u.Key("hello"), max)
74
+ isInClient := false
75
+ for p := range providersFromClient {
76
+ if bytes.Equal(p.ID, peer.ID) {
77
+ isInClient = true
78
+ }
79
+ }
80
+ if !isInClient {
81
+ t.Fatal("Despite client providing key, client didn't receive peer when finding providers")
82
+ }
83
+}
84
+
85
+func TestClientOverMax(t *testing.T) {
86
+ rs := newRoutingServer()
87
+ k := u.Key("hello")
88
+ numProvidersForHelloKey := 100
89
+ for i := 0; i < numProvidersForHelloKey; i++ {
90
+ peer := &peer.Peer{
91
+ ID: []byte(string(i)),
92
+ }
93
+ err := rs.Announce(peer, k)
94
+ if err != nil {
95
+ t.Fatal(err)
96
+ }
97
+ }
98
+ providersFromHashTable := rs.Providers(k)
99
+ if len(providersFromHashTable) != numProvidersForHelloKey {
100
+ t.Log(1 == len(providersFromHashTable))
101
+ t.Fatal("not all providers were returned")
102
+ }
103
+
104
+ max := 10
105
+ client := rs.Client(&peer.Peer{ID: []byte("TODO")})
106
+ providersFromClient := client.FindProvidersAsync(context.Background(), k, max)
107
+ i := 0
108
+ for _ = range providersFromClient {
109
+ i++
110
+ }
111
+ if i != max {
112
+ t.Fatal("Too many providers returned")
113
+ }
114
+}
115
+
116
+// TODO does dht ensure won't receive self as a provider? probably not.
117
+func TestCanceledContext(t *testing.T) {
118
+ rs := newRoutingServer()
119
+ k := u.Key("hello")
120
+
121
+ t.Log("async'ly announce infinite stream of providers for key")
122
+ i := 0
123
+ go func() { // infinite stream
124
+ for {
125
+ peer := &peer.Peer{
126
+ ID: []byte(string(i)),
127
+ }
128
+ err := rs.Announce(peer, k)
129
+ if err != nil {
130
+ t.Fatal(err)
131
+ }
132
+ i++
133
+ }
134
+ }()
135
+
136
+ client := rs.Client(&peer.Peer{ID: []byte("peer id doesn't matter")})
137
+
138
+ t.Log("warning: max is finite so this test is non-deterministic")
139
+ t.Log("context cancellation could simply take lower priority")
140
+ t.Log("and result in receiving the max number of results")
141
+ max := 1000
142
+
143
+ t.Log("cancel the context before consuming")
144
+ ctx, cancelFunc := context.WithCancel(context.Background())
145
+ cancelFunc()
146
+ providers := client.FindProvidersAsync(ctx, k, max)
147
+
148
+ numProvidersReturned := 0
149
+ for _ = range providers {
150
+ numProvidersReturned++
151
+ }
152
+ t.Log(numProvidersReturned)
153
+
154
+ if numProvidersReturned == max {
155
+ t.Fatal("Context cancel had no effect")
156
+ }
157
+}
exchange/bitswap/local_network.go
new
+174
@@ -0,0 +1,174 @@
1
+package bitswap
2
+
3
+import (
4
+ "bytes"
5
+ "errors"
6
+
7
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
+ bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
9
+ bsnet "github.com/jbenet/go-ipfs/exchange/bitswap/network"
10
+ peer "github.com/jbenet/go-ipfs/peer"
11
+ "github.com/jbenet/go-ipfs/util"
12
+)
13
+
14
+type Network interface {
15
+ Adapter(*peer.Peer) bsnet.Adapter
16
+
17
+ SendMessage(
18
+ ctx context.Context,
19
+ from *peer.Peer,
20
+ to *peer.Peer,
21
+ message bsmsg.BitSwapMessage) error
22
+
23
+ SendRequest(
24
+ ctx context.Context,
25
+ from *peer.Peer,
26
+ to *peer.Peer,
27
+ message bsmsg.BitSwapMessage) (
28
+ incoming bsmsg.BitSwapMessage, err error)
29
+}
30
+
31
+// network impl
32
+
33
+func LocalNetwork() Network {
34
+ return &network{
35
+ clients: make(map[util.Key]bsnet.Receiver),
36
+ }
37
+}
38
+
39
+type network struct {
40
+ clients map[util.Key]bsnet.Receiver
41
+}
42
+
43
+func (n *network) Adapter(p *peer.Peer) bsnet.Adapter {
44
+ client := &networkClient{
45
+ local: p,
46
+ network: n,
47
+ }
48
+ n.clients[p.Key()] = client
49
+ return client
50
+}
51
+
52
+// TODO should this be completely asynchronous?
53
+// TODO what does the network layer do with errors received from services?
54
+func (n *network) SendMessage(
55
+ ctx context.Context,
56
+ from *peer.Peer,
57
+ to *peer.Peer,
58
+ message bsmsg.BitSwapMessage) error {
59
+
60
+ receiver, ok := n.clients[to.Key()]
61
+ if !ok {
62
+ return errors.New("Cannot locate peer on network")
63
+ }
64
+
65
+ // nb: terminate the context since the context wouldn't actually be passed
66
+ // over the network in a real scenario
67
+
68
+ go n.deliver(receiver, from, message)
69
+
70
+ return nil
71
+}
72
+
73
+func (n *network) deliver(
74
+ r bsnet.Receiver, from *peer.Peer, message bsmsg.BitSwapMessage) error {
75
+ if message == nil || from == nil {
76
+ return errors.New("Invalid input")
77
+ }
78
+
79
+ nextPeer, nextMsg, err := r.ReceiveMessage(context.TODO(), from, message)
80
+ if err != nil {
81
+
82
+ // TODO should this error be returned across network boundary?
83
+
84
+ // TODO this raises an interesting question about network contract. How
85
+ // can the network be expected to behave under different failure
86
+ // conditions? What if peer is unreachable? Will we know if messages
87
+ // aren't delivered?
88
+
89
+ return err
90
+ }
91
+
92
+ if (nextPeer == nil && nextMsg != nil) || (nextMsg == nil && nextPeer != nil) {
93
+ return errors.New("Malformed client request")
94
+ }
95
+
96
+ if nextPeer == nil && nextMsg == nil {
97
+ return nil
98
+ }
99
+
100
+ nextReceiver, ok := n.clients[nextPeer.Key()]
101
+ if !ok {
102
+ return errors.New("Cannot locate peer on network")
103
+ }
104
+ go n.deliver(nextReceiver, nextPeer, nextMsg)
105
+ return nil
106
+}
107
+
108
+var NoResponse = errors.New("No response received from the receiver")
109
+
110
+// TODO
111
+func (n *network) SendRequest(
112
+ ctx context.Context,
113
+ from *peer.Peer,
114
+ to *peer.Peer,
115
+ message bsmsg.BitSwapMessage) (
116
+ incoming bsmsg.BitSwapMessage, err error) {
117
+
118
+ r, ok := n.clients[to.Key()]
119
+ if !ok {
120
+ return nil, errors.New("Cannot locate peer on network")
121
+ }
122
+ nextPeer, nextMsg, err := r.ReceiveMessage(context.TODO(), from, message)
123
+ if err != nil {
124
+ return nil, err
125
+ // TODO return nil, NoResponse
126
+ }
127
+
128
+ // TODO dedupe code
129
+ if (nextPeer == nil && nextMsg != nil) || (nextMsg == nil && nextPeer != nil) {
130
+ return nil, errors.New("Malformed client request")
131
+ }
132
+
133
+ // TODO dedupe code
134
+ if nextPeer == nil && nextMsg == nil {
135
+ return nil, nil
136
+ }
137
+
138
+ // TODO test when receiver doesn't immediately respond to the initiator of the request
139
+ if !bytes.Equal(nextPeer.ID, from.ID) {
140
+ go func() {
141
+ nextReceiver, ok := n.clients[nextPeer.Key()]
142
+ if !ok {
143
+ // TODO log the error?
144
+ }
145
+ n.deliver(nextReceiver, nextPeer, nextMsg)
146
+ }()
147
+ return nil, NoResponse
148
+ }
149
+ return nextMsg, nil
150
+}
151
+
152
+type networkClient struct {
153
+ local *peer.Peer
154
+ bsnet.Receiver
155
+ network Network
156
+}
157
+
158
+func (nc *networkClient) SendMessage(
159
+ ctx context.Context,
160
+ to *peer.Peer,
161
+ message bsmsg.BitSwapMessage) error {
162
+ return nc.network.SendMessage(ctx, nc.local, to, message)
163
+}
164
+
165
+func (nc *networkClient) SendRequest(
166
+ ctx context.Context,
167
+ to *peer.Peer,
168
+ message bsmsg.BitSwapMessage) (incoming bsmsg.BitSwapMessage, err error) {
169
+ return nc.network.SendRequest(ctx, nc.local, to, message)
170
+}
171
+
172
+func (nc *networkClient) SetDelegate(r bsnet.Receiver) {
173
+ nc.Receiver = r
174
+}
exchange/bitswap/local_network_test.go
new
+138
@@ -0,0 +1,138 @@
1
+package bitswap
2
+
3
+import (
4
+ "sync"
5
+ "testing"
6
+
7
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
+ bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
9
+ bsnet "github.com/jbenet/go-ipfs/exchange/bitswap/network"
10
+ peer "github.com/jbenet/go-ipfs/peer"
11
+ testutil "github.com/jbenet/go-ipfs/util/testutil"
12
+)
13
+
14
+func TestSendRequestToCooperativePeer(t *testing.T) {
15
+ net := LocalNetwork()
16
+
17
+ idOfRecipient := []byte("recipient")
18
+
19
+ t.Log("Get two network adapters")
20
+
21
+ initiator := net.Adapter(&peer.Peer{ID: []byte("initiator")})
22
+ recipient := net.Adapter(&peer.Peer{ID: idOfRecipient})
23
+
24
+ expectedStr := "response from recipient"
25
+ recipient.SetDelegate(lambda(func(
26
+ ctx context.Context,
27
+ from *peer.Peer,
28
+ incoming bsmsg.BitSwapMessage) (
29
+ *peer.Peer, bsmsg.BitSwapMessage, error) {
30
+
31
+ t.Log("Recipient received a message from the network")
32
+
33
+ // TODO test contents of incoming message
34
+
35
+ m := bsmsg.New()
36
+ m.AppendBlock(testutil.NewBlockOrFail(t, expectedStr))
37
+
38
+ return from, m, nil
39
+ }))
40
+
41
+ t.Log("Build a message and send a synchronous request to recipient")
42
+
43
+ message := bsmsg.New()
44
+ message.AppendBlock(testutil.NewBlockOrFail(t, "data"))
45
+ response, err := initiator.SendRequest(
46
+ context.Background(), &peer.Peer{ID: idOfRecipient}, message)
47
+ if err != nil {
48
+ t.Fatal(err)
49
+ }
50
+
51
+ t.Log("Check the contents of the response from recipient")
52
+
53
+ for _, blockFromRecipient := range response.Blocks() {
54
+ if string(blockFromRecipient.Data) == expectedStr {
55
+ return
56
+ }
57
+ }
58
+ t.Fatal("Should have returned after finding expected block data")
59
+}
60
+
61
+func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
62
+ net := LocalNetwork()
63
+ idOfResponder := []byte("responder")
64
+ waiter := net.Adapter(&peer.Peer{ID: []byte("waiter")})
65
+ responder := net.Adapter(&peer.Peer{ID: idOfResponder})
66
+
67
+ var wg sync.WaitGroup
68
+
69
+ wg.Add(1)
70
+
71
+ expectedStr := "received async"
72
+
73
+ responder.SetDelegate(lambda(func(
74
+ ctx context.Context,
75
+ fromWaiter *peer.Peer,
76
+ msgFromWaiter bsmsg.BitSwapMessage) (
77
+ *peer.Peer, bsmsg.BitSwapMessage, error) {
78
+
79
+ msgToWaiter := bsmsg.New()
80
+ msgToWaiter.AppendBlock(testutil.NewBlockOrFail(t, expectedStr))
81
+
82
+ return fromWaiter, msgToWaiter, nil
83
+ }))
84
+
85
+ waiter.SetDelegate(lambda(func(
86
+ ctx context.Context,
87
+ fromResponder *peer.Peer,
88
+ msgFromResponder bsmsg.BitSwapMessage) (
89
+ *peer.Peer, bsmsg.BitSwapMessage, error) {
90
+
91
+ // TODO assert that this came from the correct peer and that the message contents are as expected
92
+ ok := false
93
+ for _, b := range msgFromResponder.Blocks() {
94
+ if string(b.Data) == expectedStr {
95
+ wg.Done()
96
+ ok = true
97
+ }
98
+ }
99
+
100
+ if !ok {
101
+ t.Fatal("Message not received from the responder")
102
+
103
+ }
104
+ return nil, nil, nil
105
+ }))
106
+
107
+ messageSentAsync := bsmsg.New()
108
+ messageSentAsync.AppendBlock(testutil.NewBlockOrFail(t, "data"))
109
+ errSending := waiter.SendMessage(
110
+ context.Background(), &peer.Peer{ID: idOfResponder}, messageSentAsync)
111
+ if errSending != nil {
112
+ t.Fatal(errSending)
113
+ }
114
+
115
+ wg.Wait() // until waiter delegate function is executed
116
+}
117
+
118
+type receiverFunc func(ctx context.Context, p *peer.Peer,
119
+ incoming bsmsg.BitSwapMessage) (*peer.Peer, bsmsg.BitSwapMessage, error)
120
+
121
+// lambda returns a Receiver instance given a receiver function
122
+func lambda(f receiverFunc) bsnet.Receiver {
123
+ return &lambdaImpl{
124
+ f: f,
125
+ }
126
+}
127
+
128
+type lambdaImpl struct {
129
+ f func(ctx context.Context, p *peer.Peer,
130
+ incoming bsmsg.BitSwapMessage) (
131
+ *peer.Peer, bsmsg.BitSwapMessage, error)
132
+}
133
+
134
+func (lam *lambdaImpl) ReceiveMessage(ctx context.Context,
135
+ p *peer.Peer, incoming bsmsg.BitSwapMessage) (
136
+ *peer.Peer, bsmsg.BitSwapMessage, error) {
137
+ return lam.f(ctx, p, incoming)
138
+}
exchange/bitswap/strategy/strategy.go
+7
@@ -51,6 +51,13 @@ func (s *strategist) Seed(int64) {
51
}
52
53
func (s *strategist) MessageReceived(p *peer.Peer, m bsmsg.BitSwapMessage) error {
54
+ // TODO find a more elegant way to handle this check
55
+ if p == nil {
56
+ return errors.New("Strategy received nil peer")
57
+ }
58
+ if m == nil {
59
+ return errors.New("Strategy received nil message")
60
+ }
61
l := s.ledger(p)
62
for _, key := range m.Wantlist() {
63
l.Wants(key)