Special Relativity in Financial Modeling 1.0.0
Lorentz transforms, spacetime classification, and geodesic price paths for quantitative finance
Loading...
Searching...
No Matches
spsc_ring.hpp
Go to the documentation of this file.
1#pragma once
2/**
3 * @file spsc_ring.hpp
4 * @brief Lock-free Single-Producer / Single-Consumer ring buffer.
5 *
6 * Module: include/srfm/stream/
7 * Owner: AGT-10 (Builder) — 2026-03-01
8 *
9 * Responsibility
10 * --------------
11 * Provide a zero-heap-allocation, lock-free SPSC queue with sub-microsecond
12 * throughput. This is the sole inter-thread communication primitive in the
13 * streaming pipeline.
14 *
15 * Guarantees
16 * ----------
17 * • Deterministic: push()/pop() never allocate heap memory.
18 * • Lock-free: no mutexes, no condition variables, no spin-locks.
19 * • Cache-friendly: producer and consumer indices live on separate
20 * cache lines (alignas(64)) to eliminate false sharing.
21 * • Correct: all cross-thread visibility relies on C++11 acquire/release
22 * atomic ordering — no relaxed cross-thread loads.
23 * • Non-blocking: push() returns false when full; pop() returns nullopt when
24 * empty. Neither ever blocks the calling thread.
25 * • noexcept: both push() and pop() are unconditionally noexcept, provided
26 * T's move-constructor and move-assignment are noexcept.
27 *
28 * Template Parameters
29 * -------------------
30 * T — Element type. Must be noexcept move-constructible and
31 * noexcept move-assignable.
32 * SIZE — Capacity in elements. Must be an exact power of two.
33 * Valid range: [2, 2^30].
34 *
35 * Usage
36 * -----
37 * @code
38 * SPSCRing<OHLCVTick, 65536> ring;
39 *
40 * // Producer thread:
41 * OHLCVTick tick = ...;
42 * if (!ring.push(std::move(tick))) { ++drop_count; }
43 *
44 * // Consumer thread:
45 * while (auto t = ring.pop()) { process(*t); }
46 * @endcode
47 *
48 * NOT Responsible For
49 * -------------------
50 * • Multi-producer or multi-consumer safety — use one producer, one consumer.
51 * • Blocking wait — callers must implement their own spin or yield logic.
52 * • Persistence — in-memory only.
53 *
54 * Memory Layout
55 * -------------
56 * The buffer is stored inline (std::array<T, SIZE>), so the total object size
57 * is approximately SIZE * sizeof(T) + 128 bytes (two padded indices).
58 * For SIZE=65536 and OHLCVTick (48 bytes): ~3 MiB per ring.
59 */
60
61#include <array>
62#include <atomic>
63#include <cstddef>
64#include <optional>
65#include <type_traits>
66
67// Suppress C4324 (structure was padded due to alignment specifier) — padding
68// between tail_ and head_ is intentional; it eliminates false sharing.
69#ifdef _MSC_VER
70# pragma warning(push)
71# pragma warning(disable: 4324)
72#endif
73
74namespace srfm::stream {
75
76// ── SPSCRing ──────────────────────────────────────────────────────────────────
77
78/**
79 * @brief Lock-free single-producer / single-consumer ring buffer.
80 *
81 * @tparam T Element type (noexcept move-constructible + move-assignable).
82 * @tparam SIZE Capacity; must be a power of two.
83 */
84template<typename T, std::size_t SIZE>
85class SPSCRing {
86 // ── Compile-time invariants ────────────────────────────────────────────────
87 static_assert(SIZE >= 2,
88 "SPSCRing: SIZE must be at least 2");
89 static_assert((SIZE & (SIZE - 1u)) == 0u,
90 "SPSCRing: SIZE must be a power of two");
91 static_assert(SIZE <= (1u << 30u),
92 "SPSCRing: SIZE exceeds safe limit (2^30)");
93 static_assert(std::is_nothrow_move_constructible_v<T>,
94 "SPSCRing: T must be noexcept move-constructible");
95 static_assert(std::is_nothrow_move_assignable_v<T>,
96 "SPSCRing: T must be noexcept move-assignable");
97
98 static constexpr std::size_t MASK = SIZE - 1u;
99
100public:
101 // ── Capacity query ─────────────────────────────────────────────────────────
102
103 /// Maximum number of elements the ring can hold simultaneously.
104 [[nodiscard]] static constexpr std::size_t capacity() noexcept {
105 return SIZE - 1u; // One slot reserved to distinguish full from empty.
106 }
107
108 // ── Producer API (call from a single producer thread only) ─────────────────
109
110 /**
111 * @brief Try to enqueue one element.
112 *
113 * Moves @p item into the ring. Returns immediately without blocking.
114 *
115 * @param item Element to enqueue (moved from).
116 * @return true on success.
117 * @return false if the ring is full; @p item is left in a moved-from state.
118 *
119 * @note noexcept — never throws, never allocates.
120 * @note Call from the producer thread only.
121 */
122 [[nodiscard]] bool push(T&& item) noexcept {
123 const std::size_t tail = tail_.load(std::memory_order_relaxed);
124 const std::size_t next = (tail + 1u) & MASK;
125
126 // If the next slot equals head_, the ring is full.
127 if (next == head_.load(std::memory_order_acquire)) {
128 return false;
129 }
130
131 buffer_[tail] = std::move(item);
132
133 // Release: make the written element visible to the consumer.
134 tail_.store(next, std::memory_order_release);
135 return true;
136 }
137
138 /**
139 * @brief Try to enqueue a copy of one element.
140 *
141 * Convenience overload for copy-constructible types.
142 *
143 * @param item Element to copy-enqueue.
144 * @return true on success, false if full.
145 */
146 [[nodiscard]] bool push_copy(const T& item) noexcept(
147 std::is_nothrow_copy_constructible_v<T>) {
148 T copy{item};
149 return push(std::move(copy));
150 }
151
152 // ── Consumer API (call from a single consumer thread only) ─────────────────
153
154 /**
155 * @brief Try to dequeue one element.
156 *
157 * Returns immediately without blocking.
158 *
159 * @return std::optional<T> containing the dequeued element, or
160 * std::nullopt if the ring is empty.
161 *
162 * @note noexcept — never throws, never allocates.
163 * @note Call from the consumer thread only.
164 */
165 [[nodiscard]] std::optional<T> pop() noexcept {
166 const std::size_t head = head_.load(std::memory_order_relaxed);
167
168 // If head == tail_, the ring is empty.
169 if (head == tail_.load(std::memory_order_acquire)) {
170 return std::nullopt;
171 }
172
173 T item{std::move(buffer_[head])};
174
175 // Release: make the freed slot visible to the producer.
176 head_.store((head + 1u) & MASK, std::memory_order_release);
177 return item;
178 }
179
180 // ── Diagnostic queries (approximate — non-atomic snapshot) ─────────────────
181
182 /**
183 * @brief Approximate number of elements currently in the ring.
184 *
185 * Not guaranteed to be exact under concurrent use (non-atomic snapshot of
186 * two separate atomics). Suitable for monitoring / metrics only.
187 */
188 [[nodiscard]] std::size_t size_approx() const noexcept {
189 const std::size_t t = tail_.load(std::memory_order_acquire);
190 const std::size_t h = head_.load(std::memory_order_acquire);
191 return (t - h + SIZE) & MASK;
192 }
193
194 /**
195 * @brief Approximate check for emptiness.
196 *
197 * Non-atomic snapshot — use for monitoring only.
198 */
199 [[nodiscard]] bool empty_approx() const noexcept {
200 return tail_.load(std::memory_order_acquire) ==
201 head_.load(std::memory_order_acquire);
202 }
203
204 /**
205 * @brief Approximate check for fullness.
206 *
207 * Non-atomic snapshot — use for monitoring only.
208 */
209 [[nodiscard]] bool full_approx() const noexcept {
210 const std::size_t t = tail_.load(std::memory_order_acquire);
211 const std::size_t h = head_.load(std::memory_order_acquire);
212 return ((t + 1u) & MASK) == h;
213 }
214
215private:
216 // ── Storage ────────────────────────────────────────────────────────────────
217
218 // Producer index: next slot to write.
219 // alignas(64) places this on its own cache line.
220 alignas(64) std::atomic<std::size_t> tail_{0};
221
222 // Consumer index: next slot to read.
223 // alignas(64) separates from tail_ — eliminates false sharing.
224 alignas(64) std::atomic<std::size_t> head_{0};
225
226 // Element storage. Inline array — zero heap allocation.
227 std::array<T, SIZE> buffer_{};
228};
229
230} // namespace srfm::stream
231
232#ifdef _MSC_VER
233# pragma warning(pop)
234#endif
Lock-free single-producer / single-consumer ring buffer.
Definition spsc_ring.hpp:85
std::optional< T > pop() noexcept
Try to dequeue one element.
bool full_approx() const noexcept
Approximate check for fullness.
bool push_copy(const T &item) noexcept(std::is_nothrow_copy_constructible_v< T >)
Try to enqueue a copy of one element.
static constexpr std::size_t capacity() noexcept
Maximum number of elements the ring can hold simultaneously.
std::size_t size_approx() const noexcept
Approximate number of elements currently in the ring.
bool empty_approx() const noexcept
Approximate check for emptiness.
bool push(T &&item) noexcept
Try to enqueue one element.