Integrated new network into ipfs
Juan Batiz-Benet committed
Dec 16, 2014 at 08:55 UTC
41751b49384721606110b0000482b03898aa7fdb
20 files changed
+230
-301
core/bootstrap.go
+3
-2
@@ -48,10 +48,11 @@ func bootstrap(ctx context.Context,
48
ps peer.Peerstore,
49
boots []*config.BootstrapPeer) error {
50
51
- if len(n.GetConnections()) >= recoveryThreshold {
51
+ connectedPeers := n.Peers()
52
+ if len(connectedPeers) >= recoveryThreshold {
53
return nil
54
}
54
- numCxnsToCreate := recoveryThreshold - len(n.GetConnections())
55
+ numCxnsToCreate := recoveryThreshold - len(connectedPeers)
56
57
var bootstrapPeers []peer.Peer
58
for _, bootstrap := range boots {
core/commands/swarm.go
+1
-1
@@ -55,7 +55,7 @@ ipfs swarm peers lists the set of peers this node is connected to.
55
return nil, errNotOnline
56
}
57
58
- conns := n.Network.GetConnections()
58
+ conns := n.Network.Conns()
59
addrs := make([]string, len(conns))
60
for i, c := range conns {
61
pid := c.RemotePeer().ID()
core/core.go
+10
-25
@@ -6,6 +6,7 @@ import (
6
7
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
b58 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-base58"
9
+ ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
10
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
11
12
bstore "github.com/jbenet/go-ipfs/blocks/blockstore"
@@ -21,14 +22,11 @@ import (
22
namesys "github.com/jbenet/go-ipfs/namesys"
23
inet "github.com/jbenet/go-ipfs/net"
24
handshake "github.com/jbenet/go-ipfs/net/handshake"
24
- mux "github.com/jbenet/go-ipfs/net/mux"
25
- netservice "github.com/jbenet/go-ipfs/net/service"
25
path "github.com/jbenet/go-ipfs/path"
26
peer "github.com/jbenet/go-ipfs/peer"
27
pin "github.com/jbenet/go-ipfs/pin"
28
routing "github.com/jbenet/go-ipfs/routing"
29
dht "github.com/jbenet/go-ipfs/routing/dht"
31
- ctxc "github.com/jbenet/go-ipfs/util/ctxcloser"
30
ds2 "github.com/jbenet/go-ipfs/util/datastore2"
31
debugerror "github.com/jbenet/go-ipfs/util/debugerror"
32
eventlog "github.com/jbenet/go-ipfs/util/eventlog"
@@ -63,7 +61,7 @@ type IpfsNode struct {
61
Namesys namesys.NameSystem // the name system, resolves paths to hashes
62
Diagnostics *diag.Diagnostics // the diagnostics service
63
66
- ctxc.ContextCloser
64
+ ctxgroup.ContextGroup
65
}
66
67
// Mounts defines what the node's mount state is. This should
@@ -91,8 +89,8 @@ func NewIpfsNode(ctx context.Context, cfg *config.Config, online bool) (n *IpfsN
89
onlineMode: online,
90
Config: cfg,
91
}
94
- n.ContextCloser = ctxc.NewContextCloser(ctx, n.teardown)
95
- ctx = n.ContextCloser.Context()
92
+ n.ContextGroup = ctxgroup.WithContextAndTeardown(ctx, n.teardown)
93
+ ctx = n.ContextGroup.Context()
94
95
// setup datastore.
96
if n.Datastore, err = makeDatastore(cfg.Datastore); err != nil {
@@ -112,45 +110,32 @@ func NewIpfsNode(ctx context.Context, cfg *config.Config, online bool) (n *IpfsN
110
// setup online services
111
if online {
112
115
- dhtService := netservice.NewService(ctx, nil) // nil handler for now, need to patch it
116
- exchangeService := netservice.NewService(ctx, nil) // nil handler for now, need to patch it
117
- diagService := netservice.NewService(ctx, nil) // nil handler for now, need to patch it
118
-
119
- muxMap := &mux.ProtocolMap{
120
- mux.ProtocolID_Routing: dhtService,
121
- mux.ProtocolID_Exchange: exchangeService,
122
- mux.ProtocolID_Diagnostic: diagService,
123
- // add protocol services here.
124
- }
125
-
113
// setup the network
114
listenAddrs, err := listenAddresses(cfg)
115
if err != nil {
116
return nil, debugerror.Wrap(err)
117
}
118
132
- n.Network, err = inet.NewIpfsNetwork(ctx, listenAddrs, n.Identity, n.Peerstore, muxMap)
119
+ n.Network, err = inet.NewNetwork(ctx, listenAddrs, n.Identity, n.Peerstore)
120
if err != nil {
121
return nil, debugerror.Wrap(err)
122
}
136
- n.AddCloserChild(n.Network)
123
+ n.AddChildGroup(n.Network.CtxGroup())
124
125
// setup diagnostics service
139
- n.Diagnostics = diag.NewDiagnostics(n.Identity, n.Network, diagService)
140
- diagService.SetHandler(n.Diagnostics)
126
+ n.Diagnostics = diag.NewDiagnostics(n.Identity, n.Network)
127
128
// setup routing service
143
- dhtRouting := dht.NewDHT(ctx, n.Identity, n.Peerstore, n.Network, dhtService, n.Datastore)
129
+ dhtRouting := dht.NewDHT(ctx, n.Identity, n.Peerstore, n.Network, n.Datastore)
130
dhtRouting.Validators[IpnsValidatorTag] = namesys.ValidateIpnsRecord
131
132
// TODO(brian): perform this inside NewDHT factory method
147
- dhtService.SetHandler(dhtRouting) // wire the handler to the service.
133
n.Routing = dhtRouting
149
- n.AddCloserChild(dhtRouting)
134
+ n.AddChildGroup(dhtRouting)
135
136
// setup exchange service
137
const alwaysSendToPeer = true // use YesManStrategy
153
- bitswapNetwork := bsnet.NewFromIpfsNetwork(exchangeService, n.Network)
138
+ bitswapNetwork := bsnet.NewFromIpfsNetwork(n.Network)
139
140
n.Exchange = bitswap.New(ctx, n.Identity, bitswapNetwork, n.Routing, blockstore, alwaysSendToPeer)
141
diagnostics/diag.go
+40
-38
@@ -13,12 +13,12 @@ import (
13
14
"crypto/rand"
15
16
+ ggio "code.google.com/p/gogoprotobuf/io"
17
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
18
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
19
20
pb "github.com/jbenet/go-ipfs/diagnostics/internal/pb"
21
net "github.com/jbenet/go-ipfs/net"
21
- msg "github.com/jbenet/go-ipfs/net/message"
22
peer "github.com/jbenet/go-ipfs/peer"
23
util "github.com/jbenet/go-ipfs/util"
24
)
@@ -31,7 +31,6 @@ const ResponseTimeout = time.Second * 10
31
// requests
32
type Diagnostics struct {
33
network net.Network
34
- sender net.Sender
34
self peer.Peer
35
36
diagLock sync.Mutex
@@ -40,14 +39,16 @@ type Diagnostics struct {
39
}
40
41
// NewDiagnostics instantiates a new diagnostics service running on the given network
43
-func NewDiagnostics(self peer.Peer, inet net.Network, sender net.Sender) *Diagnostics {
44
- return &Diagnostics{
42
+func NewDiagnostics(self peer.Peer, inet net.Network) *Diagnostics {
43
+ d := &Diagnostics{
44
network: inet,
46
- sender: sender,
45
self: self,
48
- diagMap: make(map[string]time.Time),
46
birth: time.Now(),
47
+ diagMap: make(map[string]time.Time),
48
}
49
+
50
+ inet.SetHandler(net.ProtocolDiag, d.handleNewStream)
51
+ return d
52
}
53
54
type connDiagInfo struct {
@@ -91,7 +92,7 @@ func (di *DiagInfo) Marshal() []byte {
92
}
93
94
func (d *Diagnostics) getPeers() []peer.Peer {
94
- return d.network.GetPeerList()
95
+ return d.network.Peers()
96
}
97
98
func (d *Diagnostics) getDiagInfo() *DiagInfo {
@@ -100,7 +101,7 @@ func (d *Diagnostics) getDiagInfo() *DiagInfo {
101
di.ID = d.self.ID().Pretty()
102
di.LifeSpan = time.Since(d.birth)
103
di.Keys = nil // Currently no way to query datastore
103
- di.BwIn, di.BwOut = d.network.GetBandwidthTotals()
104
+ di.BwIn, di.BwOut = d.network.BandwidthTotals()
105
106
for _, p := range d.getPeers() {
107
d := connDiagInfo{p.GetLatency(), p.ID().Pretty()}
@@ -110,7 +111,7 @@ func (d *Diagnostics) getDiagInfo() *DiagInfo {
111
}
112
113
func newID() string {
113
- id := make([]byte, 4)
114
+ id := make([]byte, 16)
115
rand.Read(id)
116
return string(id)
117
}
@@ -196,29 +197,31 @@ func newMessage(diagID string) *pb.Message {
197
198
func (d *Diagnostics) sendRequest(ctx context.Context, p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
199
199
- mes, err := msg.FromObject(p, pmes)
200
+ s, err := d.network.NewStream(net.ProtocolDiag, p)
201
if err != nil {
202
return nil, err
203
}
204
+ defer s.Close()
205
+
206
+ r := ggio.NewDelimitedReader(s, net.MessageSizeMax)
207
+ w := ggio.NewDelimitedWriter(s)
208
209
start := time.Now()
210
206
- rmes, err := d.sender.SendRequest(ctx, mes)
207
- if err != nil {
211
+ if err := w.WriteMsg(pmes); err != nil {
212
+ return nil, err
213
+ }
214
+
215
+ var rpmes *pb.Message
216
+ if err := r.ReadMsg(rpmes); err != nil {
217
return nil, err
218
}
210
- if rmes == nil {
219
+ if rpmes == nil {
220
return nil, errors.New("no response to request")
221
}
222
223
rtt := time.Since(start)
224
log.Infof("diagnostic request took: %s", rtt.String())
216
-
217
- rpmes := new(pb.Message)
218
- if err := proto.Unmarshal(rmes.Data(), rpmes); err != nil {
219
- return nil, err
220
- }
221
-
225
return rpmes, nil
226
}
227
@@ -271,33 +274,25 @@ func (d *Diagnostics) handleDiagnostic(p peer.Peer, pmes *pb.Message) (*pb.Messa
274
return resp, nil
275
}
276
274
-func (d *Diagnostics) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.NetMessage {
275
- mData := mes.Data()
276
- if mData == nil {
277
- log.Error("message did not include Data")
278
- return nil
279
- }
277
+func (d *Diagnostics) HandleMessage(ctx context.Context, s net.Stream) error {
278
281
- mPeer := mes.Peer()
282
- if mPeer == nil {
283
- log.Error("message did not include a Peer")
284
- return nil
285
- }
279
+ r := ggio.NewDelimitedReader(s, 32768) // maxsize
280
+ w := ggio.NewDelimitedWriter(s)
281
282
// deserialize msg
283
pmes := new(pb.Message)
289
- err := proto.Unmarshal(mData, pmes)
290
- if err != nil {
284
+ if err := r.ReadMsg(pmes); err != nil {
285
log.Errorf("Failed to decode protobuf message: %v", err)
286
return nil
287
}
288
289
// Print out diagnostic
290
log.Infof("[peer: %s] Got message from [%s]\n",
297
- d.self.ID().Pretty(), mPeer.ID().Pretty())
291
+ d.self.ID().Pretty(), s.Conn().RemotePeer().ID().Pretty())
292
293
// dispatch handler.
300
- rpmes, err := d.handleDiagnostic(mPeer, pmes)
294
+ p := s.Conn().RemotePeer()
295
+ rpmes, err := d.handleDiagnostic(p, pmes)
296
if err != nil {
297
log.Errorf("handleDiagnostic error: %s", err)
298
return nil
@@ -308,12 +303,19 @@ func (d *Diagnostics) HandleMessage(ctx context.Context, mes msg.NetMessage) msg
303
return nil
304
}
305
311
- // serialize response msg
312
- rmes, err := msg.FromObject(mPeer, rpmes)
313
- if err != nil {
306
+ // serialize + send response msg
307
+ if err := w.WriteMsg(rpmes); err != nil {
308
log.Errorf("Failed to encode protobuf message: %v", err)
309
return nil
310
}
311
318
- return rmes
312
+ return nil
313
+}
314
+
315
+func (d *Diagnostics) handleNewStream(s net.Stream) {
316
+
317
+ go func() {
318
+ d.HandleMessage(context.Background(), s)
319
+ }()
320
+
321
}
exchange/bitswap/message/message.go
+18
-9
@@ -1,13 +1,14 @@
1
package message
2
3
import (
4
- proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
4
+ "io"
5
+
6
blocks "github.com/jbenet/go-ipfs/blocks"
7
pb "github.com/jbenet/go-ipfs/exchange/bitswap/message/internal/pb"
7
- netmsg "github.com/jbenet/go-ipfs/net/message"
8
- nm "github.com/jbenet/go-ipfs/net/message"
9
- peer "github.com/jbenet/go-ipfs/peer"
8
+ inet "github.com/jbenet/go-ipfs/net"
9
u "github.com/jbenet/go-ipfs/util"
10
+
11
+ ggio "code.google.com/p/gogoprotobuf/io"
12
)
13
14
// TODO move message.go into the bitswap package
@@ -38,7 +39,7 @@ type BitSwapMessage interface {
39
40
type Exportable interface {
41
ToProto() *pb.Message
41
- ToNet(p peer.Peer) (nm.NetMessage, error)
42
+ ToNet(w io.Writer) error
43
}
44
45
type impl struct {
@@ -92,11 +93,14 @@ func (m *impl) AddBlock(b *blocks.Block) {
93
m.blocks[b.Key()] = b
94
}
95
95
-func FromNet(nmsg netmsg.NetMessage) (BitSwapMessage, error) {
96
+func FromNet(r io.Reader) (BitSwapMessage, error) {
97
+ pbr := ggio.NewDelimitedReader(r, inet.MessageSizeMax)
98
+
99
pb := new(pb.Message)
97
- if err := proto.Unmarshal(nmsg.Data(), pb); err != nil {
100
+ if err := pbr.ReadMsg(pb); err != nil {
101
return nil, err
102
}
103
+
104
m := newMessageFromProto(*pb)
105
return m, nil
106
}
@@ -112,6 +116,11 @@ func (m *impl) ToProto() *pb.Message {
116
return pb
117
}
118
115
-func (m *impl) ToNet(p peer.Peer) (nm.NetMessage, error) {
116
- return nm.FromObject(p, m.ToProto())
119
+func (m *impl) ToNet(w io.Writer) error {
120
+ pbw := ggio.NewDelimitedWriter(w)
121
+
122
+ if err := pbw.WriteMsg(m.ToProto()); err != nil {
123
+ return err
124
+ }
125
+ return nil
126
}
exchange/bitswap/message/message_test.go
+6
-21
@@ -7,7 +7,6 @@ import (
7
blocks "github.com/jbenet/go-ipfs/blocks"
8
pb "github.com/jbenet/go-ipfs/exchange/bitswap/message/internal/pb"
9
u "github.com/jbenet/go-ipfs/util"
10
- testutil "github.com/jbenet/go-ipfs/util/testutil"
10
)
11
12
func TestAppendWanted(t *testing.T) {
@@ -87,18 +86,6 @@ func TestCopyProtoByValue(t *testing.T) {
86
}
87
}
88
90
-func TestToNetMethodSetsPeer(t *testing.T) {
91
- m := New()
92
- p := testutil.NewPeerWithIDString("X")
93
- netmsg, err := m.ToNet(p)
94
- if err != nil {
95
- t.Fatal(err)
96
- }
97
- if !(netmsg.Peer().Key() == p.Key()) {
98
- t.Fatal("Peer key is different")
99
- }
100
-}
101
-
89
func TestToNetFromNetPreservesWantList(t *testing.T) {
90
original := New()
91
original.AddWanted(u.Key("M"))
@@ -107,13 +94,12 @@ func TestToNetFromNetPreservesWantList(t *testing.T) {
94
original.AddWanted(u.Key("T"))
95
original.AddWanted(u.Key("F"))
96
110
- p := testutil.NewPeerWithIDString("X")
111
- netmsg, err := original.ToNet(p)
112
- if err != nil {
97
+ var buf bytes.Buffer
98
+ if err := original.ToNet(&buf); err != nil {
99
t.Fatal(err)
100
}
101
116
- copied, err := FromNet(netmsg)
102
+ copied, err := FromNet(&buf)
103
if err != nil {
104
t.Fatal(err)
105
}
@@ -138,13 +124,12 @@ func TestToAndFromNetMessage(t *testing.T) {
124
original.AddBlock(blocks.NewBlock([]byte("F")))
125
original.AddBlock(blocks.NewBlock([]byte("M")))
126
141
- p := testutil.NewPeerWithIDString("X")
142
- netmsg, err := original.ToNet(p)
143
- if err != nil {
127
+ var buf bytes.Buffer
128
+ if err := original.ToNet(&buf); err != nil {
129
t.Fatal(err)
130
}
131
147
- m2, err := FromNet(netmsg)
132
+ m2, err := FromNet(&buf)
133
if err != nil {
134
t.Fatal(err)
135
}
exchange/bitswap/network/ipfs_impl.go
+31
-26
@@ -5,7 +5,6 @@ import (
5
6
bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
7
inet "github.com/jbenet/go-ipfs/net"
8
- netmsg "github.com/jbenet/go-ipfs/net/message"
8
peer "github.com/jbenet/go-ipfs/peer"
9
util "github.com/jbenet/go-ipfs/util"
10
)
@@ -14,46 +13,48 @@ var log = util.Logger("bitswap_network")
13
14
// NewFromIpfsNetwork returns a BitSwapNetwork supported by underlying IPFS
15
// Dialer & Service
17
-func NewFromIpfsNetwork(s inet.Service, dialer inet.Dialer) BitSwapNetwork {
16
+func NewFromIpfsNetwork(n inet.Network) BitSwapNetwork {
17
bitswapNetwork := impl{
19
- service: s,
20
- dialer: dialer,
18
+ network: n,
19
}
22
- s.SetHandler(&bitswapNetwork)
20
+ n.SetHandler(inet.ProtocolBitswap, bitswapNetwork.handleNewStream)
21
return &bitswapNetwork
22
}
23
24
// impl transforms the ipfs network interface, which sends and receives
25
// NetMessage objects, into the bitswap network interface.
26
type impl struct {
29
- service inet.Service
30
- dialer inet.Dialer
27
+ network inet.Network
28
29
// inbound messages from the network are forwarded to the receiver
30
receiver Receiver
31
}
32
36
-// HandleMessage marshals and unmarshals net messages, forwarding them to the
37
-// BitSwapMessage receiver
38
-func (bsnet *impl) HandleMessage(
39
- ctx context.Context, incoming netmsg.NetMessage) netmsg.NetMessage {
33
+// handleNewStream receives a new stream from the network.
34
+func (bsnet *impl) handleNewStream(s inet.Stream) {
35
36
if bsnet.receiver == nil {
42
- return nil
37
+ return
38
}
39
45
- received, err := bsmsg.FromNet(incoming)
46
- if err != nil {
47
- go bsnet.receiver.ReceiveError(err)
48
- return nil
49
- }
40
+ go func() {
41
+ defer s.Close()
42
+
43
+ received, err := bsmsg.FromNet(s)
44
+ if err != nil {
45
+ go bsnet.receiver.ReceiveError(err)
46
+ return
47
+ }
48
+
49
+ p := s.Conn().RemotePeer()
50
+ ctx := context.Background()
51
+ bsnet.receiver.ReceiveMessage(ctx, p, received)
52
+ }()
53
51
- bsnet.receiver.ReceiveMessage(ctx, incoming.Peer(), received)
52
- return nil
54
}
55
56
func (bsnet *impl) DialPeer(ctx context.Context, p peer.Peer) error {
56
- return bsnet.dialer.DialPeer(ctx, p)
57
+ return bsnet.network.DialPeer(ctx, p)
58
}
59
60
func (bsnet *impl) SendMessage(
@@ -61,11 +62,13 @@ func (bsnet *impl) SendMessage(
62
p peer.Peer,
63
outgoing bsmsg.BitSwapMessage) error {
64
64
- nmsg, err := outgoing.ToNet(p)
65
+ s, err := bsnet.network.NewStream(inet.ProtocolBitswap, p)
66
if err != nil {
67
return err
68
}
68
- return bsnet.service.SendMessage(ctx, nmsg)
69
+ defer s.Close()
70
+
71
+ return outgoing.ToNet(s)
72
}
73
74
func (bsnet *impl) SendRequest(
@@ -73,15 +76,17 @@ func (bsnet *impl) SendRequest(
76
p peer.Peer,
77
outgoing bsmsg.BitSwapMessage) (bsmsg.BitSwapMessage, error) {
78
76
- outgoingMsg, err := outgoing.ToNet(p)
79
+ s, err := bsnet.network.NewStream(inet.ProtocolBitswap, p)
80
if err != nil {
81
return nil, err
82
}
80
- incomingMsg, err := bsnet.service.SendRequest(ctx, outgoingMsg)
81
- if err != nil {
83
+ defer s.Close()
84
+
85
+ if err := outgoing.ToNet(s); err != nil {
86
return nil, err
87
}
84
- return bsmsg.FromNet(incomingMsg)
88
+
89
+ return bsmsg.FromNet(s)
90
}
91
92
func (bsnet *impl) SetDelegate(r Receiver) {
fuse/ipns/mount_unix.go
+2
-2
@@ -33,9 +33,9 @@ func Mount(ipfs *core.IpfsNode, fpath string, ipfspath string) (mount.Mount, err
33
// assume it worked...
34
}
35
36
- // bind the mount (ContextCloser) to the node, so that when the node exits
36
+ // bind the mount (ContextGroup) to the node, so that when the node exits
37
// the fsclosers are automatically closed.
38
- ipfs.AddCloserChild(m)
38
+ ipfs.AddChildGroup(m)
39
return m, nil
40
}
41
fuse/mount/mount.go
+6
-8
@@ -6,9 +6,9 @@ import (
6
"time"
7
8
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9
+ ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
10
11
u "github.com/jbenet/go-ipfs/util"
11
- ctxc "github.com/jbenet/go-ipfs/util/ctxcloser"
12
)
13
14
var log = u.Logger("mount")
@@ -25,7 +25,7 @@ type Mount interface {
25
// Unmount calls Close.
26
Unmount() error
27
28
- ctxc.ContextCloser
28
+ ctxgroup.ContextGroup
29
}
30
31
// UnmountFunc is a function used to Unmount a mount
@@ -39,12 +39,12 @@ type MountFunc func(Mount) error
39
// in an UnmountFunc to perform the unmounting logic.
40
func New(ctx context.Context, mountpoint string) Mount {
41
m := &mount{mpoint: mountpoint}
42
- m.ContextCloser = ctxc.NewContextCloser(ctx, m.persistentUnmount)
42
+ m.ContextGroup = ctxgroup.WithContextAndTeardown(ctx, m.persistentUnmount)
43
return m
44
}
45
46
type mount struct {
47
- ctxc.ContextCloser
47
+ ctxgroup.ContextGroup
48
49
unmount UnmountFunc
50
mpoint string
@@ -80,15 +80,13 @@ func (m *mount) Unmount() error {
80
}
81
82
func (m *mount) Mount(mount MountFunc, unmount UnmountFunc) {
83
- m.Children().Add(1)
83
m.unmount = unmount
84
85
// go serve the mount
87
- go func() {
86
+ m.ContextGroup.AddChildFunc(func(parent ctxgroup.ContextGroup) {
87
if err := mount(m); err != nil {
88
log.Error("%s mount: %s", m.MountPoint(), err)
89
}
91
- m.Children().Done()
90
m.Unmount()
93
- }()
91
+ })
92
}
fuse/readonly/readonly_unix.go
+2
-2
@@ -179,9 +179,9 @@ func Mount(ipfs *core.IpfsNode, fpath string) (mount.Mount, error) {
179
// assume it worked...
180
}
181
182
- // bind the mount (ContextCloser) to the node, so that when the node exits
182
+ // bind the mount (ContextGroup) to the node, so that when the node exits
183
// the fsclosers are automatically closed.
184
- ipfs.AddCloserChild(m)
184
+ ipfs.AddChildGroup(m)
185
return m, nil
186
}
187
net/interface.go
+21
-5
@@ -4,21 +4,31 @@ import (
4
"io"
5
6
conn "github.com/jbenet/go-ipfs/net/conn"
7
- swarm "github.com/jbenet/go-ipfs/net/swarm2"
7
+ // swarm "github.com/jbenet/go-ipfs/net/swarm2"
8
peer "github.com/jbenet/go-ipfs/peer"
9
10
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
11
+ ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
12
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
13
)
14
15
+// ProtocolID is an identifier used to write protocol headers in streams
16
type ProtocolID string
17
18
+// These are the ProtocolIDs of the protocols running. It is useful
19
+// to keep them in one place.
20
const (
21
ProtocolBitswap ProtocolID = "/ipfs/bitswap"
22
ProtocolDHT ProtocolID = "/ipfs/dht"
23
ProtocolDiag ProtocolID = "/ipfs/diagnostics"
24
)
25
26
+// MessageSizeMax is a soft (recommended) maximum for network messages.
27
+// One can write more, as the interface is a stream. But it is useful
28
+// to bunch it up into multiple read/writes when the whole message is
29
+// a single, large serialized object.
30
+const MessageSizeMax = 2 << 22 // 4MB
31
+
32
// Stream represents a bidirectional channel between two agents in
33
// the IPFS network. "agent" is as granular as desired, potentially
34
// being a "request -> reply" pair, or whole protocols.
@@ -45,8 +55,8 @@ type StreamHandlerMap map[ProtocolID]StreamHandler
55
type Conn interface {
56
conn.PeerConn
57
48
- // NewStream constructs a new Stream directly connected to p.
49
- NewStream(pr ProtocolID, p peer.Peer) (Stream, error)
58
+ // NewStreamWithProtocol constructs a new Stream directly connected to p.
59
+ NewStreamWithProtocol(pr ProtocolID, p peer.Peer) (Stream, error)
60
}
61
62
// Network is the interface IPFS uses for connecting to the world.
@@ -66,8 +76,11 @@ type Network interface {
76
// If ProtocolID is "", writes no header.
77
NewStream(ProtocolID, peer.Peer) (Stream, error)
78
69
- // Swarm returns the connection Swarm
70
- Swarm() *swarm.Swarm
79
+ // Peers returns the peers connected
80
+ Peers() []peer.Peer
81
+
82
+ // Conns returns the connections in this Netowrk
83
+ Conns() []Conn
84
85
// BandwidthTotals returns the total number of bytes passed through
86
// the network since it was instantiated
@@ -80,6 +93,9 @@ type Network interface {
93
// listens. It expands "any interface" addresses (/ip4/0.0.0.0, /ip6/::) to
94
// use the known local interfaces.
95
InterfaceListenAddresses() ([]ma.Multiaddr, error)
96
+
97
+ // CtxGroup returns the network's contextGroup
98
+ CtxGroup() ctxgroup.ContextGroup
99
}
100
101
// Dialer represents a service that can dial out to peers
net/net.go
+39
-3
@@ -44,7 +44,7 @@ func (c *conn_) SwarmConn() *swarm.Conn {
44
return (*swarm.Conn)(c)
45
}
46
47
-func (c *conn_) NewStream(pr ProtocolID, p peer.Peer) (Stream, error) {
47
+func (c *conn_) NewStreamWithProtocol(pr ProtocolID, p peer.Peer) (Stream, error) {
48
s, err := (*swarm.Conn)(c).NewStream()
49
if err != nil {
50
return nil, err
@@ -91,7 +91,7 @@ type network struct {
91
92
// NewNetwork constructs a new network and starts listening on given addresses.
93
func NewNetwork(ctx context.Context, listen []ma.Multiaddr, local peer.Peer,
94
- peers peer.Peerstore) (*network, error) {
94
+ peers peer.Peerstore) (Network, error) {
95
96
s, err := swarm.NewSwarm(ctx, listen, local, peers)
97
if err != nil {
@@ -109,6 +109,7 @@ func NewNetwork(ctx context.Context, listen []ma.Multiaddr, local peer.Peer,
109
n.mux.Handle((*stream)(s))
110
})
111
112
+ n.cg.SetTeardown(n.close)
113
n.cg.AddChildGroup(s.CtxGroup())
114
return n, nil
115
}
@@ -123,16 +124,51 @@ func (n *network) DialPeer(ctx context.Context, p peer.Peer) error {
124
return err
125
}
126
127
+// CtxGroup returns the network's ContextGroup
128
+func (n *network) CtxGroup() ctxgroup.ContextGroup {
129
+ return n.cg
130
+}
131
+
132
+// Swarm returns the network's peerstream.Swarm
133
+func (n *network) Swarm() *swarm.Swarm {
134
+ return n.Swarm()
135
+}
136
+
137
// LocalPeer the network's LocalPeer
138
func (n *network) LocalPeer() peer.Peer {
139
return n.swarm.LocalPeer()
140
}
141
142
+// Peers returns the connected peers
143
+func (n *network) Peers() []peer.Peer {
144
+ return n.swarm.Peers()
145
+}
146
+
147
+// Conns returns the connected peers
148
+func (n *network) Conns() []Conn {
149
+ conns1 := n.swarm.Connections()
150
+ out := make([]Conn, len(conns1))
151
+ for i, c := range conns1 {
152
+ out[i] = (*conn_)(c)
153
+ }
154
+ return out
155
+}
156
+
157
// ClosePeer connection to peer
158
func (n *network) ClosePeer(p peer.Peer) error {
159
return n.swarm.CloseConnection(p)
160
}
161
162
+// close is the real teardown function
163
+func (n *network) close() error {
164
+ return n.swarm.Close()
165
+}
166
+
167
+// Close calls the ContextCloser func
168
+func (n *network) Close() error {
169
+ return n.cg.Close()
170
+}
171
+
172
// BandwidthTotals returns the total amount of bandwidth transferred
173
func (n *network) BandwidthTotals() (in uint64, out uint64) {
174
// need to implement this. probably best to do it in swarm this time.
@@ -165,7 +201,7 @@ func (n *network) Connectedness(p peer.Peer) Connectedness {
201
// NewStream returns a new stream to given peer p.
202
// If there is no connection to p, attempts to create one.
203
// If ProtocolID is "", writes no header.
168
-func (c *network) NewStreamWithPeer(pr ProtocolID, p peer.Peer) (Stream, error) {
204
+func (c *network) NewStream(pr ProtocolID, p peer.Peer) (Stream, error) {
205
s, err := c.swarm.NewStreamWithPeer(p)
206
if err != nil {
207
return nil, err
routing/dht/dht.go
+21
-123
@@ -11,19 +11,17 @@ import (
11
"time"
12
13
inet "github.com/jbenet/go-ipfs/net"
14
- msg "github.com/jbenet/go-ipfs/net/message"
14
peer "github.com/jbenet/go-ipfs/peer"
15
routing "github.com/jbenet/go-ipfs/routing"
16
pb "github.com/jbenet/go-ipfs/routing/dht/pb"
17
kb "github.com/jbenet/go-ipfs/routing/kbucket"
18
u "github.com/jbenet/go-ipfs/util"
20
- ctxc "github.com/jbenet/go-ipfs/util/ctxcloser"
19
"github.com/jbenet/go-ipfs/util/eventlog"
20
21
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
24
- ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
25
-
22
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
23
+ ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
24
+ ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
25
)
26
27
var log = eventlog.Logger("dht")
@@ -35,50 +33,37 @@ const doPinging = false
33
// IpfsDHT is an implementation of Kademlia with Coral and S/Kademlia modifications.
34
// It is used to implement the base IpfsRouting module.
35
type IpfsDHT struct {
38
- // Array of routing tables for differently distanced nodes
39
- // NOTE: (currently, only a single table is used)
40
- routingTable *kb.RoutingTable
41
-
42
- // the network services we need
43
- dialer inet.Dialer
44
- sender inet.Sender
36
+ network inet.Network // the network services we need
37
+ self peer.Peer // Local peer (yourself)
38
+ peerstore peer.Peerstore // Other peers
39
46
- // Local peer (yourself)
47
- self peer.Peer
48
-
49
- // Other peers
50
- peerstore peer.Peerstore
51
-
52
- // Local data
53
- datastore ds.Datastore
40
+ datastore ds.Datastore // Local data
41
dslock sync.Mutex
42
56
- providers *ProviderManager
57
-
58
- // When this peer started up
59
- birth time.Time
43
+ routingTable *kb.RoutingTable // Array of routing tables for differently distanced nodes
44
+ providers *ProviderManager
45
61
- //lock to make diagnostics work better
62
- diaglock sync.Mutex
46
+ birth time.Time // When this peer started up
47
+ diaglock sync.Mutex // lock to make diagnostics work better
48
49
// record validator funcs
50
Validators map[string]ValidatorFunc
51
67
- ctxc.ContextCloser
52
+ ctxgroup.ContextGroup
53
}
54
55
// NewDHT creates a new DHT object with the given peer as the 'local' host
71
-func NewDHT(ctx context.Context, p peer.Peer, ps peer.Peerstore, dialer inet.Dialer, sender inet.Sender, dstore ds.Datastore) *IpfsDHT {
56
+func NewDHT(ctx context.Context, p peer.Peer, ps peer.Peerstore, n inet.Network, dstore ds.Datastore) *IpfsDHT {
57
dht := new(IpfsDHT)
73
- dht.dialer = dialer
74
- dht.sender = sender
58
dht.datastore = dstore
59
dht.self = p
60
dht.peerstore = ps
78
- dht.ContextCloser = ctxc.NewContextCloser(ctx, nil)
61
+ dht.ContextGroup = ctxgroup.WithContext(ctx)
62
+ dht.network = n
63
+ n.SetHandler(inet.ProtocolDHT, dht.handleNewStream)
64
65
dht.providers = NewProviderManager(dht.Context(), p.ID())
81
- dht.AddCloserChild(dht.providers)
66
+ dht.AddChildGroup(dht.providers)
67
68
dht.routingTable = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Minute)
69
dht.birth = time.Now()
@@ -95,7 +80,7 @@ func NewDHT(ctx context.Context, p peer.Peer, ps peer.Peerstore, dialer inet.Dia
80
81
// Connect to a new peer at the given address, ping and add to the routing table
82
func (dht *IpfsDHT) Connect(ctx context.Context, npeer peer.Peer) error {
98
- err := dht.dialer.DialPeer(ctx, npeer)
83
+ err := dht.network.DialPeer(ctx, npeer)
84
if err != nil {
85
return err
86
}
@@ -113,93 +98,6 @@ func (dht *IpfsDHT) Connect(ctx context.Context, npeer peer.Peer) error {
98
return nil
99
}
100
116
-// HandleMessage implements the inet.Handler interface.
117
-func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.NetMessage {
118
-
119
- mData := mes.Data()
120
- if mData == nil {
121
- log.Error("Message contained nil data.")
122
- return nil
123
- }
124
-
125
- mPeer := mes.Peer()
126
- if mPeer == nil {
127
- log.Error("Message contained nil peer.")
128
- return nil
129
- }
130
-
131
- // deserialize msg
132
- pmes := new(pb.Message)
133
- err := proto.Unmarshal(mData, pmes)
134
- if err != nil {
135
- log.Error("Error unmarshaling data")
136
- return nil
137
- }
138
-
139
- // update the peer (on valid msgs only)
140
- dht.Update(ctx, mPeer)
141
-
142
- log.Event(ctx, "foo", dht.self, mPeer, pmes)
143
-
144
- // get handler for this msg type.
145
- handler := dht.handlerForMsgType(pmes.GetType())
146
- if handler == nil {
147
- log.Error("got back nil handler from handlerForMsgType")
148
- return nil
149
- }
150
-
151
- // dispatch handler.
152
- rpmes, err := handler(ctx, mPeer, pmes)
153
- if err != nil {
154
- log.Errorf("handle message error: %s", err)
155
- return nil
156
- }
157
-
158
- // if nil response, return it before serializing
159
- if rpmes == nil {
160
- log.Warning("Got back nil response from request.")
161
- return nil
162
- }
163
-
164
- // serialize response msg
165
- rmes, err := msg.FromObject(mPeer, rpmes)
166
- if err != nil {
167
- log.Errorf("serialze response error: %s", err)
168
- return nil
169
- }
170
-
171
- return rmes
172
-}
173
-
174
-// sendRequest sends out a request using dht.sender, but also makes sure to
175
-// measure the RTT for latency measurements.
176
-func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
177
-
178
- mes, err := msg.FromObject(p, pmes)
179
- if err != nil {
180
- return nil, err
181
- }
182
-
183
- start := time.Now()
184
-
185
- rmes, err := dht.sender.SendRequest(ctx, mes) // respect?
186
- if err != nil {
187
- return nil, err
188
- }
189
- if rmes == nil {
190
- return nil, errors.New("no response to request")
191
- }
192
- log.Event(ctx, "sentMessage", dht.self, p, pmes)
193
-
194
- rmes.Peer().SetLatency(time.Since(start))
195
-
196
- rpmes := new(pb.Message)
197
- if err := proto.Unmarshal(rmes.Data(), rpmes); err != nil {
198
- return nil, err
199
- }
200
- return rpmes, nil
201
-}
202
-
101
// putValueToNetwork stores the given key/value pair at the peer 'p'
102
func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p peer.Peer,
103
key string, rec *pb.Record) error {
@@ -224,7 +122,7 @@ func (dht *IpfsDHT) putProvider(ctx context.Context, p peer.Peer, key string) er
122
pmes := pb.NewMessage(pb.Message_ADD_PROVIDER, string(key), 0)
123
124
// add self as the provider
227
- pmes.ProviderPeers = pb.PeersToPBPeers(dht.dialer, []peer.Peer{dht.self})
125
+ pmes.ProviderPeers = pb.PeersToPBPeers(dht.network, []peer.Peer{dht.self})
126
127
rpmes, err := dht.sendRequest(ctx, p, pmes)
128
if err != nil {
@@ -478,12 +376,12 @@ func (dht *IpfsDHT) ensureConnectedToPeer(ctx context.Context, pbp *pb.Message_P
376
return nil, err
377
}
378
481
- if dht.dialer.LocalPeer().ID().Equal(p.ID()) {
379
+ if dht.self.ID().Equal(p.ID()) {
380
return nil, errors.New("attempting to ensure connection to self")
381
}
382
383
// dial connection
486
- err = dht.dialer.DialPeer(ctx, p)
384
+ err = dht.network.DialPeer(ctx, p)
385
return p, err
386
}
387
@@ -536,7 +434,7 @@ func (dht *IpfsDHT) Bootstrap(ctx context.Context) {
434
if err != nil {
435
log.Errorf("Bootstrap peer error: %s", err)
436
}
539
- err = dht.dialer.DialPeer(ctx, p)
437
+ err = dht.network.DialPeer(ctx, p)
438
if err != nil {
439
log.Errorf("Bootstrap peer error: %s", err)
440
}
routing/dht/dht_test.go
+15
-19
@@ -12,8 +12,6 @@ import (
12
13
ci "github.com/jbenet/go-ipfs/crypto"
14
inet "github.com/jbenet/go-ipfs/net"
15
- mux "github.com/jbenet/go-ipfs/net/mux"
16
- netservice "github.com/jbenet/go-ipfs/net/service"
15
peer "github.com/jbenet/go-ipfs/peer"
16
u "github.com/jbenet/go-ipfs/util"
17
testutil "github.com/jbenet/go-ipfs/util/testutil"
@@ -25,16 +23,14 @@ import (
23
func setupDHT(ctx context.Context, t *testing.T, p peer.Peer) *IpfsDHT {
24
peerstore := peer.NewPeerstore()
25
28
- dhts := netservice.NewService(ctx, nil) // nil handler for now, need to patch it
29
- net, err := inet.NewIpfsNetwork(ctx, p.Addresses(), p, peerstore, &mux.ProtocolMap{
30
- mux.ProtocolID_Routing: dhts,
31
- })
26
+ n, err := inet.NewNetwork(ctx, p.Addresses(), p, peerstore)
27
if err != nil {
28
t.Fatal(err)
29
}
30
36
- d := NewDHT(ctx, p, peerstore, net, dhts, ds.NewMapDatastore())
37
- dhts.SetHandler(d)
31
+ d := NewDHT(ctx, p, peerstore, n, ds.NewMapDatastore())
32
+ d.network.SetHandler(inet.ProtocolDHT, d.handleNewStream)
33
+
34
d.Validators["v"] = func(u.Key, []byte) error {
35
return nil
36
}
@@ -107,8 +103,8 @@ func TestPing(t *testing.T) {
103
104
defer dhtA.Close()
105
defer dhtB.Close()
110
- defer dhtA.dialer.(inet.Network).Close()
111
- defer dhtB.dialer.(inet.Network).Close()
106
+ defer dhtA.network.Close()
107
+ defer dhtB.network.Close()
108
109
err = dhtA.Connect(ctx, peerB)
110
if err != nil {
@@ -157,8 +153,8 @@ func TestValueGetSet(t *testing.T) {
153
154
defer dhtA.Close()
155
defer dhtB.Close()
160
- defer dhtA.dialer.(inet.Network).Close()
161
- defer dhtB.dialer.(inet.Network).Close()
156
+ defer dhtA.network.Close()
157
+ defer dhtB.network.Close()
158
159
err = dhtA.Connect(ctx, peerB)
160
if err != nil {
@@ -199,7 +195,7 @@ func TestProvides(t *testing.T) {
195
defer func() {
196
for i := 0; i < 4; i++ {
197
dhts[i].Close()
202
- defer dhts[i].dialer.(inet.Network).Close()
198
+ defer dhts[i].network.Close()
199
}
200
}()
201
@@ -261,7 +257,7 @@ func TestProvidesAsync(t *testing.T) {
257
defer func() {
258
for i := 0; i < 4; i++ {
259
dhts[i].Close()
264
- defer dhts[i].dialer.(inet.Network).Close()
260
+ defer dhts[i].network.Close()
261
}
262
}()
263
@@ -326,7 +322,7 @@ func TestLayeredGet(t *testing.T) {
322
defer func() {
323
for i := 0; i < 4; i++ {
324
dhts[i].Close()
329
- defer dhts[i].dialer.(inet.Network).Close()
325
+ defer dhts[i].network.Close()
326
}
327
}()
328
@@ -381,7 +377,7 @@ func TestFindPeer(t *testing.T) {
377
defer func() {
378
for i := 0; i < 4; i++ {
379
dhts[i].Close()
384
- dhts[i].dialer.(inet.Network).Close()
380
+ dhts[i].network.Close()
381
}
382
}()
383
@@ -427,7 +423,7 @@ func TestFindPeersConnectedToPeer(t *testing.T) {
423
defer func() {
424
for i := 0; i < 4; i++ {
425
dhts[i].Close()
430
- dhts[i].dialer.(inet.Network).Close()
426
+ dhts[i].network.Close()
427
}
428
}()
429
@@ -566,8 +562,8 @@ func TestConnectCollision(t *testing.T) {
562
563
dhtA.Close()
564
dhtB.Close()
569
- dhtA.dialer.(inet.Network).Close()
570
- dhtB.dialer.(inet.Network).Close()
565
+ dhtA.network.Close()
566
+ dhtB.network.Close()
567
568
<-time.After(200 * time.Millisecond)
569
}
routing/dht/ext_test.go
-2
@@ -9,8 +9,6 @@ import (
9
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
10
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
11
inet "github.com/jbenet/go-ipfs/net"
12
- msg "github.com/jbenet/go-ipfs/net/message"
13
- mux "github.com/jbenet/go-ipfs/net/mux"
12
peer "github.com/jbenet/go-ipfs/peer"
13
routing "github.com/jbenet/go-ipfs/routing"
14
pb "github.com/jbenet/go-ipfs/routing/dht/pb"
routing/dht/handlers.go
+5
-5
@@ -87,7 +87,7 @@ func (dht *IpfsDHT) handleGetValue(ctx context.Context, p peer.Peer, pmes *pb.Me
87
provs := dht.providers.GetProviders(ctx, u.Key(pmes.GetKey()))
88
if len(provs) > 0 {
89
log.Debugf("handleGetValue returning %d provider[s]", len(provs))
90
- resp.ProviderPeers = pb.PeersToPBPeers(dht.dialer, provs)
90
+ resp.ProviderPeers = pb.PeersToPBPeers(dht.network, provs)
91
}
92
93
// Find closest peer on given cluster to desired key and reply with that info
@@ -99,7 +99,7 @@ func (dht *IpfsDHT) handleGetValue(ctx context.Context, p peer.Peer, pmes *pb.Me
99
log.Critical("no addresses on peer being sent!")
100
}
101
}
102
- resp.CloserPeers = pb.PeersToPBPeers(dht.dialer, closer)
102
+ resp.CloserPeers = pb.PeersToPBPeers(dht.network, closer)
103
}
104
105
return resp, nil
@@ -160,7 +160,7 @@ func (dht *IpfsDHT) handleFindPeer(ctx context.Context, p peer.Peer, pmes *pb.Me
160
log.Debugf("handleFindPeer: sending back '%s'", p)
161
}
162
163
- resp.CloserPeers = pb.PeersToPBPeers(dht.dialer, withAddresses)
163
+ resp.CloserPeers = pb.PeersToPBPeers(dht.network, withAddresses)
164
return resp, nil
165
}
166
@@ -183,13 +183,13 @@ func (dht *IpfsDHT) handleGetProviders(ctx context.Context, p peer.Peer, pmes *p
183
}
184
185
if providers != nil && len(providers) > 0 {
186
- resp.ProviderPeers = pb.PeersToPBPeers(dht.dialer, providers)
186
+ resp.ProviderPeers = pb.PeersToPBPeers(dht.network, providers)
187
}
188
189
// Also send closer peers.
190
closer := dht.betterPeersToQuery(pmes, CloserPeerCount)
191
if closer != nil {
192
- resp.CloserPeers = pb.PeersToPBPeers(dht.dialer, closer)
192
+ resp.CloserPeers = pb.PeersToPBPeers(dht.network, closer)
193
}
194
195
return resp, nil
routing/dht/pb/message.go
+1
-1
@@ -65,7 +65,7 @@ func RawPeersToPBPeers(peers []peer.Peer) []*Message_Peer {
65
// which can be written to a message and sent out. the key thing this function
66
// does (in addition to PeersToPBPeers) is set the ConnectionType with
67
// information from the given inet.Dialer.
68
-func PeersToPBPeers(d inet.Dialer, peers []peer.Peer) []*Message_Peer {
68
+func PeersToPBPeers(d inet.Network, peers []peer.Peer) []*Message_Peer {
69
pbps := RawPeersToPBPeers(peers)
70
for i, pbp := range pbps {
71
c := ConnectionType(d.Connectedness(peers[i]))
routing/dht/providers.go
+3
-3
@@ -3,9 +3,9 @@ package dht
3
import (
4
"time"
5
6
+ ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
7
peer "github.com/jbenet/go-ipfs/peer"
8
u "github.com/jbenet/go-ipfs/util"
8
- ctxc "github.com/jbenet/go-ipfs/util/ctxcloser"
9
10
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
11
)
@@ -18,7 +18,7 @@ type ProviderManager struct {
18
newprovs chan *addProv
19
getprovs chan *getProv
20
period time.Duration
21
- ctxc.ContextCloser
21
+ ctxgroup.ContextGroup
22
}
23
24
type addProv struct {
@@ -38,7 +38,7 @@ func NewProviderManager(ctx context.Context, local peer.ID) *ProviderManager {
38
pm.providers = make(map[u.Key][]*providerInfo)
39
pm.getlocal = make(chan chan []u.Key)
40
pm.local = make(map[u.Key]struct{})
41
- pm.ContextCloser = ctxc.NewContextCloser(ctx, nil)
41
+ pm.ContextGroup = ctxgroup.WithContext(ctx)
42
43
pm.Children().Add(1)
44
go pm.run()
routing/dht/records.go
+1
-1
@@ -50,7 +50,7 @@ func (dht *IpfsDHT) getPublicKey(pid peer.ID) (ci.PubKey, error) {
50
}
51
52
log.Debug("not in peerstore, searching dht.")
53
- ctxT, _ := context.WithTimeout(dht.ContextCloser.Context(), time.Second*5)
53
+ ctxT, _ := context.WithTimeout(dht.ContextGroup.Context(), time.Second*5)
54
val, err := dht.GetValue(ctxT, u.Key("/pk/"+string(pid)))
55
if err != nil {
56
log.Warning("Failed to find requested public key.")
routing/dht/routing.go
+5
-5
@@ -40,7 +40,7 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
40
41
peers := dht.routingTable.NearestPeers(kb.ConvertKey(key), KValue)
42
43
- query := newQuery(key, dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
43
+ query := newQuery(key, dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
44
log.Debugf("%s PutValue qry part %v", dht.self, p)
45
err := dht.putValueToNetwork(ctx, p, string(key), rec)
46
if err != nil {
@@ -75,7 +75,7 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
75
}
76
77
// setup the Query
78
- query := newQuery(key, dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
78
+ query := newQuery(key, dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
79
80
val, peers, err := dht.getValueOrPeers(ctx, p, key)
81
if err != nil {
@@ -159,7 +159,7 @@ func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, co
159
}
160
161
// setup the Query
162
- query := newQuery(key, dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
162
+ query := newQuery(key, dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
163
164
pmes, err := dht.findProvidersSingle(ctx, p, key)
165
if err != nil {
@@ -262,7 +262,7 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (peer.Peer, error)
262
}
263
264
// setup the Query
265
- query := newQuery(u.Key(id), dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
265
+ query := newQuery(u.Key(id), dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
266
267
pmes, err := dht.findPeerSingle(ctx, p, id)
268
if err != nil {
@@ -316,7 +316,7 @@ func (dht *IpfsDHT) FindPeersConnectedToPeer(ctx context.Context, id peer.ID) (<
316
}
317
318
// setup the Query
319
- query := newQuery(u.Key(id), dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
319
+ query := newQuery(u.Key(id), dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
320
321
pmes, err := dht.findPeerSingle(ctx, p, id)
322
if err != nil {