fix(net) pass contexts to dial peer
Brian Tiger Chow committed
Nov 5, 2014 at 04:26 UTC
23096de3c4e77c09bfc9c88fde8e43a9c80b2f2f
8 files changed
+16
-16
exchange/bitswap/network/net_message_adapter.go
+1
-1
@@ -68,7 +68,7 @@ func (adapter *impl) HandleMessage(
68
}
69
70
func (adapter *impl) DialPeer(ctx context.Context, p peer.Peer) error {
71
- return adapter.net.DialPeer(p)
71
+ return adapter.net.DialPeer(ctx, p)
72
}
73
74
func (adapter *impl) SendMessage(
exchange/bitswap/testnet/network.go
+1
-1
@@ -163,7 +163,7 @@ 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 {
166
+func (nc *networkClient) DialPeer(ctx context.Context, 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)
net/interface.go
+3
-2
@@ -1,6 +1,7 @@
1
package net
2
3
import (
4
+ "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
msg "github.com/jbenet/go-ipfs/net/message"
6
mux "github.com/jbenet/go-ipfs/net/mux"
7
srv "github.com/jbenet/go-ipfs/net/service"
@@ -19,7 +20,7 @@ type Network interface {
20
// TODO: for now, only listen on addrs in local peer when initializing.
21
22
// DialPeer attempts to establish a connection to a given peer
22
- DialPeer(peer.Peer) error
23
+ DialPeer(context.Context, peer.Peer) error
24
25
// ClosePeer connection to peer
26
ClosePeer(peer.Peer) error
@@ -64,5 +65,5 @@ type Service srv.Service
65
type Dialer interface {
66
67
// DialPeer attempts to establish a connection to a given peer
67
- DialPeer(peer.Peer) error
68
+ DialPeer(context.Context, peer.Peer) error
69
}
net/net.go
+1
-1
@@ -57,7 +57,7 @@ func NewIpfsNetwork(ctx context.Context, listen []ma.Multiaddr, local peer.Peer,
57
// func (n *IpfsNetwork) Listen(*ma.Muliaddr) error {}
58
59
// DialPeer attempts to establish a connection to a given peer
60
-func (n *IpfsNetwork) DialPeer(p peer.Peer) error {
60
+func (n *IpfsNetwork) DialPeer(ctx context.Context, p peer.Peer) error {
61
_, err := n.swarm.Dial(p)
62
return err
63
}
routing/dht/dht.go
+5
-5
@@ -100,7 +100,7 @@ func (dht *IpfsDHT) Connect(ctx context.Context, npeer peer.Peer) (peer.Peer, er
100
//
101
// /ip4/10.20.30.40/tcp/1234/ipfs/Qxhxxchxzcncxnzcnxzcxzm
102
//
103
- err := dht.dialer.DialPeer(npeer)
103
+ err := dht.dialer.DialPeer(ctx, npeer)
104
if err != nil {
105
return nil, err
106
}
@@ -311,7 +311,7 @@ func (dht *IpfsDHT) getFromPeerList(ctx context.Context, key u.Key,
311
peerlist []*pb.Message_Peer, level int) ([]byte, error) {
312
313
for _, pinfo := range peerlist {
314
- p, err := dht.ensureConnectedToPeer(pinfo)
314
+ p, err := dht.ensureConnectedToPeer(ctx, pinfo)
315
if err != nil {
316
log.Errorf("getFromPeers error: %s", err)
317
continue
@@ -496,14 +496,14 @@ func (dht *IpfsDHT) peerFromInfo(pbp *pb.Message_Peer) (peer.Peer, error) {
496
return p, nil
497
}
498
499
-func (dht *IpfsDHT) ensureConnectedToPeer(pbp *pb.Message_Peer) (peer.Peer, error) {
499
+func (dht *IpfsDHT) ensureConnectedToPeer(ctx context.Context, pbp *pb.Message_Peer) (peer.Peer, error) {
500
p, err := dht.peerFromInfo(pbp)
501
if err != nil {
502
return nil, err
503
}
504
505
// dial connection
506
- err = dht.dialer.DialPeer(p)
506
+ err = dht.dialer.DialPeer(ctx, p)
507
return p, err
508
}
509
@@ -556,7 +556,7 @@ func (dht *IpfsDHT) Bootstrap(ctx context.Context) {
556
if err != nil {
557
log.Error("Bootstrap peer error: %s", err)
558
}
559
- err = dht.dialer.DialPeer(p)
559
+ err = dht.dialer.DialPeer(ctx, p)
560
if err != nil {
561
log.Errorf("Bootstrap peer error: %s", err)
562
}
routing/dht/ext_test.go
+1
-2
@@ -4,7 +4,6 @@ import (
4
"testing"
5
6
crand "crypto/rand"
7
-
7
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
9
@@ -82,7 +81,7 @@ type fauxNet struct {
81
}
82
83
// DialPeer attempts to establish a connection to a given peer
85
-func (f *fauxNet) DialPeer(peer.Peer) error {
84
+func (f *fauxNet) DialPeer(context.Context, peer.Peer) error {
85
return nil
86
}
87
routing/dht/query.go
+1
-1
@@ -230,7 +230,7 @@ func (r *dhtQueryRunner) queryPeer(p peer.Peer) {
230
231
// make sure we're connected to the peer.
232
// (Incidentally, this will add it to the peerstore too)
233
- err := r.query.dialer.DialPeer(p)
233
+ err := r.query.dialer.DialPeer(r.ctx, p)
234
if err != nil {
235
log.Debugf("ERROR worker for: %v -- err connecting: %v", p, err)
236
r.Lock()
routing/dht/routing.go
+3
-3
@@ -145,7 +145,7 @@ func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int
145
log.Error(err)
146
return
147
}
148
- dht.addPeerListAsync(key, pmes.GetProviderPeers(), ps, count, peerOut)
148
+ dht.addPeerListAsync(ctx, key, pmes.GetProviderPeers(), ps, count, peerOut)
149
}(pp)
150
}
151
wg.Wait()
@@ -154,13 +154,13 @@ func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int
154
return peerOut
155
}
156
157
-func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*pb.Message_Peer, ps *peerSet, count int, out chan peer.Peer) {
157
+func (dht *IpfsDHT) addPeerListAsync(ctx context.Context, k u.Key, peers []*pb.Message_Peer, ps *peerSet, count int, out chan peer.Peer) {
158
done := make(chan struct{})
159
for _, pbp := range peers {
160
go func(mp *pb.Message_Peer) {
161
defer func() { done <- struct{}{} }()
162
// construct new peer
163
- p, err := dht.ensureConnectedToPeer(mp)
163
+ p, err := dht.ensureConnectedToPeer(ctx, mp)
164
if err != nil {
165
log.Error("%s", err)
166
return