@cryptotaxi247 / kubo / commits / ae1f7688a

separate to ensure sync safety

Juan Batiz-Benet committed Sep 17, 2014 at 03:32 UTC ae1f7688aaf5da6a88aa0a18680c608a4929edf9
1 file changed +42 -34
peer/queue/distance.go
+42 -34
@@ -10,61 +10,69 @@ import (
10 u "github.com/jbenet/go-ipfs/util"
11 )
12
13 -// peerDistance tracks a peer and its distance to something else.
14 -type peerDistance struct {
13 +// peerMetric tracks a peer and its distance to something else.
14 +type peerMetric struct {
15 // the peer
16 peer *peer.Peer
17
18 // big.Int for XOR metric
19 - distance *big.Int
19 + metric *big.Int
20 }
21
22 -// distancePQ implements heap.Interface and PeerQueue
23 -type distancePQ struct {
24 - // from is the Key this PQ measures against
25 - from ks.Key
26 -
27 - // peers is a heap of peerDistance items
28 - peers []*peerDistance
29 -
30 - sync.RWMutex
31 -}
22 +// peerMetricHeap implements a heap of peerDistances
23 +type peerMetricHeap []*peerMetric
24
33 -func (pq *distancePQ) Len() int {
34 - return len(pq.peers)
25 +func (ph peerMetricHeap) Len() int {
26 + return len(ph)
27 }
28
37 -func (pq *distancePQ) Less(i, j int) bool {
38 - return -1 == pq.peers[i].distance.Cmp(pq.peers[j].distance)
29 +func (ph peerMetricHeap) Less(i, j int) bool {
30 + return -1 == ph[i].metric.Cmp(ph[j].metric)
31 }
32
41 -func (pq *distancePQ) Swap(i, j int) {
42 - p := pq.peers
43 - p[i], p[j] = p[j], p[i]
33 +func (ph peerMetricHeap) Swap(i, j int) {
34 + ph[i], ph[j] = ph[j], ph[i]
35 }
36
46 -func (pq *distancePQ) Push(x interface{}) {
47 - item := x.(*peerDistance)
48 - pq.peers = append(pq.peers, item)
37 +func (ph *peerMetricHeap) Push(x interface{}) {
38 + item := x.(*peerMetric)
39 + *ph = append(*ph, item)
40 }
41
51 -func (pq *distancePQ) Pop() interface{} {
52 - old := pq.peers
42 +func (ph *peerMetricHeap) Pop() interface{} {
43 + old := *ph
44 n := len(old)
45 item := old[n-1]
55 - pq.peers = old[0 : n-1]
46 + *ph = old[0 : n-1]
47 return item
48 }
49
50 +// distancePQ implements heap.Interface and PeerQueue
51 +type distancePQ struct {
52 + // from is the Key this PQ measures against
53 + from ks.Key
54 +
55 + // heap is a heap of peerDistance items
56 + heap peerMetricHeap
57 +
58 + sync.RWMutex
59 +}
60 +
61 +func (pq *distancePQ) Len() int {
62 + pq.Lock()
63 + defer pq.Unlock()
64 + return len(pq.heap)
65 +}
66 +
67 func (pq *distancePQ) Enqueue(p *peer.Peer) {
68 pq.Lock()
69 defer pq.Unlock()
70
71 distance := ks.XORKeySpace.Key(p.ID).Distance(pq.from)
72
65 - heap.Push(pq, &peerDistance{
66 - peer: p,
67 - distance: distance,
73 + heap.Push(&pq.heap, &peerMetric{
74 + peer: p,
75 + metric: distance,
76 })
77 }
78
@@ -72,13 +80,13 @@ func (pq *distancePQ) Dequeue() *peer.Peer {
80 pq.Lock()
81 defer pq.Unlock()
82
75 - if len(pq.peers) < 1 {
83 + if len(pq.heap) < 1 {
84 panic("called Dequeue on an empty PeerQueue")
85 // will panic internally anyway, but we can help debug here
86 }
87
80 - o := heap.Pop(pq)
81 - p := o.(*peerDistance)
88 + o := heap.Pop(&pq.heap)
89 + p := o.(*peerMetric)
90 return p.peer
91 }
92
@@ -87,7 +95,7 @@ func (pq *distancePQ) Dequeue() *peer.Peer {
95 // XOR as a metric of distance).
96 func NewXORDistancePQ(fromKey u.Key) PeerQueue {
97 return &distancePQ{
90 - from: ks.XORKeySpace.Key([]byte(fromKey)),
91 - peers: []*peerDistance{},
98 + from: ks.XORKeySpace.Key([]byte(fromKey)),
99 + heap: peerMetricHeap{},
100 }
101 }