-
Notifications
You must be signed in to change notification settings - Fork 12
Expand file tree
/
Copy pathqueue.cpp
More file actions
173 lines (157 loc) · 7.22 KB
/
Copy pathqueue.cpp
File metadata and controls
173 lines (157 loc) · 7.22 KB
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
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
// compat.concurrentqueue — exercise the C++ API the way the library's README
// shows it: enqueue/try_dequeue, bulk ops, a blocking consumer, and a real
// multi-producer/multi-consumer race. Every check asserts observable
// behaviour (order, exact-once delivery, timeout), so a package that links
// but miscompiles the lock-free paths still fails here.
#include <concurrentqueue.h>
#include <blockingconcurrentqueue.h>
#include <atomic>
#include <chrono>
#include <cstdio>
#include <memory>
#include <thread>
#include <utility>
#include <vector>
int main() {
bool ok = true;
auto check = [&](bool cond, const char* what) {
if (!cond) {
std::printf("FAIL: %s\n", what);
ok = false;
}
};
// ---- single-threaded FIFO --------------------------------------------
// The queue is MPMC, not ordered ACROSS producers — but a single producer
// hands its own items to the consumer in order, and that is what we hold
// it to here.
{
moodycamel::ConcurrentQueue<int> q;
check(q.size_approx() == 0, "fresh queue reports size 0");
for (int i = 0; i < 1000; ++i) q.enqueue(i);
check(q.size_approx() == 1000, "size_approx after 1000 enqueues");
int v = -1;
for (int i = 0; i < 1000; ++i) {
if (!q.try_dequeue(v) || v != i) {
check(false, "single-producer FIFO order preserved");
break;
}
}
check(!q.try_dequeue(v), "drained queue reports try_dequeue == false");
check(q.size_approx() == 0, "size_approx back to 0 after draining");
}
// ---- bulk ops ---------------------------------------------------------
{
moodycamel::ConcurrentQueue<int> q;
std::vector<int> in(5000);
for (int i = 0; i < 5000; ++i) in[i] = i;
// enqueue_bulk answers YES/NO (bool); only try_dequeue_bulk answers
// HOW MANY (size_t). Mixing the two up is exactly the kind of mistake
// this test exists to catch — starting with its own author's.
check(q.enqueue_bulk(in.data(), in.size()),
"enqueue_bulk reports everything enqueued");
std::vector<int> out(5000, -1);
check(q.try_dequeue_bulk(out.data(), out.size()) == out.size(),
"try_dequeue_bulk reports everything dequeued");
check(in == out, "bulk round-trip preserves contents and order");
}
// ---- producer token ----------------------------------------------------
// Tokens are the documented fast path for a dedicated producer; exercise
// at least the enqueue side of them.
{
moodycamel::ConcurrentQueue<int> q;
moodycamel::ProducerToken ptok(q);
q.enqueue(ptok, 7);
int v = 0;
check(q.try_dequeue(v) && v == 7, "token enqueue round-trips");
}
// ---- move-only elements ------------------------------------------------
{
moodycamel::ConcurrentQueue<std::unique_ptr<int>> q;
q.enqueue(std::make_unique<int>(41));
q.enqueue(std::make_unique<int>(42));
std::unique_ptr<int> p;
check(q.try_dequeue(p) && p && *p == 41, "move-only element survives the queue");
check(q.try_dequeue(p) && p && *p == 42, "second move-only element comes out in order");
check(!q.try_dequeue(p), "queue drains after move-only elements");
}
// ---- blocking consumer -------------------------------------------------
// wait_dequeue must actually block until the producer lands an item, and
// wait_dequeue_timed on an EMPTY queue must come back false after the
// timeout — the second half is what separates "semaphore posted eagerly"
// from a real wait.
{
moodycamel::BlockingConcurrentQueue<int> q;
int missed = -1;
auto t0 = std::chrono::steady_clock::now();
bool got = q.wait_dequeue_timed(missed, std::chrono::milliseconds(50));
auto elapsed = std::chrono::steady_clock::now() - t0;
check(!got, "wait_dequeue_timed on an empty queue times out");
check(elapsed >= std::chrono::milliseconds(45),
"the timeout was actually waited, not spun through");
std::thread producer([&q] {
for (int i = 0; i < 2000; ++i) q.enqueue(i);
});
int prev = -1;
for (int i = 0; i < 2000; ++i) {
int v = -1;
// wait_dequeue returns void in 1.0.5 (only the _timed spelling
// reports success), so the assertion is on the VALUE alone: if
// the semaphore path were broken this call would either hang
// (test timeout) or hand back a value out of sequence.
q.wait_dequeue(v);
if (v != prev + 1) {
check(false, "blocking consumer receives every item in order");
break;
}
prev = v;
}
producer.join();
int drained = -1;
check(!q.wait_dequeue_timed(drained, std::chrono::milliseconds(50)),
"nothing left after the consumer caught up");
}
// ---- MPMC: exact-once delivery under contention -------------------------
// 4 producers x 10000 unique values, 4 consumers racing them out. Every
// value must arrive EXACTLY once: a lost slot, a duplicated slot, or an
// ABA-shaped torn read all show up as a counter != 1. This is the check
// the single-threaded assertions cannot make.
{
constexpr int kProducers = 4, kConsumers = 4, kPerProducer = 10000;
constexpr int kTotal = kProducers * kPerProducer;
moodycamel::ConcurrentQueue<int> q;
std::vector<std::atomic<char>> seen(kTotal);
for (auto& s : seen) s.store(0, std::memory_order_relaxed);
std::atomic<int> dequeued{0};
std::vector<std::thread> producers, consumers;
for (int p = 0; p < kProducers; ++p) {
producers.emplace_back([&q, p] {
for (int i = 0; i < kPerProducer; ++i) q.enqueue(p * kPerProducer + i);
});
}
for (int c = 0; c < kConsumers; ++c) {
consumers.emplace_back([&] {
int v;
// try_dequeue until the producers are done AND the queue is
// empty — the consumer that would otherwise spin on an empty
// queue while producers are still working is the whole point
// of a concurrent queue.
for (;;) {
while (q.try_dequeue(v)) {
seen[v].store(1, std::memory_order_relaxed);
dequeued.fetch_add(1, std::memory_order_relaxed);
}
if (dequeued.load(std::memory_order_relaxed) >= kTotal) return;
std::this_thread::yield();
}
});
}
for (auto& t : producers) t.join();
for (auto& t : consumers) t.join();
check(dequeued.load() == kTotal, "every enqueued value was dequeued");
int exactly_once = 0;
for (const auto& s : seen) exactly_once += (s.load() == 1) ? 1 : 0;
check(exactly_once == kTotal, "every value delivered exactly once (no loss, no duplicate)");
}
std::printf("%s\n", ok ? "concurrentqueue: all checks passed" : "concurrentqueue: FAILED");
return ok ? 0 : 1;
}