@cryptotaxi247 / kubo / commits / 014157cac

refac(bitswap) simply network interfaces

Brian Tiger Chow committed Sep 16, 2014 at 02:05 UTC 014157cac6ab3cc249e0a39c156a27b7a2e08612
5 files changed +85 -49
bitswap/bitswap.go
+7 -12
@@ -8,10 +8,9 @@ import (
8 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
9
10 bsmsg "github.com/jbenet/go-ipfs/bitswap/message"
11 + bsnet "github.com/jbenet/go-ipfs/bitswap/network"
12 notifications "github.com/jbenet/go-ipfs/bitswap/notifications"
12 - tx "github.com/jbenet/go-ipfs/bitswap/transmission"
13 blocks "github.com/jbenet/go-ipfs/blocks"
14 - net "github.com/jbenet/go-ipfs/net"
14 peer "github.com/jbenet/go-ipfs/peer"
15 routing "github.com/jbenet/go-ipfs/routing"
16 u "github.com/jbenet/go-ipfs/util"
@@ -34,7 +33,7 @@ type BitSwap struct {
33 peer *peer.Peer
34
35 // sender delivers messages on behalf of the session
37 - sender tx.Sender
36 + sender bsnet.NetworkAdapter
37
38 // datastore is the local database // Ledgers of known
39 datastore ds.Datastore
@@ -62,21 +61,16 @@ type BitSwap struct {
61 }
62
63 // NewSession initializes a bitswap session.
65 -func NewSession(parent context.Context, s net.Sender, p *peer.Peer, d ds.Datastore, r routing.IpfsRouting) *BitSwap {
64 +func NewSession(parent context.Context, s bsnet.NetworkService, p *peer.Peer, d ds.Datastore, r routing.IpfsRouting) *BitSwap {
65
67 - // TODO(brian): define a contract for management of async operations that
68 - // fall under bitswap's purview
69 - // ctx, _ := context.WithCancel(parent)
70 -
71 - receiver := tx.Forwarder{}
72 - sender := tx.NewSender(s)
66 + receiver := bsnet.Forwarder{}
67 bs := &BitSwap{
68 peer: p,
69 datastore: d,
70 partners: LedgerMap{},
71 wantList: KeySet{},
72 routing: r,
79 - sender: sender,
73 + sender: bsnet.NewNetworkAdapter(s, &receiver),
74 haltChan: make(chan struct{}),
75 notifications: notifications.New(),
76 strategy: YesManStrategy,
@@ -246,7 +240,7 @@ func (bs *BitSwap) Halt() {
240
241 func (bs *BitSwap) ReceiveMessage(
242 ctx context.Context, sender *peer.Peer, incoming bsmsg.BitSwapMessage) (
249 - bsmsg.BitSwapMessage, *peer.Peer, error) {
243 + *peer.Peer, bsmsg.BitSwapMessage, error) {
244 if incoming.Blocks() != nil {
245 for _, block := range incoming.Blocks() {
246 go bs.blockReceive(sender, block)
@@ -255,6 +249,7 @@ func (bs *BitSwap) ReceiveMessage(
249
250 if incoming.Wantlist() != nil {
251 for _, want := range incoming.Wantlist() {
252 + // TODO(brian): return the block synchronously
253 go bs.peerWantsBlock(sender, want)
254 }
255 }
bitswap/network/forwarder.go
+6 -6
@@ -1,4 +1,4 @@
1 -package transmission
1 +package network
2
3 import (
4 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
@@ -6,17 +6,17 @@ import (
6 peer "github.com/jbenet/go-ipfs/peer"
7 )
8
9 -// Forwarder breaks the circular dependency between bitswap and its sender
10 -// NB: A sender is instantiated with a handler and this sender is then passed
11 -// as a constructor argument to BitSwap. However, the handler is BitSwap!
12 -// Hence, this receiver.
9 +// Forwarder receives messages and forwards them to the delegate.
10 +//
11 +// Forwarder breaks the circular dependency between the BitSwap Session and the
12 +// Network Service.
13 type Forwarder struct {
14 delegate Receiver
15 }
16
17 func (r *Forwarder) ReceiveMessage(
18 ctx context.Context, sender *peer.Peer, incoming bsmsg.BitSwapMessage) (
19 - bsmsg.BitSwapMessage, *peer.Peer, error) {
19 + *peer.Peer, bsmsg.BitSwapMessage, error) {
20 if r.delegate == nil {
21 return nil, nil, nil
22 }
bitswap/network/forwarder_test.go
+12 -2
@@ -1,4 +1,4 @@
1 -package transmission
1 +package network
2
3 import (
4 "testing"
@@ -13,4 +13,14 @@ func TestDoesntPanicIfDelegateNotPresent(t *testing.T) {
13 fwdr.ReceiveMessage(context.Background(), &peer.Peer{}, bsmsg.New())
14 }
15
16 -// TODO(brian): func TestForwardsMessageToDelegate(t *testing.T)
16 +func TestForwardsMessageToDelegate(t *testing.T) {
17 + fwdr := Forwarder{delegate: &EchoDelegate{}}
18 + fwdr.ReceiveMessage(context.Background(), &peer.Peer{}, bsmsg.New())
19 +}
20 +
21 +type EchoDelegate struct{}
22 +
23 +func (d *EchoDelegate) ReceiveMessage(ctx context.Context, p *peer.Peer,
24 + incoming bsmsg.BitSwapMessage) (*peer.Peer, bsmsg.BitSwapMessage, error) {
25 + return p, incoming, nil
26 +}
bitswap/network/interface.go
+22 -7
@@ -1,23 +1,38 @@
1 -package transmission
1 +package network
2
3 import (
4 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5 + netservice "github.com/jbenet/go-ipfs/net/service"
6
7 bsmsg "github.com/jbenet/go-ipfs/bitswap/message"
8 + netmsg "github.com/jbenet/go-ipfs/net/message"
9 peer "github.com/jbenet/go-ipfs/peer"
10 )
11
10 -type Sender interface {
11 - SendMessage(ctx context.Context, destination *peer.Peer, message bsmsg.Exportable) error
12 - SendRequest(ctx context.Context, destination *peer.Peer, outgoing bsmsg.Exportable) (
13 - incoming bsmsg.BitSwapMessage, err error)
12 +// NetworkAdapter mediates the exchange's communication with the network.
13 +type NetworkAdapter interface {
14 +
15 + // SendMessage sends a BitSwap message to a peer.
16 + SendMessage(
17 + context.Context,
18 + *peer.Peer,
19 + bsmsg.BitSwapMessage) error
20 +
21 + // SendRequest sends a BitSwap message to a peer and waits for a response.
22 + SendRequest(
23 + context.Context,
24 + *peer.Peer,
25 + bsmsg.BitSwapMessage) (incoming bsmsg.BitSwapMessage, err error)
26 +
27 + // SetDelegate registers the Reciver to handle messages received from the
28 + // network.
29 + SetDelegate(Receiver)
30 }
31
16 -// TODO(brian): consider returning a NetMessage
32 type Receiver interface {
33 ReceiveMessage(
34 ctx context.Context, sender *peer.Peer, incoming bsmsg.BitSwapMessage) (
20 - outgoing bsmsg.BitSwapMessage, destination *peer.Peer, err error)
35 + destination *peer.Peer, outgoing bsmsg.BitSwapMessage, err error)
36 }
37
38 // TODO(brian): move this to go-ipfs/net package
bitswap/network/network_adapter.go
+38 -22
@@ -1,43 +1,54 @@
1 -package transmission
1 +package network
2
3 import (
4 + "errors"
5 +
6 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
8 bsmsg "github.com/jbenet/go-ipfs/bitswap/message"
7 - net "github.com/jbenet/go-ipfs/net"
9 netmsg "github.com/jbenet/go-ipfs/net/message"
10 peer "github.com/jbenet/go-ipfs/peer"
11 )
12
12 -// NewSender wraps the net.service.Sender to perform translation between
13 +// NewSender wraps a network Service to perform translation between
14 // BitSwapMessage and NetMessage formats. This allows the BitSwap session to
15 // ignore these details.
15 -func NewSender(s net.Sender) Sender {
16 - return &senderWrapper{s}
16 +func NewNetworkAdapter(s NetworkService, r Receiver) NetworkAdapter {
17 + adapter := networkAdapter{
18 + networkService: s,
19 + receiver: r,
20 + }
21 + s.SetHandler(&adapter)
22 + return &adapter
23 }
24
19 -// handlerWrapper implements the net.service.Handler interface. It is
20 -// responsible for converting between
21 -// delegates calls to the BitSwap delegate.
22 -type handlerWrapper struct {
23 - bitswapDelegate Receiver
25 +// networkAdapter implements NetworkAdapter
26 +type networkAdapter struct {
27 + networkService NetworkService
28 + receiver Receiver
29 }
30
31 // HandleMessage marshals and unmarshals net messages, forwarding them to the
32 // BitSwapMessage receiver
28 -func (wrapper *handlerWrapper) HandleMessage(
33 +func (adapter *networkAdapter) HandleMessage(
34 ctx context.Context, incoming netmsg.NetMessage) (netmsg.NetMessage, error) {
35
36 + if adapter.receiver == nil {
37 + return nil, errors.New("No receiver. NetMessage dropped")
38 + }
39 +
40 received, err := bsmsg.FromNet(incoming)
41 if err != nil {
42 return nil, err
43 }
44
36 - bsmsg, p, err := wrapper.bitswapDelegate.ReceiveMessage(ctx, incoming.Peer(), received)
45 + p, bsmsg, err := adapter.receiver.ReceiveMessage(ctx, incoming.Peer(), received)
46 if err != nil {
47 return nil, err
48 }
40 - if bsmsg == nil {
49 +
50 + // TODO(brian): put this in a helper function
51 + if bsmsg == nil || p == nil {
52 return nil, nil
53 }
54
@@ -49,29 +60,34 @@ func (wrapper *handlerWrapper) HandleMessage(
60 return outgoing, nil
61 }
62
52 -type senderWrapper struct {
53 - serviceDelegate net.Sender
54 -}
63 +func (adapter *networkAdapter) SendMessage(
64 + ctx context.Context,
65 + p *peer.Peer,
66 + outgoing bsmsg.BitSwapMessage) error {
67
56 -func (wrapper *senderWrapper) SendMessage(
57 - ctx context.Context, p *peer.Peer, outgoing bsmsg.Exportable) error {
68 nmsg, err := outgoing.ToNet(p)
69 if err != nil {
70 return err
71 }
62 - return wrapper.serviceDelegate.SendMessage(ctx, nmsg)
72 + return adapter.networkService.SendMessage(ctx, nmsg)
73 }
74
65 -func (wrapper *senderWrapper) SendRequest(ctx context.Context,
66 - p *peer.Peer, outgoing bsmsg.Exportable) (bsmsg.BitSwapMessage, error) {
75 +func (adapter *networkAdapter) SendRequest(
76 + ctx context.Context,
77 + p *peer.Peer,
78 + outgoing bsmsg.BitSwapMessage) (bsmsg.BitSwapMessage, error) {
79
80 outgoingMsg, err := outgoing.ToNet(p)
81 if err != nil {
82 return nil, err
83 }
72 - incomingMsg, err := wrapper.serviceDelegate.SendRequest(ctx, outgoingMsg)
84 + incomingMsg, err := adapter.networkService.SendRequest(ctx, outgoingMsg)
85 if err != nil {
86 return nil, err
87 }
88 return bsmsg.FromNet(incomingMsg)
89 }
90 +
91 +func (adapter *networkAdapter) SetDelegate(r Receiver) {
92 + adapter.receiver = r
93 +}