chan queue
Juan Batiz-Benet committed
Sep 17, 2014 at 05:09 UTC
551c40930ec28dd967b37df5973b5cf2450335a1
2 files changed
+117
-1
peer/queue/queue_test.go
+59
-1
@@ -1,17 +1,21 @@
1
package queue
2
3
import (
4
+ "fmt"
5
"testing"
6
+ "time"
7
8
peer "github.com/jbenet/go-ipfs/peer"
9
u "github.com/jbenet/go-ipfs/util"
10
+
11
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
12
)
13
14
func newPeer(id string) *peer.Peer {
15
return &peer.Peer{ID: peer.ID(id)}
16
}
17
14
-func TestPeerstore(t *testing.T) {
18
+func TestQueue(t *testing.T) {
19
20
p1 := newPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31")
21
p2 := newPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a32")
@@ -60,3 +64,57 @@ func TestPeerstore(t *testing.T) {
64
}
65
66
}
67
+
68
+func newPeerTime(t time.Time) *peer.Peer {
69
+ s := fmt.Sprintf("hmmm time: %v", t)
70
+ h, _ := u.Hash([]byte(s))
71
+ return &peer.Peer{ID: peer.ID(h)}
72
+}
73
+
74
+func TestSyncQueue(t *testing.T) {
75
+ ctx, _ := context.WithTimeout(context.Background(), time.Second*2)
76
+
77
+ pq := NewXORDistancePQ(u.Key("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31"))
78
+ cq := NewChanQueue(ctx, pq)
79
+ countIn := 0
80
+ countOut := 0
81
+
82
+ produce := func() {
83
+ tick := time.Tick(time.Millisecond)
84
+ for {
85
+ select {
86
+ case tim := <-tick:
87
+ countIn++
88
+ cq.EnqChan <- newPeerTime(tim)
89
+ case <-ctx.Done():
90
+ return
91
+ }
92
+ }
93
+ }
94
+
95
+ consume := func() {
96
+ for {
97
+ select {
98
+ case <-cq.DeqChan:
99
+ countOut++
100
+ case <-ctx.Done():
101
+ return
102
+ }
103
+ }
104
+ }
105
+
106
+ for i := 0; i < 10; i++ {
107
+ go produce()
108
+ go produce()
109
+ go consume()
110
+ }
111
+
112
+ select {
113
+ case <-ctx.Done():
114
+ }
115
+
116
+ if countIn != countOut {
117
+ t.Errorf("didnt get them all out: %d/%d", countOut, countIn)
118
+ }
119
+
120
+}
peer/queue/sync.go
new
+58
@@ -0,0 +1,58 @@
1
+package queue
2
+
3
+import (
4
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
+
6
+ peer "github.com/jbenet/go-ipfs/peer"
7
+)
8
+
9
+// ChanQueue makes any PeerQueue synchronizable through channels.
10
+type ChanQueue struct {
11
+ Queue PeerQueue
12
+ EnqChan chan *peer.Peer
13
+ DeqChan chan *peer.Peer
14
+}
15
+
16
+// NewChanQueue creates a ChanQueue by wrapping pq.
17
+func NewChanQueue(ctx context.Context, pq PeerQueue) *ChanQueue {
18
+ cq := &ChanQueue{
19
+ Queue: pq,
20
+ EnqChan: make(chan *peer.Peer, 10),
21
+ DeqChan: make(chan *peer.Peer, 10),
22
+ }
23
+ go cq.process(ctx)
24
+ return cq
25
+}
26
+
27
+func (cq *ChanQueue) process(ctx context.Context) {
28
+ var next *peer.Peer
29
+
30
+ for {
31
+
32
+ if cq.Queue.Len() == 0 {
33
+ select {
34
+ case next = <-cq.EnqChan:
35
+ case <-ctx.Done():
36
+ close(cq.DeqChan)
37
+ return
38
+ }
39
+
40
+ } else {
41
+ next = cq.Queue.Dequeue()
42
+ }
43
+
44
+ select {
45
+ case item := <-cq.EnqChan:
46
+ cq.Queue.Enqueue(item)
47
+ cq.Queue.Enqueue(next)
48
+ next = nil
49
+
50
+ case cq.DeqChan <- next:
51
+ next = nil
52
+
53
+ case <-ctx.Done():
54
+ close(cq.DeqChan)
55
+ return
56
+ }
57
+ }
58
+}