28#include "../../include/srfm/stream/signal_consumer.hpp"
38 if (running_.load(std::memory_order_acquire))
return;
40 stop_requested_.store(
false, std::memory_order_release);
41 running_.store(
true, std::memory_order_release);
43 thread_ = std::thread([
this]()
noexcept { run_loop(); });
49 stop_requested_.store(
true, std::memory_order_release);
50 if (thread_.joinable()) thread_.join();
53 running_.store(
false, std::memory_order_release);
60 "{\"bar\": %lld, \"beta\": %.4f, \"gamma\": %.4f,"
61 " \"regime\": \"%s\", \"signal\": %.4f}\n",
62 static_cast<long long>(sig.bar),
71void SignalConsumer::run_loop() noexcept {
72 std::size_t since_flush = 0;
74 while (!stop_requested_.load(std::memory_order_acquire)) {
76 auto maybe_sig = ring_.
pop();
78 if (!maybe_sig.has_value()) {
80 if (since_flush > 0) {
85 std::this_thread::yield();
102 auto maybe_sig = ring_.
pop();
103 if (!maybe_sig.has_value())
break;
110 if (since_flush > 0) {
115 running_.store(
false, std::memory_order_release);
std::optional< T > pop() noexcept
Try to dequeue one element.
void start() noexcept
Launch the consumer thread. Idempotent.
void stop() noexcept
Stop the consumer thread and perform a final flush. Idempotent.
void write_one(const StreamRelativisticSignal &sig) noexcept
Serialise one signal to the sink immediately.
void flush() noexcept
Flush the output sink.
static constexpr std::size_t FLUSH_INTERVAL
Flush every N signals.
const char * regime_to_str(Regime r) noexcept
Convert Regime to a null-terminated ASCII string.
std::uint64_t signals_consumed
Total signals written to output.
std::uint64_t flush_count
Number of explicit flushes performed.
Fully-characterised relativistic signal for one processed tick.