Top-level streaming pipeline owner. Non-copyable, non-movable.
More...
#include <stream_engine.hpp>
|
| static constexpr std::size_t | RING1_SIZE = 65536 |
| | Tick ring capacity.
|
| |
| static constexpr std::size_t | RING2_SIZE = 65536 |
| | Signal ring capacity.
|
| |
Top-level streaming pipeline owner. Non-copyable, non-movable.
Definition at line 64 of file stream_engine.hpp.
◆ Ring1
◆ Ring2
◆ StreamEngine() [1/3]
| srfm::engine::StreamEngine::StreamEngine |
( |
stream::TickSource & |
source, |
|
|
FILE * |
sink = stdout |
|
) |
| |
|
inlineexplicitnoexcept |
Construct the engine with a tick source and optional output sink.
- Parameters
-
| source | Source of OHLCVTick data (non-owning ref). Callers may pass a QueueTickSource for testing or a PipeTickSource for live operation. |
| sink | Output sink for JSON lines (default: stdout). |
Definition at line 80 of file stream_engine.hpp.
◆ ~StreamEngine()
| srfm::engine::StreamEngine::~StreamEngine |
( |
| ) |
|
|
inlinenoexcept |
◆ StreamEngine() [2/3]
| srfm::engine::StreamEngine::StreamEngine |
( |
const StreamEngine & |
| ) |
|
|
delete |
◆ StreamEngine() [3/3]
| srfm::engine::StreamEngine::StreamEngine |
( |
StreamEngine && |
| ) |
|
|
delete |
◆ consumer_counters()
◆ ingester_counters()
◆ inject()
Inject a tick directly into Ring1, bypassing the TickSource.
The tick is validated by tick_is_valid() before pushing. Invalid ticks are silently dropped; the drop count is tracked in inject_counters_.
- Parameters
-
- Returns
- true if the tick was pushed successfully.
-
false if the tick was invalid or Ring1 was full.
- Note
- Thread constraint: only one external thread may call inject() concurrently (SPSC constraint on Ring1 producer side). When the TickIngester thread is also running and reading from a TickSource, do not mix inject() calls with that source producing ticks.
Definition at line 59 of file stream_engine.cpp.
◆ operator=() [1/2]
◆ operator=() [2/2]
◆ process_sync()
Process one tick synchronously (no threads, no ring buffers).
Calls SignalProcessor::process_one() directly. Useful for deterministic unit tests that do not need the full thread machinery.
- Parameters
-
| tick | Validated tick to process. |
| bar_index | Sequence number assigned to this tick. |
Definition at line 92 of file stream_engine.cpp.
◆ processor_counters()
◆ reset_processor_state()
| void srfm::engine::StreamEngine::reset_processor_state |
( |
| ) |
|
|
noexcept |
Reset all signal-processor state (for test reuse).
Must only be called when the engine is stopped.
Definition at line 99 of file stream_engine.cpp.
◆ ring1_size_approx()
| std::size_t srfm::engine::StreamEngine::ring1_size_approx |
( |
| ) |
const |
|
inlinenoexcept |
◆ ring2_size_approx()
| std::size_t srfm::engine::StreamEngine::ring2_size_approx |
( |
| ) |
const |
|
inlinenoexcept |
◆ running()
| bool srfm::engine::StreamEngine::running |
( |
| ) |
const |
|
inlinenoexcept |
◆ start()
| void srfm::engine::StreamEngine::start |
( |
| ) |
|
|
noexcept |
Start all three pipeline threads.
Idempotent. Order: ingester → processor → consumer.
Definition at line 34 of file stream_engine.cpp.
◆ stop()
| void srfm::engine::StreamEngine::stop |
( |
| ) |
|
|
noexcept |
Stop all three pipeline threads.
Idempotent. Order: ingester → processor → consumer (drain order). Blocks until all threads have joined.
Definition at line 46 of file stream_engine.cpp.
◆ RING1_SIZE
| constexpr std::size_t srfm::engine::StreamEngine::RING1_SIZE = 65536 |
|
staticconstexpr |
◆ RING2_SIZE
| constexpr std::size_t srfm::engine::StreamEngine::RING2_SIZE = 65536 |
|
staticconstexpr |
The documentation for this class was generated from the following files: