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.hpp
Go to the documentation of this file.
1#pragma once
2/**
3 * @file tick_ingester.hpp
4 * @brief TickIngester — ingestion thread that reads, validates, and ring-pushes ticks.
5 *
6 * Module: include/srfm/stream/
7 * Owner: AGT-10 (Builder) — 2026-03-01
8 *
9 * Responsibility
10 * --------------
11 * Run the ingestion thread: repeatedly call source.read(), validate each tick
12 * with tick_is_valid(), push valid ticks into SPSCRing<OHLCVTick, 65536>, and
13 * count dropped ticks by rejection category.
14 *
15 * Guarantees
16 * ----------
17 * • Never blocks the calling thread (start() launches a std::thread).
18 * • Never throws from the hot path (noexcept read/validate/push loop).
19 * • Counts all dropped ticks; drops are never silently lost.
20 * • Graceful shutdown: stop() signals the thread and join()s within the
21 * destructor if the caller forgets to call stop().
22 *
23 * NOT Responsible For
24 * -------------------
25 * • Sourcing ticks (TickSource)
26 * • Signal processing (SignalProcessor)
27 */
28
29#include "spsc_ring.hpp"
30#include "tick.hpp"
31#include "tick_source.hpp"
32
33#include <atomic>
34#include <cstdint>
35#include <thread>
36
37namespace srfm::stream {
38
39// ── TickIngesterCounters ──────────────────────────────────────────────────────
40
41/**
42 * @brief Diagnostic counters for the ingestion thread.
43 *
44 * All fields are updated by the ingestion thread only. Reads from external
45 * threads are approximate (no synchronisation guarantee beyond coherence).
46 */
48 std::uint64_t ticks_received{0}; ///< Raw ticks read from source.
49 std::uint64_t ticks_pushed{0}; ///< Valid ticks pushed to ring.
50 std::uint64_t ticks_dropped_invalid{0}; ///< Dropped: failed tick_is_valid().
51 std::uint64_t ticks_dropped_ring_full{0}; ///< Dropped: ring was full.
52};
53
54// ── TickIngester ──────────────────────────────────────────────────────────────
55
56/**
57 * @brief Ingestion thread owner. Non-copyable, non-movable.
58 *
59 * @code
60 * SPSCRing<OHLCVTick, 65536> ring;
61 * QueueTickSource src;
62 * TickIngester ingester{src, ring};
63 * ingester.start();
64 * src.push(make_valid_tick());
65 * // ... let it run ...
66 * ingester.stop();
67 * @endcode
68 */
70public:
71 static constexpr std::size_t RING_SIZE = 65536;
73
74 /**
75 * @brief Construct bound to a source and a ring buffer.
76 *
77 * Both @p source and @p ring must outlive this TickIngester.
78 *
79 * @param source Tick source to read from (non-owning reference).
80 * @param ring Ring buffer to push valid ticks into (non-owning reference).
81 */
82 TickIngester(TickSource& source, Ring& ring) noexcept
83 : source_{source}, ring_{ring}
84 {}
85
86 /// Destructor: ensures the thread is stopped and joined.
87 ~TickIngester() noexcept { stop(); }
88
89 // Not copyable or movable — owns a thread handle.
90 TickIngester(const TickIngester&) = delete;
94
95 // ── Lifecycle ──────────────────────────────────────────────────────────────
96
97 /**
98 * @brief Launch the ingestion thread.
99 *
100 * Idempotent: a second call while running is a no-op.
101 */
102 void start() noexcept;
103
104 /**
105 * @brief Signal the ingestion thread to stop and wait for it to finish.
106 *
107 * Idempotent: safe to call multiple times.
108 */
109 void stop() noexcept;
110
111 /// Whether the ingestion thread is currently running.
112 [[nodiscard]] bool running() const noexcept {
113 return running_.load(std::memory_order_acquire);
114 }
115
116 // ── Diagnostic accessors ───────────────────────────────────────────────────
117
118 /**
119 * @brief Snapshot of ingestion counters.
120 *
121 * Non-atomic read — approximate under concurrency. Safe for monitoring.
122 */
123 [[nodiscard]] TickIngesterCounters counters() const noexcept {
124 return counters_;
125 }
126
127private:
128 void run_loop() noexcept;
129
130 TickSource& source_;
131 Ring& ring_;
132
133 std::atomic<bool> running_{false};
134 std::atomic<bool> stop_requested_{false};
135 std::thread thread_;
136 TickIngesterCounters counters_{};
137};
138
139} // namespace srfm::stream
Lock-free single-producer / single-consumer ring buffer.
Definition spsc_ring.hpp:85
Ingestion thread owner. Non-copyable, non-movable.
TickIngester & operator=(const TickIngester &)=delete
bool running() const noexcept
Whether the ingestion thread is currently running.
TickIngester & operator=(TickIngester &&)=delete
SPSCRing< OHLCVTick, RING_SIZE > Ring
TickIngester(const TickIngester &)=delete
~TickIngester() noexcept
Destructor: ensures the thread is stopped and joined.
TickIngester(TickIngester &&)=delete
void stop() noexcept
Signal the ingestion thread to stop and wait for it to finish.
TickIngesterCounters counters() const noexcept
Snapshot of ingestion counters.
static constexpr std::size_t RING_SIZE
void start() noexcept
Launch the ingestion thread.
TickIngester(TickSource &source, Ring &ring) noexcept
Construct bound to a source and a ring buffer.
Abstract source of raw OHLCVTick values.
Lock-free Single-Producer / Single-Consumer ring buffer.
Diagnostic counters for the ingestion thread.
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.
OHLCVTick — atomic market data unit for the lock-free streaming pipeline.
Abstract tick source interface + concrete implementations.