feat(bs:net) impl service wrapper
Brian Tiger Chow committed
Sep 12, 2014 at 20:51 UTC
7fcb5d3a4bb01d7098081c4ea0dfd3088eee80b2
7 files changed
+173
-6
bitswap/bitswap.go
+20
-6
@@ -1,12 +1,14 @@
1
package bitswap
2
3
import (
4
+ "errors"
5
"time"
6
7
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
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"
13
blocks "github.com/jbenet/go-ipfs/blocks"
14
swarm "github.com/jbenet/go-ipfs/net/swarm"
@@ -31,6 +33,7 @@ type BitSwap struct {
33
peer *peer.Peer
34
35
// net holds the connections to all peers.
36
+ sender bsnet.Sender
37
net swarm.Network
38
meschan *swarm.Chan
39
@@ -61,17 +64,22 @@ type BitSwap struct {
64
65
// NewBitSwap creates a new BitSwap instance. It does not check its parameters.
66
func NewBitSwap(p *peer.Peer, net swarm.Network, d ds.Datastore, r routing.IpfsRouting) *BitSwap {
67
+ receiver := receiver{}
68
+ sender := bsnet.NewBSNetService(context.Background(), &receiver)
69
bs := &BitSwap{
65
- peer: p,
66
- net: net,
67
- datastore: d,
68
- partners: LedgerMap{},
69
- wantList: KeySet{},
70
- routing: r.(*dht.IpfsDHT),
70
+ peer: p,
71
+ net: net,
72
+ datastore: d,
73
+ partners: LedgerMap{},
74
+ wantList: KeySet{},
75
+ routing: r.(*dht.IpfsDHT),
76
+ // TODO(brian): replace |meschan| with |sender| in BitSwap impl
77
meschan: net.GetChannel(swarm.PBWrapper_BITSWAP),
78
+ sender: sender,
79
haltChan: make(chan struct{}),
80
notifications: notifications.New(),
81
}
82
+ receiver.Delegate(bs)
83
84
go bs.handleMessages()
85
return bs
@@ -274,3 +282,9 @@ func (bs *BitSwap) SetStrategy(sf StrategyFunc) {
282
ledger.Strategy = sf
283
}
284
}
285
+
286
+func (r *BitSwap) ReceiveMessage(
287
+ ctx context.Context, incoming bsmsg.BitSwapMessage) (
288
+ bsmsg.BitSwapMessage, *peer.Peer, error) {
289
+ return nil, nil, errors.New("TODO implement")
290
+}
bitswap/message/Makefile
new
+8
@@ -0,0 +1,8 @@
1
+# TODO(brian): add proto tasks
2
+all: message.pb.go
3
+
4
+message.pb.go: message.proto
5
+ protoc --gogo_out=. --proto_path=../../../../../:/usr/local/opt/protobuf/include:. $<
6
+
7
+clean:
8
+ rm message.pb.go
bitswap/message/message.go
+7
@@ -1,7 +1,10 @@
1
package message
2
3
import (
4
+ "errors"
5
+
6
proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
7
+ netmsg "github.com/jbenet/go-ipfs/net/message"
8
9
blocks "github.com/jbenet/go-ipfs/blocks"
10
nm "github.com/jbenet/go-ipfs/net/message"
@@ -55,6 +58,10 @@ func (m *message) AppendBlock(b *blocks.Block) {
58
m.pb.Blocks = append(m.pb.Blocks, b.Data)
59
}
60
61
+func FromNet(nmsg netmsg.NetMessage) (BitSwapMessage, error) {
62
+ return nil, errors.New("TODO implement")
63
+}
64
+
65
func FromSwarm(sms swarm.Message) (BitSwapMessage, error) {
66
var protoMsg PBMessage
67
err := proto.Unmarshal(sms.Data, &protoMsg)
bitswap/network/interface.go
new
+20
@@ -0,0 +1,20 @@
1
+package network
2
+
3
+import (
4
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
+
6
+ bsmsg "github.com/jbenet/go-ipfs/bitswap/message"
7
+ peer "github.com/jbenet/go-ipfs/peer"
8
+)
9
+
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)
14
+}
15
+
16
+// TODO(brian): consider returning a NetMessage
17
+type Receiver interface {
18
+ ReceiveMessage(ctx context.Context, incoming bsmsg.BitSwapMessage) (
19
+ outgoing bsmsg.BitSwapMessage, destination *peer.Peer, err error)
20
+}
bitswap/network/service_wrapper.go
new
+76
@@ -0,0 +1,76 @@
1
+package network
2
+
3
+import (
4
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
+
6
+ bsmsg "github.com/jbenet/go-ipfs/bitswap/message"
7
+ netmsg "github.com/jbenet/go-ipfs/net/message"
8
+ netservice "github.com/jbenet/go-ipfs/net/service"
9
+ peer "github.com/jbenet/go-ipfs/peer"
10
+)
11
+
12
+func NewBSNetService(ctx context.Context, r Receiver) Sender {
13
+ h := &handlerWrapper{r}
14
+ s := netservice.NewService(ctx, h)
15
+ return &serviceWrapper{*s}
16
+}
17
+
18
+// handlerWrapper is responsible for marshaling/unmarshaling NetMessages. It
19
+// delegates calls to the BitSwap delegate.
20
+type handlerWrapper struct {
21
+ bitswapDelegate Receiver
22
+}
23
+
24
+// HandleMessage marshals and unmarshals net messages, forwarding them to the
25
+// BitSwapMessage receiver
26
+func (wrapper *handlerWrapper) HandleMessage(
27
+ ctx context.Context, incoming netmsg.NetMessage) (netmsg.NetMessage, error) {
28
+
29
+ received, err := bsmsg.FromNet(incoming)
30
+ if err != nil {
31
+ return nil, err
32
+ }
33
+
34
+ bsmsg, p, err := wrapper.bitswapDelegate.ReceiveMessage(ctx, received)
35
+ if err != nil {
36
+ return nil, err
37
+ }
38
+ if bsmsg == nil {
39
+ return nil, nil
40
+ }
41
+
42
+ outgoing, err := bsmsg.ToNet(p)
43
+ if err != nil {
44
+ return nil, err
45
+ }
46
+
47
+ return outgoing, nil
48
+}
49
+
50
+type serviceWrapper struct {
51
+ serviceDelegate netservice.Service
52
+}
53
+
54
+func (wrapper *serviceWrapper) SendMessage(
55
+ ctx context.Context, p *peer.Peer, outgoing bsmsg.Exportable) error {
56
+ nmsg, err := outgoing.ToNet(p)
57
+ if err != nil {
58
+ return err
59
+ }
60
+ req, err := netservice.NewRequest(p.ID)
61
+ return wrapper.serviceDelegate.SendMessage(ctx, nmsg, req.ID)
62
+}
63
+
64
+func (wrapper *serviceWrapper) SendRequest(ctx context.Context,
65
+ p *peer.Peer, outgoing bsmsg.Exportable) (bsmsg.BitSwapMessage, error) {
66
+
67
+ outgoingMsg, err := outgoing.ToNet(p)
68
+ if err != nil {
69
+ return nil, err
70
+ }
71
+ incomingMsg, err := wrapper.serviceDelegate.SendRequest(ctx, outgoingMsg)
72
+ if err != nil {
73
+ return nil, err
74
+ }
75
+ return bsmsg.FromNet(incomingMsg)
76
+}
bitswap/receiver.go
new
+29
@@ -0,0 +1,29 @@
1
+package bitswap
2
+
3
+import (
4
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
+ bsmsg "github.com/jbenet/go-ipfs/bitswap/message"
6
+ bsnet "github.com/jbenet/go-ipfs/bitswap/network"
7
+ peer "github.com/jbenet/go-ipfs/peer"
8
+)
9
+
10
+// receiver breaks the circular dependency between bitswap and its sender
11
+// NB: A sender is instantiated with a handler and this sender is then passed
12
+// as a constructor argument to BitSwap. However, the handler is BitSwap!
13
+// Hence, this receiver.
14
+type receiver struct {
15
+ delegate bsnet.Receiver
16
+}
17
+
18
+func (r *receiver) ReceiveMessage(
19
+ ctx context.Context, incoming bsmsg.BitSwapMessage) (
20
+ bsmsg.BitSwapMessage, *peer.Peer, error) {
21
+ if r.delegate == nil {
22
+ return nil, nil, nil
23
+ }
24
+ return r.delegate.ReceiveMessage(ctx, incoming)
25
+}
26
+
27
+func (r *receiver) Delegate(delegate bsnet.Receiver) {
28
+ r.delegate = delegate
29
+}
bitswap/receiver_test.go
new
+13
@@ -0,0 +1,13 @@
1
+package bitswap
2
+
3
+import (
4
+ "testing"
5
+
6
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
+ bsmsg "github.com/jbenet/go-ipfs/bitswap/message"
8
+)
9
+
10
+func TestDoesntPanicIfDelegateNotPresent(t *testing.T) {
11
+ r := receiver{}
12
+ r.ReceiveMessage(context.Background(), bsmsg.New())
13
+}