gstreamer
Loading...
Searching...
No Matches
gSROCrate.cc
Go to the documentation of this file.
1#include "gSROCrate.h"
2
3#include <gemc/guts/gthreads.h>
4
5#include <condition_variable>
6#include <deque>
7#include <exception>
8#include <map>
9#include <mutex>
10#include <stdexcept>
11#include <tuple>
12#include <utility>
13#include <variant>
14
16{
17public:
19 DeliveryCallback delivered, FailureCallback failed)
20 : crate_id(id), make_resources(std::move(factory)), limits(bounds),
21 on_delivery(std::move(delivered)), on_failure(std::move(failed)) {
22 if (!make_resources || !limits.queue_messages || !limits.queue_bytes ||
23 !limits.pending_payloads || !limits.pending_bytes) {
24 throw std::invalid_argument("SRO crate requires a resource factory and positive buffer limits");
25 }
26 thread = jthread_alias([this] { run(); });
27 }
28
30 if (!payload.data || payload.crate_id != crate_id) {
31 throw std::invalid_argument("SRO payload is null or addressed to another crate");
32 }
33 const auto bytes = payload.data->size_bytes();
34 if (limits.queue_bytes < sizeof(GSROPayload) || bytes > limits.queue_bytes - sizeof(GSROPayload)) {
35 throw std::length_error("SRO payload exceeds the crate input queue byte limit");
36 }
37 const auto total_bytes = bytes + sizeof(GSROPayload);
38 std::unique_lock lock(mutex);
39 space_available.wait(lock, [&] {
40 return closing || failure ||
41 (queue.size() < limits.queue_messages && total_bytes <= limits.queue_bytes - queue_bytes);
42 });
43 check_open();
44 if (submitted_safe_time && payload.time < *submitted_safe_time) {
45 throw std::invalid_argument("SRO payload precedes an already submitted safe-time boundary");
46 }
47 queue.emplace_back(QueuedPayload{std::move(payload), total_bytes});
48 queue_bytes += total_bytes;
49 lock.unlock();
50 input_available.notify_one();
51 }
52
53 void advance_time(GSROTime safe_time) {
54 std::unique_lock lock(mutex);
55 space_available.wait(lock, [&] { return closing || failure || queue.size() < limits.queue_messages; });
56 check_open();
57 if (submitted_safe_time && safe_time < *submitted_safe_time) {
58 throw std::invalid_argument("SRO safe-time boundaries must not decrease");
59 }
60 queue.emplace_back(safe_time);
61 submitted_safe_time = safe_time;
62 lock.unlock();
63 input_available.notify_one();
64 }
65
67 {
68 std::lock_guard lock(mutex);
69 if (!closing && !failure) {
70 if (context.safe_time && submitted_safe_time && *context.safe_time < *submitted_safe_time) {
71 throw std::invalid_argument("SRO final boundary precedes an already submitted boundary");
72 }
73 if (!context.safe_time) { context.safe_time = submitted_safe_time; }
74 end_context = context;
75 closing = true;
76 }
77 }
78 input_available.notify_one();
79 space_available.notify_all();
80 }
81
82 void cancel(std::exception_ptr error) {
83 if (!error) { throw std::invalid_argument("SRO cancellation requires an exception"); }
84 {
85 std::lock_guard lock(mutex);
86 if (!failure) { failure = error; }
87 closing = true;
88 }
89 input_available.notify_one();
90 space_available.notify_all();
91 }
92
94 request_finish(context);
95 if (thread.joinable()) { thread.join(); }
97 }
98
99 void rethrow_if_failed() const {
100 std::lock_guard lock(mutex);
101 if (failure) { std::rethrow_exception(failure); }
102 }
103
104private:
105 struct QueuedPayload {
106 GSROPayload payload;
107 std::size_t bytes;
108 };
109 using Message = std::variant<QueuedPayload, GSROTime>;
110 using Key = std::tuple<GSROTime, GSROEventId, std::uint64_t>;
111 using Pending = std::map<Key, QueuedPayload>;
112
113 void check_open() const {
114 if (failure) { std::rethrow_exception(failure); }
115 if (closing) { throw std::logic_error("SRO crate is closed to input"); }
116 }
117
118 void consume_before(Pending& pending, std::size_t& bytes, GSROCratePlugin& plugin,
119 std::optional<GSROTime> boundary) {
120 while (!pending.empty() && (!boundary || std::get<0>(pending.begin()->first) < *boundary)) {
122 auto entry = pending.extract(pending.begin());
123 bytes -= entry.mapped().bytes;
124 plugin.consume_payload(std::move(entry.mapped().payload));
125 }
126 }
127
128 void run() noexcept {
129 try {
130 // Plugin and sink are local so destruction happens here, in dependency order, even on failure.
131 auto resources = make_resources();
132 if (!resources.sink || !resources.plugin) {
133 throw std::invalid_argument("SRO resource factory returned a null plugin or sink");
134 }
135 Pending pending;
136 std::size_t pending_bytes = 0;
137 std::optional<GSROTime> processed_safe_time;
138 GSROEndContext context{GSROEndReason::interrupted, std::nullopt};
139 for (;;) {
140 Message message;
141 {
142 std::unique_lock lock(mutex);
143 input_available.wait(lock, [&] { return closing || !queue.empty(); });
144 if (failure) { std::rethrow_exception(failure); }
145 if (queue.empty()) {
146 context = *end_context;
147 break;
148 }
149 message = std::move(queue.front());
150 queue.pop_front();
151 if (const auto* item = std::get_if<QueuedPayload>(&message)) { queue_bytes -= item->bytes; }
152 }
153 space_available.notify_all();
154 if (auto* item = std::get_if<QueuedPayload>(&message)) {
155 if (pending.size() >= limits.pending_payloads ||
156 item->bytes > limits.pending_bytes - pending_bytes) {
157 throw std::length_error("SRO pending input limit exceeded while waiting for safe time");
158 }
159 const auto event = item->payload.event_id;
160 const auto sequence = item->payload.sequence;
161 const Key key{item->payload.time, event, sequence};
162 const auto bytes = item->bytes;
163 if (!pending.emplace(key, std::move(*item)).second) {
164 throw std::invalid_argument("Duplicate SRO payload ordering key");
165 }
166 pending_bytes += bytes;
167 if (on_delivery) { on_delivery(event, sequence); }
168 }
169 else {
170 const auto boundary = std::get<GSROTime>(message);
171 consume_before(pending, pending_bytes, *resources.plugin, boundary);
172 resources.plugin->advance_time(boundary);
173 processed_safe_time = boundary;
174 }
175 }
176 if (context.safe_time && (!processed_safe_time || *context.safe_time > *processed_safe_time)) {
177 consume_before(pending, pending_bytes, *resources.plugin, context.safe_time);
178 resources.plugin->advance_time(*context.safe_time);
179 }
180 consume_before(pending, pending_bytes, *resources.plugin, std::nullopt);
182 resources.plugin->finish_run(context);
184 resources.sink->finish_output();
185 }
186 catch (...) {
187 // Destroy rejected payloads outside the mutex; data destructors belong to external implementations.
188 std::deque<Message> rejected;
189 std::exception_ptr error;
190 {
191 std::lock_guard lock(mutex);
192 if (!failure) { failure = std::current_exception(); }
193 error = failure;
194 closing = true;
195 rejected.swap(queue);
196 queue_bytes = 0;
197 }
198 space_available.notify_all();
199 if (on_failure) {
200 try { on_failure(error); }
201 catch (...) { /* Preserve the original background failure. */ }
202 }
203 }
204 }
205
206 const GSROCrateId crate_id;
207 ResourceFactory make_resources;
208 const GSROCrateLimits limits;
209 DeliveryCallback on_delivery;
210 FailureCallback on_failure;
211 mutable std::mutex mutex;
212 std::condition_variable input_available;
213 std::condition_variable space_available;
214 std::deque<Message> queue;
215 std::size_t queue_bytes = 0;
216 std::optional<GSROTime> submitted_safe_time;
217 std::optional<GSROEndContext> end_context;
218 std::exception_ptr failure;
219 bool closing = false;
220 jthread_alias thread;
221};
222
224 GSROCrateLimits limits, DeliveryCallback on_delivery, FailureCallback on_failure)
225 : impl(std::make_unique<Impl>(crate_id, std::move(make_resources), limits,
226 std::move(on_delivery), std::move(on_failure))) {}
227
229 try { impl->finish_and_join({GSROEndReason::interrupted, std::nullopt}); }
230 catch (...) { /* Explicit finish_and_join is the error-reporting path. */ }
231}
232
233void GSROCrate::enqueue_payload(GSROPayload payload) { impl->enqueue_payload(std::move(payload)); }
234void GSROCrate::advance_time(GSROTime safe_time) { impl->advance_time(safe_time); }
235void GSROCrate::request_finish(GSROEndContext context) { impl->request_finish(context); }
236void GSROCrate::cancel(std::exception_ptr error) { impl->cancel(error); }
237void GSROCrate::finish_and_join(GSROEndContext context) { impl->finish_and_join(context); }
238void GSROCrate::rethrow_if_failed() const { impl->rethrow_if_failed(); }
virtual void consume_payload(GSROPayload payload)=0
void advance_time(GSROTime safe_time)
Definition gSROCrate.cc:53
void enqueue_payload(GSROPayload payload)
Definition gSROCrate.cc:29
void request_finish(GSROEndContext context)
Definition gSROCrate.cc:66
void cancel(std::exception_ptr error)
Definition gSROCrate.cc:82
void finish_and_join(GSROEndContext context)
Definition gSROCrate.cc:93
Impl(GSROCrateId id, ResourceFactory factory, GSROCrateLimits bounds, DeliveryCallback delivered, FailureCallback failed)
Definition gSROCrate.cc:18
void rethrow_if_failed() const
Definition gSROCrate.cc:99
void cancel(std::exception_ptr error)
Fail input and wake waiters. Cleanup runs on the crate thread after any executing callback returns.
Definition gSROCrate.cc:236
void enqueue_payload(GSROPayload payload)
Transfers ownership. Rejects null data, wrong crates, oversized input, and timestamps before a bounda...
Definition gSROCrate.cc:233
std::function< void(std::exception_ptr)> FailureCallback
Definition gSROCrate.h:76
void advance_time(GSROTime safe_time)
Queue a nondecreasing boundary after all earlier input has been submitted. Equality is permitted.
Definition gSROCrate.cc:234
GSROCrate(GSROCrateId crate_id, ResourceFactory make_resources, GSROCrateLimits limits={}, DeliveryCallback on_delivery={}, FailureCallback on_failure={})
Definition gSROCrate.cc:223
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 request_finish(GSROEndContext context)
Close input and wake producers without joining; permits a service to stop every crate before waiting.
Definition gSROCrate.cc:235
void finish_and_join(GSROEndContext context)
Definition gSROCrate.cc:237
std::function< void(GSROEventId, std::uint64_t)> DeliveryCallback
Definition gSROCrate.h:75
std::chrono::duration< std::int64_t, std::nano > GSROTime
Definition gSROData.h:20
std::uint32_t GSROCrateId
Definition gSROData.h:18
event
Per-crate bounds preventing asynchronous delivery from consuming unlimited memory.
Definition gSROCrate.h:29
Input-delivery state passed to the implementation at orderly shutdown.
Definition gSROData.h:94
std::optional< GSROTime > safe_time
Definition gSROData.h:96
One worker-produced contribution, transferred by move to its destination crate.
Definition gSROData.h:58
GSROTime time
Definition gSROData.h:62
std::unique_ptr< const GSROData > data
Definition gSROData.h:63
GSROCrateId crate_id
Definition gSROData.h:59