PeerQueue (based on XOR distance metric)
Juan Batiz-Benet committed
Sep 17, 2014 at 02:58 UTC
a21c1b6b6229e8fc01defca2f57141194d38aa1a
3 files changed
+161
peer/queue/distance.go
new
+84
@@ -0,0 +1,84 @@
1
+package queue
2
+
3
+import (
4
+ "container/heap"
5
+ "math/big"
6
+
7
+ peer "github.com/jbenet/go-ipfs/peer"
8
+ ks "github.com/jbenet/go-ipfs/routing/keyspace"
9
+ u "github.com/jbenet/go-ipfs/util"
10
+)
11
+
12
+// peerDistance tracks a peer and its distance to something else.
13
+type peerDistance struct {
14
+ // the peer
15
+ peer *peer.Peer
16
+
17
+ // big.Int for XOR metric
18
+ distance *big.Int
19
+}
20
+
21
+// distancePQ implements heap.Interface and PeerQueue
22
+type distancePQ struct {
23
+ // from is the Key this PQ measures against
24
+ from ks.Key
25
+
26
+ // peers is a heap of peerDistance items
27
+ peers []*peerDistance
28
+}
29
+
30
+func (pq distancePQ) Len() int {
31
+ return len(pq.peers)
32
+}
33
+
34
+func (pq distancePQ) Less(i, j int) bool {
35
+ return -1 == pq.peers[i].distance.Cmp(pq.peers[j].distance)
36
+}
37
+
38
+func (pq distancePQ) Swap(i, j int) {
39
+ p := pq.peers
40
+ p[i], p[j] = p[j], p[i]
41
+}
42
+
43
+func (pq *distancePQ) Push(x interface{}) {
44
+ item := x.(*peerDistance)
45
+ pq.peers = append(pq.peers, item)
46
+}
47
+
48
+func (pq *distancePQ) Pop() interface{} {
49
+ old := pq.peers
50
+ n := len(old)
51
+ item := old[n-1]
52
+ pq.peers = old[0 : n-1]
53
+ return item
54
+}
55
+
56
+func (pq *distancePQ) Enqueue(p *peer.Peer) {
57
+ distance := ks.XORKeySpace.Key(p.ID).Distance(pq.from)
58
+
59
+ heap.Push(pq, &peerDistance{
60
+ peer: p,
61
+ distance: distance,
62
+ })
63
+}
64
+
65
+func (pq *distancePQ) Dequeue() *peer.Peer {
66
+ if len(pq.peers) < 1 {
67
+ panic("called Dequeue on an empty PeerQueue")
68
+ // will panic internally anyway, but we can help debug here
69
+ }
70
+
71
+ o := heap.Pop(pq)
72
+ p := o.(*peerDistance)
73
+ return p.peer
74
+}
75
+
76
+// NewXORDistancePQ returns a PeerQueue which maintains its peers sorted
77
+// in terms of their distances to each other in an XORKeySpace (i.e. using
78
+// XOR as a metric of distance).
79
+func NewXORDistancePQ(fromKey u.Key) PeerQueue {
80
+ return &distancePQ{
81
+ from: ks.XORKeySpace.Key([]byte(fromKey)),
82
+ peers: []*peerDistance{},
83
+ }
84
+}
peer/queue/interface.go
new
+15
@@ -0,0 +1,15 @@
1
+package queue
2
+
3
+import peer "github.com/jbenet/go-ipfs/peer"
4
+
5
+// PeerQueue maintains a set of peers ordered according to a metric.
6
+// Implementations of PeerQueue could order peers based on distances along
7
+// a KeySpace, latency measurements, trustworthiness, reputation, etc.
8
+type PeerQueue interface {
9
+
10
+ // Enqueue adds this node to the queue.
11
+ Enqueue(*peer.Peer)
12
+
13
+ // Dequeue retrieves the highest (smallest int) priority node
14
+ Dequeue() *peer.Peer
15
+}
peer/queue/queue_test.go
new
+62
@@ -0,0 +1,62 @@
1
+package queue
2
+
3
+import (
4
+ "testing"
5
+
6
+ peer "github.com/jbenet/go-ipfs/peer"
7
+ u "github.com/jbenet/go-ipfs/util"
8
+)
9
+
10
+func newPeer(id string) *peer.Peer {
11
+ return &peer.Peer{ID: peer.ID(id)}
12
+}
13
+
14
+func TestPeerstore(t *testing.T) {
15
+
16
+ p1 := newPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31")
17
+ p2 := newPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a32")
18
+ p3 := newPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a33")
19
+ p4 := newPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a34")
20
+ p5 := newPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31")
21
+
22
+ // these are the peer.IDs' XORKeySpace Key values:
23
+ // [228 47 151 130 156 102 222 232 218 31 132 94 170 208 80 253 120 103 55 35 91 237 48 157 81 245 57 247 66 150 9 40]
24
+ // [26 249 85 75 54 49 25 30 21 86 117 62 85 145 48 175 155 194 210 216 58 14 241 143 28 209 129 144 122 28 163 6]
25
+ // [78 135 26 216 178 181 224 181 234 117 2 248 152 115 255 103 244 34 4 152 193 88 9 225 8 127 216 158 226 8 236 246]
26
+ // [125 135 124 6 226 160 101 94 192 57 39 12 18 79 121 140 190 154 147 55 44 83 101 151 63 255 94 179 51 203 241 51]
27
+
28
+ pq := NewXORDistancePQ(u.Key("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31"))
29
+ pq.Enqueue(p3)
30
+ pq.Enqueue(p1)
31
+ pq.Enqueue(p2)
32
+ pq.Enqueue(p4)
33
+ pq.Enqueue(p5)
34
+ pq.Enqueue(p1)
35
+
36
+ // should come out as: p1, p4, p3, p2
37
+
38
+ if d := pq.Dequeue(); d != p1 && d != p5 {
39
+ t.Error("ordering failed")
40
+ }
41
+
42
+ if d := pq.Dequeue(); d != p1 && d != p5 {
43
+ t.Error("ordering failed")
44
+ }
45
+
46
+ if d := pq.Dequeue(); d != p1 && d != p5 {
47
+ t.Error("ordering failed")
48
+ }
49
+
50
+ if pq.Dequeue() != p4 {
51
+ t.Error("ordering failed")
52
+ }
53
+
54
+ if pq.Dequeue() != p3 {
55
+ t.Error("ordering failed")
56
+ }
57
+
58
+ if pq.Dequeue() != p2 {
59
+ t.Error("ordering failed")
60
+ }
61
+
62
+}