blob: 7d56fe0e47e2ee0e7bac5cff44553bef99bc4e36 (
plain) (
blame)
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
|
#include <doctest/doctest.h>
#include <atomic>
#include <chrono>
#include <thread>
#include "packeteer/capture_queue.hpp"
using namespace packeteer;
TEST_CASE("try_push/pop returns packets in FIFO order") {
CaptureQueue queue(4);
for (std::uint32_t i = 0; i < 3; ++i) {
CapturedPacket p;
p.ts_sec = i;
p.data = {static_cast<unsigned char>(i)};
CHECK(queue.try_push(std::move(p)));
}
for (std::uint32_t i = 0; i < 3; ++i) {
auto p = queue.pop();
REQUIRE(p.has_value());
CHECK(p->ts_sec == i);
}
CHECK(queue.dropped() == 0);
}
TEST_CASE("try_push drops and counts once the queue is full") {
CaptureQueue queue(2);
CapturedPacket a, b, c;
CHECK(queue.try_push(std::move(a)));
CHECK(queue.try_push(std::move(b)));
CHECK_FALSE(queue.try_push(std::move(c))); // full: dropped, not blocked
CHECK(queue.dropped() == 1);
}
TEST_CASE("stop() drains items already queued before pop() returns nullopt") {
CaptureQueue queue(4);
CapturedPacket a, b;
queue.try_push(std::move(a));
queue.try_push(std::move(b));
queue.stop();
CHECK(queue.pop().has_value());
CHECK(queue.pop().has_value());
CHECK_FALSE(queue.pop().has_value()); // drained and stopped
}
TEST_CASE("pop() blocks until a packet is pushed") {
CaptureQueue queue(4);
std::thread producer([&] {
std::this_thread::sleep_for(std::chrono::milliseconds(50));
CapturedPacket p;
p.data = {0x42};
queue.try_push(std::move(p));
});
auto p = queue.pop();
REQUIRE(p.has_value());
CHECK(p->data[0] == 0x42);
producer.join();
}
TEST_CASE("push() succeeds immediately when there's room") {
CaptureQueue queue(4);
CapturedPacket p;
p.data = {0x01};
CHECK(queue.push(std::move(p)));
CHECK(queue.dropped() == 0);
}
TEST_CASE("push() blocks for room instead of dropping, unlike try_push()") {
CaptureQueue queue(1);
CapturedPacket a;
a.data = {0xAA};
CHECK(queue.try_push(std::move(a))); // fills the only slot
std::atomic<bool> pushed{false};
std::thread producer([&] {
CapturedPacket b;
b.data = {0xBB};
CHECK(queue.push(std::move(b))); // must block until the pop() below frees room
pushed.store(true);
});
std::this_thread::sleep_for(std::chrono::milliseconds(50));
CHECK_FALSE(pushed.load()); // still blocked: queue was full this whole time
auto first = queue.pop(); // frees a slot
REQUIRE(first.has_value());
CHECK(first->data[0] == 0xAA);
producer.join();
CHECK(pushed.load());
CHECK(queue.dropped() == 0); // never dropped - it waited instead
auto second = queue.pop();
REQUIRE(second.has_value());
CHECK(second->data[0] == 0xBB);
}
TEST_CASE("push() returns false without pushing if stop() is called while it's waiting") {
CaptureQueue queue(1);
CapturedPacket a;
a.data = {0xAA};
queue.try_push(std::move(a)); // fill the only slot
std::atomic<bool> result_ready{false};
std::atomic<bool> push_result{true};
std::thread producer([&] {
CapturedPacket b;
b.data = {0xBB};
push_result.store(queue.push(std::move(b)));
result_ready.store(true);
});
std::this_thread::sleep_for(std::chrono::milliseconds(50));
queue.stop();
producer.join();
CHECK(result_ready.load());
CHECK_FALSE(push_result.load());
}
|