Special Relativity in Financial Modeling 1.0.0
Lorentz transforms, spacetime classification, and geodesic price paths for quantitative finance
Loading...
Searching...
No Matches
stream_engine.hpp
Go to the documentation of this file.
1#pragma once
2/**
3 * @file stream_engine.hpp
4 * @brief StreamEngine — top-level owner of the lock-free streaming pipeline.
5 *
6 * Module: include/srfm/engine/
7 * Owner: AGT-10 (Builder) — 2026-03-01
8 *
9 * Responsibility
10 * --------------
11 * Own and coordinate the three threads and two ring buffers that form the
12 * streaming pipeline:
13 *
14 * TickSource → [TickIngester] → Ring1<OHLCVTick>
15 * → [SignalProcessor] → Ring2<StreamRelativisticSignal>
16 * → [SignalConsumer] → stdout (or configured sink)
17 *
18 * Provide a test-injection path (inject()) that bypasses the TickSource and
19 * pushes ticks directly into Ring1 from the calling thread.
20 *
21 * Guarantees
22 * ----------
23 * • Owns Ring1 and Ring2 (inline storage — no heap allocation for rings).
24 * • start() / stop() are idempotent.
25 * • Graceful shutdown: stop() signals all threads and joins them in order.
26 * • inject() is safe to call from any thread while the engine is running,
27 * provided only one external thread calls inject() (SPSC constraint on Ring1).
28 *
29 * NOT Responsible For
30 * -------------------
31 * • Tick validation (TickIngester)
32 * • Signal computation (SignalProcessor)
33 * • JSON formatting (SignalConsumer)
34 *
35 * @code
36 * QueueTickSource src;
37 * StreamEngine engine{src};
38 * engine.start();
39 * engine.inject(make_valid_tick());
40 * // let it process...
41 * engine.stop();
42 * @endcode
43 */
44
45#include "../stream/signal_consumer.hpp"
46#include "../stream/signal_processor.hpp"
47#include "../stream/spsc_ring.hpp"
48#include "../stream/stream_signal.hpp"
49#include "../stream/tick.hpp"
50#include "../stream/tick_ingester.hpp"
51#include "../stream/tick_source.hpp"
52
53#include <atomic>
54#include <cstdint>
55#include <cstdio>
56
57namespace srfm::engine {
58
59// ── StreamEngine ──────────────────────────────────────────────────────────────
60
61/**
62 * @brief Top-level streaming pipeline owner. Non-copyable, non-movable.
63 */
65public:
66 static constexpr std::size_t RING1_SIZE = 65536; ///< Tick ring capacity.
67 static constexpr std::size_t RING2_SIZE = 65536; ///< Signal ring capacity.
68
71
72 /**
73 * @brief Construct the engine with a tick source and optional output sink.
74 *
75 * @param source Source of OHLCVTick data (non-owning ref). Callers may
76 * pass a QueueTickSource for testing or a PipeTickSource for
77 * live operation.
78 * @param sink Output sink for JSON lines (default: stdout).
79 */
81 FILE* sink = stdout) noexcept
82 : ingester_{source, ring1_}
83 , processor_{ring1_, ring2_}
84 , consumer_{ring2_, sink}
85 {}
86
87 ~StreamEngine() noexcept { stop(); }
88
89 StreamEngine(const StreamEngine&) = delete;
93
94 // ── Lifecycle ──────────────────────────────────────────────────────────────
95
96 /**
97 * @brief Start all three pipeline threads.
98 *
99 * Idempotent. Order: ingester → processor → consumer.
100 */
101 void start() noexcept;
102
103 /**
104 * @brief Stop all three pipeline threads.
105 *
106 * Idempotent. Order: ingester → processor → consumer (drain order).
107 * Blocks until all threads have joined.
108 */
109 void stop() noexcept;
110
111 /// Whether the engine is currently running.
112 [[nodiscard]] bool running() const noexcept {
113 return running_.load(std::memory_order_acquire);
114 }
115
116 // ── Test injection path ────────────────────────────────────────────────────
117
118 /**
119 * @brief Inject a tick directly into Ring1, bypassing the TickSource.
120 *
121 * The tick is validated by tick_is_valid() before pushing. Invalid ticks
122 * are silently dropped; the drop count is tracked in inject_counters_.
123 *
124 * @param tick Tick to inject.
125 * @return true if the tick was pushed successfully.
126 * @return false if the tick was invalid or Ring1 was full.
127 *
128 * @note Thread constraint: only one external thread may call inject()
129 * concurrently (SPSC constraint on Ring1 producer side). When the
130 * TickIngester thread is also running and reading from a TickSource,
131 * do not mix inject() calls with that source producing ticks.
132 */
133 [[nodiscard]] bool inject(stream::OHLCVTick tick) noexcept;
134
135 // ── Accessors ──────────────────────────────────────────────────────────────
136
137 [[nodiscard]] stream::TickIngesterCounters ingester_counters() const noexcept;
138 [[nodiscard]] stream::SignalProcessorCounters processor_counters() const noexcept;
139 [[nodiscard]] stream::SignalConsumerCounters consumer_counters() const noexcept;
140
141 /// Approximate number of ticks waiting in Ring1.
142 [[nodiscard]] std::size_t ring1_size_approx() const noexcept {
143 return ring1_.size_approx();
144 }
145
146 /// Approximate number of signals waiting in Ring2.
147 [[nodiscard]] std::size_t ring2_size_approx() const noexcept {
148 return ring2_.size_approx();
149 }
150
151 // ── Direct processor access (for synchronous test use) ────────────────────
152
153 /**
154 * @brief Process one tick synchronously (no threads, no ring buffers).
155 *
156 * Calls SignalProcessor::process_one() directly. Useful for deterministic
157 * unit tests that do not need the full thread machinery.
158 *
159 * @param tick Validated tick to process.
160 * @param bar_index Sequence number assigned to this tick.
161 */
164 std::int64_t bar_index) noexcept;
165
166 /**
167 * @brief Reset all signal-processor state (for test reuse).
168 *
169 * Must only be called when the engine is stopped.
170 */
171 void reset_processor_state() noexcept;
172
173private:
174 Ring1 ring1_; ///< Tick ring: TickIngester → SignalProcessor.
175 Ring2 ring2_; ///< Signal ring: SignalProcessor → SignalConsumer.
176
177 stream::TickIngester ingester_;
178 stream::SignalProcessor processor_;
179 stream::SignalConsumer consumer_;
180
181 std::atomic<bool> running_{false};
182
183 // inject() diagnostic counters.
184 std::uint64_t inject_dropped_invalid_{0};
185 std::uint64_t inject_dropped_ring_full_{0};
186 std::uint64_t inject_pushed_{0};
187};
188
189} // namespace srfm::engine
Top-level streaming pipeline owner. Non-copyable, non-movable.
stream::SPSCRing< stream::OHLCVTick, RING1_SIZE > Ring1
stream::SignalConsumerCounters consumer_counters() const noexcept
StreamEngine(const StreamEngine &)=delete
StreamEngine(StreamEngine &&)=delete
static constexpr std::size_t RING2_SIZE
Signal ring capacity.
bool inject(stream::OHLCVTick tick) noexcept
Inject a tick directly into Ring1, bypassing the TickSource.
bool running() const noexcept
Whether the engine is currently running.
stream::TickIngesterCounters ingester_counters() const noexcept
StreamEngine & operator=(StreamEngine &&)=delete
StreamEngine(stream::TickSource &source, FILE *sink=stdout) noexcept
Construct the engine with a tick source and optional output sink.
stream::SPSCRing< stream::StreamRelativisticSignal, RING2_SIZE > Ring2
stream::SignalProcessorCounters processor_counters() const noexcept
void stop() noexcept
Stop all three pipeline threads.
std::size_t ring2_size_approx() const noexcept
Approximate number of signals waiting in Ring2.
void start() noexcept
Start all three pipeline threads.
stream::StreamRelativisticSignal process_sync(const stream::OHLCVTick &tick, std::int64_t bar_index) noexcept
Process one tick synchronously (no threads, no ring buffers).
void reset_processor_state() noexcept
Reset all signal-processor state (for test reuse).
std::size_t ring1_size_approx() const noexcept
Approximate number of ticks waiting in Ring1.
static constexpr std::size_t RING1_SIZE
Tick ring capacity.
StreamEngine & operator=(const StreamEngine &)=delete
std::size_t size_approx() const noexcept
Approximate number of elements currently in the ring.
Abstract source of raw OHLCVTick values.
Single OHLCV bar tick from the market data feed.
Definition tick.hpp:50
Fully-characterised relativistic signal for one processed tick.
Diagnostic counters for the ingestion thread.