@cryptotaxi247 / kubo / commits / 1555ce7c4

bitswap dials peers

Important bugfix. Otherwise bitswap cannot message peers the node has not connected to yet :(

Juan Batiz-Benet committed Oct 11, 2014 at 06:31 UTC 1555ce7c48bc5ddadadadadf5e66e8548caa48f9
4 files changed +35 -6
exchange/bitswap/bitswap.go
+9 -5
@@ -24,13 +24,12 @@ func NetMessageSession(parent context.Context, p *peer.Peer,
24 net inet.Network, srv inet.Service, directory bsnet.Routing,
25 d ds.Datastore, nice bool) exchange.Interface {
26
27 - networkAdapter := bsnet.NetMessageAdapter(srv, nil)
27 + networkAdapter := bsnet.NetMessageAdapter(srv, net, nil)
28 bs := &bitswap{
29 blockstore: blockstore.NewBlockstore(d),
30 notifications: notifications.New(),
31 strategy: strategy.New(nice),
32 routing: directory,
33 - network: net,
33 sender: networkAdapter,
34 wantlist: u.NewKeySet(),
35 }
@@ -42,9 +41,6 @@ func NetMessageSession(parent context.Context, p *peer.Peer,
41 // bitswap instances implement the bitswap protocol.
42 type bitswap struct {
43
45 - // network maintains connections to the outside world.
46 - network inet.Network
47 -
44 // sender delivers messages on behalf of the session
45 sender bsnet.Adapter
46
@@ -88,8 +84,16 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
84 for iiiii := range peersToQuery {
85 log.Debug("bitswap got peersToQuery: %s", iiiii)
86 go func(p *peer.Peer) {
87 +
88 + err := bs.sender.DialPeer(p)
89 + if err != nil {
90 + log.Error("Error sender.DialPeer(%s)", p)
91 + return
92 + }
93 +
94 response, err := bs.sender.SendRequest(ctx, p, message)
95 if err != nil {
96 + log.Error("Error sender.SendRequest(%s)", p)
97 return
98 }
99 // FIXME ensure accounting is handled correctly when
exchange/bitswap/network/interface.go
+3
@@ -11,6 +11,9 @@ import (
11 // Adapter provides network connectivity for BitSwap sessions
12 type Adapter interface {
13
14 + // DialPeer ensures there is a connection to peer.
15 + DialPeer(*peer.Peer) error
16 +
17 // SendMessage sends a BitSwap message to a peer.
18 SendMessage(
19 context.Context,
exchange/bitswap/network/net_message_adapter.go
+7 -1
@@ -10,9 +10,10 @@ import (
10 )
11
12 // NetMessageAdapter wraps a NetMessage network service
13 -func NetMessageAdapter(s inet.Service, r Receiver) Adapter {
13 +func NetMessageAdapter(s inet.Service, n inet.Network, r Receiver) Adapter {
14 adapter := impl{
15 nms: s,
16 + net: n,
17 receiver: r,
18 }
19 s.SetHandler(&adapter)
@@ -22,6 +23,7 @@ func NetMessageAdapter(s inet.Service, r Receiver) Adapter {
23 // implements an Adapter that integrates with a NetMessage network service
24 type impl struct {
25 nms inet.Service
26 + net inet.Network
27
28 // inbound messages from the network are forwarded to the receiver
29 receiver Receiver
@@ -58,6 +60,10 @@ func (adapter *impl) HandleMessage(
60 return outgoing
61 }
62
63 +func (adapter *impl) DialPeer(p *peer.Peer) error {
64 + return adapter.DialPeer(p)
65 +}
66 +
67 func (adapter *impl) SendMessage(
68 ctx context.Context,
69 p *peer.Peer,
exchange/bitswap/testnet/network.go
+16
@@ -3,6 +3,7 @@ package bitswap
3 import (
4 "bytes"
5 "errors"
6 + "fmt"
7
8 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9 bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
@@ -14,6 +15,8 @@ import (
15 type Network interface {
16 Adapter(*peer.Peer) bsnet.Adapter
17
18 + HasPeer(*peer.Peer) bool
19 +
20 SendMessage(
21 ctx context.Context,
22 from *peer.Peer,
@@ -49,6 +52,11 @@ func (n *network) Adapter(p *peer.Peer) bsnet.Adapter {
52 return client
53 }
54
55 +func (n *network) HasPeer(p *peer.Peer) bool {
56 + _, found := n.clients[p.Key()]
57 + return found
58 +}
59 +
60 // TODO should this be completely asynchronous?
61 // TODO what does the network layer do with errors received from services?
62 func (n *network) SendMessage(
@@ -155,6 +163,14 @@ func (nc *networkClient) SendRequest(
163 return nc.network.SendRequest(ctx, nc.local, to, message)
164 }
165
166 +func (nc *networkClient) DialPeer(p *peer.Peer) error {
167 + // no need to do anything because dialing isn't a thing in this test net.
168 + if !nc.network.HasPeer(p) {
169 + return fmt.Errorf("Peer not in network: %s", p)
170 + }
171 + return nil
172 +}
173 +
174 func (nc *networkClient) SetDelegate(r bsnet.Receiver) {
175 nc.Receiver = r
176 }