Special Relativity in Financial Modeling 1.0.0
Lorentz transforms, spacetime classification, and geodesic price paths for quantitative finance
Loading...
Searching...
No Matches
tick_source.hpp
Go to the documentation of this file.
1#pragma once
2/**
3 * @file tick_source.hpp
4 * @brief Abstract tick source interface + concrete implementations.
5 *
6 * Module: include/srfm/stream/
7 * Owner: AGT-10 (Builder) — 2026-03-01
8 *
9 * Responsibility
10 * --------------
11 * Abstract the byte-level transport from the TickIngester's validation and
12 * ring-push logic. The TickIngester accepts any TickSource* and calls read()
13 * in a tight loop.
14 *
15 * Concrete implementations provided here
16 * ---------------------------------------
17 * PipeTickSource — reads raw OHLCVTick bytes from a named pipe (Windows)
18 * or POSIX FIFO.
19 * QueueTickSource — in-process queue for unit-testing without I/O.
20 *
21 * NOT Responsible For
22 * -------------------
23 * • Tick validation (TickIngester::push_tick)
24 * • Ring-buffer push (TickIngester)
25 * • Connection retry (caller manages reconnect policy)
26 */
27
28#include "tick.hpp"
29
30#include <chrono>
31#include <deque>
32#include <optional>
33#include <string>
34#include <thread>
35
36#ifdef _WIN32
37# define WIN32_LEAN_AND_MEAN
38# define NOMINMAX
39# include <windows.h>
40#else
41# include <cerrno>
42# include <cstring>
43# include <fcntl.h>
44# include <unistd.h>
45#endif
46
47namespace srfm::stream {
48
49// ── TickSource (abstract) ─────────────────────────────────────────────────────
50
51/**
52 * @brief Abstract source of raw OHLCVTick values.
53 *
54 * TickIngester calls read() in its hot loop. Implementors must be thread-safe
55 * if called from the ingestion thread while other threads interact with the
56 * underlying resource (this is not required by the default implementations).
57 */
59public:
60 virtual ~TickSource() noexcept = default;
61
62 /**
63 * @brief Read the next OHLCVTick from the source.
64 *
65 * @return std::optional<OHLCVTick>:
66 * - has_value() → a tick was read; it may or may not pass validation.
67 * - nullopt → no data available right now (transient) or source
68 * closed (permanent — caller checks is_open()).
69 *
70 * @note Implementations must not throw.
71 */
72 [[nodiscard]] virtual std::optional<OHLCVTick> read() noexcept = 0;
73
74 /**
75 * @brief Whether the source is still usable.
76 *
77 * Returns false after a permanent error or after close() is called.
78 */
79 [[nodiscard]] virtual bool is_open() const noexcept = 0;
80
81 /**
82 * @brief Close the source and release any OS resources.
83 *
84 * Safe to call multiple times.
85 */
86 virtual void close() noexcept = 0;
87};
88
89// ── QueueTickSource ───────────────────────────────────────────────────────────
90
91#ifdef SRFM_TESTING
92// QueueTickSource is for unit tests only. Include this class only when
93// SRFM_TESTING is defined (set by CMake for test targets).
94
95/**
96 * @brief In-process tick source backed by a std::deque.
97 *
98 * Designed for unit tests and the StreamEngine::inject() path. Thread-safe
99 * enough for single-producer / single-consumer usage when the producer uses
100 * inject() from the test thread and the consumer calls read() from the
101 * ingestion thread. For concurrent push/pop correctness in tests the caller
102 * is responsible for ensuring only one side runs at a time.
103 */
104class QueueTickSource : public TickSource {
105public:
106 QueueTickSource() noexcept = default;
107
108 /**
109 * @brief Push one tick for future reads.
110 *
111 * Called from the test / injection thread.
112 */
113 void push(OHLCVTick tick) noexcept {
114 queue_.push_back(tick);
115 }
116
117 /**
118 * @brief Push multiple ticks at once.
119 */
120 template<typename InputIt>
121 void push_range(InputIt first, InputIt last) noexcept {
122 for (auto it = first; it != last; ++it) {
123 queue_.push_back(*it);
124 }
125 }
126
127 [[nodiscard]] std::optional<OHLCVTick> read() noexcept override {
128 if (queue_.empty()) return std::nullopt;
129 OHLCVTick t = queue_.front();
130 queue_.pop_front();
131 return t;
132 }
133
134 [[nodiscard]] bool is_open() const noexcept override { return open_; }
135
136 void close() noexcept override { open_ = false; }
137
138 /// Number of ticks waiting in the queue.
139 [[nodiscard]] std::size_t pending() const noexcept { return queue_.size(); }
140
141private:
142 std::deque<OHLCVTick> queue_;
143 bool open_{true};
144};
145
146#endif // SRFM_TESTING
147
148// ── PipeTickSource ────────────────────────────────────────────────────────────
149
150/**
151 * @brief Tick source that reads raw OHLCVTick bytes from a named pipe / FIFO.
152 *
153 * Wire protocol: little-endian binary, exactly sizeof(OHLCVTick) bytes per tick.
154 * No framing header — the ingester detects malformed ticks via tick_is_valid().
155 *
156 * On Windows: uses a named pipe (\\\\.\\pipe\<name>).
157 * On POSIX: uses a named FIFO at path <name>.
158 */
160public:
161 /**
162 * @brief Open the named pipe / FIFO for reading.
163 *
164 * @param name Pipe name (Windows) or FIFO path (POSIX).
165 * Windows: omit the \\.\pipe\ prefix — it is added automatically.
166 * POSIX: full path, e.g. "/tmp/srfm_ticks".
167 *
168 * On Windows, waits up to @p timeout_ms for the server to create the pipe.
169 * On POSIX, the FIFO must already exist (call mkfifo before constructing).
170 */
171 explicit PipeTickSource(const std::string& name,
172 unsigned int timeout_ms = 5000) noexcept {
173 open_pipe(name, timeout_ms);
174 }
175
176 ~PipeTickSource() noexcept override { close(); }
177
178 [[nodiscard]] std::optional<OHLCVTick> read() noexcept override {
179 if (!is_open()) return std::nullopt;
180
181 OHLCVTick tick{};
182 const bool ok = read_exact(reinterpret_cast<char*>(&tick),
183 sizeof(OHLCVTick));
184 if (!ok) {
185 close();
186 return std::nullopt;
187 }
188 return tick;
189 }
190
191 [[nodiscard]] bool is_open() const noexcept override {
192#ifdef _WIN32
193 return handle_ != INVALID_HANDLE_VALUE;
194#else
195 return fd_ >= 0;
196#endif
197 }
198
199 void close() noexcept override {
200#ifdef _WIN32
201 if (handle_ != INVALID_HANDLE_VALUE) {
202 CloseHandle(handle_);
203 handle_ = INVALID_HANDLE_VALUE;
204 }
205#else
206 if (fd_ >= 0) {
207 ::close(fd_);
208 fd_ = -1;
209 }
210#endif
211 }
212
213private:
214 void open_pipe(const std::string& name, unsigned int timeout_ms) noexcept {
215#ifdef _WIN32
216 const std::string full = "\\\\.\\pipe\\" + name;
217 const DWORD deadline = GetTickCount() + timeout_ms;
218 while (true) {
219 handle_ = CreateFileA(full.c_str(), GENERIC_READ, 0, nullptr,
220 OPEN_EXISTING, 0, nullptr);
221 if (handle_ != INVALID_HANDLE_VALUE) break;
222 if (GetLastError() != ERROR_PIPE_BUSY) break;
223 if (GetTickCount() >= deadline) break;
224 WaitNamedPipeA(full.c_str(), 100);
225 }
226#else
227 fd_ = ::open(name.c_str(), O_RDONLY | O_NONBLOCK);
228 (void)timeout_ms;
229#endif
230 }
231
232 bool read_exact(char* buf, std::size_t n) noexcept {
233#ifdef _WIN32
234 DWORD total = 0;
235 while (total < static_cast<DWORD>(n)) {
236 DWORD got = 0;
237 if (!ReadFile(handle_, buf + total,
238 static_cast<DWORD>(n) - total, &got, nullptr)) {
239 return false;
240 }
241 total += got;
242 }
243 return true;
244#else
245 std::size_t total = 0;
246 while (total < n) {
247 ssize_t got = ::read(fd_, buf + total, n - total);
248 if (got <= 0) return false;
249 total += static_cast<std::size_t>(got);
250 }
251 return true;
252#endif
253 }
254
255#ifdef _WIN32
256 HANDLE handle_{INVALID_HANDLE_VALUE};
257#else
258 int fd_{-1};
259#endif
260};
261
262// ── ReconnectingTickSource ────────────────────────────────────────────────────
263
264/**
265 * @brief Wraps any TickSource with automatic reconnection on read failure.
266 *
267 * Uses exponential backoff starting at @p initial_delay (default 100 ms),
268 * doubling each attempt up to @p max_delay (default 30 000 ms = 30 s).
269 * After a successful read the delay resets to @p initial_delay.
270 *
271 * Thread safety: same as the wrapped Source type.
272 *
273 * @par Example
274 * @code{.cpp}
275 * // Create a QueueTickSource wrapped with reconnect logic.
276 * // (Useful in tests to simulate transient read failures.)
277 * QueueTickSource inner;
278 * inner.push(OHLCVTick{...});
279 *
280 * ReconnectingTickSource<QueueTickSource> src(
281 * std::move(inner),
282 * std::chrono::milliseconds{50}, // initial backoff
283 * std::chrono::milliseconds{5000} // max backoff
284 * );
285 *
286 * // Returns nullopt (and sleeps briefly) when the queue is empty;
287 * // returns the tick and resets backoff on success.
288 * while (true) {
289 * auto tick = src.read();
290 * if (tick) process(*tick);
291 * }
292 * @endcode
293 */
294template<typename Source>
296public:
297 explicit ReconnectingTickSource(Source source,
298 std::chrono::milliseconds initial_delay = std::chrono::milliseconds{100},
299 std::chrono::milliseconds max_delay = std::chrono::milliseconds{30'000})
300 : source_(std::move(source))
301 , initial_delay_(initial_delay)
302 , max_delay_(max_delay)
303 , current_delay_(initial_delay)
304 {}
305
306 /// Read the next tick. Returns std::nullopt on transient failure (after
307 /// sleeping current_delay_). Reconnect is attempted automatically.
308 std::optional<OHLCVTick> read() {
309 auto result = source_.read();
310 if (result) {
311 current_delay_ = initial_delay_; // reset on success
312 return result;
313 }
314 // Back off and signal caller to retry
315 std::this_thread::sleep_for(current_delay_);
316 current_delay_ = std::min(current_delay_ * 2, max_delay_);
317 return std::nullopt;
318 }
319
320 /// Reset the backoff state without reconnecting the underlying source.
321 void reset_backoff() noexcept { current_delay_ = initial_delay_; }
322
323 /// Access the underlying source (e.g. to re-open a pipe).
324 Source& source() noexcept { return source_; }
325 const Source& source() const noexcept { return source_; }
326
327private:
328 Source source_;
329 std::chrono::milliseconds initial_delay_;
330 std::chrono::milliseconds max_delay_;
331 std::chrono::milliseconds current_delay_;
332};
333
334} // namespace srfm::stream
Tick source that reads raw OHLCVTick bytes from a named pipe / FIFO.
bool is_open() const noexcept override
Whether the source is still usable.
PipeTickSource(const std::string &name, unsigned int timeout_ms=5000) noexcept
Open the named pipe / FIFO for reading.
void close() noexcept override
Close the source and release any OS resources.
~PipeTickSource() noexcept override
std::optional< OHLCVTick > read() noexcept override
Read the next OHLCVTick from the source.
Wraps any TickSource with automatic reconnection on read failure.
void reset_backoff() noexcept
Reset the backoff state without reconnecting the underlying source.
Source & source() noexcept
Access the underlying source (e.g. to re-open a pipe).
ReconnectingTickSource(Source source, std::chrono::milliseconds initial_delay=std::chrono::milliseconds{100}, std::chrono::milliseconds max_delay=std::chrono::milliseconds{30 '000})
std::optional< OHLCVTick > read()
Read the next tick. Returns std::nullopt on transient failure (after sleeping current_delay_)....
const Source & source() const noexcept
Abstract source of raw OHLCVTick values.
virtual bool is_open() const noexcept=0
Whether the source is still usable.
virtual ~TickSource() noexcept=default
virtual void close() noexcept=0
Close the source and release any OS resources.
virtual std::optional< OHLCVTick > read() noexcept=0
Read the next OHLCVTick from the source.
Single OHLCV bar tick from the market data feed.
Definition tick.hpp:50
OHLCVTick — atomic market data unit for the lock-free streaming pipeline.