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.hpp
Go to the documentation of this file.
1#pragma once
2/**
3 * @file signal_consumer.hpp
4 * @brief SignalConsumer — JSON serialisation thread for relativistic signals.
5 *
6 * Module: include/srfm/stream/
7 * Owner: AGT-10 (Builder) — 2026-03-01
8 *
9 * Responsibility
10 * --------------
11 * Pop StreamRelativisticSignal objects from the output ring and write one
12 * JSON line per signal to a configurable output stream (default: stdout).
13 * Flush the output every FLUSH_INTERVAL signals to amortise syscall cost.
14 *
15 * Output format (one line per signal):
16 * @code
17 * {"bar": 42, "beta": 0.4200, "gamma": 1.2309, "regime": "TIMELIKE", "signal": 0.8700}
18 * @endcode
19 *
20 * Guarantees
21 * ----------
22 * • Batched I/O: flushes every FLUSH_INTERVAL (default 100) signals, not per signal.
23 * • Non-throwing: the serialisation loop never throws.
24 * • Final flush: stop() performs a final flush before joining the thread.
25 * • Configurable sink: constructor accepts any FILE* for testability.
26 *
27 * NOT Responsible For
28 * -------------------
29 * • Signal computation (SignalProcessor)
30 * • Routing or fanout (one sink only)
31 */
32
33#include "spsc_ring.hpp"
34#include "stream_signal.hpp"
35
36#include <atomic>
37#include <cstdint>
38#include <cstdio>
39#include <thread>
40
41namespace srfm::stream {
42
43// ── SignalConsumerCounters ────────────────────────────────────────────────────
44
45/**
46 * @brief Diagnostic counters for the signal consumer thread.
47 */
49 std::uint64_t signals_consumed{0}; ///< Total signals written to output.
50 std::uint64_t flush_count{0}; ///< Number of explicit flushes performed.
51};
52
53// ── SignalConsumer ────────────────────────────────────────────────────────────
54
55/**
56 * @brief JSON output thread owner. Non-copyable, non-movable.
57 *
58 * @code
59 * SPSCRing<StreamRelativisticSignal, 65536> out_ring;
60 * SignalConsumer consumer{out_ring}; // writes to stdout
61 * consumer.start();
62 * // ...
63 * consumer.stop();
64 * @endcode
65 */
67public:
68 static constexpr std::size_t RING_SIZE = 65536;
69 static constexpr std::size_t FLUSH_INTERVAL = 100; ///< Flush every N signals.
70
72
73 /**
74 * @brief Construct bound to the output ring and an I/O sink.
75 *
76 * @param ring Ring populated by SignalProcessor (non-owning ref).
77 * @param sink Output stream (default: stdout). Must remain valid for the
78 * lifetime of this SignalConsumer.
79 */
80 explicit SignalConsumer(Ring& ring, FILE* sink = stdout) noexcept
81 : ring_{ring}, sink_{sink}
82 {}
83
84 ~SignalConsumer() noexcept { stop(); }
85
90
91 // ── Lifecycle ──────────────────────────────────────────────────────────────
92
93 /// Launch the consumer thread. Idempotent.
94 void start() noexcept;
95
96 /// Stop the consumer thread and perform a final flush. Idempotent.
97 void stop() noexcept;
98
99 /// Whether the consumer thread is currently running.
100 [[nodiscard]] bool running() const noexcept {
101 return running_.load(std::memory_order_acquire);
102 }
103
104 // ── Single-signal serialisation (usable from test without threading) ───────
105
106 /**
107 * @brief Serialise one signal to the sink immediately.
108 *
109 * Does not flush. Useful in unit tests.
110 *
111 * @param sig Signal to serialise.
112 */
113 void write_one(const StreamRelativisticSignal& sig) noexcept;
114
115 /**
116 * @brief Flush the output sink.
117 *
118 * Calls std::fflush(sink_). noexcept.
119 */
120 void flush() noexcept { std::fflush(sink_); }
121
122 // ── Accessors ──────────────────────────────────────────────────────────────
123
124 [[nodiscard]] SignalConsumerCounters counters() const noexcept {
125 return counters_;
126 }
127
128 [[nodiscard]] FILE* sink() const noexcept { return sink_; }
129
130private:
131 void run_loop() noexcept;
132
133 Ring& ring_;
134 FILE* sink_;
135
136 std::atomic<bool> running_{false};
137 std::atomic<bool> stop_requested_{false};
138 std::thread thread_;
139 SignalConsumerCounters counters_{};
140};
141
142} // namespace srfm::stream
Lock-free single-producer / single-consumer ring buffer.
Definition spsc_ring.hpp:85
JSON output thread owner. Non-copyable, non-movable.
SignalConsumer(Ring &ring, FILE *sink=stdout) noexcept
Construct bound to the output ring and an I/O sink.
SignalConsumer(const SignalConsumer &)=delete
SignalConsumer & operator=(const SignalConsumer &)=delete
SignalConsumerCounters counters() const noexcept
bool running() const noexcept
Whether the consumer thread is currently running.
SPSCRing< StreamRelativisticSignal, RING_SIZE > Ring
void start() noexcept
Launch the consumer thread. Idempotent.
void stop() noexcept
Stop the consumer thread and perform a final flush. Idempotent.
FILE * sink() const noexcept
static constexpr std::size_t RING_SIZE
SignalConsumer(SignalConsumer &&)=delete
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.
SignalConsumer & operator=(SignalConsumer &&)=delete
Lock-free Single-Producer / Single-Consumer ring buffer.
StreamRelativisticSignal — output unit of the signal-processing pipeline.
Diagnostic counters for the signal consumer thread.
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.