@cryptotaxi247 / kubo / commits / c5c0e7e8f

dht: changed msgs, include multiple addrs + conn type

See https://github.com/jbenet/go-ipfs/issues/153#issuecomment-63350535

Juan Batiz-Benet committed Nov 20, 2014 at 10:46 UTC c5c0e7e8f35cf090be39012e769974390371385d
6 files changed +124 -28
routing/dht/dht.go
+5 -2
@@ -517,11 +517,14 @@ func (dht *IpfsDHT) peerFromInfo(pbp *pb.Message_Peer) (peer.Peer, error) {
517 return nil, err
518 }
519
520 - maddr, err := pbp.Address()
520 + // add addresses we've just discovered
521 + maddrs, err := pbp.Addresses()
522 if err != nil {
523 return nil, err
524 }
524 - p.AddAddress(maddr)
525 + for _, maddr := range maddrs {
526 + p.AddAddress(maddr)
527 + }
528 return p, nil
529 }
530
routing/dht/handlers.go
+6 -4
@@ -210,14 +210,16 @@ func (dht *IpfsDHT) handleAddProvider(ctx context.Context, p peer.Peer, pmes *pb
210 pid := peer.ID(pb.GetId())
211 if pid.Equal(p.ID()) {
212
213 - addr, err := pb.Address()
213 + maddrs, err := pb.Addresses()
214 if err != nil {
215 - log.Errorf("provider %s error with address %s", p, *pb.Addr)
215 + log.Errorf("provider %s error with addresses %s", p, pb.Addrs)
216 continue
217 }
218
219 - log.Infof("received provider %s %s for %s", p, addr, key)
220 - p.AddAddress(addr)
219 + log.Infof("received provider %s %s for %s", p, maddrs, key)
220 + for _, maddr := range maddrs {
221 + p.AddAddress(maddr)
222 + }
223 dht.providers.AddProvider(key, p)
224
225 } else {
routing/dht/pb/dht.pb.go
+66 -8
@@ -15,10 +15,12 @@ It has these top-level messages:
15 package dht_pb
16
17 import proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/gogoprotobuf/proto"
18 +import json "encoding/json"
19 import math "math"
20
20 -// Reference imports to suppress errors if they are not otherwise used.
21 +// Reference proto, json, and math imports to suppress error if they are not otherwise used.
22 var _ = proto.Marshal
23 +var _ = &json.SyntaxError{}
24 var _ = math.Inf
25
26 type Message_MessageType int32
@@ -66,6 +68,50 @@ func (x *Message_MessageType) UnmarshalJSON(data []byte) error {
68 return nil
69 }
70
71 +type Message_ConnectionType int32
72 +
73 +const (
74 + // sender does not have a connection to peer, and no extra information (default)
75 + Message_NOT_CONNECTED Message_ConnectionType = 0
76 + // sender has a live connection to peer
77 + Message_CONNECTED Message_ConnectionType = 1
78 + // sender recently connected to peer
79 + Message_CAN_CONNECT Message_ConnectionType = 2
80 + // sender recently tried to connect to peer repeatedly but failed to connect
81 + // ("try" here is loose, but this should signal "made strong effort, failed")
82 + Message_CANNOT_CONNECT Message_ConnectionType = 3
83 +)
84 +
85 +var Message_ConnectionType_name = map[int32]string{
86 + 0: "NOT_CONNECTED",
87 + 1: "CONNECTED",
88 + 2: "CAN_CONNECT",
89 + 3: "CANNOT_CONNECT",
90 +}
91 +var Message_ConnectionType_value = map[string]int32{
92 + "NOT_CONNECTED": 0,
93 + "CONNECTED": 1,
94 + "CAN_CONNECT": 2,
95 + "CANNOT_CONNECT": 3,
96 +}
97 +
98 +func (x Message_ConnectionType) Enum() *Message_ConnectionType {
99 + p := new(Message_ConnectionType)
100 + *p = x
101 + return p
102 +}
103 +func (x Message_ConnectionType) String() string {
104 + return proto.EnumName(Message_ConnectionType_name, int32(x))
105 +}
106 +func (x *Message_ConnectionType) UnmarshalJSON(data []byte) error {
107 + value, err := proto.UnmarshalJSONEnum(Message_ConnectionType_value, data, "Message_ConnectionType")
108 + if err != nil {
109 + return err
110 + }
111 + *x = Message_ConnectionType(value)
112 + return nil
113 +}
114 +
115 type Message struct {
116 // defines what type of message it is.
117 Type *Message_MessageType `protobuf:"varint,1,opt,name=type,enum=dht.pb.Message_MessageType" json:"type,omitempty"`
@@ -133,9 +179,13 @@ func (m *Message) GetProviderPeers() []*Message_Peer {
179 }
180
181 type Message_Peer struct {
136 - Id *string `protobuf:"bytes,1,opt,name=id" json:"id,omitempty"`
137 - Addr *string `protobuf:"bytes,2,opt,name=addr" json:"addr,omitempty"`
138 - XXX_unrecognized []byte `json:"-"`
182 + // ID of a given peer.
183 + Id *string `protobuf:"bytes,1,opt,name=id" json:"id,omitempty"`
184 + // multiaddrs for a given peer
185 + Addrs []string `protobuf:"bytes,2,rep,name=addrs" json:"addrs,omitempty"`
186 + // used to signal the sender's connection capabilities to the peer
187 + Connection *Message_ConnectionType `protobuf:"varint,3,opt,name=connection,enum=dht.pb.Message_ConnectionType" json:"connection,omitempty"`
188 + XXX_unrecognized []byte `json:"-"`
189 }
190
191 func (m *Message_Peer) Reset() { *m = Message_Peer{} }
@@ -149,11 +199,18 @@ func (m *Message_Peer) GetId() string {
199 return ""
200 }
201
152 -func (m *Message_Peer) GetAddr() string {
153 - if m != nil && m.Addr != nil {
154 - return *m.Addr
202 +func (m *Message_Peer) GetAddrs() []string {
203 + if m != nil {
204 + return m.Addrs
205 }
156 - return ""
206 + return nil
207 +}
208 +
209 +func (m *Message_Peer) GetConnection() Message_ConnectionType {
210 + if m != nil && m.Connection != nil {
211 + return *m.Connection
212 + }
213 + return Message_NOT_CONNECTED
214 }
215
216 // Record represents a dht record that contains a value
@@ -204,4 +261,5 @@ func (m *Record) GetSignature() []byte {
261
262 func init() {
263 proto.RegisterEnum("dht.pb.Message_MessageType", Message_MessageType_name, Message_MessageType_value)
264 + proto.RegisterEnum("dht.pb.Message_ConnectionType", Message_ConnectionType_name, Message_ConnectionType_value)
265 }
routing/dht/pb/dht.proto
+22 -1
@@ -12,9 +12,30 @@ message Message {
12 PING = 5;
13 }
14
15 + enum ConnectionType {
16 + // sender does not have a connection to peer, and no extra information (default)
17 + NOT_CONNECTED = 0;
18 +
19 + // sender has a live connection to peer
20 + CONNECTED = 1;
21 +
22 + // sender recently connected to peer
23 + CAN_CONNECT = 2;
24 +
25 + // sender recently tried to connect to peer repeatedly but failed to connect
26 + // ("try" here is loose, but this should signal "made strong effort, failed")
27 + CANNOT_CONNECT = 3;
28 + }
29 +
30 message Peer {
31 + // ID of a given peer.
32 optional string id = 1;
17 - optional string addr = 2;
33 +
34 + // multiaddrs for a given peer
35 + repeated string addrs = 2;
36 +
37 + // used to signal the sender's connection capabilities to the peer
38 + optional ConnectionType connection = 3;
39 }
40
41 // defines what type of message it is.
routing/dht/pb/message.go
+17 -10
@@ -3,7 +3,6 @@ package dht_pb
3 import (
4 "errors"
5
6 - "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
6 ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
7 peer "github.com/jbenet/go-ipfs/peer"
8 )
@@ -19,12 +18,11 @@ func NewMessage(typ Message_MessageType, key string, level int) *Message {
18
19 func peerToPBPeer(p peer.Peer) *Message_Peer {
20 pbp := new(Message_Peer)
22 - addrs := p.Addresses()
23 - if len(addrs) == 0 || addrs[0] == nil {
24 - pbp.Addr = proto.String("")
25 - } else {
26 - addr := addrs[0].String()
27 - pbp.Addr = &addr
21 +
22 + maddrs := p.Addresses()
23 + pbp.Addrs = make([]string, len(maddrs))
24 + for i, maddr := range maddrs {
25 + pbp.Addrs[i] = maddr.String()
26 }
27 pid := string(p.ID())
28 pbp.Id = &pid
@@ -41,12 +39,21 @@ func PeersToPBPeers(peers []peer.Peer) []*Message_Peer {
39 return pbpeers
40 }
41
44 -// Address returns a multiaddr associated with the Message_Peer entry
45 -func (m *Message_Peer) Address() (ma.Multiaddr, error) {
42 +// Addresses returns a multiaddr associated with the Message_Peer entry
43 +func (m *Message_Peer) Addresses() ([]ma.Multiaddr, error) {
44 if m == nil {
45 return nil, errors.New("MessagePeer is nil")
46 }
49 - return ma.NewMultiaddr(*m.Addr)
47 +
48 + var err error
49 + maddrs := make([]ma.Multiaddr, len(m.Addrs))
50 + for i, addr := range m.Addrs {
51 + maddrs[i], err = ma.NewMultiaddr(addr)
52 + if err != nil {
53 + return nil, err
54 + }
55 + }
56 + return maddrs, nil
57 }
58
59 // GetClusterLevel gets and adjusts the cluster level on the message.
routing/dht/routing.go
+8 -3
@@ -241,12 +241,17 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (peer.Peer, error)
241 log.Warningf("Received invalid peer from query: %v", err)
242 continue
243 }
244 - ma, err := pbp.Address()
244 +
245 + // add addresses
246 + maddrs, err := pbp.Addresses()
247 if err != nil {
246 - log.Warning("Received peer with bad or missing address.")
248 + log.Warning("Received peer with bad or missing addresses: %s", pbp.Addrs)
249 continue
250 }
249 - np.AddAddress(ma)
251 + for _, maddr := range maddrs {
252 + np.AddAddress(maddr)
253 + }
254 +
255 if pbp.GetId() == string(id) {
256 return &dhtQueryResult{
257 peer: np,