gstreamer
Loading...
Searching...
No Matches
sro_service_example.cc
Go to the documentation of this file.
1#include "sro/gSROService.h"
2
3#include <algorithm>
4#include <atomic>
5#include <chrono>
6#include <condition_variable>
7#include <future>
8#include <iostream>
9#include <limits>
10#include <map>
11#include <mutex>
12#include <stdexcept>
13#include <string>
14#include <thread>
15#include <tuple>
16#include <utility>
17#include <vector>
18
19using namespace std::chrono_literals;
20
21namespace {
22void require(bool condition, const char* message) {
23 if (!condition) { throw std::runtime_error(message); }
24}
25
26template<class F>
27void expect_error(F&& action, const std::string& message) {
28 try { action(); }
29 catch (const std::exception& error) {
30 require(std::string(error.what()).find(message) != std::string::npos, "Unexpected service error");
31 return;
32 }
33 throw std::runtime_error("Expected an exception: " + message);
34}
35
36std::atomic<int> live_data{0};
37struct Data final : GSROData {
38 Data() { ++live_data; }
39 ~Data() override { --live_data; }
40 std::size_t size_bytes() const noexcept override { return sizeof(*this); }
41};
42
43GSROPayload payload(GSROCrateId crate, GSROEventId event, std::int64_t time = 1) {
44 return {crate, event, 0, GSROTime{time}, std::make_unique<Data>()};
45}
46
47using Key = std::tuple<GSROTime, GSROEventId, std::uint64_t>;
48struct CrateState {
49 int created = 0;
50 int destroyed = 0;
51 int closed = 0;
52 int delivered = 0;
53 std::thread::id owner;
54 std::vector<Key> consumed;
55 std::vector<GSROTime> boundaries;
56 std::optional<GSROEndContext> end;
57};
58
59struct State {
60 std::mutex mutex;
61 std::condition_variable changed;
62 std::map<GSROCrateId, CrateState> crates;
63 std::optional<GSROCrateId> fail_close;
64 bool wrong_thread = false;
65
66 template<class F>
67 void update(GSROCrateId id, F&& action) {
68 std::lock_guard lock(mutex);
69 auto& crate = crates.at(id);
70 wrong_thread |= crate.owner != std::this_thread::get_id();
71 action(crate);
72 changed.notify_all();
73 }
74
75 template<class F>
76 void wait_for(F&& condition) {
77 std::unique_lock lock(mutex);
78 require(changed.wait_for(lock, 5s, condition), "Timed out waiting for crate progress");
79 }
80};
81
82struct Sink final : GSROFrameSink {
83 Sink(GSROCrateId crate, std::shared_ptr<State> s) : id(crate), state(std::move(s)) {}
84 ~Sink() override { state->update(id, [](auto& c) { ++c.destroyed; }); }
85 void write_frame(GSROFrame) override {}
86 void finish_output() override {
87 state->update(id, [](auto& c) { ++c.closed; });
88 if (state->fail_close == id) { throw std::runtime_error("synthetic close failure"); }
89 }
90 GSROCrateId id;
91 std::shared_ptr<State> state;
92};
93
94struct Plugin final : GSROCratePlugin {
95 Plugin(GSROFrameSink& sink, GSROCrateId crate, std::shared_ptr<State> s)
96 : GSROCratePlugin(sink), id(crate), state(std::move(s)) {}
97 void consume_payload(GSROPayload item) override {
98 require(item.crate_id == id, "Service misrouted a payload");
99 state->update(id, [&](auto& c) { c.consumed.emplace_back(item.time, item.event_id, item.sequence); });
100 }
101 void advance_time(GSROTime time) override {
102 state->update(id, [&](auto& c) { c.boundaries.push_back(time); });
103 }
104 void finish_run(const GSROEndContext& context) override {
105 state->update(id, [&](auto& c) { c.end = context; });
106 }
107 GSROCrateId id;
108 std::shared_ptr<State> state;
109};
110
111GSROService::ResourceFactory factory(const std::shared_ptr<State>& state,
112 std::function<void(GSROCrateId)> before = {}) {
113 return [state, before](GSROCrateId id) {
114 {
115 std::lock_guard lock(state->mutex);
116 auto& crate = state->crates[id];
117 ++crate.created;
118 crate.owner = std::this_thread::get_id();
119 }
120 state->changed.notify_all();
121 if (before) { before(id); }
122 auto sink = std::make_unique<Sink>(id, state);
123 auto plugin = std::make_unique<Plugin>(*sink, id, state);
124 return GSROCrateResources{std::move(sink), std::move(plugin)};
125 };
126}
127
128void test_routing_and_runs() {
129 GSROService service;
130 expect_error([&] { service.create_crate_thread_if_needed(1); }, "not accepting");
131 expect_error([&] { service.begin_run({}); }, "resource factory");
132 service.finish_run({GSROEndReason::completed, std::nullopt});
133 auto state = std::make_shared<State>();
134 service.begin_run(factory(state), {}, [state](GSROCrateId id, GSROEventId, std::uint64_t) {
135 state->update(id, [](auto& c) { ++c.delivered; });
136 });
137 require(state->crates.empty(), "begin_run eagerly created crates");
138 expect_error([&] { service.begin_run(factory(state)); }, "previous SRO run");
139 std::vector<std::future<void>> workers;
140 for (GSROEventId worker = 0; worker < 8; ++worker) {
141 workers.push_back(std::async(std::launch::async, [&, worker] {
142 for (int n = 31; n >= 0; --n) {
144 service.dispatch_payload_to_crate(payload(7 + n % 2, worker * 32 + n, n % 4));
145 }
146 }));
147 }
148 for (auto& worker : workers) { worker.get(); }
149 service.advance_time(GSROTime{10});
150 // A crate appearing after progress must inherit the boundary before receiving any payload.
151 service.dispatch_payload_to_crate(payload(9, 0, 10));
152 service.create_crate_thread_if_needed(10); // Explicit empty-crate output.
153 expect_error([&] { service.advance_time(GSROTime{9}); }, "must not decrease");
154 expect_error([&] { service.finish_run({GSROEndReason::completed, GSROTime{9}}); }, "precedes");
156 service.finish_run({GSROEndReason::completed, std::nullopt});
157 require(state->crates.size() == 4 && !state->wrong_thread, "Unexpected crate ownership");
158 for (const auto& [id, c] : state->crates) {
159 require(c.created == 1 && c.destroyed == 1 && c.closed == 1, "Crate created or closed more than once");
160 require(c.owner != std::this_thread::get_id(), "Resource factory ran on the caller");
161 require(c.end->safe_time == GSROTime{11}, "Crates did not share final progress");
162 require(std::is_sorted(c.consumed.begin(), c.consumed.end()), "Crate delivery is out of order");
163 require(c.boundaries == std::vector<GSROTime>{GSROTime{10}, GSROTime{11}},
164 "Lost inherited boundary");
165 require(c.delivered == (id < 9 ? 128 : id == 9 ? 1 : 0), "Incorrect delivery acknowledgments");
166 require(c.consumed.size() == static_cast<std::size_t>(c.delivered), "Payload lost during routing");
167 }
168 expect_error([&] { service.dispatch_payload_to_crate(payload(7, 0)); }, "not accepting");
169 state = std::make_shared<State>();
170 service.begin_run(factory(state));
171 // Event IDs and simulation time restart per invocation.
172 service.dispatch_payload_to_crate(payload(7, 0, -5));
173 service.finish_run({GSROEndReason::completed, std::nullopt});
174 require(state->crates.at(7).consumed.size() == 1 && !state->crates.at(7).end->safe_time,
175 "A new run inherited stale crate state");
176}
177
178void test_slow_crate(bool fail_other_crate) {
179 GSROService service;
180 auto state = std::make_shared<State>();
181 std::promise<void> release;
182 auto gate = release.get_future().share();
183 GSROCrateLimits limits;
184 limits.queue_messages = 1;
185 service.begin_run(factory(state, [=](GSROCrateId id) {
186 if (id == 1) { gate.wait(); }
187 if (id == 2 && fail_other_crate) { throw std::runtime_error("synthetic setup failure"); }
188 }), limits);
189 service.dispatch_payload_to_crate(payload(1, 0));
190 std::promise<void> started;
191 auto ready = started.get_future();
192 auto producer = std::async(std::launch::async, [&] {
193 started.set_value();
194 try { service.dispatch_payload_to_crate(payload(1, 1)); }
195 catch (const std::exception& error) { return std::string(error.what()); }
196 return std::string{};
197 });
198 ready.wait();
199 const bool blocked = producer.wait_for(30ms) == std::future_status::timeout;
200 // This dispatch must proceed while crate 1's producer is waiting for space.
201 try { service.dispatch_payload_to_crate(payload(2, 2)); }
202 catch (const std::runtime_error&) { /* Setup may fail before the initial dispatch returns. */ }
203 std::future<void> finishing;
204 if (!fail_other_crate) {
205 finishing = std::async(std::launch::async, [&] {
206 service.finish_run({GSROEndReason::interrupted, std::nullopt});
207 });
208 }
209 const bool woke = producer.wait_for(5s) == std::future_status::ready;
210 release.set_value();
211 const auto error = producer.get();
212 require(blocked && woke, "Cross-crate progress, cancellation, or shutdown deadlocked");
213 if (fail_other_crate) {
214 require(error.find("synthetic setup failure") != std::string::npos,
215 "Peer failure did not wake producer");
216 expect_error([&] { service.finish_run({GSROEndReason::completed, std::nullopt}); }, "setup failure");
217 expect_error([&] { service.rethrow_if_failed(); }, "setup failure");
218 expect_error([&] { service.create_crate_thread_if_needed(3); }, "setup failure");
219 // Recovery requires joining the failed run, then starting a fresh invocation.
220 state = std::make_shared<State>();
221 service.begin_run(factory(state));
222 service.dispatch_payload_to_crate(payload(3, 0));
223 service.finish_run({GSROEndReason::completed, std::nullopt});
224 require(state->crates.at(3).closed == 1, "Failed run poisoned the following run");
225 }
226 else {
227 finishing.get();
228 require(error.find("closed") != std::string::npos, "Shutdown did not reject blocked submission");
229 require(state->crates.at(1).consumed.size() == 1 && state->crates.at(2).consumed.size() == 1,
230 "Normal shutdown cancelled previously accepted data");
231 }
232}
233
234void test_close_failure_and_destructor() {
235 auto state = std::make_shared<State>();
236 state->fail_close = 1;
237 {
238 GSROService service;
239 service.begin_run(factory(state));
240 service.dispatch_payload_to_crate(payload(1, 0));
241 service.dispatch_payload_to_crate(payload(2, 0));
242 expect_error([&] { service.finish_run({GSROEndReason::completed, std::nullopt}); }, "close failure");
243 require(state->crates.at(1).destroyed == 1 && state->crates.at(2).destroyed == 1,
244 "One close failure prevented another crate from joining");
245 }
246 state = std::make_shared<State>();
247 {
248 GSROService service;
249 service.begin_run(factory(state));
250 service.dispatch_payload_to_crate(payload(4, 0));
251 }
252 require(state->crates.at(4).end->reason == GSROEndReason::interrupted && state->crates.at(4).destroyed == 1,
253 "Service destructor did not finalize and join");
254}
255
256struct Timing final : GSROTiming {
257 std::function<std::optional<GSROTime>(GSROEventId)> bound = [](GSROEventId id) {
258 return GSROTime{static_cast<std::int64_t>(id) * 10 - 10};
259 };
260 std::optional<GSROTime> earliest_remaining_time(GSROEventId id) const override { return bound(id); }
261};
262
263void test_automatic_delivery() {
264 GSROService service;
265 auto state = std::make_shared<State>();
266 auto timing = std::make_shared<Timing>();
267 std::promise<void> release;
268 auto gate = release.get_future().share();
269 service.begin_run(factory(state, [=](GSROCrateId id) { if (id == 2) { gate.wait(); } }), timing);
270 service.dispatch_payload_to_crate(payload(1, 0, -5));
271 auto second = payload(2, 0, -10);
272 second.sequence = 1;
273 service.dispatch_payload_to_crate(std::move(second));
274 service.complete_event(1); // Empty event finishes before event zero.
275 service.complete_event(0); // Closed, but crate 2 has not retained its payload yet.
276 state->wait_for([&] { return state->crates.count(1) && !state->crates.at(1).boundaries.empty(); });
277 bool held;
278 {
279 std::lock_guard lock(state->mutex);
280 held = state->crates.at(1).boundaries.back() == GSROTime{-10} &&
281 state->crates.at(1).consumed.empty();
282 }
283 release.set_value();
284 require(held, "Closed event advanced before every crate acknowledged delivery");
285 state->wait_for([&] { return state->crates.at(1).boundaries.back() == GSROTime{10}; });
286 // A new crate inherits progress. Equality stays pending until another complete event advances time.
287 service.dispatch_payload_to_crate(payload(3, 2, 10));
288 state->wait_for([&] {
289 const auto found = state->crates.find(3);
290 return found != state->crates.end() && !found->second.boundaries.empty();
291 });
292 {
293 std::lock_guard lock(state->mutex);
294 require(state->crates.at(3).consumed.empty(), "Equal-boundary payload released too early");
295 }
296 expect_error([&] { service.complete_event(1); }, "already complete");
297 expect_error([&] { service.dispatch_payload_to_crate(payload(1, 0)); }, "already complete");
298 expect_error([&] { service.advance_time(GSROTime{100}); }, "automatic timing");
299 expect_error([&] { service.finish_run({GSROEndReason::completed, GSROTime{100}}); }, "automatic timing");
300 service.complete_event(2);
301 service.finish_run({GSROEndReason::completed, std::nullopt});
302 for (const auto& [id, crate] : state->crates) {
303 require(crate.consumed.size() == 1 && crate.end->safe_time == GSROTime{20},
304 "Final acknowledgment or automatic boundary was lost");
305 }
306}
307
308void test_automatic_workers() {
309 GSROService service;
310 auto state = std::make_shared<State>();
311 GSROCrateLimits limits;
312 limits.queue_messages = 2; // Exercise progress broadcasts while producers wait for queue space.
313 service.begin_run(factory(state), std::make_shared<Timing>(), limits);
314 std::vector<std::future<void>> workers;
315 for (GSROEventId worker = 0; worker < 8; ++worker) {
316 workers.push_back(std::async(std::launch::async, [&, worker] {
317 for (int n = 31; n >= 0; --n) {
318 const auto event = worker * 32 + n;
319 if (event % 4) {
320 service.dispatch_payload_to_crate(payload(7, event, event * 10 - 5));
321 auto second = payload(8, event, event * 10 - 10);
322 second.sequence = 1;
323 service.dispatch_payload_to_crate(std::move(second));
324 }
325 service.complete_event(event);
326 }
327 }));
328 }
329 for (auto& worker : workers) { worker.get(); }
330 service.finish_run({GSROEndReason::completed, std::nullopt});
331 for (const auto& [id, crate] : state->crates) {
332 require(crate.consumed.size() == 192, "Concurrent automatic dispatch lost payloads");
333 require(std::is_sorted(crate.consumed.begin(), crate.consumed.end()),
334 "Automatic delivery out of order");
335 require(crate.end->safe_time == GSROTime{2550}, "Empty or reordered events left incorrect progress");
336 }
337 require(!state->wrong_thread, "Automatic progress bypassed crate thread ownership");
338 // A subsequent invocation starts at event zero, with no inherited completion or timing state.
339 state = std::make_shared<State>();
340 service.begin_run(factory(state), std::make_shared<Timing>());
341 service.dispatch_payload_to_crate(payload(1, 0, -10));
342 service.complete_event(0);
343 service.finish_run({GSROEndReason::completed, std::nullopt});
344 require(state->crates.at(1).end->safe_time == GSROTime{0}, "Automatic state leaked across runs");
345}
346
347void test_missing_events_and_unknown_time() {
348 for (int mode = 0; mode < 3; ++mode) {
349 auto state = std::make_shared<State>();
350 {
351 GSROService service;
352 auto timing = std::make_shared<Timing>();
353 if (mode == 2) { timing->bound = [](GSROEventId) { return std::nullopt; }; }
354 service.begin_run(factory(state), timing);
355 service.dispatch_payload_to_crate(payload(1, 1, 15));
356 service.complete_event(1);
357 if (mode == 1) {
358 service.dispatch_payload_to_crate(payload(1, 0, -5)); // Delivered but never closed.
359 }
360 // Destructor drains accepted payloads but must not complete missing/unfinished event zero.
361 }
362 const auto& crate = state->crates.at(1);
363 require(crate.end->reason == GSROEndReason::interrupted, "Interruption was lost");
364 require(crate.end->safe_time == (mode == 2 ? std::nullopt : std::optional<GSROTime>{GSROTime{-10}}),
365 "Shutdown invented progress across a missing event or unknown time");
366 require(crate.consumed.size() == (mode == 1 ? 2u : 1u), "Interrupted shutdown lost accepted payloads");
367 }
368 // All-empty runs still query the timing implementation through the final contiguous prefix.
369 GSROService service;
370 auto state = std::make_shared<State>();
371 auto timing = std::make_shared<Timing>();
372 std::atomic<GSROEventId> last{0};
373 timing->bound = [&](GSROEventId id) { last = id; return std::nullopt; };
374 service.begin_run(factory(state), timing);
375 service.complete_event(2);
376 service.complete_event(0);
377 service.complete_event(1);
378 service.finish_run({GSROEndReason::completed, std::nullopt});
379 require(last == 3 && state->crates.empty(), "All-empty run lost completion or created an output crate");
380}
381
382void test_automatic_failures() {
383 GSROService service;
384 auto state = std::make_shared<State>();
385 auto timing = std::make_shared<Timing>();
386 expect_error([&] { service.begin_run(factory(state), std::shared_ptr<const GSROTiming>{}); }, "null");
387 expect_error([&] { service.begin_run(factory(state), timing, {}, 0); }, "positive");
388 service.begin_run(factory(state), timing, {}, 2);
389 expect_error([&] { service.complete_event(std::numeric_limits<GSROEventId>::max()); }, "overflow");
390 service.complete_event(1);
391 service.complete_event(2);
392 expect_error([&] { service.complete_event(3); }, "pending event limit");
393 expect_error([&] { service.finish_run({GSROEndReason::completed, std::nullopt}); }, "pending event limit");
394 for (bool throws : {false, true}) {
395 state = std::make_shared<State>();
396 timing = std::make_shared<Timing>();
397 timing->bound = [=](GSROEventId id) {
398 if (id && throws) { throw std::runtime_error("synthetic timing failure"); }
399 return GSROTime{id ? -20 : -10};
400 };
401 service.begin_run(factory(state), timing);
403 state->wait_for([&] { return state->crates.count(1) && !state->crates.at(1).boundaries.empty(); });
404 service.complete_event(0);
405 expect_error([&] { service.finish_run({GSROEndReason::completed, std::nullopt}); },
406 throws ? "timing failure" : "must not decrease");
407 require(state->crates.at(1).destroyed == 1, "Timing failure skipped crate cleanup");
408 }
409 // A failed enqueue was registered before submission; shutdown must fail, not wait for a missing ack.
410 state = std::make_shared<State>();
411 GSROCrateLimits limits;
412 limits.queue_bytes = sizeof(GSROPayload);
413 service.begin_run(factory(state), std::make_shared<Timing>(), limits);
414 expect_error([&] { service.dispatch_payload_to_crate(payload(1, 0)); }, "byte limit");
415 expect_error([&] { service.finish_run({GSROEndReason::interrupted, std::nullopt}); }, "byte limit");
416}
417
418void test_automatic_waiters(bool fail_timing) {
419 GSROService service;
420 auto state = std::make_shared<State>();
421 auto timing = std::make_shared<Timing>();
422 timing->bound = [=](GSROEventId id) -> std::optional<GSROTime> {
423 if (!id) { return std::nullopt; } // No initial control message occupies the gated input queue.
424 if (fail_timing) { throw std::runtime_error("synthetic timing failure"); }
425 return GSROTime{static_cast<std::int64_t>(id) * 10};
426 };
427 std::promise<void> release;
428 auto gate = release.get_future().share();
429 GSROCrateLimits limits;
430 limits.queue_messages = 1;
431 service.begin_run(factory(state, [=](GSROCrateId) { gate.wait(); }), timing, limits);
432 service.dispatch_payload_to_crate(payload(1, fail_timing ? 1 : 0));
433 if (fail_timing) {
434 std::promise<void> started;
435 auto ready = started.get_future();
436 auto producer = std::async(std::launch::async, [&] {
437 started.set_value();
438 try { service.dispatch_payload_to_crate(payload(1, 2)); }
439 catch (const std::exception& error) { return std::string(error.what()); }
440 return std::string{};
441 });
442 ready.wait();
443 const bool blocked = producer.wait_for(30ms) == std::future_status::timeout;
444 service.complete_event(0); // Timing failure must cancel the crate and wake its blocked worker.
445 const bool woke = producer.wait_for(5s) == std::future_status::ready;
446 release.set_value();
447 const auto error = producer.get();
448 require(blocked && woke && error.find("timing failure") != std::string::npos,
449 "Timing failure did not wake a worker blocked on a full crate queue");
450 expect_error([&] { service.finish_run({GSROEndReason::interrupted, std::nullopt}); }, "timing failure");
451 }
452 else {
453 service.complete_event(0);
454 auto finishing = std::async(std::launch::async, [&] {
455 service.finish_run({GSROEndReason::completed, std::nullopt});
456 });
457 const bool waiting = finishing.wait_for(30ms) == std::future_status::timeout;
458 release.set_value();
459 finishing.get();
460 require(waiting && state->crates.at(1).end->safe_time == GSROTime{10},
461 "Shutdown failed to wait for the last delivery and derive its final bound");
462 }
463}
464} // namespace
465
466int main() {
467 try {
468 test_routing_and_runs();
469 test_slow_crate(false);
470 test_slow_crate(true);
471 test_close_failure_and_destructor();
472 test_automatic_delivery();
473 test_automatic_workers();
474 test_missing_events_and_unknown_time();
475 test_automatic_failures();
476 test_automatic_waiters(false);
477 test_automatic_waiters(true);
478 require(live_data == 0, "Service leaked payloads");
479 std::cout << "SRO service routing, automatic progress, shutdown, and failure checks passed\n";
480 return 0;
481 }
482 catch (const std::exception& error) {
483 std::cerr << error.what() << '\n';
484 return 1;
485 }
486}
Implementation-owned framing and optional electronics response for one crate.
Definition gSROPlugin.h:79
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
Shared, run-scoped registry routing worker payloads directly to crate threads.
Definition gSROService.h:28
void begin_run(ResourceFactory make_resources, GSROCrateLimits limits={}, DeliveryCallback on_delivery={})
Begin a fresh invocation after the previous one was joined; clears its progress and failure state.
void finish_run(GSROEndContext context)
void advance_time(GSROTime safe_time)
Manual mode only: broadcast a proven bound after all earlier dispatches return. New crates inherit it...
void rethrow_if_failed() const
void complete_event(GSROEventId event_id)
Automatic mode only. Close once, after all dispatch calls for this event have returned.
void create_crate_thread_if_needed(GSROCrateId crate_id)
Also permits implementations to request output for empty crates. Concurrent requests create once.
std::function< GSROCrateResources(GSROCrateId)> ResourceFactory
Definition gSROService.h:30
void dispatch_payload_to_crate(GSROPayload payload)
Transfers ownership directly to the destination queue, creating its crate on first use.
Implementation-defined bound used by GEMC to establish safe input progress.
Definition gSROPlugin.h:24
std::chrono::duration< std::int64_t, std::nano > GSROTime
Definition gSROData.h:20
std::uint32_t GSROCrateId
Definition gSROData.h:18
std::uint64_t GSROEventId
Definition gSROData.h:19
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
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
GSROCrateId crate_id
Definition gSROData.h:59