@cryptotaxi247 / kubo / commits / cc5c181ae

Dialer for dht

dht doesn't need the whole network interface, only needs a Dialer. (much reduced surface of possible errors)

Juan Batiz-Benet committed Oct 21, 2014 at 03:13 UTC cc5c181ae018d2ee1df4a0bf4cc22f10e6850417
4 files changed +22 -17
net/interface.go
+9
@@ -48,3 +48,12 @@ type Handler srv.Handler
48
49 // Service interface for network resources.
50 type Service srv.Service
51 +
52 +// Dialer service that can dial to peers
53 +// (this is usually just a Network, but other services may not need the whole
54 +// thing, and thus it becomes easier to mock)
55 +type Dialer interface {
56 +
57 + // DialPeer attempts to establish a connection to a given peer
58 + DialPeer(peer.Peer) error
59 +}
routing/dht/dht.go
+7 -7
@@ -33,9 +33,9 @@ type IpfsDHT struct {
33 // NOTE: (currently, only a single table is used)
34 routingTables []*kb.RoutingTable
35
36 - // the network interface. service
37 - network inet.Network
38 - sender inet.Sender
36 + // the network services we need
37 + dialer inet.Dialer
38 + sender inet.Sender
39
40 // Local peer (yourself)
41 self peer.Peer
@@ -59,9 +59,9 @@ type IpfsDHT struct {
59 }
60
61 // NewDHT creates a new DHT object with the given peer as the 'local' host
62 -func NewDHT(ctx context.Context, p peer.Peer, ps peer.Peerstore, net inet.Network, sender inet.Sender, dstore ds.Datastore) *IpfsDHT {
62 +func NewDHT(ctx context.Context, p peer.Peer, ps peer.Peerstore, dialer inet.Dialer, sender inet.Sender, dstore ds.Datastore) *IpfsDHT {
63 dht := new(IpfsDHT)
64 - dht.network = net
64 + dht.dialer = dialer
65 dht.sender = sender
66 dht.datastore = dstore
67 dht.self = p
@@ -95,7 +95,7 @@ func (dht *IpfsDHT) Connect(ctx context.Context, npeer peer.Peer) (peer.Peer, er
95 //
96 // /ip4/10.20.30.40/tcp/1234/ipfs/Qxhxxchxzcncxnzcnxzcxzm
97 //
98 - err := dht.network.DialPeer(npeer)
98 + err := dht.dialer.DialPeer(npeer)
99 if err != nil {
100 return nil, err
101 }
@@ -499,7 +499,7 @@ func (dht *IpfsDHT) ensureConnectedToPeer(pbp *Message_Peer) (peer.Peer, error)
499 }
500
501 // dial connection
502 - err = dht.network.DialPeer(p)
502 + err = dht.dialer.DialPeer(p)
503 return p, err
504 }
505
routing/dht/query.go
+3 -7
@@ -3,6 +3,7 @@ package dht
3 import (
4 "sync"
5
6 + inet "github.com/jbenet/go-ipfs/net"
7 peer "github.com/jbenet/go-ipfs/peer"
8 queue "github.com/jbenet/go-ipfs/peer/queue"
9 kb "github.com/jbenet/go-ipfs/routing/kbucket"
@@ -14,17 +15,12 @@ import (
15
16 var maxQueryConcurrency = AlphaValue
17
17 -type dhtDialer interface {
18 - // DialPeer attempts to establish a connection to a given peer
19 - DialPeer(peer.Peer) error
20 -}
21 -
18 type dhtQuery struct {
19 // the key we're querying for
20 key u.Key
21
22 // dialer used to ensure we're connected to peers
27 - dialer dhtDialer
23 + dialer inet.Dialer
24
25 // the function to execute per peer
26 qfunc queryFunc
@@ -42,7 +38,7 @@ type dhtQueryResult struct {
38 }
39
40 // constructs query
45 -func newQuery(k u.Key, d dhtDialer, f queryFunc) *dhtQuery {
41 +func newQuery(k u.Key, d inet.Dialer, f queryFunc) *dhtQuery {
42 return &dhtQuery{
43 key: k,
44 dialer: d,
routing/dht/routing.go
+3 -3
@@ -29,7 +29,7 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
29 peers = append(peers, npeers...)
30 }
31
32 - query := newQuery(key, dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
32 + query := newQuery(key, dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
33 log.Debug("%s PutValue qry part %v", dht.self, p)
34 err := dht.putValueToNetwork(ctx, p, string(key), value)
35 if err != nil {
@@ -65,7 +65,7 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
65 }
66
67 // setup the Query
68 - query := newQuery(key, dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
68 + query := newQuery(key, dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
69
70 val, peers, err := dht.getValueOrPeers(ctx, p, key, routeLevel)
71 if err != nil {
@@ -230,7 +230,7 @@ func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (peer.Peer
230 }
231
232 // setup query function
233 - query := newQuery(u.Key(id), dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
233 + query := newQuery(u.Key(id), dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
234 pmes, err := dht.findPeerSingle(ctx, p, id, routeLevel)
235 if err != nil {
236 log.Error("%s getPeer error: %v", dht.self, err)