sync safety to pq
Juan Batiz-Benet committed
Sep 17, 2014 at 03:16 UTC
51eeec1a79973ed38faa91875f1ea0ea78377129
1 file changed
+12
-3
peer/queue/distance.go
+12
-3
@@ -3,6 +3,7 @@ package queue
3
import (
4
"container/heap"
5
"math/big"
6
+ "sync"
7
8
peer "github.com/jbenet/go-ipfs/peer"
9
ks "github.com/jbenet/go-ipfs/routing/keyspace"
@@ -25,17 +26,19 @@ type distancePQ struct {
26
27
// peers is a heap of peerDistance items
28
peers []*peerDistance
29
+
30
+ sync.RWMutex
31
}
32
30
-func (pq distancePQ) Len() int {
33
+func (pq *distancePQ) Len() int {
34
return len(pq.peers)
35
}
36
34
-func (pq distancePQ) Less(i, j int) bool {
37
+func (pq *distancePQ) Less(i, j int) bool {
38
return -1 == pq.peers[i].distance.Cmp(pq.peers[j].distance)
39
}
40
38
-func (pq distancePQ) Swap(i, j int) {
41
+func (pq *distancePQ) Swap(i, j int) {
42
p := pq.peers
43
p[i], p[j] = p[j], p[i]
44
}
@@ -54,6 +57,9 @@ func (pq *distancePQ) Pop() interface{} {
57
}
58
59
func (pq *distancePQ) Enqueue(p *peer.Peer) {
60
+ pq.Lock()
61
+ defer pq.Unlock()
62
+
63
distance := ks.XORKeySpace.Key(p.ID).Distance(pq.from)
64
65
heap.Push(pq, &peerDistance{
@@ -63,6 +69,9 @@ func (pq *distancePQ) Enqueue(p *peer.Peer) {
69
}
70
71
func (pq *distancePQ) Dequeue() *peer.Peer {
72
+ pq.Lock()
73
+ defer pq.Unlock()
74
+
75
if len(pq.peers) < 1 {
76
panic("called Dequeue on an empty PeerQueue")
77
// will panic internally anyway, but we can help debug here