@cryptotaxi247 / kubo / commits / bd9fc2b78

fix bug in routing table lookups

Jeromy committed Aug 6, 2014 at 10:02 UTC bd9fc2b782192ff37c7d9e15d3ecc90a55ba4c66
10 files changed +81 -31
routing/dht/dht.go
+13 -1
@@ -106,6 +106,14 @@ func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) (*peer.Peer, error) {
106 if removed != nil {
107 panic("need to remove this peer.")
108 }
109 +
110 + // Ping new peer to register in their routing table
111 + // NOTE: this should be done better...
112 + err = dht.Ping(peer, time.Second * 2)
113 + if err != nil {
114 + panic("Failed to ping new peer.")
115 + }
116 +
117 return peer, nil
118 }
119
@@ -149,7 +157,7 @@ func (dht *IpfsDHT) handleMessages() {
157 }
158 //
159
152 - u.DOut("Got message type: '%s' [id = %x]", mesNames[pmes.GetType()], pmes.GetId())
160 + u.DOut("Got message type: '%s' [id = %x]", DHTMessage_MessageType_name[int32(pmes.GetType())], pmes.GetId())
161 switch pmes.GetType() {
162 case DHTMessage_GET_VALUE:
163 dht.handleGetValue(mes.Peer, pmes)
@@ -215,14 +223,18 @@ func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *DHTMessage) {
223 }
224
225 func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *DHTMessage) {
226 + u.POut("handleFindPeer: searching for '%s'", peer.ID(pmes.GetKey()).Pretty())
227 closest := dht.routes.NearestPeer(convertKey(u.Key(pmes.GetKey())))
228 if closest == nil {
229 + panic("could not find anything.")
230 }
231
232 if len(closest.Addresses) == 0 {
233 panic("no addresses for connected peer...")
234 }
235
236 + u.POut("handleFindPeer: sending back '%s'", closest.ID.Pretty())
237 +
238 addr,err := closest.Addresses[0].String()
239 if err != nil {
240 panic(err)
routing/dht/dht_test.go
+2 -2
@@ -45,7 +45,7 @@ func TestPing(t *testing.T) {
45 dht_a.Start()
46 dht_b.Start()
47
48 - err = dht_a.Connect(addr_b)
48 + _,err = dht_a.Connect(addr_b)
49 if err != nil {
50 t.Fatal(err)
51 }
@@ -92,7 +92,7 @@ func TestValueGetSet(t *testing.T) {
92 dht_a.Start()
93 dht_b.Start()
94
95 - err = dht_a.Connect(addr_b)
95 + _,err = dht_a.Connect(addr_b)
96 if err != nil {
97 t.Fatal(err)
98 }
routing/dht/messages.pb.go
+9 -8
@@ -29,6 +29,7 @@ const (
29 DHTMessage_GET_PROVIDERS DHTMessage_MessageType = 3
30 DHTMessage_FIND_NODE DHTMessage_MessageType = 4
31 DHTMessage_PING DHTMessage_MessageType = 5
32 + DHTMessage_DIAGNOSTIC DHTMessage_MessageType = 6
33 )
34
35 var DHTMessage_MessageType_name = map[int32]string{
@@ -38,6 +39,7 @@ var DHTMessage_MessageType_name = map[int32]string{
39 3: "GET_PROVIDERS",
40 4: "FIND_NODE",
41 5: "PING",
42 + 6: "DIAGNOSTIC",
43 }
44 var DHTMessage_MessageType_value = map[string]int32{
45 "PUT_VALUE": 0,
@@ -46,6 +48,7 @@ var DHTMessage_MessageType_value = map[string]int32{
48 "GET_PROVIDERS": 3,
49 "FIND_NODE": 4,
50 "PING": 5,
51 + "DIAGNOSTIC": 6,
52 }
53
54 func (x DHTMessage_MessageType) Enum() *DHTMessage_MessageType {
@@ -66,14 +69,12 @@ func (x *DHTMessage_MessageType) UnmarshalJSON(data []byte) error {
69 }
70
71 type DHTMessage struct {
69 - Type *DHTMessage_MessageType `protobuf:"varint,1,req,name=type,enum=dht.DHTMessage_MessageType" json:"type,omitempty"`
70 - Key *string `protobuf:"bytes,2,opt,name=key" json:"key,omitempty"`
71 - Value []byte `protobuf:"bytes,3,opt,name=value" json:"value,omitempty"`
72 - // Unique ID of this message, used to match queries with responses
73 - Id *uint64 `protobuf:"varint,4,req,name=id" json:"id,omitempty"`
74 - // Signals whether or not this message is a response to another message
75 - Response *bool `protobuf:"varint,5,opt,name=response" json:"response,omitempty"`
76 - XXX_unrecognized []byte `json:"-"`
72 + Type *DHTMessage_MessageType `protobuf:"varint,1,req,name=type,enum=dht.DHTMessage_MessageType" json:"type,omitempty"`
73 + Key *string `protobuf:"bytes,2,opt,name=key" json:"key,omitempty"`
74 + Value []byte `protobuf:"bytes,3,opt,name=value" json:"value,omitempty"`
75 + Id *uint64 `protobuf:"varint,4,req,name=id" json:"id,omitempty"`
76 + Response *bool `protobuf:"varint,5,opt,name=response" json:"response,omitempty"`
77 + XXX_unrecognized []byte `json:"-"`
78 }
79
80 func (m *DHTMessage) Reset() { *m = DHTMessage{} }
routing/dht/messages.proto
+1
@@ -10,6 +10,7 @@ message DHTMessage {
10 GET_PROVIDERS = 3;
11 FIND_NODE = 4;
12 PING = 5;
13 + DIAGNOSTIC = 6;
14 }
15
16 required MessageType type = 1;
routing/dht/pDHTMessage.go
-11
@@ -9,17 +9,6 @@ type pDHTMessage struct {
9 Id uint64
10 }
11
12 -var mesNames [10]string
13 -
14 -func init() {
15 - mesNames[DHTMessage_ADD_PROVIDER] = "add provider"
16 - mesNames[DHTMessage_FIND_NODE] = "find node"
17 - mesNames[DHTMessage_GET_PROVIDERS] = "get providers"
18 - mesNames[DHTMessage_GET_VALUE] = "get value"
19 - mesNames[DHTMessage_PUT_VALUE] = "put value"
20 - mesNames[DHTMessage_PING] = "ping"
21 -}
22 -
12 func (m *pDHTMessage) ToProtobuf() *DHTMessage {
13 pmes := new(DHTMessage)
14 if m.Value != nil {
routing/dht/routing.go
+12 -1
@@ -194,6 +194,17 @@ func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error
194 return nil, err
195 }
196
197 - return s.Connect(maddr)
197 + found_peer, err := s.Connect(maddr)
198 + if err != nil {
199 + u.POut("Found peer but couldnt connect.")
200 + return nil, err
201 + }
202 +
203 + if !found_peer.ID.Equal(id) {
204 + u.POut("FindPeer: searching for '%s' but found '%s'", id.Pretty(), found_peer.ID.Pretty())
205 + return found_peer, u.ErrSearchIncomplete
206 + }
207 +
208 + return found_peer, nil
209 }
210 }
routing/dht/table.go
+19 -5
@@ -1,10 +1,12 @@
1 package dht
2
3 import (
4 + "encoding/hex"
5 "container/list"
6 "sort"
7
8 peer "github.com/jbenet/go-ipfs/peer"
9 + u "github.com/jbenet/go-ipfs/util"
10 )
11
12 // RoutingTable defines the routing table.
@@ -87,13 +89,13 @@ func (p peerSorterArr) Less(a, b int) bool {
89 }
90 //
91
90 -func (rt *RoutingTable) copyPeersFromList(peerArr peerSorterArr, peerList *list.List) peerSorterArr {
92 +func copyPeersFromList(target ID, peerArr peerSorterArr, peerList *list.List) peerSorterArr {
93 for e := peerList.Front(); e != nil; e = e.Next() {
94 p := e.Value.(*peer.Peer)
95 p_id := convertPeerID(p.ID)
96 pd := peerDistance{
97 p: p,
96 - distance: xor(rt.local, p_id),
98 + distance: xor(target, p_id),
99 }
100 peerArr = append(peerArr, &pd)
101 }
@@ -112,6 +114,7 @@ func (rt *RoutingTable) NearestPeer(id ID) *peer.Peer {
114
115 // Returns a list of the 'count' closest peers to the given ID
116 func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
117 + u.POut("Searching table, size = %d", rt.Size())
118 cpl := xor(id, rt.local).commonPrefixLen()
119
120 // Get bucket at cpl index or last bucket
@@ -127,16 +130,16 @@ func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
130 // if this happens, search both surrounding buckets for nearest peer
131 if cpl > 0 {
132 plist := (*list.List)(rt.Buckets[cpl - 1])
130 - peerArr = rt.copyPeersFromList(peerArr, plist)
133 + peerArr = copyPeersFromList(id, peerArr, plist)
134 }
135
136 if cpl < len(rt.Buckets) - 1 {
137 plist := (*list.List)(rt.Buckets[cpl + 1])
135 - peerArr = rt.copyPeersFromList(peerArr, plist)
138 + peerArr = copyPeersFromList(id, peerArr, plist)
139 }
140 } else {
141 plist := (*list.List)(bucket)
139 - peerArr = rt.copyPeersFromList(peerArr, plist)
142 + peerArr = copyPeersFromList(id, peerArr, plist)
143 }
144
145 // Sort by distance to local peer
@@ -145,7 +148,18 @@ func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
148 var out []*peer.Peer
149 for i := 0; i < count && i < peerArr.Len(); i++ {
150 out = append(out, peerArr[i].p)
151 + u.POut("peer out: %s - %s", peerArr[i].p.ID.Pretty(),
152 + hex.EncodeToString(xor(id, convertPeerID(peerArr[i].p.ID))))
153 }
154
155 return out
156 }
157 +
158 +// Returns the total number of peers in the routing table
159 +func (rt *RoutingTable) Size() int {
160 + var tot int
161 + for _,buck := range rt.Buckets {
162 + tot += buck.Len()
163 + }
164 + return tot
165 +}
routing/dht/table_test.go
+17
@@ -90,3 +90,20 @@ func TestTableUpdate(t *testing.T) {
90 }
91 }
92 }
93 +
94 +func TestTableFind(t *testing.T) {
95 + local := _randPeer()
96 + rt := NewRoutingTable(10, convertPeerID(local.ID))
97 +
98 + peers := make([]*peer.Peer, 100)
99 + for i := 0; i < 5; i++ {
100 + peers[i] = _randPeer()
101 + rt.Update(peers[i])
102 + }
103 +
104 + t.Logf("Searching for peer: '%s'", peers[2].ID.Pretty())
105 + found := rt.NearestPeer(convertPeerID(peers[2].ID))
106 + if !found.ID.Equal(peers[2].ID) {
107 + t.Fatalf("Failed to lookup known node...")
108 + }
109 +}
swarm/swarm.go
+1 -1
@@ -245,7 +245,7 @@ func (s *Swarm) fanOut() {
245 if !ok {
246 return
247 }
248 - u.DOut("fanOut: outgoing message for: '%s'", msg.Peer.Key().Pretty())
248 + //u.DOut("fanOut: outgoing message for: '%s'", msg.Peer.Key().Pretty())
249
250 s.connsLock.RLock()
251 conn, found := s.conns[msg.Peer.Key()]
util/util.go
+7 -2
@@ -2,6 +2,7 @@ package util
2
3 import (
4 "fmt"
5 + "errors"
6 mh "github.com/jbenet/go-multihash"
7 "os"
8 "os/user"
@@ -13,10 +14,14 @@ import (
14 var Debug bool
15
16 // ErrNotImplemented signifies a function has not been implemented yet.
16 -var ErrNotImplemented = fmt.Errorf("Error: not implemented yet.")
17 +var ErrNotImplemented = errors.New("Error: not implemented yet.")
18
19 // ErrTimeout implies that a timeout has been triggered
19 -var ErrTimeout = fmt.Errorf("Error: Call timed out.")
20 +var ErrTimeout = errors.New("Error: Call timed out.")
21 +
22 +// ErrSeErrSearchIncomplete implies that a search type operation didnt
23 +// find the expected node, but did find 'a' node.
24 +var ErrSearchIncomplete = errors.New("Error: Search Incomplete.")
25
26 // Key is a string representation of multihash for use with maps.
27 type Key string