Special Relativity in Financial Modeling 1.0.0
Lorentz transforms, spacetime classification, and geodesic price paths for quantitative finance
Loading...
Searching...
No Matches
signal_consumer.cpp
Go to the documentation of this file.
1/**
2 * @file signal_consumer.cpp
3 * @brief SignalConsumer implementation — JSON serialisation and batched I/O.
4 *
5 * See include/srfm/stream/signal_consumer.hpp for the full module contract.
6 *
7 * JSON format (one line per signal)
8 * ----------------------------------
9 * {"bar": 0, "beta": 0.0000, "gamma": 1.0000, "regime": "LIGHTLIKE", "signal": 0.0000}
10 *
11 * - All floating-point fields printed with 4 decimal places.
12 * - No trailing comma.
13 * - Fields always appear in the order: bar, beta, gamma, regime, signal.
14 *
15 * I/O batching
16 * ------------
17 * std::fflush() is expensive. We accumulate FLUSH_INTERVAL signals in the
18 * internal buffer (backed by the OS I/O buffer) before flushing. On stop(),
19 * we drain the ring and perform a final flush to ensure no signals are lost.
20 *
21 * Output correctness
22 * ------------------
23 * std::fprintf() is used rather than std::cout to avoid locale effects and
24 * C++ stream synchronisation overhead. The format string is compile-time
25 * constant; no dynamic allocation occurs in write_one().
26 */
27
28#include "../../include/srfm/stream/signal_consumer.hpp"
29
30#include <cstdio>
31#include <thread>
32
33namespace srfm::stream {
34
35// ── SignalConsumer::start ─────────────────────────────────────────────────────
36
37void SignalConsumer::start() noexcept {
38 if (running_.load(std::memory_order_acquire)) return;
39
40 stop_requested_.store(false, std::memory_order_release);
41 running_.store(true, std::memory_order_release);
42
43 thread_ = std::thread([this]() noexcept { run_loop(); });
44}
45
46// ── SignalConsumer::stop ──────────────────────────────────────────────────────
47
48void SignalConsumer::stop() noexcept {
49 stop_requested_.store(true, std::memory_order_release);
50 if (thread_.joinable()) thread_.join();
51 // Final flush to ensure no buffered output is lost.
52 flush();
53 running_.store(false, std::memory_order_release);
54}
55
56// ── SignalConsumer::write_one ─────────────────────────────────────────────────
57
59 std::fprintf(sink_,
60 "{\"bar\": %lld, \"beta\": %.4f, \"gamma\": %.4f,"
61 " \"regime\": \"%s\", \"signal\": %.4f}\n",
62 static_cast<long long>(sig.bar),
63 sig.beta,
64 sig.gamma,
65 regime_to_str(sig.regime),
66 sig.signal);
67}
68
69// ── SignalConsumer::run_loop ──────────────────────────────────────────────────
70
71void SignalConsumer::run_loop() noexcept {
72 std::size_t since_flush = 0;
73
74 while (!stop_requested_.load(std::memory_order_acquire)) {
75
76 auto maybe_sig = ring_.pop();
77
78 if (!maybe_sig.has_value()) {
79 // Flush pending output if we're idle and have buffered data.
80 if (since_flush > 0) {
81 flush();
82 ++counters_.flush_count;
83 since_flush = 0;
84 }
85 std::this_thread::yield();
86 continue;
87 }
88
89 write_one(*maybe_sig);
90 ++counters_.signals_consumed;
91 ++since_flush;
92
93 if (since_flush >= FLUSH_INTERVAL) {
94 flush();
95 ++counters_.flush_count;
96 since_flush = 0;
97 }
98 }
99
100 // Drain any remaining signals that arrived before stop was signalled.
101 while (true) {
102 auto maybe_sig = ring_.pop();
103 if (!maybe_sig.has_value()) break;
104 write_one(*maybe_sig);
105 ++counters_.signals_consumed;
106 ++since_flush;
107 }
108
109 // Final flush.
110 if (since_flush > 0) {
111 flush();
112 ++counters_.flush_count;
113 }
114
115 running_.store(false, std::memory_order_release);
116}
117
118} // namespace srfm::stream
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.