Special Relativity in Financial Modeling 1.0.0
Lorentz transforms, spacetime classification, and geodesic price paths for quantitative finance
Loading...
Searching...
No Matches
signal_processor.hpp
Go to the documentation of this file.
1#pragma once
2/**
3 * @file signal_processor.hpp
4 * @brief SignalProcessor — per-tick relativistic signal computation thread.
5 *
6 * Module: include/srfm/stream/
7 * Owner: AGT-10 (Builder) — 2026-03-01
8 *
9 * Responsibility
10 * --------------
11 * Pop OHLCVTick objects from the ingestion ring, run each tick through the
12 * four-stage signal chain (normalise → β → Lorentz → manifold), compute a
13 * StreamRelativisticSignal, and push it into the output ring.
14 *
15 * Signal chain (per tick)
16 * -----------------------
17 * 1. CoordinateNormalizer (window=20): z-score(close) → norm_close
18 * 2. BetaCalculatorFix3: log-returns(close) → β
19 * 3. LorentzTransform: (bar_idx, norm_close, β) → (t', x', γ)
20 * 4. SpacetimeManifold: (t', x') → Δs², regime, signal
21 *
22 * Output StreamRelativisticSignal:
23 * bar = tick sequence number
24 * beta = β from BetaCalculatorFix3
25 * gamma = γ from LorentzTransform
26 * regime = TIMELIKE / LIGHTLIKE / SPACELIKE from SpacetimeManifold
27 * signal = γ * manifold_signal (Lorentz-scaled)
28 * ds2 = raw Δs² from SpacetimeManifold
29 *
30 * Guarantees
31 * ----------
32 * • No heap allocation in the hot path — all state pre-allocated in ctor.
33 * • Never throws from the processing loop.
34 * • Graceful shutdown via stop().
35 *
36 * NOT Responsible For
37 * -------------------
38 * • Tick validation (TickIngester)
39 * • JSON serialisation (SignalConsumer)
40 */
41
42#include "beta_calculator.hpp"
44#include "lorentz_transform.hpp"
46#include "spsc_ring.hpp"
47#include "stream_signal.hpp"
48#include "tick.hpp"
49
50#include <atomic>
51#include <cstdint>
52#include <thread>
53
54namespace srfm::stream {
55
56// ── SignalProcessorCounters ───────────────────────────────────────────────────
57
58/**
59 * @brief Diagnostic counters for the signal-processing thread.
60 */
62 std::uint64_t ticks_processed{0}; ///< Ticks popped from input ring.
63 std::uint64_t signals_emitted{0}; ///< Signals pushed to output ring.
64 std::uint64_t signals_dropped_ring_full{0}; ///< Dropped: output ring full.
65 std::uint64_t ticks_rejected_invalid{0}; ///< Dropped: NaN/Inf/non-positive field.
66 std::uint64_t timelike_count{0}; ///< TIMELIKE regime events.
67 std::uint64_t lightlike_count{0}; ///< LIGHTLIKE regime events.
68 std::uint64_t spacelike_count{0}; ///< SPACELIKE regime events.
69};
70
71// ── SignalProcessor ───────────────────────────────────────────────────────────
72
73/**
74 * @brief Signal processing thread owner. Non-copyable, non-movable.
75 *
76 * @code
77 * SPSCRing<OHLCVTick, 65536> in_ring;
78 * SPSCRing<StreamRelativisticSignal, 65536> out_ring;
79 * SignalProcessor proc{in_ring, out_ring};
80 * proc.start();
81 * // ...
82 * proc.stop();
83 * @endcode
84 */
86public:
87 static constexpr std::size_t IN_RING_SIZE = 65536;
88 static constexpr std::size_t OUT_RING_SIZE = 65536;
89
92
93 static constexpr std::size_t NORMALIZER_WINDOW = 20;
94
95 /**
96 * @brief Construct bound to the input and output ring buffers.
97 *
98 * All signal-chain objects are default-constructed here; no allocation
99 * occurs in the processing loop.
100 *
101 * @param in_ring Ring populated by TickIngester (non-owning ref).
102 * @param out_ring Ring consumed by SignalConsumer (non-owning ref).
103 */
104 SignalProcessor(InRing& in_ring, OutRing& out_ring) noexcept
105 : in_ring_{in_ring}
106 , out_ring_{out_ring}
107 , normalizer_{NORMALIZER_WINDOW}
108 {}
109
110 ~SignalProcessor() noexcept { stop(); }
111
116
117 // ── Lifecycle ──────────────────────────────────────────────────────────────
118
119 /// Launch the processing thread. Idempotent.
120 void start() noexcept;
121
122 /// Stop and join the processing thread. Idempotent.
123 void stop() noexcept;
124
125 /// Whether the processing thread is currently running.
126 [[nodiscard]] bool running() const noexcept {
127 return running_.load(std::memory_order_acquire);
128 }
129
130 // ── Single-tick processing (usable from test without threading) ────────────
131
132 /**
133 * @brief Process one tick synchronously and return the resulting signal.
134 *
135 * Does not push to any ring. Useful for unit tests.
136 *
137 * @param tick Validated tick to process.
138 * @param bar_index Sequence number for this tick.
139 * @return The computed StreamRelativisticSignal.
140 */
141 [[nodiscard]] StreamRelativisticSignal process_one(const OHLCVTick& tick,
142 std::int64_t bar_index) noexcept;
143
144 // ── Diagnostic accessors ───────────────────────────────────────────────────
145
146 [[nodiscard]] SignalProcessorCounters counters() const noexcept {
147 return counters_;
148 }
149
150 // ── State reset (for test reuse) ──────────────────────────────────────────
151
152 /**
153 * @brief Reset all signal-chain component state.
154 *
155 * Must only be called when the processing thread is not running.
156 */
157 void reset_state() noexcept {
158 normalizer_.reset();
159 beta_calc_.reset();
160 manifold_.reset();
161 bar_counter_ = 0;
162 counters_ = {};
163 }
164
165private:
166 void run_loop() noexcept;
167
168 InRing& in_ring_;
169 OutRing& out_ring_;
170
171 // Signal-chain components — pre-allocated, zero per-tick allocation.
172 CoordinateNormalizer normalizer_;
173 BetaCalculatorFix3 beta_calc_{};
174 LorentzTransform lorentz_{};
175 SpacetimeManifold manifold_{};
176
177 std::int64_t bar_counter_{0};
178 std::atomic<bool> running_{false};
179 std::atomic<bool> stop_requested_{false};
180 std::thread thread_;
181 SignalProcessorCounters counters_{};
182};
183
184} // namespace srfm::stream
void reset() noexcept
Reset all state as if no ticks have been seen.
Rolling z-score normaliser using Welford's online algorithm.
void reset() noexcept
Reset all state as if no ticks have been seen.
Lock-free single-producer / single-consumer ring buffer.
Definition spsc_ring.hpp:85
Signal processing thread owner. Non-copyable, non-movable.
SignalProcessorCounters counters() const noexcept
void stop() noexcept
Stop and join the processing thread. Idempotent.
SPSCRing< StreamRelativisticSignal, OUT_RING_SIZE > OutRing
static constexpr std::size_t IN_RING_SIZE
void start() noexcept
Launch the processing thread. Idempotent.
static constexpr std::size_t NORMALIZER_WINDOW
void reset_state() noexcept
Reset all signal-chain component state.
static constexpr std::size_t OUT_RING_SIZE
SignalProcessor(const SignalProcessor &)=delete
StreamRelativisticSignal process_one(const OHLCVTick &tick, std::int64_t bar_index) noexcept
Process one tick synchronously and return the resulting signal.
SPSCRing< OHLCVTick, IN_RING_SIZE > InRing
SignalProcessor(SignalProcessor &&)=delete
SignalProcessor(InRing &in_ring, OutRing &out_ring) noexcept
Construct bound to the input and output ring buffers.
SignalProcessor & operator=(SignalProcessor &&)=delete
SignalProcessor & operator=(const SignalProcessor &)=delete
bool running() const noexcept
Whether the processing thread is currently running.
void reset() noexcept
Reset all state as if no events have been seen.
Rolling-window z-score normaliser for tick close prices.
Lock-free Single-Producer / Single-Consumer ring buffer.
BetaCalculator — financial-to-physics velocity mapping (AGT-01).
Lorentz Transform Engine — AGT-01 public header.
Spacetime manifold processor with Christoffel symbols (AGT-13 / SRFM)
StreamRelativisticSignal — output unit of the signal-processing pipeline.
Single OHLCV bar tick from the market data feed.
Definition tick.hpp:50
Diagnostic counters for the signal-processing thread.
std::uint64_t timelike_count
TIMELIKE regime events.
std::uint64_t signals_emitted
Signals pushed to output ring.
std::uint64_t ticks_rejected_invalid
Dropped: NaN/Inf/non-positive field.
std::uint64_t lightlike_count
LIGHTLIKE regime events.
std::uint64_t ticks_processed
Ticks popped from input ring.
std::uint64_t spacelike_count
SPACELIKE regime events.
std::uint64_t signals_dropped_ring_full
Dropped: output ring full.
Fully-characterised relativistic signal for one processed tick.
OHLCVTick — atomic market data unit for the lock-free streaming pipeline.