moved routing table code into its own package
Jeromy committed
Aug 8, 2014 at 19:58 UTC
9f7604378c7fa1d3535186f7a70d62ace5dc6d50
7 files changed
+30
-28
routing/dht/dht.go
+7
-6
@@ -10,6 +10,7 @@ import (
10
peer "github.com/jbenet/go-ipfs/peer"
11
swarm "github.com/jbenet/go-ipfs/swarm"
12
u "github.com/jbenet/go-ipfs/util"
13
+ kb "github.com/jbenet/go-ipfs/routing/kbucket"
14
15
ma "github.com/jbenet/go-multiaddr"
16
@@ -25,7 +26,7 @@ import (
26
type IpfsDHT struct {
27
// Array of routing tables for differently distanced nodes
28
// NOTE: (currently, only a single table is used)
28
- routes []*RoutingTable
29
+ routes []*kb.RoutingTable
30
31
network *swarm.Swarm
32
@@ -84,8 +85,8 @@ func NewDHT(p *peer.Peer) (*IpfsDHT, error) {
85
dht.listeners = make(map[uint64]*listenInfo)
86
dht.providers = make(map[u.Key][]*providerInfo)
87
dht.shutdown = make(chan struct{})
87
- dht.routes = make([]*RoutingTable, 1)
88
- dht.routes[0] = NewRoutingTable(20, convertPeerID(p.ID))
88
+ dht.routes = make([]*kb.RoutingTable, 1)
89
+ dht.routes[0] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID))
90
dht.birth = time.Now()
91
return dht, nil
92
}
@@ -253,7 +254,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
254
}
255
} else if err == ds.ErrNotFound {
256
// Find closest peer(s) to desired key and reply with that info
256
- closer := dht.routes[0].NearestPeer(convertKey(u.Key(pmes.GetKey())))
257
+ closer := dht.routes[0].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
258
resp = &DHTMessage{
259
Response: true,
260
Id: *pmes.Id,
@@ -290,7 +291,7 @@ func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *PBDHTMessage) {
291
func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *PBDHTMessage) {
292
success := true
293
u.POut("handleFindPeer: searching for '%s'", peer.ID(pmes.GetKey()).Pretty())
293
- closest := dht.routes[0].NearestPeer(convertKey(u.Key(pmes.GetKey())))
294
+ closest := dht.routes[0].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
295
if closest == nil {
296
u.PErr("handleFindPeer: could not find anything.")
297
success = false
@@ -432,7 +433,7 @@ func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *PBDHTMessage) {
433
}
434
dht.diaglock.Unlock()
435
435
- seq := dht.routes[0].NearestPeers(convertPeerID(dht.self.ID), 10)
436
+ seq := dht.routes[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
437
listen_chan := dht.ListenFor(pmes.GetId(), len(seq), time.Second*30)
438
439
for _, ps := range seq {
routing/dht/diag.go
+1
-1
@@ -37,7 +37,7 @@ func (dht *IpfsDHT) getDiagInfo() *diagInfo {
37
di.LifeSpan = time.Since(dht.birth)
38
di.Keys = nil // Currently no way to query datastore
39
40
- for _,p := range dht.routes[0].listpeers() {
40
+ for _,p := range dht.routes[0].Listpeers() {
41
di.Connections = append(di.Connections, connDiagInfo{p.GetLatency(), p.ID})
42
}
43
return di
routing/dht/routing.go
+7
-6
@@ -12,6 +12,7 @@ import (
12
ma "github.com/jbenet/go-multiaddr"
13
14
peer "github.com/jbenet/go-ipfs/peer"
15
+ kb "github.com/jbenet/go-ipfs/routing/kbucket"
16
swarm "github.com/jbenet/go-ipfs/swarm"
17
u "github.com/jbenet/go-ipfs/util"
18
)
@@ -33,7 +34,7 @@ func GenerateMessageID() uint64 {
34
// This is the top level "Store" operation of the DHT
35
func (s *IpfsDHT) PutValue(key u.Key, value []byte) error {
36
var p *peer.Peer
36
- p = s.routes[0].NearestPeer(convertKey(key))
37
+ p = s.routes[0].NearestPeer(kb.ConvertKey(key))
38
if p == nil {
39
return errors.New("Table returned nil peer!")
40
}
@@ -46,7 +47,7 @@ func (s *IpfsDHT) PutValue(key u.Key, value []byte) error {
47
// returned along with util.ErrSearchIncomplete
48
func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
49
var p *peer.Peer
49
- p = s.routes[0].NearestPeer(convertKey(key))
50
+ p = s.routes[0].NearestPeer(kb.ConvertKey(key))
51
if p == nil {
52
return nil, errors.New("Table returned nil peer!")
53
}
@@ -90,7 +91,7 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
91
92
// Announce that this node can provide value for given key
93
func (s *IpfsDHT) Provide(key u.Key) error {
93
- peers := s.routes[0].NearestPeers(convertKey(key), PoolSize)
94
+ peers := s.routes[0].NearestPeers(kb.ConvertKey(key), PoolSize)
95
if len(peers) == 0 {
96
//return an error
97
}
@@ -110,7 +111,7 @@ func (s *IpfsDHT) Provide(key u.Key) error {
111
112
// FindProviders searches for peers who can provide the value for given key.
113
func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer, error) {
113
- p := s.routes[0].NearestPeer(convertKey(key))
114
+ p := s.routes[0].NearestPeer(kb.ConvertKey(key))
115
116
pmes := DHTMessage{
117
Type: PBDHTMessage_GET_PROVIDERS,
@@ -168,7 +169,7 @@ func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer,
169
170
// FindPeer searches for a peer with given ID.
171
func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error) {
171
- p := s.routes[0].NearestPeer(convertPeerID(id))
172
+ p := s.routes[0].NearestPeer(kb.ConvertPeerID(id))
173
174
pmes := DHTMessage{
175
Type: PBDHTMessage_FIND_NODE,
@@ -242,7 +243,7 @@ func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
243
func (dht *IpfsDHT) GetDiagnostic(timeout time.Duration) ([]*diagInfo, error) {
244
u.DOut("Begin Diagnostic")
245
//Send to N closest peers
245
- targets := dht.routes[0].NearestPeers(convertPeerID(dht.self.ID), 10)
246
+ targets := dht.routes[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
247
248
// TODO: Add timeout to this struct so nodes know when to return
249
pmes := DHTMessage{
routing/kbucket/bucket.go
renamed
+1
-1
@@ -48,7 +48,7 @@ func (b *Bucket) Split(cpl int, target ID) *Bucket {
48
out := list.New()
49
e := bucket_list.Front()
50
for e != nil {
51
- peer_id := convertPeerID(e.Value.(*peer.Peer).ID)
51
+ peer_id := ConvertPeerID(e.Value.(*peer.Peer).ID)
52
peer_cpl := prefLen(peer_id, target)
53
if peer_cpl > cpl {
54
cur := e
routing/kbucket/table.go
renamed
+3
-3
@@ -36,7 +36,7 @@ func NewRoutingTable(bucketsize int, local_id ID) *RoutingTable {
36
func (rt *RoutingTable) Update(p *peer.Peer) *peer.Peer {
37
rt.tabLock.Lock()
38
defer rt.tabLock.Unlock()
39
- peer_id := convertPeerID(p.ID)
39
+ peer_id := ConvertPeerID(p.ID)
40
cpl := xor(peer_id, rt.local).commonPrefixLen()
41
42
b_id := cpl
@@ -97,7 +97,7 @@ func (p peerSorterArr) Less(a, b int) bool {
97
func copyPeersFromList(target ID, peerArr peerSorterArr, peerList *list.List) peerSorterArr {
98
for e := peerList.Front(); e != nil; e = e.Next() {
99
p := e.Value.(*peer.Peer)
100
- p_id := convertPeerID(p.ID)
100
+ p_id := ConvertPeerID(p.ID)
101
pd := peerDistance{
102
p: p,
103
distance: xor(target, p_id),
@@ -173,7 +173,7 @@ func (rt *RoutingTable) Size() int {
173
}
174
175
// NOTE: This is potentially unsafe... use at your own risk
176
-func (rt *RoutingTable) listpeers() []*peer.Peer {
176
+func (rt *RoutingTable) Listpeers() []*peer.Peer {
177
var peers []*peer.Peer
178
for _,buck := range rt.Buckets {
179
for e := buck.getIter(); e != nil; e = e.Next() {
routing/kbucket/table_test.go
renamed
+9
-9
@@ -36,7 +36,7 @@ func TestBucket(t *testing.T) {
36
}
37
38
local := _randPeer()
39
- local_id := convertPeerID(local.ID)
39
+ local_id := ConvertPeerID(local.ID)
40
41
i := rand.Intn(len(peers))
42
e := b.Find(peers[i].ID)
@@ -44,10 +44,10 @@ func TestBucket(t *testing.T) {
44
t.Errorf("Failed to find peer: %v", peers[i])
45
}
46
47
- spl := b.Split(0, convertPeerID(local.ID))
47
+ spl := b.Split(0, ConvertPeerID(local.ID))
48
llist := (*list.List)(b)
49
for e := llist.Front(); e != nil; e = e.Next() {
50
- p := convertPeerID(e.Value.(*peer.Peer).ID)
50
+ p := ConvertPeerID(e.Value.(*peer.Peer).ID)
51
cpl := xor(p, local_id).commonPrefixLen()
52
if cpl > 0 {
53
t.Fatalf("Split failed. found id with cpl > 0 in 0 bucket")
@@ -56,7 +56,7 @@ func TestBucket(t *testing.T) {
56
57
rlist := (*list.List)(spl)
58
for e := rlist.Front(); e != nil; e = e.Next() {
59
- p := convertPeerID(e.Value.(*peer.Peer).ID)
59
+ p := ConvertPeerID(e.Value.(*peer.Peer).ID)
60
cpl := xor(p, local_id).commonPrefixLen()
61
if cpl == 0 {
62
t.Fatalf("Split failed. found id with cpl == 0 in non 0 bucket")
@@ -67,7 +67,7 @@ func TestBucket(t *testing.T) {
67
// Right now, this just makes sure that it doesnt hang or crash
68
func TestTableUpdate(t *testing.T) {
69
local := _randPeer()
70
- rt := NewRoutingTable(10, convertPeerID(local.ID))
70
+ rt := NewRoutingTable(10, ConvertPeerID(local.ID))
71
72
peers := make([]*peer.Peer, 100)
73
for i := 0; i < 100; i++ {
@@ -93,7 +93,7 @@ func TestTableUpdate(t *testing.T) {
93
94
func TestTableFind(t *testing.T) {
95
local := _randPeer()
96
- rt := NewRoutingTable(10, convertPeerID(local.ID))
96
+ rt := NewRoutingTable(10, ConvertPeerID(local.ID))
97
98
peers := make([]*peer.Peer, 100)
99
for i := 0; i < 5; i++ {
@@ -102,7 +102,7 @@ func TestTableFind(t *testing.T) {
102
}
103
104
t.Logf("Searching for peer: '%s'", peers[2].ID.Pretty())
105
- found := rt.NearestPeer(convertPeerID(peers[2].ID))
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
}
@@ -110,7 +110,7 @@ func TestTableFind(t *testing.T) {
110
111
func TestTableFindMultiple(t *testing.T) {
112
local := _randPeer()
113
- rt := NewRoutingTable(20, convertPeerID(local.ID))
113
+ rt := NewRoutingTable(20, ConvertPeerID(local.ID))
114
115
peers := make([]*peer.Peer, 100)
116
for i := 0; i < 18; i++ {
@@ -119,7 +119,7 @@ func TestTableFindMultiple(t *testing.T) {
119
}
120
121
t.Logf("Searching for peer: '%s'", peers[2].ID.Pretty())
122
- found := rt.NearestPeers(convertPeerID(peers[2].ID), 15)
122
+ found := rt.NearestPeers(ConvertPeerID(peers[2].ID), 15)
123
if len(found) != 15 {
124
t.Fatalf("Got back different number of peers than we expected.")
125
}
routing/kbucket/util.go
renamed
+2
-2
@@ -71,12 +71,12 @@ func equalizeSizes(a, b ID) (ID, ID) {
71
return a, b
72
}
73
74
-func convertPeerID(id peer.ID) ID {
74
+func ConvertPeerID(id peer.ID) ID {
75
hash := sha256.Sum256(id)
76
return hash[:]
77
}
78
79
-func convertKey(id u.Key) ID {
79
+func ConvertKey(id u.Key) ID {
80
hash := sha256.Sum256([]byte(id))
81
return hash[:]
82
}