@cryptotaxi247 / kubo / commits / 0324b4b28

mild refactor of bitswap

Jeromy Johnson committed May 5, 2015 at 12:28 UTC 0324b4b28356059151d6215ae8d58880f5354e20
4 files changed +23 -155
exchange/bitswap/bitswap.go
+11 -23
@@ -4,6 +4,7 @@ package bitswap
4
5 import (
6 "errors"
7 + "fmt"
8 "math"
9 "sync"
10 "time"
@@ -324,47 +325,31 @@ func (bs *Bitswap) sendWantlistToProviders(ctx context.Context, entries []wantli
325 }
326
327 // TODO(brian): handle errors
327 -func (bs *Bitswap) ReceiveMessage(ctx context.Context, p peer.ID, incoming bsmsg.BitSwapMessage) (
328 - peer.ID, bsmsg.BitSwapMessage) {
328 +func (bs *Bitswap) ReceiveMessage(ctx context.Context, p peer.ID, incoming bsmsg.BitSwapMessage) error {
329 defer log.EventBegin(ctx, "receiveMessage", p, incoming).Done()
330
331 - if p == "" {
332 - log.Debug("Received message from nil peer!")
333 - // TODO propagate the error upward
334 - return "", nil
335 - }
336 - if incoming == nil {
337 - log.Debug("Got nil bitswap message!")
338 - // TODO propagate the error upward
339 - return "", nil
340 - }
341 -
331 // This call records changes to wantlists, blocks received,
332 // and number of bytes transfered.
333 bs.engine.MessageReceived(p, incoming)
334 // TODO: this is bad, and could be easily abused.
335 // Should only track *useful* messages in ledger
336
337 + var keys []u.Key
338 for _, block := range incoming.Blocks() {
339 bs.blocksRecvd++
340 if has, err := bs.blockstore.Has(block.Key()); err == nil && has {
341 bs.dupBlocksRecvd++
342 }
343 + log.Debugf("got block %s from %s", block, p)
344 hasBlockCtx, cancel := context.WithTimeout(ctx, hasBlockTimeout)
345 if err := bs.HasBlock(hasBlockCtx, block); err != nil {
355 - log.Debug(err)
346 + return fmt.Errorf("ReceiveMessage HasBlock error: %s", err)
347 }
348 cancel()
358 - }
359 -
360 - var keys []u.Key
361 - for _, block := range incoming.Blocks() {
349 keys = append(keys, block.Key())
350 }
364 - bs.cancelBlocks(ctx, keys)
351
366 - // TODO: consider changing this function to not return anything
367 - return "", nil
352 + return bs.cancelBlocks(ctx, keys)
353 }
354
355 // Connected/Disconnected warns bitswap about peer connections
@@ -384,21 +369,24 @@ func (bs *Bitswap) PeerDisconnected(p peer.ID) {
369 bs.engine.PeerDisconnected(p)
370 }
371
387 -func (bs *Bitswap) cancelBlocks(ctx context.Context, bkeys []u.Key) {
372 +func (bs *Bitswap) cancelBlocks(ctx context.Context, bkeys []u.Key) error {
373 if len(bkeys) < 1 {
389 - return
374 + return nil
375 }
376 message := bsmsg.New()
377 message.SetFull(false)
378 for _, k := range bkeys {
379 + log.Debug("cancel block: %s", k)
380 message.Cancel(k)
381 }
382 for _, p := range bs.engine.Peers() {
383 err := bs.send(ctx, p, message)
384 if err != nil {
385 log.Debugf("Error sending message: %s", err)
386 + return err
387 }
388 }
389 + return nil
390 }
391
392 func (bs *Bitswap) wantNewBlocks(ctx context.Context, bkeys []u.Key) {
exchange/bitswap/network/interface.go
+3 -8
@@ -19,12 +19,6 @@ type BitSwapNetwork interface {
19 peer.ID,
20 bsmsg.BitSwapMessage) error
21
22 - // SendRequest sends a BitSwap message to a peer and waits for a response.
23 - SendRequest(
24 - context.Context,
25 - peer.ID,
26 - bsmsg.BitSwapMessage) (incoming bsmsg.BitSwapMessage, err error)
27 -
22 // SetDelegate registers the Reciver to handle messages received from the
23 // network.
24 SetDelegate(Receiver)
@@ -35,8 +29,9 @@ type BitSwapNetwork interface {
29 // Implement Receiver to receive messages from the BitSwapNetwork
30 type Receiver interface {
31 ReceiveMessage(
38 - ctx context.Context, sender peer.ID, incoming bsmsg.BitSwapMessage) (
39 - destination peer.ID, outgoing bsmsg.BitSwapMessage)
32 + ctx context.Context,
33 + sender peer.ID,
34 + incoming bsmsg.BitSwapMessage) error
35
36 ReceiveError(error)
37
exchange/bitswap/testnet/network_test.go
+8 -62
@@ -14,57 +14,6 @@ import (
14 testutil "github.com/ipfs/go-ipfs/util/testutil"
15 )
16
17 -func TestSendRequestToCooperativePeer(t *testing.T) {
18 - net := VirtualNetwork(mockrouting.NewServer(), delay.Fixed(0))
19 -
20 - recipientPeer := testutil.RandIdentityOrFatal(t)
21 -
22 - t.Log("Get two network adapters")
23 -
24 - initiator := net.Adapter(testutil.RandIdentityOrFatal(t))
25 - recipient := net.Adapter(recipientPeer)
26 -
27 - expectedStr := "response from recipient"
28 - recipient.SetDelegate(lambda(func(
29 - ctx context.Context,
30 - from peer.ID,
31 - incoming bsmsg.BitSwapMessage) (
32 - peer.ID, bsmsg.BitSwapMessage) {
33 -
34 - t.Log("Recipient received a message from the network")
35 -
36 - // TODO test contents of incoming message
37 -
38 - m := bsmsg.New()
39 - m.AddBlock(blocks.NewBlock([]byte(expectedStr)))
40 -
41 - return from, m
42 - }))
43 -
44 - t.Log("Build a message and send a synchronous request to recipient")
45 -
46 - message := bsmsg.New()
47 - message.AddBlock(blocks.NewBlock([]byte("data")))
48 - response, err := initiator.SendRequest(
49 - context.Background(), recipientPeer.ID(), message)
50 - if err != nil {
51 - t.Fatal(err)
52 - }
53 -
54 - t.Log("Check the contents of the response from recipient")
55 -
56 - if response == nil {
57 - t.Fatal("Should have received a response")
58 - }
59 -
60 - for _, blockFromRecipient := range response.Blocks() {
61 - if string(blockFromRecipient.Data) == expectedStr {
62 - return
63 - }
64 - }
65 - t.Fatal("Should have returned after finding expected block data")
66 -}
67 -
17 func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
18 net := VirtualNetwork(mockrouting.NewServer(), delay.Fixed(0))
19 responderPeer := testutil.RandIdentityOrFatal(t)
@@ -80,20 +29,19 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
29 responder.SetDelegate(lambda(func(
30 ctx context.Context,
31 fromWaiter peer.ID,
83 - msgFromWaiter bsmsg.BitSwapMessage) (
84 - peer.ID, bsmsg.BitSwapMessage) {
32 + msgFromWaiter bsmsg.BitSwapMessage) error {
33
34 msgToWaiter := bsmsg.New()
35 msgToWaiter.AddBlock(blocks.NewBlock([]byte(expectedStr)))
36 + waiter.SendMessage(ctx, fromWaiter, msgToWaiter)
37
89 - return fromWaiter, msgToWaiter
38 + return nil
39 }))
40
41 waiter.SetDelegate(lambda(func(
42 ctx context.Context,
43 fromResponder peer.ID,
95 - msgFromResponder bsmsg.BitSwapMessage) (
96 - peer.ID, bsmsg.BitSwapMessage) {
44 + msgFromResponder bsmsg.BitSwapMessage) error {
45
46 // TODO assert that this came from the correct peer and that the message contents are as expected
47 ok := false
@@ -108,7 +56,7 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
56 t.Fatal("Message not received from the responder")
57
58 }
111 - return "", nil
59 + return nil
60 }))
61
62 messageSentAsync := bsmsg.New()
@@ -123,7 +71,7 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
71 }
72
73 type receiverFunc func(ctx context.Context, p peer.ID,
126 - incoming bsmsg.BitSwapMessage) (peer.ID, bsmsg.BitSwapMessage)
74 + incoming bsmsg.BitSwapMessage) error
75
76 // lambda returns a Receiver instance given a receiver function
77 func lambda(f receiverFunc) bsnet.Receiver {
@@ -133,13 +81,11 @@ func lambda(f receiverFunc) bsnet.Receiver {
81 }
82
83 type lambdaImpl struct {
136 - f func(ctx context.Context, p peer.ID, incoming bsmsg.BitSwapMessage) (
137 - peer.ID, bsmsg.BitSwapMessage)
84 + f func(ctx context.Context, p peer.ID, incoming bsmsg.BitSwapMessage) error
85 }
86
87 func (lam *lambdaImpl) ReceiveMessage(ctx context.Context,
141 - p peer.ID, incoming bsmsg.BitSwapMessage) (
142 - peer.ID, bsmsg.BitSwapMessage) {
88 + p peer.ID, incoming bsmsg.BitSwapMessage) error {
89 return lam.f(ctx, p, incoming)
90 }
91
exchange/bitswap/testnet/virtual.go
+1 -62
@@ -72,61 +72,7 @@ func (n *network) deliver(
72
73 n.delay.Wait()
74
75 - nextPeer, nextMsg := r.ReceiveMessage(context.TODO(), from, message)
76 -
77 - if (nextPeer == "" && nextMsg != nil) || (nextMsg == nil && nextPeer != "") {
78 - return errors.New("Malformed client request")
79 - }
80 -
81 - if nextPeer == "" && nextMsg == nil { // no response to send
82 - return nil
83 - }
84 -
85 - nextReceiver, ok := n.clients[nextPeer]
86 - if !ok {
87 - return errors.New("Cannot locate peer on network")
88 - }
89 - go n.deliver(nextReceiver, nextPeer, nextMsg)
90 - return nil
91 -}
92 -
93 -// TODO
94 -func (n *network) SendRequest(
95 - ctx context.Context,
96 - from peer.ID,
97 - to peer.ID,
98 - message bsmsg.BitSwapMessage) (
99 - incoming bsmsg.BitSwapMessage, err error) {
100 -
101 - r, ok := n.clients[to]
102 - if !ok {
103 - return nil, errors.New("Cannot locate peer on network")
104 - }
105 - nextPeer, nextMsg := r.ReceiveMessage(context.TODO(), from, message)
106 -
107 - // TODO dedupe code
108 - if (nextPeer == "" && nextMsg != nil) || (nextMsg == nil && nextPeer != "") {
109 - r.ReceiveError(errors.New("Malformed client request"))
110 - return nil, nil
111 - }
112 -
113 - // TODO dedupe code
114 - if nextPeer == "" && nextMsg == nil {
115 - return nil, nil
116 - }
117 -
118 - // TODO test when receiver doesn't immediately respond to the initiator of the request
119 - if nextPeer != from {
120 - go func() {
121 - nextReceiver, ok := n.clients[nextPeer]
122 - if !ok {
123 - // TODO log the error?
124 - }
125 - n.deliver(nextReceiver, nextPeer, nextMsg)
126 - }()
127 - return nil, nil
128 - }
129 - return nextMsg, nil
75 + return r.ReceiveMessage(context.TODO(), from, message)
76 }
77
78 type networkClient struct {
@@ -143,13 +89,6 @@ func (nc *networkClient) SendMessage(
89 return nc.network.SendMessage(ctx, nc.local, to, message)
90 }
91
146 -func (nc *networkClient) SendRequest(
147 - ctx context.Context,
148 - to peer.ID,
149 - message bsmsg.BitSwapMessage) (incoming bsmsg.BitSwapMessage, err error) {
150 - return nc.network.SendRequest(ctx, nc.local, to, message)
151 -}
152 -
92 // FindProvidersAsync returns a channel of providers for the given key
93 func (nc *networkClient) FindProvidersAsync(ctx context.Context, k util.Key, max int) <-chan peer.ID {
94