6#include <condition_variable>
19using namespace std::chrono_literals;
22void require(
bool condition,
const char* message) {
23 if (!condition) {
throw std::runtime_error(message); }
27void expect_error(F&& action,
const std::string& message) {
29 catch (
const std::exception& error) {
30 require(std::string(error.what()).find(message) != std::string::npos,
"Unexpected service error");
33 throw std::runtime_error(
"Expected an exception: " + message);
36std::atomic<int> live_data{0};
38 Data() { ++live_data; }
39 ~Data()
override { --live_data; }
40 std::size_t size_bytes() const noexcept
override {
return sizeof(*this); }
44 return {crate,
event, 0,
GSROTime{time}, std::make_unique<Data>()};
47using Key = std::tuple<GSROTime, GSROEventId, std::uint64_t>;
53 std::thread::id owner;
54 std::vector<Key> consumed;
55 std::vector<GSROTime> boundaries;
56 std::optional<GSROEndContext> end;
61 std::condition_variable changed;
62 std::map<GSROCrateId, CrateState> crates;
63 std::optional<GSROCrateId> fail_close;
64 bool wrong_thread =
false;
68 std::lock_guard lock(mutex);
69 auto& crate = crates.at(
id);
70 wrong_thread |= crate.owner != std::this_thread::get_id();
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");
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"); }
91 std::shared_ptr<State> state;
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); });
101 void advance_time(
GSROTime time)
override {
102 state->update(
id, [&](
auto& c) { c.boundaries.push_back(time); });
104 void finish_run(
const GSROEndContext& context)
override {
105 state->update(
id, [&](
auto& c) { c.end = context; });
108 std::shared_ptr<State> state;
115 std::lock_guard lock(state->mutex);
116 auto& crate = state->crates[id];
118 crate.owner = std::this_thread::get_id();
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);
128void test_routing_and_runs() {
131 expect_error([&] { service.
begin_run({}); },
"resource factory");
133 auto state = std::make_shared<State>();
135 state->update(
id, [](
auto& c) { ++c.delivered; });
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) {
148 for (
auto& worker : workers) { worker.get(); }
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");
169 state = std::make_shared<State>();
174 require(state->crates.at(7).consumed.size() == 1 && !state->crates.at(7).end->safe_time,
175 "A new run inherited stale crate state");
178void test_slow_crate(
bool fail_other_crate) {
180 auto state = std::make_shared<State>();
181 std::promise<void> release;
182 auto gate = release.get_future().share();
186 if (
id == 1) { gate.wait(); }
187 if (
id == 2 && fail_other_crate) {
throw std::runtime_error(
"synthetic setup failure"); }
190 std::promise<void> started;
191 auto ready = started.get_future();
192 auto producer = std::async(std::launch::async, [&] {
195 catch (
const std::exception& error) {
return std::string(error.what()); }
196 return std::string{};
199 const bool blocked = producer.wait_for(30ms) == std::future_status::timeout;
202 catch (
const std::runtime_error&) { }
203 std::future<void> finishing;
204 if (!fail_other_crate) {
205 finishing = std::async(std::launch::async, [&] {
209 const bool woke = producer.wait_for(5s) == std::future_status::ready;
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");
220 state = std::make_shared<State>();
224 require(state->crates.at(3).closed == 1,
"Failed run poisoned the following run");
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");
234void test_close_failure_and_destructor() {
235 auto state = std::make_shared<State>();
236 state->fail_close = 1;
243 require(state->crates.at(1).destroyed == 1 && state->crates.at(2).destroyed == 1,
244 "One close failure prevented another crate from joining");
246 state = std::make_shared<State>();
253 "Service destructor did not finalize and join");
258 return GSROTime{
static_cast<std::int64_t
>(id) * 10 - 10};
260 std::optional<GSROTime> earliest_remaining_time(
GSROEventId id)
const override {
return bound(
id); }
263void test_automatic_delivery() {
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);
271 auto second = payload(2, 0, -10);
276 state->wait_for([&] {
return state->crates.count(1) && !state->crates.at(1).boundaries.empty(); });
279 std::lock_guard lock(state->mutex);
280 held = state->crates.at(1).boundaries.back() ==
GSROTime{-10} &&
281 state->crates.at(1).consumed.empty();
284 require(held,
"Closed event advanced before every crate acknowledged delivery");
285 state->wait_for([&] {
return state->crates.at(1).boundaries.back() ==
GSROTime{10}; });
288 state->wait_for([&] {
289 const auto found = state->crates.find(3);
290 return found != state->crates.end() && !found->second.boundaries.empty();
293 std::lock_guard lock(state->mutex);
294 require(state->crates.at(3).consumed.empty(),
"Equal-boundary payload released too early");
296 expect_error([&] { service.
complete_event(1); },
"already complete");
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");
308void test_automatic_workers() {
310 auto state = std::make_shared<State>();
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;
321 auto second = payload(8,
event,
event * 10 - 10);
329 for (
auto& worker : workers) { worker.get(); }
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");
337 require(!state->wrong_thread,
"Automatic progress bypassed crate thread ownership");
339 state = std::make_shared<State>();
340 service.
begin_run(factory(state), std::make_shared<Timing>());
344 require(state->crates.at(1).end->safe_time ==
GSROTime{0},
"Automatic state leaked across runs");
347void test_missing_events_and_unknown_time() {
348 for (
int mode = 0; mode < 3; ++mode) {
349 auto state = std::make_shared<State>();
352 auto timing = std::make_shared<Timing>();
353 if (mode == 2) { timing->bound = [](
GSROEventId) {
return std::nullopt; }; }
354 service.
begin_run(factory(state), timing);
362 const auto& crate = state->crates.at(1);
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");
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);
379 require(last == 3 && state->crates.empty(),
"All-empty run lost completion or created an output crate");
382void test_automatic_failures() {
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");
392 expect_error([&] { service.
complete_event(3); },
"pending event limit");
394 for (
bool throws : {
false,
true}) {
395 state = std::make_shared<State>();
396 timing = std::make_shared<Timing>();
398 if (
id &&
throws) {
throw std::runtime_error(
"synthetic timing failure"); }
401 service.
begin_run(factory(state), timing);
403 state->wait_for([&] {
return state->crates.count(1) && !state->crates.at(1).boundaries.empty(); });
406 throws ?
"timing failure" :
"must not decrease");
407 require(state->crates.at(1).destroyed == 1,
"Timing failure skipped crate cleanup");
410 state = std::make_shared<State>();
413 service.
begin_run(factory(state), std::make_shared<Timing>(), limits);
418void test_automatic_waiters(
bool fail_timing) {
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; }
424 if (fail_timing) {
throw std::runtime_error(
"synthetic timing failure"); }
425 return GSROTime{
static_cast<std::int64_t
>(id) * 10};
427 std::promise<void> release;
428 auto gate = release.get_future().share();
434 std::promise<void> started;
435 auto ready = started.get_future();
436 auto producer = std::async(std::launch::async, [&] {
439 catch (
const std::exception& error) {
return std::string(error.what()); }
440 return std::string{};
443 const bool blocked = producer.wait_for(30ms) == std::future_status::timeout;
445 const bool woke = producer.wait_for(5s) == std::future_status::ready;
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");
454 auto finishing = std::async(std::launch::async, [&] {
457 const bool waiting = finishing.wait_for(30ms) == std::future_status::timeout;
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");
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";
482 catch (
const std::exception& error) {
483 std::cerr << error.what() <<
'\n';
Implementation-owned framing and optional electronics response for one crate.
Base for implementation-defined payload and frame contents.
Output destination called exclusively by its owning crate thread.
Shared, run-scoped registry routing worker payloads directly to crate threads.
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
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.
std::chrono::duration< std::int64_t, std::nano > GSROTime
std::uint32_t GSROCrateId
std::uint64_t GSROEventId
Per-crate bounds preventing asynchronous delivery from consuming unlimited memory.
std::size_t queue_messages
Maximum queued payload/progress messages; producers wait when the queue is full.
std::size_t queue_bytes
Maximum accounted payload bytes in the input queue; producers wait until their payload fits.
Declare the sink first so the plugin referencing it is destroyed first.
One worker-produced contribution, transferred by move to its destination crate.