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.cpp
Go to the documentation of this file.
1/**
2 * @file stream_engine.cpp
3 * @brief StreamEngine implementation — pipeline lifecycle and inject() path.
4 *
5 * See include/srfm/engine/stream_engine.hpp for the full module contract.
6 *
7 * Startup / shutdown order
8 * ------------------------
9 * start() launches threads in data-flow order:
10 * 1. TickIngester (produces to Ring1)
11 * 2. SignalProcessor (consumes Ring1, produces Ring2)
12 * 3. SignalConsumer (consumes Ring2)
13 *
14 * stop() signals threads in the same order and joins them. This ensures
15 * in-flight ticks are processed before the downstream consumers exit.
16 *
17 * inject() path
18 * -------------
19 * inject() validates the tick (tick_is_valid) and pushes directly into Ring1.
20 * It does not interact with the TickIngester's loop or its counters — it is
21 * a separate producer path intended for tests.
22 *
23 * When the TickIngester is also running with a live TickSource, the caller
24 * MUST NOT mix inject() calls with TickSource production (SPSC constraint:
25 * only one producer thread at a time).
26 */
27
28#include "../../include/srfm/engine/stream_engine.hpp"
29
30namespace srfm::engine {
31
32// ── StreamEngine::start ───────────────────────────────────────────────────────
33
34void StreamEngine::start() noexcept {
35 if (running_.load(std::memory_order_acquire)) return;
36
37 ingester_.start();
38 processor_.start();
39 consumer_.start();
40
41 running_.store(true, std::memory_order_release);
42}
43
44// ── StreamEngine::stop ────────────────────────────────────────────────────────
45
46void StreamEngine::stop() noexcept {
47 if (!running_.load(std::memory_order_acquire)) return;
48
49 // Stop in data-flow order: cut off input first, then let downstream drain.
50 ingester_.stop();
51 processor_.stop();
52 consumer_.stop();
53
54 running_.store(false, std::memory_order_release);
55}
56
57// ── StreamEngine::inject ──────────────────────────────────────────────────────
58
60 if (!stream::tick_is_valid(tick)) {
61 ++inject_dropped_invalid_;
62 return false;
63 }
64 if (!ring1_.push(std::move(tick))) {
65 ++inject_dropped_ring_full_;
66 return false;
67 }
68 ++inject_pushed_;
69 return true;
70}
71
72// ── StreamEngine::counters ────────────────────────────────────────────────────
73
76 return ingester_.counters();
77}
78
81 return processor_.counters();
82}
83
86 return consumer_.counters();
87}
88
89// ── StreamEngine::process_sync ────────────────────────────────────────────────
90
93 std::int64_t bar_index) noexcept {
94 return processor_.process_one(tick, bar_index);
95}
96
97// ── StreamEngine::reset_processor_state ──────────────────────────────────────
98
100 processor_.reset_state();
101}
102
103} // namespace srfm::engine
stream::SignalConsumerCounters consumer_counters() const noexcept
bool inject(stream::OHLCVTick tick) noexcept
Inject a tick directly into Ring1, bypassing the TickSource.
stream::TickIngesterCounters ingester_counters() const noexcept
stream::SignalProcessorCounters processor_counters() const noexcept
void stop() noexcept
Stop all three pipeline threads.
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).
SignalConsumerCounters counters() const noexcept
void start() noexcept
Launch the consumer thread. Idempotent.
void stop() noexcept
Stop the consumer thread and perform a final flush. Idempotent.
SignalProcessorCounters counters() const noexcept
void stop() noexcept
Stop and join the processing thread. Idempotent.
void start() noexcept
Launch the processing thread. Idempotent.
void reset_state() noexcept
Reset all signal-chain component state.
void stop() noexcept
Signal the ingestion thread to stop and wait for it to finish.
TickIngesterCounters counters() const noexcept
Snapshot of ingestion counters.
void start() noexcept
Launch the ingestion thread.
bool tick_is_valid(const OHLCVTick &t) noexcept
Validate all fields of an OHLCVTick.
Definition tick.hpp:73
Single OHLCV bar tick from the market data feed.
Definition tick.hpp:50
Diagnostic counters for the signal consumer thread.
Diagnostic counters for the signal-processing thread.
Fully-characterised relativistic signal for one processed tick.
Diagnostic counters for the ingestion thread.