gstreamer
Loading...
Searching...
No Matches
sro_crate_example.cc
Go to the documentation of this file.
1#include "sro/gSROCrate.h"
2
3#include <algorithm>
4#include <atomic>
5#include <chrono>
6#include <condition_variable>
7#include <future>
8#include <iostream>
9#include <mutex>
10#include <stdexcept>
11#include <string>
12#include <thread>
13#include <tuple>
14#include <utility>
15#include <vector>
16
17using namespace std::chrono_literals;
18
19namespace {
20void require(bool condition, const char* message) {
21 if (!condition) { throw std::runtime_error(message); }
22}
23
24template<class F>
25void expect_error(F&& action, const std::string& message) {
26 try { action(); }
27 catch (const std::exception& error) {
28 require(std::string(error.what()).find(message) != std::string::npos, "Unexpected exception");
29 return;
30 }
31 throw std::runtime_error("Expected an exception: " + message);
32}
33
34std::atomic<int> live_data{0};
35struct Data final : GSROData {
36 Data() { ++live_data; }
37 ~Data() override { --live_data; }
38 std::size_t size_bytes() const noexcept override { return sizeof(*this); }
39};
40
41GSROPayload payload(std::int64_t time, GSROEventId event = 0, std::uint64_t sequence = 0) {
42 return {7, event, sequence, GSROTime{time}, std::make_unique<Data>()};
43}
44
45using Key = std::tuple<GSROTime, GSROEventId, std::uint64_t>;
46struct State {
47 std::mutex mutex;
48 std::condition_variable changed;
49 std::thread::id owner;
50 bool wrong_thread = false;
51 std::vector<Key> consumed;
52 std::vector<std::string> lifecycle;
53 std::optional<GSROTime> boundary;
54 std::optional<GSROEndContext> end;
55 std::atomic<int> delivered{0};
56 std::string fail_at;
57
58 void record(const std::string& operation) {
59 std::lock_guard lock(mutex);
60 wrong_thread |= owner != std::this_thread::get_id();
61 lifecycle.push_back(operation);
62 }
63 void maybe_fail(const std::string& operation) const {
64 if (fail_at == operation) { throw std::runtime_error("synthetic " + operation + " failure"); }
65 }
66};
67
68struct Sink final : GSROFrameSink {
69 explicit Sink(std::shared_ptr<State> s) : state(std::move(s)) { state->record("sink-created"); }
70 ~Sink() override { state->record("sink-destroyed"); }
71 void write_frame(GSROFrame frame) override {
72 require(frame.crate_id == 7 && frame.data != nullptr, "Invalid emitted frame");
73 state->record("write");
74 state->maybe_fail("write");
75 }
76 void finish_output() override {
77 state->record("close");
78 state->maybe_fail("close");
79 }
80 std::shared_ptr<State> state;
81};
82
83// Deliberately simple implementation: retain contributions until advance_time; discard its final partial data.
84struct Plugin final : GSROCratePlugin {
85 Plugin(GSROFrameSink& sink, std::shared_ptr<State> s) : GSROCratePlugin(sink), state(std::move(s)) {
86 state->record("plugin-created");
87 }
88 ~Plugin() override { state->record("plugin-destroyed"); }
89 void consume_payload(GSROPayload item) override {
90 state->record("consume");
91 state->maybe_fail("consume");
92 {
93 std::lock_guard lock(state->mutex);
94 state->consumed.emplace_back(item.time, item.event_id, item.sequence);
95 }
96 pending.push_back(std::move(item));
97 }
98 void advance_time(GSROTime time) override {
99 state->record("advance");
100 state->maybe_fail("advance");
101 for (auto& item : pending) {
102 output.write_frame({item.crate_id, item.sequence, item.time, time, std::move(item.data)});
103 }
104 pending.clear();
105 {
106 std::lock_guard lock(state->mutex);
107 state->boundary = time;
108 }
109 state->changed.notify_all();
110 }
111 void finish_run(const GSROEndContext& context) override {
112 state->record("finish");
113 state->maybe_fail("finish");
114 state->end = context;
115 pending.clear();
116 }
117 std::shared_ptr<State> state;
118 std::vector<GSROPayload> pending;
119};
120
121GSROCrate::ResourceFactory factory(const std::shared_ptr<State>& state) {
122 return [state] {
123 state->owner = std::this_thread::get_id();
124 state->maybe_fail("setup");
125 auto sink = std::make_unique<Sink>(state);
126 auto plugin = std::make_unique<Plugin>(*sink, state);
127 return GSROCrateResources{std::move(sink), std::move(plugin)};
128 };
129}
130
131void test_ordering_and_lifecycle() {
132 auto state = std::make_shared<State>();
133 GSROCrate crate(7, factory(state), {}, [state](GSROEventId, std::uint64_t) { ++state->delivered; });
134 std::vector<std::future<void>> workers;
135 std::vector<Key> expected;
136 for (GSROEventId event = 0; event < 4; ++event) {
137 for (int seq = 0; seq < 32; ++seq) { expected.emplace_back(GSROTime{seq % 5 - 2}, event, seq); }
138 workers.push_back(std::async(std::launch::async, [&, event] {
139 for (int seq = 31; seq >= 0; --seq) { crate.enqueue_payload(payload(seq % 5 - 2, event, seq)); }
140 }));
141 }
142 for (auto& worker : workers) { worker.get(); }
143 crate.advance_time(GSROTime{0});
144 {
145 std::unique_lock lock(state->mutex);
146 require(state->changed.wait_for(lock, 5s, [&] { return state->boundary == GSROTime{0}; }),
147 "Crate did not process its boundary");
148 require(state->consumed.size() == 56, "Boundary equality or negative-time ordering is incorrect");
149 for (const auto& key : state->consumed) { require(std::get<0>(key) < GSROTime{0}, "Early release"); }
150 }
151 crate.advance_time(GSROTime{0});
152 expect_error([&] { crate.advance_time(GSROTime{-1}); }, "must not decrease");
153 expect_error([&] { crate.enqueue_payload(payload(-1)); }, "precedes");
154 expect_error([&] { crate.finish_and_join({GSROEndReason::completed, GSROTime{-1}}); }, "precedes");
155 crate.finish_and_join({GSROEndReason::completed, GSROTime{2}});
156 crate.finish_and_join({GSROEndReason::completed, std::nullopt});
157 std::sort(expected.begin(), expected.end());
158 require(state->consumed == expected && state->delivered == 128, "Lost, duplicated, or reordered payload");
159 require(state->end->safe_time == GSROTime{2}, "Final progress was lost");
160 require(!state->wrong_thread && state->owner != std::this_thread::get_id(), "Wrong callback thread");
161 const std::vector<std::string> ending{"finish", "close", "plugin-destroyed", "sink-destroyed"};
162 require(std::equal(ending.begin(), ending.end(), state->lifecycle.end() - 4), "Wrong shutdown order");
163 require(std::count(state->lifecycle.begin(), state->lifecycle.end(), "write") == 104,
164 "Implementation did not retain and discard the final partial input");
165 expect_error([&] { crate.enqueue_payload(payload(3)); }, "closed");
166}
167
168void test_blocked_producer(bool fail_setup, bool close_while_blocked, bool bound_by_bytes = false) {
169 auto state = std::make_shared<State>();
170 std::promise<void> release;
171 auto gate = release.get_future().share();
172 auto setup = factory(state);
173 GSROCrateLimits limits;
174 if (bound_by_bytes) { limits.queue_bytes = sizeof(GSROPayload) + sizeof(Data); }
175 else { limits.queue_messages = 1; }
176 GSROCrate crate(7, [=] {
177 gate.wait();
178 if (fail_setup) { throw std::runtime_error("synthetic setup failure"); }
179 return setup();
180 }, limits);
181 crate.enqueue_payload(payload(1));
182 std::promise<void> started;
183 auto started_future = started.get_future();
184 auto producer = std::async(std::launch::async, [&] {
185 started.set_value();
186 try { crate.enqueue_payload(payload(2)); }
187 catch (const std::exception& error) { return std::string(error.what()); }
188 return std::string{};
189 });
190 started_future.wait();
191 const bool blocked = producer.wait_for(30ms) == std::future_status::timeout;
192 std::future<void> finishing;
193 if (close_while_blocked) {
194 finishing = std::async(std::launch::async, [&] {
195 crate.finish_and_join({GSROEndReason::interrupted, std::nullopt});
196 });
197 }
198 const bool woke_on_close = !close_while_blocked || producer.wait_for(5s) == std::future_status::ready;
199 release.set_value();
200 const auto error = producer.get();
201 if (finishing.valid()) { finishing.get(); }
202 if (fail_setup) {
203 expect_error([&] { crate.finish_and_join({GSROEndReason::completed, std::nullopt}); }, "setup failure");
204 require(error.find("setup failure") != std::string::npos, "Blocked producer missed writer failure");
205 }
206 else {
207 crate.finish_and_join({GSROEndReason::completed, std::nullopt});
208 require(close_while_blocked ? error.find("closed") != std::string::npos : error.empty(),
209 "Unexpected blocked producer result");
210 require(state->consumed.size() == (close_while_blocked ? 1u : 2u),
211 "Shutdown lost accepted input");
212 }
213 require(blocked && woke_on_close, "Queue backpressure or shutdown wakeup failed");
214}
215
216void test_failures_and_limits() {
217 for (const std::string stage : {"setup", "consume", "advance", "write", "finish", "close", "delivery"}) {
218 auto state = std::make_shared<State>();
219 state->fail_at = stage;
220 GSROCrate crate(7, factory(state), {}, [state](GSROEventId, std::uint64_t) {
221 state->maybe_fail("delivery");
222 });
223 try {
224 crate.enqueue_payload(payload(0));
225 crate.advance_time(GSROTime{1});
226 }
227 catch (const std::runtime_error&) { /* An asynchronous failure can already reach a producer. */ }
228 expect_error([&] { crate.finish_and_join({GSROEndReason::completed, std::nullopt}); },
229 "synthetic " + stage + " failure");
230 expect_error([&] { crate.rethrow_if_failed(); }, "synthetic " + stage + " failure");
231 }
232 for (const bool byte_limit : {false, true}) {
233 auto state = std::make_shared<State>();
234 GSROCrateLimits limits;
235 if (byte_limit) { limits.pending_bytes = sizeof(GSROPayload) + sizeof(Data); }
236 else { limits.pending_payloads = 1; }
237 GSROCrate crate(7, factory(state), limits);
238 crate.enqueue_payload(payload(1));
239 crate.enqueue_payload(payload(2));
240 expect_error([&] { crate.finish_and_join({GSROEndReason::completed, std::nullopt}); },
241 "pending input limit");
242 }
243 auto state = std::make_shared<State>();
244 GSROCrateLimits limits;
245 limits.queue_bytes = 1;
246 GSROCrate crate(7, factory(state), limits);
247 expect_error([&] { crate.enqueue_payload(payload(0)); }, "byte limit");
248 auto wrong_crate = payload(0);
249 wrong_crate.crate_id = 8;
250 expect_error([&] { crate.enqueue_payload(std::move(wrong_crate)); }, "another crate");
251 expect_error([&] { crate.enqueue_payload({7, 0, 0, GSROTime{0}, nullptr}); }, "null");
252 crate.finish_and_join({GSROEndReason::completed, std::nullopt});
253 limits.queue_messages = 0;
254 expect_error([&] { GSROCrate invalid(7, factory(state), limits); }, "positive buffer limits");
255 GSROCrate invalid_resources(7, [] { return GSROCrateResources{}; });
256 expect_error([&] { invalid_resources.finish_and_join({GSROEndReason::completed, std::nullopt}); },
257 "null plugin or sink");
258 GSROCrate duplicates(7, factory(state));
259 duplicates.enqueue_payload(payload(1));
260 duplicates.enqueue_payload(payload(1));
261 expect_error([&] { duplicates.finish_and_join({GSROEndReason::completed, std::nullopt}); }, "Duplicate");
262}
263
264void test_destructor_shutdown() {
265 auto state = std::make_shared<State>();
266 {
267 GSROCrate crate(7, factory(state));
268 crate.enqueue_payload(payload(1));
269 crate.advance_time(GSROTime{1});
270 }
271 require(state->consumed.size() == 1 && state->end->reason == GSROEndReason::interrupted,
272 "Destructor did not drain and finalize");
273 require(state->end->safe_time == GSROTime{1}, "Destructor lost the last safe boundary");
274 require(!state->wrong_thread, "Destructor changed plugin ownership thread");
275 state = std::make_shared<State>();
276 state->fail_at = "close";
277 {
278 GSROCrate crate(7, factory(state));
279 }
280 require(state->lifecycle.back() == "sink-destroyed", "Destructor failed to clean up after close failure");
281}
282} // namespace
283
284int main() {
285 try {
286 test_ordering_and_lifecycle();
287 test_blocked_producer(false, false);
288 test_blocked_producer(false, false, true);
289 test_blocked_producer(false, true);
290 test_blocked_producer(true, false);
291 test_failures_and_limits();
292 test_destructor_shutdown();
293 require(live_data == 0, "Payload ownership leaked");
294 std::cout << "SRO crate threading, ordering, limits, and lifecycle checks passed\n";
295 return 0;
296 }
297 catch (const std::exception& error) {
298 std::cerr << error.what() << '\n';
299 return 1;
300 }
301}
Implementation-owned framing and optional electronics response for one crate.
Definition gSROPlugin.h:79
One crate's bounded input queue, ordering buffer, and background processing thread.
Definition gSROCrate.h:72
void enqueue_payload(GSROPayload payload)
Transfers ownership. Rejects null data, wrong crates, oversized input, and timestamps before a bounda...
Definition gSROCrate.cc:233
void advance_time(GSROTime safe_time)
Queue a nondecreasing boundary after all earlier input has been submitted. Equality is permitted.
Definition gSROCrate.cc:234
void rethrow_if_failed() const
Check for a background failure without waiting. Failure also wakes and rejects queued producers.
Definition gSROCrate.cc:238
std::function< GSROCrateResources()> ResourceFactory
Definition gSROCrate.h:74
void finish_and_join(GSROEndContext context)
Definition gSROCrate.cc:237
Base for implementation-defined payload and frame contents.
Definition gSROData.h:31
Output destination called exclusively by its owning crate thread.
Definition gSROPlugin.h:48
std::chrono::duration< std::int64_t, std::nano > GSROTime
Definition gSROData.h:20
std::uint64_t GSROEventId
Definition gSROData.h:19
frame
event
int main()
Per-crate bounds preventing asynchronous delivery from consuming unlimited memory.
Definition gSROCrate.h:29
std::size_t queue_messages
Maximum queued payload/progress messages; producers wait when the queue is full.
Definition gSROCrate.h:31
std::size_t queue_bytes
Maximum accounted payload bytes in the input queue; producers wait until their payload fits.
Definition gSROCrate.h:33
std::size_t pending_bytes
Maximum accounted payload bytes awaiting time ordering; exceeding this fails the crate.
Definition gSROCrate.h:37
std::size_t pending_payloads
Maximum payloads awaiting time ordering; exceeding this fails the crate.
Definition gSROCrate.h:35
Declare the sink first so the plugin referencing it is destroyed first.
Definition gSROCrate.h:42
One worker-produced contribution, transferred by move to its destination crate.
Definition gSROData.h:58
GSROTime time
Definition gSROData.h:62
std::uint64_t sequence
Definition gSROData.h:61
GSROEventId event_id
Definition gSROData.h:60
std::unique_ptr< const GSROData > data
Definition gSROData.h:63
GSROCrateId crate_id
Definition gSROData.h:59