@cryptotaxi247 / kubo / commits / 98c3afeec

clean up channel use

Juan Batiz-Benet committed Sep 18, 2014 at 06:30 UTC 98c3afeecf6d01000957cabf9583c5797b191a4e
1 file changed +46 -30
peer/queue/sync.go
+46 -30
@@ -9,50 +9,66 @@ import (
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
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)
18 + cq := &ChanQueue{Queue: pq}
19 + cq.process(ctx)
20 return cq
21 }
22
23 func (cq *ChanQueue) process(ctx context.Context) {
28 - var next *peer.Peer
24
30 - for {
25 + // construct the channels here to be able to use them bidirectionally
26 + enqChan := make(chan *peer.Peer, 10)
27 + deqChan := make(chan *peer.Peer, 10)
28
32 - if cq.Queue.Len() == 0 {
33 - select {
34 - case next = <-cq.EnqChan:
35 - case <-ctx.Done():
36 - close(cq.DeqChan)
37 - return
29 + cq.EnqChan = enqChan
30 + cq.DeqChan = deqChan
31 +
32 + go func() {
33 + defer close(deqChan)
34 +
35 + var next *peer.Peer
36 + var item *peer.Peer
37 + var more bool
38 +
39 + for {
40 + if cq.Queue.Len() == 0 {
41 + select {
42 + case next, more = <-enqChan:
43 + if !more {
44 + return
45 + }
46 +
47 + case <-ctx.Done():
48 + return
49 + }
50 +
51 + } else {
52 + next = cq.Queue.Dequeue()
53 }
54
40 - } else {
41 - next = cq.Queue.Dequeue()
42 - }
55 + select {
56 + case item, more = <-enqChan:
57 + if !more {
58 + return
59 + }
60
44 - select {
45 - case item := <-cq.EnqChan:
46 - cq.Queue.Enqueue(item)
47 - cq.Queue.Enqueue(next)
48 - next = nil
61 + cq.Queue.Enqueue(item)
62 + cq.Queue.Enqueue(next)
63 + next = nil
64
50 - case cq.DeqChan <- next:
51 - next = nil
65 + case deqChan <- next:
66 + next = nil
67
53 - case <-ctx.Done():
54 - close(cq.DeqChan)
55 - return
68 + case <-ctx.Done():
69 + return
70 + }
71 }
57 - }
72 +
73 + }()
74 }