Special Relativity in Financial Modeling 1.0.0
Lorentz transforms, spacetime classification, and geodesic price paths for quantitative finance
Loading...
Searching...
No Matches
tick_ingester.cpp
Go to the documentation of this file.
1/**
2 * @file tick_ingester.cpp
3 * @brief TickIngester implementation — ingestion thread body.
4 *
5 * See include/srfm/stream/tick_ingester.hpp for the full module contract.
6 *
7 * Hot-path design
8 * ---------------
9 * The inner loop in run_loop() is:
10 * 1. Call source_.read() — returns optional<OHLCVTick> or nullopt
11 * 2. If nullopt and source closed → exit loop
12 * 3. If nullopt and source open → yield and continue (spin with backoff)
13 * 4. Validate via tick_is_valid()
14 * 5. Push to ring1_
15 * 6. Update counters
16 *
17 * All operations are noexcept. No heap allocation occurs in the loop.
18 *
19 * Backoff strategy
20 * ----------------
21 * When source_.read() returns nullopt (no data available), we call
22 * std::this_thread::yield() to be cooperative with the OS scheduler rather
23 * than burning a core in a tight spin. A more aggressive variant could use
24 * _mm_pause() on x86 or std::this_thread::sleep_for(0ns), but yield() gives
25 * a reasonable balance between latency and CPU utilisation for this use-case.
26 */
27
28#include "../../include/srfm/stream/tick_ingester.hpp"
29
30#include <thread>
31
32namespace srfm::stream {
33
34// ── TickIngester::start ───────────────────────────────────────────────────────
35
36void TickIngester::start() noexcept {
37 // Idempotent: do nothing if already running.
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// ── TickIngester::stop ────────────────────────────────────────────────────────
47
48void TickIngester::stop() noexcept {
49 // Signal the loop to exit.
50 stop_requested_.store(true, std::memory_order_release);
51
52 // Join if the thread is joinable.
53 if (thread_.joinable()) {
54 thread_.join();
55 }
56
57 running_.store(false, std::memory_order_release);
58}
59
60// ── TickIngester::run_loop ────────────────────────────────────────────────────
61
62void TickIngester::run_loop() noexcept {
63 while (!stop_requested_.load(std::memory_order_acquire)) {
64
65 // ── Read one tick from the source ──────────────────────────────────────
66 auto maybe_tick = source_.read();
67
68 if (!maybe_tick.has_value()) {
69 if (!source_.is_open()) {
70 // Source closed permanently — exit the loop.
71 break;
72 }
73 // Transient: no data yet. Yield and retry.
74 std::this_thread::yield();
75 continue;
76 }
77
78 ++counters_.ticks_received;
79
80 // ── Validate ──────────────────────────────────────────────────────────
81 if (!tick_is_valid(*maybe_tick)) {
82 ++counters_.ticks_dropped_invalid;
83 continue;
84 }
85
86 // ── Push to ring ──────────────────────────────────────────────────────
87 if (!ring_.push(std::move(*maybe_tick))) {
88 ++counters_.ticks_dropped_ring_full;
89 continue;
90 }
91
92 ++counters_.ticks_pushed;
93 }
94
95 running_.store(false, std::memory_order_release);
96}
97
98} // namespace srfm::stream
bool push(T &&item) noexcept
Try to enqueue one element.
void stop() noexcept
Signal the ingestion thread to stop and wait for it to finish.
void start() noexcept
Launch the ingestion thread.
virtual bool is_open() const noexcept=0
Whether the source is still usable.
virtual std::optional< OHLCVTick > read() noexcept=0
Read the next OHLCVTick from the source.
bool tick_is_valid(const OHLCVTick &t) noexcept
Validate all fields of an OHLCVTick.
Definition tick.hpp:73
std::uint64_t ticks_dropped_ring_full
Dropped: ring was full.
std::uint64_t ticks_received
Raw ticks read from source.
std::uint64_t ticks_dropped_invalid
Dropped: failed tick_is_valid().
std::uint64_t ticks_pushed
Valid ticks pushed to ring.