6#include <condition_variable>
17using namespace std::chrono_literals;
20void require(
bool condition,
const char* message) {
21 if (!condition) {
throw std::runtime_error(message); }
25void expect_error(F&& action,
const std::string& message) {
27 catch (
const std::exception& error) {
28 require(std::string(error.what()).find(message) != std::string::npos,
"Unexpected exception");
31 throw std::runtime_error(
"Expected an exception: " + message);
34std::atomic<int> live_data{0};
36 Data() { ++live_data; }
37 ~Data()
override { --live_data; }
38 std::size_t size_bytes() const noexcept
override {
return sizeof(*this); }
42 return {7,
event, sequence,
GSROTime{time}, std::make_unique<Data>()};
45using Key = std::tuple<GSROTime, GSROEventId, std::uint64_t>;
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};
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);
63 void maybe_fail(
const std::string& operation)
const {
64 if (fail_at == operation) {
throw std::runtime_error(
"synthetic " + operation +
" failure"); }
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");
76 void finish_output()
override {
77 state->record(
"close");
78 state->maybe_fail(
"close");
80 std::shared_ptr<State> state;
85 Plugin(GSROFrameSink& sink, std::shared_ptr<State> s) : GSROCratePlugin(sink), state(std::move(s)) {
86 state->record(
"plugin-created");
88 ~Plugin()
override { state->record(
"plugin-destroyed"); }
89 void consume_payload(GSROPayload item)
override {
90 state->record(
"consume");
91 state->maybe_fail(
"consume");
93 std::lock_guard lock(state->mutex);
96 pending.push_back(std::move(item));
98 void advance_time(
GSROTime time)
override {
99 state->record(
"advance");
100 state->maybe_fail(
"advance");
101 for (
auto& item : pending) {
106 std::lock_guard lock(state->mutex);
107 state->boundary = time;
109 state->changed.notify_all();
111 void finish_run(
const GSROEndContext& context)
override {
112 state->record(
"finish");
113 state->maybe_fail(
"finish");
114 state->end = context;
117 std::shared_ptr<State> state;
118 std::vector<GSROPayload> pending;
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);
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;
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)); }
142 for (
auto& worker : workers) { worker.get(); }
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"); }
152 expect_error([&] { crate.advance_time(
GSROTime{-1}); },
"must not decrease");
153 expect_error([&] { crate.enqueue_payload(payload(-1)); },
"precedes");
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");
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);
178 if (fail_setup) {
throw std::runtime_error(
"synthetic setup failure"); }
182 std::promise<void> started;
183 auto started_future = started.get_future();
184 auto producer = std::async(std::launch::async, [&] {
186 try { crate.enqueue_payload(payload(2)); }
187 catch (
const std::exception& error) {
return std::string(error.what()); }
188 return std::string{};
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, [&] {
198 const bool woke_on_close = !close_while_blocked || producer.wait_for(5s) == std::future_status::ready;
200 const auto error = producer.get();
201 if (finishing.valid()) { finishing.get(); }
204 require(error.find(
"setup failure") != std::string::npos,
"Blocked producer missed writer failure");
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");
213 require(blocked && woke_on_close,
"Queue backpressure or shutdown wakeup failed");
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;
221 state->maybe_fail(
"delivery");
227 catch (
const std::runtime_error&) { }
229 "synthetic " + stage +
" failure");
230 expect_error([&] { crate.
rethrow_if_failed(); },
"synthetic " + stage +
" failure");
232 for (
const bool byte_limit : {
false,
true}) {
233 auto state = std::make_shared<State>();
237 GSROCrate crate(7, factory(state), limits);
241 "pending input limit");
243 auto state = std::make_shared<State>();
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");
254 expect_error([&] {
GSROCrate invalid(7, factory(state), limits); },
"positive buffer limits");
257 "null plugin or sink");
259 duplicates.enqueue_payload(payload(1));
260 duplicates.enqueue_payload(payload(1));
264void test_destructor_shutdown() {
265 auto state = std::make_shared<State>();
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";
280 require(state->lifecycle.back() ==
"sink-destroyed",
"Destructor failed to clean up after close failure");
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";
297 catch (
const std::exception& error) {
298 std::cerr << error.what() <<
'\n';
Implementation-owned framing and optional electronics response for one crate.
One crate's bounded input queue, ordering buffer, and background processing thread.
void enqueue_payload(GSROPayload payload)
Transfers ownership. Rejects null data, wrong crates, oversized input, and timestamps before a bounda...
void advance_time(GSROTime safe_time)
Queue a nondecreasing boundary after all earlier input has been submitted. Equality is permitted.
void rethrow_if_failed() const
Check for a background failure without waiting. Failure also wakes and rejects queued producers.
std::function< GSROCrateResources()> ResourceFactory
void finish_and_join(GSROEndContext context)
Base for implementation-defined payload and frame contents.
Output destination called exclusively by its owning crate thread.
std::chrono::duration< std::int64_t, std::nano > GSROTime
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.
std::size_t pending_bytes
Maximum accounted payload bytes awaiting time ordering; exceeding this fails the crate.
std::size_t pending_payloads
Maximum payloads awaiting time ordering; exceeding this fails the crate.
Declare the sink first so the plugin referencing it is destroyed first.
One worker-produced contribution, transferred by move to its destination crate.
std::unique_ptr< const GSROData > data