-
Notifications
You must be signed in to change notification settings - Fork 682
/
fake_queue.go
105 lines (95 loc) · 2.43 KB
/
fake_queue.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
package entrypoint
import (
"context"
"encoding/json"
"fmt"
"strings"
"sync"
"testing"
"time"
"github.com/stretchr/testify/assert"
)
// The Queue struct implements a multi-writer/multi-reader concurrent queue where the dequeue
// operation (the Get() method) takes a predicate that allows it to skip past queue entries until it
// finds one that satisfies the specified predicate.
type Queue struct {
T *testing.T
timeout time.Duration
cond *sync.Cond
entries []interface{}
offset int
}
// NewQueue constructs a new queue with the supplied timeout.
func NewQueue(t *testing.T, timeout time.Duration) *Queue {
q := &Queue{
T: t,
timeout: timeout,
cond: sync.NewCond(&sync.Mutex{}),
}
ctx, cancel := context.WithCancel(context.Background())
t.Cleanup(cancel)
// Broadcast on Queue.cond every three seconds so that anyone waiting on the condition has a
// chance to timeout. (Go doesn't support timed wait on conditions.)
go func() {
ticker := time.NewTicker(3 * time.Second)
for {
select {
case <-ticker.C:
q.cond.Broadcast()
case <-ctx.Done():
return
}
}
}()
return q
}
// Add an entry to the queue.
func (q *Queue) Add(obj interface{}) {
q.cond.L.Lock()
defer q.cond.L.Unlock()
q.entries = append(q.entries, obj)
q.cond.Broadcast()
}
// Get will return the next entry that satisfies the supplied predicate.
func (q *Queue) Get(predicate func(interface{}) bool) interface{} {
q.T.Helper()
start := time.Now()
q.cond.L.Lock()
defer q.cond.L.Unlock()
for {
for idx, obj := range q.entries[q.offset:] {
if predicate(obj) {
q.offset += idx + 1
return obj
}
}
if time.Since(start) > q.timeout {
msg := &strings.Builder{}
for idx, entry := range q.entries {
bytes, err := json.MarshalIndent(entry, "", " ")
if err != nil {
panic(err)
}
var extra string
if idx < q.offset {
extra = "(Before Offset)"
} else if idx == q.offset {
extra = "(Offset Here)"
} else {
extra = "(After Offset)"
}
msg.WriteString(fmt.Sprintf("\n--- Queue Entry[%d] %s---\n%s\n", idx, extra, string(bytes)))
}
q.T.Fatal(fmt.Sprintf("Get timed out!\n%s", msg))
}
q.cond.Wait()
}
}
// AssertEmpty will check that the queue remains empty for the supplied duration.
func (q *Queue) AssertEmpty(timeout time.Duration, msg string) {
q.T.Helper()
time.Sleep(timeout)
q.cond.L.Lock()
defer q.cond.L.Unlock()
assert.Empty(q.T, q.entries, msg)
}