Made the DHT module pass golint
Chas Leichner committed
Aug 16, 2014 at 23:03 UTC
87bfdbc599e2b055508d9b3bfd50c38ce33086a1
13 files changed
+312
-309
routing/dht/Message.go
+6
-5
@@ -4,13 +4,13 @@ import (
4
peer "github.com/jbenet/go-ipfs/peer"
5
)
6
7
-// A helper struct to make working with protbuf types easier
8
-type DHTMessage struct {
7
+// Message is a a helper struct which makes working with protbuf types easier
8
+type Message struct {
9
Type PBDHTMessage_MessageType
10
Key string
11
Value []byte
12
Response bool
13
- Id uint64
13
+ ID uint64
14
Success bool
15
Peers []*peer.Peer
16
}
@@ -28,9 +28,10 @@ func peerInfo(p *peer.Peer) *PBDHTMessage_PBPeer {
28
return pbp
29
}
30
31
+// ToProtobuf takes a Message and produces a protobuf with it.
32
// TODO: building the protobuf message this way is a little wasteful
33
// Unused fields wont be omitted, find a better way to do this
33
-func (m *DHTMessage) ToProtobuf() *PBDHTMessage {
34
+func (m *Message) ToProtobuf() *PBDHTMessage {
35
pmes := new(PBDHTMessage)
36
if m.Value != nil {
37
pmes.Value = m.Value
@@ -39,7 +40,7 @@ func (m *DHTMessage) ToProtobuf() *PBDHTMessage {
40
pmes.Type = &m.Type
41
pmes.Key = &m.Key
42
pmes.Response = &m.Response
42
- pmes.Id = &m.Id
43
+ pmes.Id = &m.ID
44
pmes.Success = &m.Success
45
for _, p := range m.Peers {
46
pmes.Peers = append(pmes.Peers, peerInfo(p))
routing/dht/dht.go
+81
-81
@@ -25,7 +25,7 @@ import (
25
type IpfsDHT struct {
26
// Array of routing tables for differently distanced nodes
27
// NOTE: (currently, only a single table is used)
28
- routes []*kb.RoutingTable
28
+ routingTables []*kb.RoutingTable
29
30
network swarm.Network
31
@@ -49,7 +49,7 @@ type IpfsDHT struct {
49
diaglock sync.Mutex
50
51
// listener is a server to register to listen for responses to messages
52
- listener *MesListener
52
+ listener *mesListener
53
}
54
55
// NewDHT creates a new DHT object with the given peer as the 'local' host
@@ -61,12 +61,11 @@ func NewDHT(p *peer.Peer, net swarm.Network) *IpfsDHT {
61
dht.providers = make(map[u.Key][]*providerInfo)
62
dht.shutdown = make(chan struct{})
63
64
- dht.routes = make([]*kb.RoutingTable, 3)
65
- dht.routes[0] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Millisecond*30)
66
- dht.routes[1] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Millisecond*100)
67
- dht.routes[2] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Hour)
68
-
69
- dht.listener = NewMesListener()
64
+ dht.routingTables = make([]*kb.RoutingTable, 3)
65
+ dht.routingTables[0] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Millisecond*30)
66
+ dht.routingTables[1] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Millisecond*100)
67
+ dht.routingTables[2] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Hour)
68
+ dht.listener = newMesListener()
69
dht.birth = time.Now()
70
return dht
71
}
@@ -175,11 +174,11 @@ func (dht *IpfsDHT) cleanExpiredProviders() {
174
}
175
176
func (dht *IpfsDHT) putValueToNetwork(p *peer.Peer, key string, value []byte) error {
178
- pmes := DHTMessage{
177
+ pmes := Message{
178
Type: PBDHTMessage_PUT_VALUE,
179
Key: key,
180
Value: value,
182
- Id: GenerateMessageID(),
181
+ ID: GenerateMessageID(),
182
}
183
184
mes := swarm.NewMessage(p, pmes.ToProtobuf())
@@ -190,9 +189,9 @@ func (dht *IpfsDHT) putValueToNetwork(p *peer.Peer, key string, value []byte) er
189
func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
190
u.DOut("handleGetValue for key: %s", pmes.GetKey())
191
dskey := ds.NewKey(pmes.GetKey())
193
- resp := &DHTMessage{
192
+ resp := &Message{
193
Response: true,
195
- Id: pmes.GetId(),
194
+ ID: pmes.GetId(),
195
Key: pmes.GetKey(),
196
}
197
iVal, err := dht.datastore.Get(dskey)
@@ -222,7 +221,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
221
}
222
u.DOut("handleGetValue searching level %d clusters", level)
223
225
- closer := dht.routes[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
224
+ closer := dht.routingTables[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
225
226
if closer.ID.Equal(dht.self.ID) {
227
u.DOut("Attempted to return self! this shouldnt happen...")
@@ -259,19 +258,19 @@ func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *PBDHTMessage) {
258
}
259
260
func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *PBDHTMessage) {
262
- resp := DHTMessage{
261
+ resp := Message{
262
Type: pmes.GetType(),
263
Response: true,
265
- Id: pmes.GetId(),
264
+ ID: pmes.GetId(),
265
}
266
267
dht.network.Send(swarm.NewMessage(p, resp.ToProtobuf()))
268
}
269
270
func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *PBDHTMessage) {
272
- resp := DHTMessage{
271
+ resp := Message{
272
Type: pmes.GetType(),
274
- Id: pmes.GetId(),
273
+ ID: pmes.GetId(),
274
Response: true,
275
}
276
defer func() {
@@ -280,7 +279,7 @@ func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *PBDHTMessage) {
279
}()
280
level := pmes.GetValue()[0]
281
u.DOut("handleFindPeer: searching for '%s'", peer.ID(pmes.GetKey()).Pretty())
283
- closest := dht.routes[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
282
+ closest := dht.routingTables[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
283
if closest == nil {
284
u.PErr("handleFindPeer: could not find anything.")
285
return
@@ -302,10 +301,10 @@ func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *PBDHTMessage) {
301
}
302
303
func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *PBDHTMessage) {
305
- resp := DHTMessage{
304
+ resp := Message{
305
Type: PBDHTMessage_GET_PROVIDERS,
306
Key: pmes.GetKey(),
308
- Id: pmes.GetId(),
307
+ ID: pmes.GetId(),
308
Response: true,
309
}
310
@@ -318,7 +317,7 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *PBDHTMessage) {
317
level = int(pmes.GetValue()[0])
318
}
319
321
- closer := dht.routes[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
320
+ closer := dht.routingTables[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
321
if kb.Closer(dht.self.ID, closer.ID, u.Key(pmes.GetKey())) {
322
resp.Peers = nil
323
} else {
@@ -346,7 +345,7 @@ func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *PBDHTMessage) {
345
dht.addProviderEntry(key, p)
346
}
347
349
-// Stop all communications from this peer and shut down
348
+// Halt stops all communications from this peer and shut down
349
func (dht *IpfsDHT) Halt() {
350
dht.shutdown <- struct{}{}
351
dht.network.Close()
@@ -362,7 +361,7 @@ func (dht *IpfsDHT) addProviderEntry(key u.Key, p *peer.Peer) {
361
362
// NOTE: not yet finished, low priority
363
func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *PBDHTMessage) {
365
- seq := dht.routes[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
364
+ seq := dht.routingTables[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
365
listenChan := dht.listener.Listen(pmes.GetId(), len(seq), time.Second*30)
366
367
for _, ps := range seq {
@@ -382,22 +381,22 @@ func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *PBDHTMessage) {
381
case <-after:
382
//Timeout, return what we have
383
goto out
385
- case req_resp := <-listenChan:
386
- pmes_out := new(PBDHTMessage)
387
- err := proto.Unmarshal(req_resp.Data, pmes_out)
384
+ case reqResp := <-listenChan:
385
+ pmesOut := new(PBDHTMessage)
386
+ err := proto.Unmarshal(reqResp.Data, pmesOut)
387
if err != nil {
388
// It broke? eh, whatever, keep going
389
continue
390
}
392
- buf.Write(req_resp.Data)
391
+ buf.Write(reqResp.Data)
392
count--
393
}
394
}
395
396
out:
398
- resp := DHTMessage{
397
+ resp := Message{
398
Type: PBDHTMessage_DIAGNOSTIC,
400
- Id: pmes.GetId(),
399
+ ID: pmes.GetId(),
400
Value: buf.Bytes(),
401
Response: true,
402
}
@@ -423,40 +422,40 @@ func (dht *IpfsDHT) getValueOrPeers(p *peer.Peer, key u.Key, timeout time.Durati
422
423
// Success! We were given the value
424
return pmes.GetValue(), nil, nil
426
- } else {
427
- // We were given a closer node
428
- var peers []*peer.Peer
429
- for _, pb := range pmes.GetPeers() {
430
- if peer.ID(pb.GetId()).Equal(dht.self.ID) {
431
- continue
432
- }
433
- addr, err := ma.NewMultiaddr(pb.GetAddr())
434
- if err != nil {
435
- u.PErr(err.Error())
436
- continue
437
- }
425
+ }
426
439
- np, err := dht.network.GetConnection(peer.ID(pb.GetId()), addr)
440
- if err != nil {
441
- u.PErr(err.Error())
442
- continue
443
- }
427
+ // We were given a closer node
428
+ var peers []*peer.Peer
429
+ for _, pb := range pmes.GetPeers() {
430
+ if peer.ID(pb.GetId()).Equal(dht.self.ID) {
431
+ continue
432
+ }
433
+ addr, err := ma.NewMultiaddr(pb.GetAddr())
434
+ if err != nil {
435
+ u.PErr(err.Error())
436
+ continue
437
+ }
438
445
- peers = append(peers, np)
439
+ np, err := dht.network.GetConnection(peer.ID(pb.GetId()), addr)
440
+ if err != nil {
441
+ u.PErr(err.Error())
442
+ continue
443
}
447
- return nil, peers, nil
444
+
445
+ peers = append(peers, np)
446
}
447
+ return nil, peers, nil
448
}
449
450
// getValueSingle simply performs the get value RPC with the given parameters
451
func (dht *IpfsDHT) getValueSingle(p *peer.Peer, key u.Key, timeout time.Duration, level int) (*PBDHTMessage, error) {
453
- pmes := DHTMessage{
452
+ pmes := Message{
453
Type: PBDHTMessage_GET_VALUE,
454
Key: string(key),
455
Value: []byte{byte(level)},
457
- Id: GenerateMessageID(),
456
+ ID: GenerateMessageID(),
457
}
459
- response_chan := dht.listener.Listen(pmes.Id, 1, time.Minute)
458
+ responseChan := dht.listener.Listen(pmes.ID, 1, time.Minute)
459
460
mes := swarm.NewMessage(p, pmes.ToProtobuf())
461
t := time.Now()
@@ -466,21 +465,21 @@ func (dht *IpfsDHT) getValueSingle(p *peer.Peer, key u.Key, timeout time.Duratio
465
timeup := time.After(timeout)
466
select {
467
case <-timeup:
469
- dht.listener.Unlisten(pmes.Id)
468
+ dht.listener.Unlisten(pmes.ID)
469
return nil, u.ErrTimeout
471
- case resp, ok := <-response_chan:
470
+ case resp, ok := <-responseChan:
471
if !ok {
472
u.PErr("response channel closed before timeout, please investigate.")
473
return nil, u.ErrTimeout
474
}
475
roundtrip := time.Since(t)
476
resp.Peer.SetLatency(roundtrip)
478
- pmes_out := new(PBDHTMessage)
479
- err := proto.Unmarshal(resp.Data, pmes_out)
477
+ pmesOut := new(PBDHTMessage)
478
+ err := proto.Unmarshal(resp.Data, pmesOut)
479
if err != nil {
480
return nil, err
481
}
483
- return pmes_out, nil
482
+ return pmesOut, nil
483
}
484
}
485
@@ -520,7 +519,7 @@ func (dht *IpfsDHT) getFromPeerList(key u.Key, timeout time.Duration,
519
return nil, u.ErrNotFound
520
}
521
523
-func (dht *IpfsDHT) GetLocal(key u.Key) ([]byte, error) {
522
+func (dht *IpfsDHT) getLocal(key u.Key) ([]byte, error) {
523
v, err := dht.datastore.Get(ds.NewKey(string(key)))
524
if err != nil {
525
return nil, err
@@ -528,17 +527,18 @@ func (dht *IpfsDHT) GetLocal(key u.Key) ([]byte, error) {
527
return v.([]byte), nil
528
}
529
531
-func (dht *IpfsDHT) PutLocal(key u.Key, value []byte) error {
530
+func (dht *IpfsDHT) putLocal(key u.Key, value []byte) error {
531
return dht.datastore.Put(ds.NewKey(string(key)), value)
532
}
533
534
+// Update TODO(chas) Document this function
535
func (dht *IpfsDHT) Update(p *peer.Peer) {
536
- for _, route := range dht.routes {
536
+ for _, route := range dht.routingTables {
537
removed := route.Update(p)
538
// Only drop the connection if no tables refer to this peer
539
if removed != nil {
540
found := false
541
- for _, r := range dht.routes {
541
+ for _, r := range dht.routingTables {
542
if r.Find(removed.ID) != nil {
543
found = true
544
break
@@ -551,9 +551,9 @@ func (dht *IpfsDHT) Update(p *peer.Peer) {
551
}
552
}
553
554
-// Look for a peer with a given ID connected to this dht
554
+// Find looks for a peer with a given ID connected to this dht and returns the peer and the table it was found in.
555
func (dht *IpfsDHT) Find(id peer.ID) (*peer.Peer, *kb.RoutingTable) {
556
- for _, table := range dht.routes {
556
+ for _, table := range dht.routingTables {
557
p := table.Find(id)
558
if p != nil {
559
return p, table
@@ -563,72 +563,72 @@ func (dht *IpfsDHT) Find(id peer.ID) (*peer.Peer, *kb.RoutingTable) {
563
}
564
565
func (dht *IpfsDHT) findPeerSingle(p *peer.Peer, id peer.ID, timeout time.Duration, level int) (*PBDHTMessage, error) {
566
- pmes := DHTMessage{
566
+ pmes := Message{
567
Type: PBDHTMessage_FIND_NODE,
568
Key: string(id),
569
- Id: GenerateMessageID(),
569
+ ID: GenerateMessageID(),
570
Value: []byte{byte(level)},
571
}
572
573
mes := swarm.NewMessage(p, pmes.ToProtobuf())
574
- listenChan := dht.listener.Listen(pmes.Id, 1, time.Minute)
574
+ listenChan := dht.listener.Listen(pmes.ID, 1, time.Minute)
575
t := time.Now()
576
dht.network.Send(mes)
577
after := time.After(timeout)
578
select {
579
case <-after:
580
- dht.listener.Unlisten(pmes.Id)
580
+ dht.listener.Unlisten(pmes.ID)
581
return nil, u.ErrTimeout
582
case resp := <-listenChan:
583
roundtrip := time.Since(t)
584
resp.Peer.SetLatency(roundtrip)
585
- pmes_out := new(PBDHTMessage)
586
- err := proto.Unmarshal(resp.Data, pmes_out)
585
+ pmesOut := new(PBDHTMessage)
586
+ err := proto.Unmarshal(resp.Data, pmesOut)
587
if err != nil {
588
return nil, err
589
}
590
591
- return pmes_out, nil
591
+ return pmesOut, nil
592
}
593
}
594
595
-func (dht *IpfsDHT) PrintTables() {
596
- for _, route := range dht.routes {
595
+func (dht *IpfsDHT) printTables() {
596
+ for _, route := range dht.routingTables {
597
route.Print()
598
}
599
}
600
601
func (dht *IpfsDHT) findProvidersSingle(p *peer.Peer, key u.Key, level int, timeout time.Duration) (*PBDHTMessage, error) {
602
- pmes := DHTMessage{
602
+ pmes := Message{
603
Type: PBDHTMessage_GET_PROVIDERS,
604
Key: string(key),
605
- Id: GenerateMessageID(),
605
+ ID: GenerateMessageID(),
606
Value: []byte{byte(level)},
607
}
608
609
mes := swarm.NewMessage(p, pmes.ToProtobuf())
610
611
- listenChan := dht.listener.Listen(pmes.Id, 1, time.Minute)
611
+ listenChan := dht.listener.Listen(pmes.ID, 1, time.Minute)
612
dht.network.Send(mes)
613
after := time.After(timeout)
614
select {
615
case <-after:
616
- dht.listener.Unlisten(pmes.Id)
616
+ dht.listener.Unlisten(pmes.ID)
617
return nil, u.ErrTimeout
618
case resp := <-listenChan:
619
u.DOut("FindProviders: got response.")
620
- pmes_out := new(PBDHTMessage)
621
- err := proto.Unmarshal(resp.Data, pmes_out)
620
+ pmesOut := new(PBDHTMessage)
621
+ err := proto.Unmarshal(resp.Data, pmesOut)
622
if err != nil {
623
return nil, err
624
}
625
626
- return pmes_out, nil
626
+ return pmesOut, nil
627
}
628
}
629
630
func (dht *IpfsDHT) addPeerList(key u.Key, peers []*PBDHTMessage_PBPeer) []*peer.Peer {
631
- var prov_arr []*peer.Peer
631
+ var provArr []*peer.Peer
632
for _, prov := range peers {
633
// Dont add outselves to the list
634
if peer.ID(prov.GetId()).Equal(dht.self.ID) {
@@ -650,7 +650,7 @@ func (dht *IpfsDHT) addPeerList(key u.Key, peers []*PBDHTMessage_PBPeer) []*peer
650
}
651
}
652
dht.addProviderEntry(key, p)
653
- prov_arr = append(prov_arr, p)
653
+ provArr = append(provArr, p)
654
}
655
- return prov_arr
655
+ return provArr
656
}
routing/dht/dht_logger.go
+6
-6
@@ -7,28 +7,28 @@ import (
7
u "github.com/jbenet/go-ipfs/util"
8
)
9
10
-type logDhtRpc struct {
10
+type logDhtRPC struct {
11
Type string
12
Start time.Time
13
End time.Time
14
Duration time.Duration
15
- RpcCount int
15
+ RPCCount int
16
Success bool
17
}
18
19
-func startNewRpc(name string) *logDhtRpc {
20
- r := new(logDhtRpc)
19
+func startNewRPC(name string) *logDhtRPC {
20
+ r := new(logDhtRPC)
21
r.Type = name
22
r.Start = time.Now()
23
return r
24
}
25
26
-func (l *logDhtRpc) EndLog() {
26
+func (l *logDhtRPC) EndLog() {
27
l.End = time.Now()
28
l.Duration = l.End.Sub(l.Start)
29
}
30
31
-func (l *logDhtRpc) Print() {
31
+func (l *logDhtRPC) Print() {
32
b, err := json.Marshal(l)
33
if err != nil {
34
u.DOut(err.Error())
routing/dht/dht_test.go
+37
-37
@@ -47,93 +47,93 @@ func setupDHTS(n int, t *testing.T) ([]*ma.Multiaddr, []*peer.Peer, []*IpfsDHT)
47
48
func TestPing(t *testing.T) {
49
u.Debug = false
50
- addr_a, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/2222")
50
+ addrA, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/2222")
51
if err != nil {
52
t.Fatal(err)
53
}
54
- addr_b, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/5678")
54
+ addrB, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/5678")
55
if err != nil {
56
t.Fatal(err)
57
}
58
59
- peer_a := new(peer.Peer)
60
- peer_a.AddAddress(addr_a)
61
- peer_a.ID = peer.ID([]byte("peer_a"))
59
+ peerA := new(peer.Peer)
60
+ peerA.AddAddress(addrA)
61
+ peerA.ID = peer.ID([]byte("peerA"))
62
63
- peer_b := new(peer.Peer)
64
- peer_b.AddAddress(addr_b)
65
- peer_b.ID = peer.ID([]byte("peer_b"))
63
+ peerB := new(peer.Peer)
64
+ peerB.AddAddress(addrB)
65
+ peerB.ID = peer.ID([]byte("peerB"))
66
67
- neta := swarm.NewSwarm(peer_a)
67
+ neta := swarm.NewSwarm(peerA)
68
err = neta.Listen()
69
if err != nil {
70
t.Fatal(err)
71
}
72
- dht_a := NewDHT(peer_a, neta)
72
+ dhtA := NewDHT(peerA, neta)
73
74
- netb := swarm.NewSwarm(peer_b)
74
+ netb := swarm.NewSwarm(peerB)
75
err = netb.Listen()
76
if err != nil {
77
t.Fatal(err)
78
}
79
- dht_b := NewDHT(peer_b, netb)
79
+ dhtB := NewDHT(peerB, netb)
80
81
- dht_a.Start()
82
- dht_b.Start()
81
+ dhtA.Start()
82
+ dhtB.Start()
83
84
- _, err = dht_a.Connect(addr_b)
84
+ _, err = dhtA.Connect(addrB)
85
if err != nil {
86
t.Fatal(err)
87
}
88
89
//Test that we can ping the node
90
- err = dht_a.Ping(peer_b, time.Second*2)
90
+ err = dhtA.Ping(peerB, time.Second*2)
91
if err != nil {
92
t.Fatal(err)
93
}
94
95
- dht_a.Halt()
96
- dht_b.Halt()
95
+ dhtA.Halt()
96
+ dhtB.Halt()
97
}
98
99
func TestValueGetSet(t *testing.T) {
100
u.Debug = false
101
- addr_a, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/1235")
101
+ addrA, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/1235")
102
if err != nil {
103
t.Fatal(err)
104
}
105
- addr_b, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/5679")
105
+ addrB, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/5679")
106
if err != nil {
107
t.Fatal(err)
108
}
109
110
- peer_a := new(peer.Peer)
111
- peer_a.AddAddress(addr_a)
112
- peer_a.ID = peer.ID([]byte("peer_a"))
110
+ peerA := new(peer.Peer)
111
+ peerA.AddAddress(addrA)
112
+ peerA.ID = peer.ID([]byte("peerA"))
113
114
- peer_b := new(peer.Peer)
115
- peer_b.AddAddress(addr_b)
116
- peer_b.ID = peer.ID([]byte("peer_b"))
114
+ peerB := new(peer.Peer)
115
+ peerB.AddAddress(addrB)
116
+ peerB.ID = peer.ID([]byte("peerB"))
117
118
- neta := swarm.NewSwarm(peer_a)
118
+ neta := swarm.NewSwarm(peerA)
119
err = neta.Listen()
120
if err != nil {
121
t.Fatal(err)
122
}
123
- dht_a := NewDHT(peer_a, neta)
123
+ dhtA := NewDHT(peerA, neta)
124
125
- netb := swarm.NewSwarm(peer_b)
125
+ netb := swarm.NewSwarm(peerB)
126
err = netb.Listen()
127
if err != nil {
128
t.Fatal(err)
129
}
130
- dht_b := NewDHT(peer_b, netb)
130
+ dhtB := NewDHT(peerB, netb)
131
132
- dht_a.Start()
133
- dht_b.Start()
132
+ dhtA.Start()
133
+ dhtB.Start()
134
135
- errsa := dht_a.network.GetChan().Errors
136
- errsb := dht_b.network.GetChan().Errors
135
+ errsa := dhtA.network.GetChan().Errors
136
+ errsb := dhtB.network.GetChan().Errors
137
go func() {
138
select {
139
case err := <-errsa:
@@ -143,14 +143,14 @@ func TestValueGetSet(t *testing.T) {
143
}
144
}()
145
146
- _, err = dht_a.Connect(addr_b)
146
+ _, err = dhtA.Connect(addrB)
147
if err != nil {
148
t.Fatal(err)
149
}
150
151
- dht_a.PutValue("hello", []byte("world"))
151
+ dhtA.PutValue("hello", []byte("world"))
152
153
- val, err := dht_a.GetValue("hello", time.Second*2)
153
+ val, err := dhtA.GetValue("hello", time.Second*2)
154
if err != nil {
155
t.Fatal(err)
156
}
routing/dht/diag.go
+4
-4
@@ -9,11 +9,11 @@ import (
9
10
type connDiagInfo struct {
11
Latency time.Duration
12
- Id peer.ID
12
+ ID peer.ID
13
}
14
15
type diagInfo struct {
16
- Id peer.ID
16
+ ID peer.ID
17
Connections []connDiagInfo
18
Keys []string
19
LifeSpan time.Duration
@@ -32,11 +32,11 @@ func (di *diagInfo) Marshal() []byte {
32
func (dht *IpfsDHT) getDiagInfo() *diagInfo {
33
di := new(diagInfo)
34
di.CodeVersion = "github.com/jbenet/go-ipfs"
35
- di.Id = dht.self.ID
35
+ di.ID = dht.self.ID
36
di.LifeSpan = time.Since(dht.birth)
37
di.Keys = nil // Currently no way to query datastore
38
39
- for _, p := range dht.routes[0].Listpeers() {
39
+ for _, p := range dht.routingTables[0].Listpeers() {
40
di.Connections = append(di.Connections, connDiagInfo{p.GetLatency(), p.ID})
41
}
42
return di
routing/dht/mes_listener.go
+8
-8
@@ -8,7 +8,7 @@ import (
8
u "github.com/jbenet/go-ipfs/util"
9
)
10
11
-type MesListener struct {
11
+type mesListener struct {
12
listeners map[uint64]*listenInfo
13
haltchan chan struct{}
14
unlist chan uint64
@@ -36,8 +36,8 @@ type listenInfo struct {
36
id uint64
37
}
38
39
-func NewMesListener() *MesListener {
40
- ml := new(MesListener)
39
+func newMesListener() *mesListener {
40
+ ml := new(mesListener)
41
ml.haltchan = make(chan struct{})
42
ml.listeners = make(map[uint64]*listenInfo)
43
ml.nlist = make(chan *listenInfo, 16)
@@ -47,7 +47,7 @@ func NewMesListener() *MesListener {
47
return ml
48
}
49
50
-func (ml *MesListener) Listen(id uint64, count int, timeout time.Duration) <-chan *swarm.Message {
50
+func (ml *mesListener) Listen(id uint64, count int, timeout time.Duration) <-chan *swarm.Message {
51
li := new(listenInfo)
52
li.count = count
53
li.eol = time.Now().Add(timeout)
@@ -57,7 +57,7 @@ func (ml *MesListener) Listen(id uint64, count int, timeout time.Duration) <-cha
57
return li.resp
58
}
59
60
-func (ml *MesListener) Unlisten(id uint64) {
60
+func (ml *mesListener) Unlisten(id uint64) {
61
ml.unlist <- id
62
}
63
@@ -66,18 +66,18 @@ type respMes struct {
66
mes *swarm.Message
67
}
68
69
-func (ml *MesListener) Respond(id uint64, mes *swarm.Message) {
69
+func (ml *mesListener) Respond(id uint64, mes *swarm.Message) {
70
ml.send <- &respMes{
71
id: id,
72
mes: mes,
73
}
74
}
75
76
-func (ml *MesListener) Halt() {
76
+func (ml *mesListener) Halt() {
77
ml.haltchan <- struct{}{}
78
}
79
80
-func (ml *MesListener) run() {
80
+func (ml *mesListener) run() {
81
for {
82
select {
83
case <-ml.haltchan:
routing/dht/messages.pb.go
+13
-10
@@ -1,4 +1,4 @@
1
-// Code generated by protoc-gen-go.
1
+// Code generated by protoc-gen-gogo.
2
// source: messages.proto
3
// DO NOT EDIT!
4
@@ -13,7 +13,7 @@ It has these top-level messages:
13
*/
14
package dht
15
16
-import proto "code.google.com/p/goprotobuf/proto"
16
+import proto "code.google.com/p/gogoprotobuf/proto"
17
import math "math"
18
19
// Reference imports to suppress errors if they are not otherwise used.
@@ -69,14 +69,17 @@ func (x *PBDHTMessage_MessageType) UnmarshalJSON(data []byte) error {
69
}
70
71
type PBDHTMessage struct {
72
- Type *PBDHTMessage_MessageType `protobuf:"varint,1,req,name=type,enum=dht.PBDHTMessage_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
- Success *bool `protobuf:"varint,6,opt,name=success" json:"success,omitempty"`
78
- Peers []*PBDHTMessage_PBPeer `protobuf:"bytes,7,rep,name=peers" json:"peers,omitempty"`
79
- XXX_unrecognized []byte `json:"-"`
72
+ Type *PBDHTMessage_MessageType `protobuf:"varint,1,req,name=type,enum=dht.PBDHTMessage_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
+ // Unique ID of this message, used to match queries with responses
76
+ Id *uint64 `protobuf:"varint,4,req,name=id" json:"id,omitempty"`
77
+ // Signals whether or not this message is a response to another message
78
+ Response *bool `protobuf:"varint,5,opt,name=response" json:"response,omitempty"`
79
+ Success *bool `protobuf:"varint,6,opt,name=success" json:"success,omitempty"`
80
+ // Used for returning peers from queries (normally, peers closer to X)
81
+ Peers []*PBDHTMessage_PBPeer `protobuf:"bytes,7,rep,name=peers" json:"peers,omitempty"`
82
+ XXX_unrecognized []byte `json:"-"`
83
}
84
85
func (m *PBDHTMessage) Reset() { *m = PBDHTMessage{} }
routing/dht/routing.go
+79
-79
@@ -27,6 +27,7 @@ var KValue = 10
27
// Its in the paper, i swear
28
var AlphaValue = 3
29
30
+// GenerateMessageID creates and returns a new message ID
31
// TODO: determine a way of creating and managing message IDs
32
func GenerateMessageID() uint64 {
33
//return (uint64(rand.Uint32()) << 32) & uint64(rand.Uint32())
@@ -39,21 +40,21 @@ func GenerateMessageID() uint64 {
40
41
// PutValue adds value corresponding to given Key.
42
// This is the top level "Store" operation of the DHT
42
-func (s *IpfsDHT) PutValue(key u.Key, value []byte) {
43
+func (dht *IpfsDHT) PutValue(key u.Key, value []byte) {
44
complete := make(chan struct{})
45
count := 0
45
- for _, route := range s.routes {
46
+ for _, route := range dht.routingTables {
47
peers := route.NearestPeers(kb.ConvertKey(key), KValue)
48
for _, p := range peers {
49
if p == nil {
49
- s.network.Error(kb.ErrLookupFailure)
50
+ dht.network.Error(kb.ErrLookupFailure)
51
continue
52
}
53
count++
54
go func(sp *peer.Peer) {
54
- err := s.putValueToNetwork(sp, string(key), value)
55
+ err := dht.putValueToNetwork(sp, string(key), value)
56
if err != nil {
56
- s.network.Error(err)
57
+ dht.network.Error(err)
58
}
59
complete <- struct{}{}
60
}(p)
@@ -121,8 +122,8 @@ func (ps *peerSet) Size() int {
122
// GetValue searches for the value corresponding to given Key.
123
// If the search does not succeed, a multiaddr string of a closer peer is
124
// returned along with util.ErrSearchIncomplete
124
-func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
125
- ll := startNewRpc("GET")
125
+func (dht *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
126
+ ll := startNewRPC("GET")
127
defer func() {
128
ll.EndLog()
129
ll.Print()
@@ -130,29 +131,29 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
131
132
// If we have it local, dont bother doing an RPC!
133
// NOTE: this might not be what we want to do...
133
- val, err := s.GetLocal(key)
134
+ val, err := dht.getLocal(key)
135
if err == nil {
136
ll.Success = true
137
u.DOut("Found local, returning.")
138
return val, nil
139
}
140
140
- route_level := 0
141
- closest := s.routes[route_level].NearestPeers(kb.ConvertKey(key), PoolSize)
141
+ routeLevel := 0
142
+ closest := dht.routingTables[routeLevel].NearestPeers(kb.ConvertKey(key), PoolSize)
143
if closest == nil || len(closest) == 0 {
144
return nil, kb.ErrLookupFailure
145
}
146
146
- val_chan := make(chan []byte)
147
- npeer_chan := make(chan *peer.Peer, 30)
148
- proc_peer := make(chan *peer.Peer, 30)
149
- err_chan := make(chan error)
147
+ valChan := make(chan []byte)
148
+ npeerChan := make(chan *peer.Peer, 30)
149
+ procPeer := make(chan *peer.Peer, 30)
150
+ errChan := make(chan error)
151
after := time.After(timeout)
152
pset := newPeerSet()
153
154
for _, p := range closest {
155
pset.Add(p)
155
- npeer_chan <- p
156
+ npeerChan <- p
157
}
158
159
c := counter{}
@@ -161,17 +162,17 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
162
go func() {
163
for {
164
select {
164
- case p := <-npeer_chan:
165
+ case p := <-npeerChan:
166
count++
167
if count >= KValue {
168
break
169
}
170
c.Increment()
171
171
- proc_peer <- p
172
+ procPeer <- p
173
default:
174
if c.Size() == 0 {
174
- err_chan <- u.ErrNotFound
175
+ errChan <- u.ErrNotFound
176
}
177
}
178
}
@@ -180,19 +181,19 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
181
process := func() {
182
for {
183
select {
183
- case p, ok := <-proc_peer:
184
+ case p, ok := <-procPeer:
185
if !ok || p == nil {
186
c.Decrement()
187
return
188
}
188
- val, peers, err := s.getValueOrPeers(p, key, timeout/4, route_level)
189
+ val, peers, err := dht.getValueOrPeers(p, key, timeout/4, routeLevel)
190
if err != nil {
191
u.DErr(err.Error())
192
c.Decrement()
193
continue
194
}
195
if val != nil {
195
- val_chan <- val
196
+ valChan <- val
197
c.Decrement()
198
return
199
}
@@ -201,7 +202,7 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
202
// TODO: filter out peers that arent closer
203
if !pset.Contains(np) && pset.Size() < KValue {
204
pset.Add(np) //This is racey... make a single function to do operation
204
- npeer_chan <- np
205
+ npeerChan <- np
206
}
207
}
208
c.Decrement()
@@ -214,9 +215,9 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
215
}
216
217
select {
217
- case val := <-val_chan:
218
+ case val := <-valChan:
219
return val, nil
219
- case err := <-err_chan:
220
+ case err := <-errChan:
221
return nil, err
222
case <-after:
223
return nil, u.ErrTimeout
@@ -226,14 +227,14 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
227
// Value provider layer of indirection.
228
// This is what DSHTs (Coral and MainlineDHT) do to store large values in a DHT.
229
229
-// Announce that this node can provide value for given key
230
-func (s *IpfsDHT) Provide(key u.Key) error {
231
- peers := s.routes[0].NearestPeers(kb.ConvertKey(key), PoolSize)
230
+// Provide makes this node announce that it can provide a value for the given key
231
+func (dht *IpfsDHT) Provide(key u.Key) error {
232
+ peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(key), PoolSize)
233
if len(peers) == 0 {
234
return kb.ErrLookupFailure
235
}
236
236
- pmes := DHTMessage{
237
+ pmes := Message{
238
Type: PBDHTMessage_ADD_PROVIDER,
239
Key: string(key),
240
}
@@ -241,57 +242,57 @@ func (s *IpfsDHT) Provide(key u.Key) error {
242
243
for _, p := range peers {
244
mes := swarm.NewMessage(p, pbmes)
244
- s.network.Send(mes)
245
+ dht.network.Send(mes)
246
}
247
return nil
248
}
249
250
// FindProviders searches for peers who can provide the value for given key.
250
-func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer, error) {
251
- ll := startNewRpc("FindProviders")
251
+func (dht *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer, error) {
252
+ ll := startNewRPC("FindProviders")
253
defer func() {
254
ll.EndLog()
255
ll.Print()
256
}()
257
u.DOut("Find providers for: '%s'", key)
257
- p := s.routes[0].NearestPeer(kb.ConvertKey(key))
258
+ p := dht.routingTables[0].NearestPeer(kb.ConvertKey(key))
259
if p == nil {
260
return nil, kb.ErrLookupFailure
261
}
262
262
- for level := 0; level < len(s.routes); {
263
- pmes, err := s.findProvidersSingle(p, key, level, timeout)
263
+ for level := 0; level < len(dht.routingTables); {
264
+ pmes, err := dht.findProvidersSingle(p, key, level, timeout)
265
if err != nil {
266
return nil, err
267
}
268
if pmes.GetSuccess() {
268
- provs := s.addPeerList(key, pmes.GetPeers())
269
+ provs := dht.addPeerList(key, pmes.GetPeers())
270
ll.Success = true
271
return provs, nil
271
- } else {
272
- closer := pmes.GetPeers()
273
- if len(closer) == 0 {
274
- level++
275
- continue
276
- }
277
- if peer.ID(closer[0].GetId()).Equal(s.self.ID) {
278
- u.DOut("Got myself back as a closer peer.")
279
- return nil, u.ErrNotFound
280
- }
281
- maddr, err := ma.NewMultiaddr(closer[0].GetAddr())
282
- if err != nil {
283
- // ??? Move up route level???
284
- panic("not yet implemented")
285
- }
272
+ }
273
287
- np, err := s.network.GetConnection(peer.ID(closer[0].GetId()), maddr)
288
- if err != nil {
289
- u.PErr("[%s] Failed to connect to: %s", s.self.ID.Pretty(), closer[0].GetAddr())
290
- level++
291
- continue
292
- }
293
- p = np
274
+ closer := pmes.GetPeers()
275
+ if len(closer) == 0 {
276
+ level++
277
+ continue
278
+ }
279
+ if peer.ID(closer[0].GetId()).Equal(dht.self.ID) {
280
+ u.DOut("Got myself back as a closer peer.")
281
+ return nil, u.ErrNotFound
282
+ }
283
+ maddr, err := ma.NewMultiaddr(closer[0].GetAddr())
284
+ if err != nil {
285
+ // ??? Move up route level???
286
+ panic("not yet implemented")
287
+ }
288
+
289
+ np, err := dht.network.GetConnection(peer.ID(closer[0].GetId()), maddr)
290
+ if err != nil {
291
+ u.PErr("[%s] Failed to connect to: %s", dht.self.ID.Pretty(), closer[0].GetAddr())
292
+ level++
293
+ continue
294
}
295
+ p = np
296
}
297
return nil, u.ErrNotFound
298
}
@@ -299,15 +300,15 @@ func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer,
300
// Find specific Peer
301
302
// FindPeer searches for a peer with given ID.
302
-func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error) {
303
+func (dht *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error) {
304
// Check if were already connected to them
304
- p, _ := s.Find(id)
305
+ p, _ := dht.Find(id)
306
if p != nil {
307
return p, nil
308
}
309
309
- route_level := 0
310
- p = s.routes[route_level].NearestPeer(kb.ConvertPeerID(id))
310
+ routeLevel := 0
311
+ p = dht.routingTables[routeLevel].NearestPeer(kb.ConvertPeerID(id))
312
if p == nil {
313
return nil, kb.ErrLookupFailure
314
}
@@ -315,11 +316,11 @@ func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error
316
return p, nil
317
}
318
318
- for route_level < len(s.routes) {
319
- pmes, err := s.findPeerSingle(p, id, timeout, route_level)
319
+ for routeLevel < len(dht.routingTables) {
320
+ pmes, err := dht.findPeerSingle(p, id, timeout, routeLevel)
321
plist := pmes.GetPeers()
322
if len(plist) == 0 {
322
- route_level++
323
+ routeLevel++
324
}
325
found := plist[0]
326
@@ -328,7 +329,7 @@ func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error
329
return nil, err
330
}
331
331
- nxtPeer, err := s.network.GetConnection(peer.ID(found.GetId()), addr)
332
+ nxtPeer, err := dht.network.GetConnection(peer.ID(found.GetId()), addr)
333
if err != nil {
334
return nil, err
335
}
@@ -337,9 +338,8 @@ func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error
338
return nil, errors.New("got back invalid peer from 'successful' response")
339
}
340
return nxtPeer, nil
340
- } else {
341
- p = nxtPeer
341
}
342
+ p = nxtPeer
343
}
344
return nil, u.ErrNotFound
345
}
@@ -349,16 +349,16 @@ func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
349
// Thoughts: maybe this should accept an ID and do a peer lookup?
350
u.DOut("Enter Ping.")
351
352
- pmes := DHTMessage{Id: GenerateMessageID(), Type: PBDHTMessage_PING}
352
+ pmes := Message{ID: GenerateMessageID(), Type: PBDHTMessage_PING}
353
mes := swarm.NewMessage(p, pmes.ToProtobuf())
354
355
before := time.Now()
356
- response_chan := dht.listener.Listen(pmes.Id, 1, time.Minute)
356
+ responseChan := dht.listener.Listen(pmes.ID, 1, time.Minute)
357
dht.network.Send(mes)
358
359
tout := time.After(timeout)
360
select {
361
- case <-response_chan:
361
+ case <-responseChan:
362
roundtrip := time.Since(before)
363
p.SetLatency(roundtrip)
364
u.DOut("Ping took %s.", roundtrip.String())
@@ -366,23 +366,23 @@ func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
366
case <-tout:
367
// Timed out, think about removing peer from network
368
u.DOut("Ping peer timed out.")
369
- dht.listener.Unlisten(pmes.Id)
369
+ dht.listener.Unlisten(pmes.ID)
370
return u.ErrTimeout
371
}
372
}
373
374
-func (dht *IpfsDHT) GetDiagnostic(timeout time.Duration) ([]*diagInfo, error) {
374
+func (dht *IpfsDHT) getDiagnostic(timeout time.Duration) ([]*diagInfo, error) {
375
u.DOut("Begin Diagnostic")
376
//Send to N closest peers
377
- targets := dht.routes[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
377
+ targets := dht.routingTables[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
378
379
// TODO: Add timeout to this struct so nodes know when to return
380
- pmes := DHTMessage{
380
+ pmes := Message{
381
Type: PBDHTMessage_DIAGNOSTIC,
382
- Id: GenerateMessageID(),
382
+ ID: GenerateMessageID(),
383
}
384
385
- listenChan := dht.listener.Listen(pmes.Id, len(targets), time.Minute*2)
385
+ listenChan := dht.listener.Listen(pmes.ID, len(targets), time.Minute*2)
386
387
pbmes := pmes.ToProtobuf()
388
for _, p := range targets {
@@ -398,15 +398,15 @@ func (dht *IpfsDHT) GetDiagnostic(timeout time.Duration) ([]*diagInfo, error) {
398
u.DOut("Diagnostic request timed out.")
399
return out, u.ErrTimeout
400
case resp := <-listenChan:
401
- pmes_out := new(PBDHTMessage)
402
- err := proto.Unmarshal(resp.Data, pmes_out)
401
+ pmesOut := new(PBDHTMessage)
402
+ err := proto.Unmarshal(resp.Data, pmesOut)
403
if err != nil {
404
// NOTE: here and elsewhere, need to audit error handling,
405
// some errors should be continued on from
406
return out, err
407
}
408
409
- dec := json.NewDecoder(bytes.NewBuffer(pmes_out.GetValue()))
409
+ dec := json.NewDecoder(bytes.NewBuffer(pmesOut.GetValue()))
410
for {
411
di := new(diagInfo)
412
err := dec.Decode(di)
routing/kbucket/bucket.go
+11
-11
@@ -13,13 +13,13 @@ type Bucket struct {
13
list *list.List
14
}
15
16
-func NewBucket() *Bucket {
16
+func newBucket() *Bucket {
17
b := new(Bucket)
18
b.list = list.New()
19
return b
20
}
21
22
-func (b *Bucket) Find(id peer.ID) *list.Element {
22
+func (b *Bucket) find(id peer.ID) *list.Element {
23
b.lk.RLock()
24
defer b.lk.RUnlock()
25
for e := b.list.Front(); e != nil; e = e.Next() {
@@ -30,19 +30,19 @@ func (b *Bucket) Find(id peer.ID) *list.Element {
30
return nil
31
}
32
33
-func (b *Bucket) MoveToFront(e *list.Element) {
33
+func (b *Bucket) moveToFront(e *list.Element) {
34
b.lk.Lock()
35
b.list.MoveToFront(e)
36
b.lk.Unlock()
37
}
38
39
-func (b *Bucket) PushFront(p *peer.Peer) {
39
+func (b *Bucket) pushFront(p *peer.Peer) {
40
b.lk.Lock()
41
b.list.PushFront(p)
42
b.lk.Unlock()
43
}
44
45
-func (b *Bucket) PopBack() *peer.Peer {
45
+func (b *Bucket) popBack() *peer.Peer {
46
b.lk.Lock()
47
defer b.lk.Unlock()
48
last := b.list.Back()
@@ -50,13 +50,13 @@ func (b *Bucket) PopBack() *peer.Peer {
50
return last.Value.(*peer.Peer)
51
}
52
53
-func (b *Bucket) Len() int {
53
+func (b *Bucket) len() int {
54
b.lk.RLock()
55
defer b.lk.RUnlock()
56
return b.list.Len()
57
}
58
59
-// Splits a buckets peers into two buckets, the methods receiver will have
59
+// Split splits a buckets peers into two buckets, the methods receiver will have
60
// peers with CPL equal to cpl, the returned bucket will have peers with CPL
61
// greater than cpl (returned bucket has closer peers)
62
func (b *Bucket) Split(cpl int, target ID) *Bucket {
@@ -64,13 +64,13 @@ func (b *Bucket) Split(cpl int, target ID) *Bucket {
64
defer b.lk.Unlock()
65
66
out := list.New()
67
- newbuck := NewBucket()
67
+ newbuck := newBucket()
68
newbuck.list = out
69
e := b.list.Front()
70
for e != nil {
71
- peer_id := ConvertPeerID(e.Value.(*peer.Peer).ID)
72
- peer_cpl := prefLen(peer_id, target)
73
- if peer_cpl > cpl {
71
+ peerID := convertPeerID(e.Value.(*peer.Peer).ID)
72
+ peerCPL := prefLen(peerID, target)
73
+ if peerCPL > cpl {
74
cur := e
75
out.PushBack(e.Value)
76
e = e.Next()
routing/kbucket/table.go
+37
-39
@@ -28,11 +28,11 @@ type RoutingTable struct {
28
bucketsize int
29
}
30
31
-func NewRoutingTable(bucketsize int, local_id ID, latency time.Duration) *RoutingTable {
31
+func newRoutingTable(bucketsize int, localID ID, latency time.Duration) *RoutingTable {
32
rt := new(RoutingTable)
33
- rt.Buckets = []*Bucket{NewBucket()}
33
+ rt.Buckets = []*Bucket{newBucket()}
34
rt.bucketsize = bucketsize
35
- rt.local = local_id
35
+ rt.local = localID
36
rt.maxLatency = latency
37
return rt
38
}
@@ -42,51 +42,50 @@ func NewRoutingTable(bucketsize int, local_id ID, latency time.Duration) *Routin
42
func (rt *RoutingTable) Update(p *peer.Peer) *peer.Peer {
43
rt.tabLock.Lock()
44
defer rt.tabLock.Unlock()
45
- peer_id := ConvertPeerID(p.ID)
46
- cpl := xor(peer_id, rt.local).commonPrefixLen()
45
+ peerID := convertPeerID(p.ID)
46
+ cpl := xor(peerID, rt.local).commonPrefixLen()
47
48
- b_id := cpl
49
- if b_id >= len(rt.Buckets) {
50
- b_id = len(rt.Buckets) - 1
48
+ bucketID := cpl
49
+ if bucketID >= len(rt.Buckets) {
50
+ bucketID = len(rt.Buckets) - 1
51
}
52
53
- bucket := rt.Buckets[b_id]
54
- e := bucket.Find(p.ID)
53
+ bucket := rt.Buckets[bucketID]
54
+ e := bucket.find(p.ID)
55
if e == nil {
56
// New peer, add to bucket
57
if p.GetLatency() > rt.maxLatency {
58
// Connection doesnt meet requirements, skip!
59
return nil
60
}
61
- bucket.PushFront(p)
61
+ bucket.pushFront(p)
62
63
// Are we past the max bucket size?
64
- if bucket.Len() > rt.bucketsize {
65
- if b_id == len(rt.Buckets)-1 {
66
- new_bucket := bucket.Split(b_id, rt.local)
67
- rt.Buckets = append(rt.Buckets, new_bucket)
68
- if new_bucket.Len() > rt.bucketsize {
64
+ if bucket.len() > rt.bucketsize {
65
+ if bucketID == len(rt.Buckets)-1 {
66
+ newBucket := bucket.Split(bucketID, rt.local)
67
+ rt.Buckets = append(rt.Buckets, newBucket)
68
+ if newBucket.len() > rt.bucketsize {
69
// TODO: This is a very rare and annoying case
70
panic("Case not handled.")
71
}
72
73
// If all elements were on left side of split...
74
- if bucket.Len() > rt.bucketsize {
75
- return bucket.PopBack()
74
+ if bucket.len() > rt.bucketsize {
75
+ return bucket.popBack()
76
}
77
} else {
78
// If the bucket cant split kick out least active node
79
- return bucket.PopBack()
79
+ return bucket.popBack()
80
}
81
}
82
return nil
83
- } else {
84
- // If the peer is already in the table, move it to the front.
85
- // This signifies that it it "more active" and the less active nodes
86
- // Will as a result tend towards the back of the list
87
- bucket.MoveToFront(e)
88
- return nil
83
}
84
+ // If the peer is already in the table, move it to the front.
85
+ // This signifies that it it "more active" and the less active nodes
86
+ // Will as a result tend towards the back of the list
87
+ bucket.moveToFront(e)
88
+ return nil
89
}
90
91
// A helper struct to sort peers by their distance to the local node
@@ -101,7 +100,7 @@ type peerSorterArr []*peerDistance
100
func (p peerSorterArr) Len() int { return len(p) }
101
func (p peerSorterArr) Swap(a, b int) { p[a], p[b] = p[b], p[a] }
102
func (p peerSorterArr) Less(a, b int) bool {
104
- return p[a].distance.Less(p[b].distance)
103
+ return p[a].distance.less(p[b].distance)
104
}
105
106
//
@@ -109,10 +108,10 @@ func (p peerSorterArr) Less(a, b int) bool {
108
func copyPeersFromList(target ID, peerArr peerSorterArr, peerList *list.List) peerSorterArr {
109
for e := peerList.Front(); e != nil; e = e.Next() {
110
p := e.Value.(*peer.Peer)
112
- p_id := ConvertPeerID(p.ID)
111
+ pID := convertPeerID(p.ID)
112
pd := peerDistance{
113
p: p,
115
- distance: xor(target, p_id),
114
+ distance: xor(target, pID),
115
}
116
peerArr = append(peerArr, &pd)
117
if e == nil {
@@ -125,24 +124,23 @@ func copyPeersFromList(target ID, peerArr peerSorterArr, peerList *list.List) pe
124
125
// Find a specific peer by ID or return nil
126
func (rt *RoutingTable) Find(id peer.ID) *peer.Peer {
128
- srch := rt.NearestPeers(ConvertPeerID(id), 1)
127
+ srch := rt.NearestPeers(convertPeerID(id), 1)
128
if len(srch) == 0 || !srch[0].ID.Equal(id) {
129
return nil
130
}
131
return srch[0]
132
}
133
135
-// Returns a single peer that is nearest to the given ID
134
+// NearestPeer returns a single peer that is nearest to the given ID
135
func (rt *RoutingTable) NearestPeer(id ID) *peer.Peer {
136
peers := rt.NearestPeers(id, 1)
137
if len(peers) > 0 {
138
return peers[0]
140
- } else {
141
- return nil
139
}
140
+ return nil
141
}
142
145
-// Returns a list of the 'count' closest peers to the given ID
143
+// NearestPeers returns a list of the 'count' closest peers to the given ID
144
func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
145
rt.tabLock.RLock()
146
defer rt.tabLock.RUnlock()
@@ -156,7 +154,7 @@ func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
154
bucket = rt.Buckets[cpl]
155
156
var peerArr peerSorterArr
159
- if bucket.Len() == 0 {
157
+ if bucket.len() == 0 {
158
// In the case of an unusual split, one bucket may be empty.
159
// if this happens, search both surrounding buckets for nearest peer
160
if cpl > 0 {
@@ -183,17 +181,17 @@ func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
181
return out
182
}
183
186
-// Returns the total number of peers in the routing table
184
+// Size returns the total number of peers in the routing table
185
func (rt *RoutingTable) Size() int {
186
var tot int
187
for _, buck := range rt.Buckets {
190
- tot += buck.Len()
188
+ tot += buck.len()
189
}
190
return tot
191
}
192
193
// NOTE: This is potentially unsafe... use at your own risk
196
-func (rt *RoutingTable) Listpeers() []*peer.Peer {
194
+func (rt *RoutingTable) listPeers() []*peer.Peer {
195
var peers []*peer.Peer
196
for _, buck := range rt.Buckets {
197
for e := buck.getIter(); e != nil; e = e.Next() {
@@ -203,10 +201,10 @@ func (rt *RoutingTable) Listpeers() []*peer.Peer {
201
return peers
202
}
203
206
-func (rt *RoutingTable) Print() {
204
+func (rt *RoutingTable) print() {
205
fmt.Printf("Routing Table, bs = %d, Max latency = %d\n", rt.bucketsize, rt.maxLatency)
206
rt.tabLock.RLock()
209
- peers := rt.Listpeers()
207
+ peers := rt.listPeers()
208
for i, p := range peers {
209
fmt.Printf("%d) %s %s\n", i, p.ID.Pretty(), p.GetLatency().String())
210
}
routing/kbucket/table_test.go
+19
-19
@@ -27,28 +27,28 @@ func _randID() ID {
27
28
// Test basic features of the bucket struct
29
func TestBucket(t *testing.T) {
30
- b := NewBucket()
30
+ b := newBucket()
31
32
peers := make([]*peer.Peer, 100)
33
for i := 0; i < 100; i++ {
34
peers[i] = _randPeer()
35
- b.PushFront(peers[i])
35
+ b.pushFront(peers[i])
36
}
37
38
local := _randPeer()
39
- local_id := ConvertPeerID(local.ID)
39
+ localID := convertPeerID(local.ID)
40
41
i := rand.Intn(len(peers))
42
- e := b.Find(peers[i].ID)
42
+ e := b.find(peers[i].ID)
43
if e == nil {
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 := b.list
49
for e := llist.Front(); e != nil; e = e.Next() {
50
- p := ConvertPeerID(e.Value.(*peer.Peer).ID)
51
- cpl := xor(p, local_id).commonPrefixLen()
50
+ p := convertPeerID(e.Value.(*peer.Peer).ID)
51
+ cpl := xor(p, localID).commonPrefixLen()
52
if cpl > 0 {
53
t.Fatalf("Split failed. found id with cpl > 0 in 0 bucket")
54
}
@@ -56,8 +56,8 @@ func TestBucket(t *testing.T) {
56
57
rlist := spl.list
58
for e := rlist.Front(); e != nil; e = e.Next() {
59
- p := ConvertPeerID(e.Value.(*peer.Peer).ID)
60
- cpl := xor(p, local_id).commonPrefixLen()
59
+ p := convertPeerID(e.Value.(*peer.Peer).ID)
60
+ cpl := xor(p, localID).commonPrefixLen()
61
if cpl == 0 {
62
t.Fatalf("Split failed. found id with cpl == 0 in non 0 bucket")
63
}
@@ -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), time.Hour)
70
+ rt := newRoutingTable(10, convertPeerID(local.ID), time.Hour)
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), time.Hour)
96
+ rt := newRoutingTable(10, convertPeerID(local.ID), time.Hour)
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), time.Hour)
113
+ rt := newRoutingTable(20, convertPeerID(local.ID), time.Hour)
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
}
@@ -130,7 +130,7 @@ func TestTableFindMultiple(t *testing.T) {
130
// and set GOMAXPROCS above 1
131
func TestTableMultithreaded(t *testing.T) {
132
local := peer.ID("localPeer")
133
- tab := NewRoutingTable(20, ConvertPeerID(local), time.Hour)
133
+ tab := newRoutingTable(20, convertPeerID(local), time.Hour)
134
var peers []*peer.Peer
135
for i := 0; i < 500; i++ {
136
peers = append(peers, _randPeer())
@@ -167,8 +167,8 @@ func TestTableMultithreaded(t *testing.T) {
167
168
func BenchmarkUpdates(b *testing.B) {
169
b.StopTimer()
170
- local := ConvertKey("localKey")
171
- tab := NewRoutingTable(20, local, time.Hour)
170
+ local := convertKey("localKey")
171
+ tab := newRoutingTable(20, local, time.Hour)
172
173
var peers []*peer.Peer
174
for i := 0; i < b.N; i++ {
@@ -183,8 +183,8 @@ func BenchmarkUpdates(b *testing.B) {
183
184
func BenchmarkFinds(b *testing.B) {
185
b.StopTimer()
186
- local := ConvertKey("localKey")
187
- tab := NewRoutingTable(20, local, time.Hour)
186
+ local := convertKey("localKey")
187
+ tab := newRoutingTable(20, local, time.Hour)
188
189
var peers []*peer.Peer
190
for i := 0; i < b.N; i++ {
routing/kbucket/util.go
+9
-9
@@ -20,11 +20,11 @@ var ErrLookupFailure = errors.New("failed to find any peer in table")
20
// peer.ID or a util.Key. This unifies the keyspace
21
type ID []byte
22
23
-func (id ID) Equal(other ID) bool {
23
+func (id ID) equal(other ID) bool {
24
return bytes.Equal(id, other)
25
}
26
27
-func (id ID) Less(other ID) bool {
27
+func (id ID) less(other ID) bool {
28
a, b := equalizeSizes(id, other)
29
for i := 0; i < len(a); i++ {
30
if a[i] != b[i] {
@@ -76,23 +76,23 @@ func equalizeSizes(a, b ID) (ID, ID) {
76
return a, b
77
}
78
79
-func ConvertPeerID(id peer.ID) ID {
79
+func convertPeerID(id peer.ID) ID {
80
hash := sha256.Sum256(id)
81
return hash[:]
82
}
83
84
-func ConvertKey(id u.Key) ID {
84
+func convertKey(id u.Key) ID {
85
hash := sha256.Sum256([]byte(id))
86
return hash[:]
87
}
88
89
-// Returns true if a is closer to key than b is
89
+// Closer returns true if a is closer to key than b is
90
func Closer(a, b peer.ID, key u.Key) bool {
91
- aid := ConvertPeerID(a)
92
- bid := ConvertPeerID(b)
93
- tgt := ConvertKey(key)
91
+ aid := convertPeerID(a)
92
+ bid := convertPeerID(b)
93
+ tgt := convertKey(key)
94
adist := xor(aid, tgt)
95
bdist := xor(bid, tgt)
96
97
- return adist.Less(bdist)
97
+ return adist.less(bdist)
98
}
routing/routing.go
+2
-1
@@ -1,9 +1,10 @@
1
package routing
2
3
import (
4
+ "time"
5
+
6
peer "github.com/jbenet/go-ipfs/peer"
7
u "github.com/jbenet/go-ipfs/util"
6
- "time"
8
)
9
10
// IpfsRouting is the routing module interface