dht/pb: changed PeersToPBPeers to set ConnectionType
Uses an inet.Dialer
Juan Batiz-Benet committed
Nov 21, 2014 at 03:45 UTC
ff1e672d5a69de87f3d8683f4821041d474aded0
5 files changed
+64
-13
net/interface.go
+3
-3
@@ -88,12 +88,12 @@ const (
88
NotConnected Connectedness = 0
89
90
// Connected means has an open, live connection to peer
91
- Connected
91
+ Connected = 1
92
93
// CanConnect means recently connected to peer, terminated gracefully
94
- CanConnect
94
+ CanConnect = 2
95
96
// CannotConnect means recently attempted connecting but failed to connect.
97
// (should signal "made effort, failed")
98
- CannotConnect
98
+ CannotConnect = 3
99
)
routing/dht/dht.go
+1
-1
@@ -227,7 +227,7 @@ func (dht *IpfsDHT) putProvider(ctx context.Context, p peer.Peer, key string) er
227
pmes := pb.NewMessage(pb.Message_ADD_PROVIDER, string(key), 0)
228
229
// add self as the provider
230
- pmes.ProviderPeers = pb.PeersToPBPeers([]peer.Peer{dht.self})
230
+ pmes.ProviderPeers = pb.PeersToPBPeers(dht.dialer, []peer.Peer{dht.self})
231
232
rpmes, err := dht.sendRequest(ctx, p, pmes)
233
if err != nil {
routing/dht/ext_test.go
+2
-2
@@ -262,7 +262,7 @@ func TestNotFound(t *testing.T) {
262
for i := 0; i < 7; i++ {
263
peers = append(peers, _randPeer())
264
}
265
- resp.CloserPeers = pb.PeersToPBPeers(peers)
265
+ resp.CloserPeers = pb.PeersToPBPeers(d.dialer, peers)
266
mes, err := msg.FromObject(mes.Peer(), resp)
267
if err != nil {
268
t.Error(err)
@@ -326,7 +326,7 @@ func TestLessThanKResponses(t *testing.T) {
326
case pb.Message_GET_VALUE:
327
resp := &pb.Message{
328
Type: pmes.Type,
329
- CloserPeers: pb.PeersToPBPeers([]peer.Peer{other}),
329
+ CloserPeers: pb.PeersToPBPeers(d.dialer, []peer.Peer{other}),
330
}
331
332
mes, err := msg.FromObject(mes.Peer(), resp)
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(provs)
90
+ resp.ProviderPeers = pb.PeersToPBPeers(dht.dialer, 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(closer)
102
+ resp.CloserPeers = pb.PeersToPBPeers(dht.dialer, closer)
103
}
104
105
return resp, nil
@@ -159,7 +159,7 @@ func (dht *IpfsDHT) handleFindPeer(ctx context.Context, p peer.Peer, pmes *pb.Me
159
for _, p := range withAddresses {
160
log.Debugf("handleFindPeer: sending back '%s'", p)
161
}
162
- resp.CloserPeers = pb.PeersToPBPeers(withAddresses)
162
+ resp.CloserPeers = pb.PeersToPBPeers(dht.dialer, withAddresses)
163
return resp, nil
164
}
165
@@ -183,13 +183,13 @@ func (dht *IpfsDHT) handleGetProviders(ctx context.Context, p peer.Peer, pmes *p
183
184
// if we've got providers, send thos those.
185
if providers != nil && len(providers) > 0 {
186
- resp.ProviderPeers = pb.PeersToPBPeers(providers)
186
+ resp.ProviderPeers = pb.PeersToPBPeers(dht.dialer, providers)
187
}
188
189
// Also send closer peers.
190
closer := dht.betterPeersToQuery(pmes, CloserPeerCount)
191
if closer != nil {
192
- resp.CloserPeers = pb.PeersToPBPeers(closer)
192
+ resp.CloserPeers = pb.PeersToPBPeers(dht.dialer, closer)
193
}
194
195
return resp, nil
routing/dht/pb/message.go
+53
-2
@@ -4,9 +4,12 @@ import (
4
"errors"
5
6
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
7
+
8
+ inet "github.com/jbenet/go-ipfs/net"
9
peer "github.com/jbenet/go-ipfs/peer"
10
)
11
12
+// NewMessage constructs a new dht message with given type, key, and level
13
func NewMessage(typ Message_MessageType, key string, level int) *Message {
14
m := &Message{
15
Type: &typ,
@@ -29,9 +32,9 @@ func peerToPBPeer(p peer.Peer) *Message_Peer {
32
return pbp
33
}
34
32
-// PeersToPBPeers converts a slice of Peers into a slice of *Message_Peers,
35
+// RawPeersToPBPeers converts a slice of Peers into a slice of *Message_Peers,
36
// ready to go out on the wire.
34
-func PeersToPBPeers(peers []peer.Peer) []*Message_Peer {
37
+func RawPeersToPBPeers(peers []peer.Peer) []*Message_Peer {
38
pbpeers := make([]*Message_Peer, len(peers))
39
for i, p := range peers {
40
pbpeers[i] = peerToPBPeer(p)
@@ -39,6 +42,19 @@ func PeersToPBPeers(peers []peer.Peer) []*Message_Peer {
42
return pbpeers
43
}
44
45
+// PeersToPBPeers converts given []peer.Peer into a set of []*Message_Peer,
46
+// which can be written to a message and sent out. the key thing this function
47
+// does (in addition to PeersToPBPeers) is set the ConnectionType with
48
+// information from the given inet.Dialer.
49
+func PeersToPBPeers(d inet.Dialer, peers []peer.Peer) []*Message_Peer {
50
+ pbps := RawPeersToPBPeers(peers)
51
+ for i, pbp := range pbps {
52
+ c := ConnectionType(d.Connectedness(peers[i]))
53
+ pbp.Connection = &c
54
+ }
55
+ return pbps
56
+}
57
+
58
// Addresses returns a multiaddr associated with the Message_Peer entry
59
func (m *Message_Peer) Addresses() ([]ma.Multiaddr, error) {
60
if m == nil {
@@ -75,6 +91,7 @@ func (m *Message) SetClusterLevel(level int) {
91
m.ClusterLevelRaw = &lvl
92
}
93
94
+// Loggable turns a Message into machine-readable log output
95
func (m *Message) Loggable() map[string]interface{} {
96
return map[string]interface{}{
97
"message": map[string]string{
@@ -82,3 +99,37 @@ func (m *Message) Loggable() map[string]interface{} {
99
},
100
}
101
}
102
+
103
+// ConnectionType returns a Message_ConnectionType associated with the
104
+// inet.Connectedness.
105
+func ConnectionType(c inet.Connectedness) Message_ConnectionType {
106
+ switch c {
107
+ default:
108
+ return Message_NOT_CONNECTED
109
+ case inet.NotConnected:
110
+ return Message_NOT_CONNECTED
111
+ case inet.Connected:
112
+ return Message_CONNECTED
113
+ case inet.CanConnect:
114
+ return Message_CAN_CONNECT
115
+ case inet.CannotConnect:
116
+ return Message_CANNOT_CONNECT
117
+ }
118
+}
119
+
120
+// Connectedness returns an inet.Connectedness associated with the
121
+// Message_ConnectionType.
122
+func Connectedness(c Message_ConnectionType) inet.Connectedness {
123
+ switch c {
124
+ default:
125
+ return inet.NotConnected
126
+ case Message_NOT_CONNECTED:
127
+ return inet.NotConnected
128
+ case Message_CONNECTED:
129
+ return inet.Connected
130
+ case Message_CAN_CONNECT:
131
+ return inet.CanConnect
132
+ case Message_CANNOT_CONNECT:
133
+ return inet.CannotConnect
134
+ }
135
+}