Special Relativity in Financial Modeling 1.0.0
Lorentz transforms, spacetime classification, and geodesic price paths for quantitative finance
Loading...
Searching...
No Matches
Public Types | Public Member Functions | Static Public Attributes | List of all members
srfm::engine::StreamEngine Class Reference

Top-level streaming pipeline owner. Non-copyable, non-movable. More...

#include <stream_engine.hpp>

Public Types

using Ring1 = stream::SPSCRing< stream::OHLCVTick, RING1_SIZE >
 
using Ring2 = stream::SPSCRing< stream::StreamRelativisticSignal, RING2_SIZE >
 

Public Member Functions

 StreamEngine (stream::TickSource &source, FILE *sink=stdout) noexcept
 Construct the engine with a tick source and optional output sink.
 
 ~StreamEngine () noexcept
 
 StreamEngine (const StreamEngine &)=delete
 
StreamEngine & operator= (const StreamEngine &)=delete
 
 StreamEngine (StreamEngine &&)=delete
 
StreamEngine & operator= (StreamEngine &&)=delete
 
void start () noexcept
 Start all three pipeline threads.
 
void stop () noexcept
 Stop all three pipeline threads.
 
bool running () const noexcept
 Whether the engine is currently running.
 
bool inject (stream::OHLCVTick tick) noexcept
 Inject a tick directly into Ring1, bypassing the TickSource.
 
stream::TickIngesterCounters ingester_counters () const noexcept
 
stream::SignalProcessorCounters processor_counters () const noexcept
 
stream::SignalConsumerCounters consumer_counters () const noexcept
 
std::size_t ring1_size_approx () const noexcept
 Approximate number of ticks waiting in Ring1.
 
std::size_t ring2_size_approx () const noexcept
 Approximate number of signals waiting in Ring2.
 
stream::StreamRelativisticSignal process_sync (const stream::OHLCVTick &tick, std::int64_t bar_index) noexcept
 Process one tick synchronously (no threads, no ring buffers).
 
void reset_processor_state () noexcept
 Reset all signal-processor state (for test reuse).
 

Static Public Attributes

static constexpr std::size_t RING1_SIZE = 65536
 Tick ring capacity.
 
static constexpr std::size_t RING2_SIZE = 65536
 Signal ring capacity.
 

Detailed Description

Top-level streaming pipeline owner. Non-copyable, non-movable.

Definition at line 64 of file stream_engine.hpp.

Member Typedef Documentation

◆ Ring1

Definition at line 69 of file stream_engine.hpp.

◆ Ring2

Definition at line 70 of file stream_engine.hpp.

Constructor & Destructor Documentation

◆ 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
sourceSource of OHLCVTick data (non-owning ref). Callers may pass a QueueTickSource for testing or a PipeTickSource for live operation.
sinkOutput sink for JSON lines (default: stdout).

Definition at line 80 of file stream_engine.hpp.

◆ ~StreamEngine()

srfm::engine::StreamEngine::~StreamEngine ( )
inlinenoexcept

Definition at line 87 of file stream_engine.hpp.

◆ StreamEngine() [2/3]

srfm::engine::StreamEngine::StreamEngine ( const StreamEngine &  )
delete

◆ StreamEngine() [3/3]

srfm::engine::StreamEngine::StreamEngine ( StreamEngine &&  )
delete

Member Function Documentation

◆ consumer_counters()

stream::SignalConsumerCounters srfm::engine::StreamEngine::consumer_counters ( ) const
noexcept

Definition at line 85 of file stream_engine.cpp.

◆ ingester_counters()

stream::TickIngesterCounters srfm::engine::StreamEngine::ingester_counters ( ) const
noexcept

Definition at line 75 of file stream_engine.cpp.

◆ inject()

bool srfm::engine::StreamEngine::inject ( stream::OHLCVTick  tick)
noexcept

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
tickTick to inject.
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]

StreamEngine & srfm::engine::StreamEngine::operator= ( const StreamEngine &  )
delete

◆ operator=() [2/2]

StreamEngine & srfm::engine::StreamEngine::operator= ( StreamEngine &&  )
delete

◆ process_sync()

stream::StreamRelativisticSignal srfm::engine::StreamEngine::process_sync ( const stream::OHLCVTick &  tick,
std::int64_t  bar_index 
)
noexcept

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
tickValidated tick to process.
bar_indexSequence number assigned to this tick.

Definition at line 92 of file stream_engine.cpp.

◆ processor_counters()

stream::SignalProcessorCounters srfm::engine::StreamEngine::processor_counters ( ) const
noexcept

Definition at line 80 of file stream_engine.cpp.

◆ 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

Approximate number of ticks waiting in Ring1.

Definition at line 142 of file stream_engine.hpp.

◆ ring2_size_approx()

std::size_t srfm::engine::StreamEngine::ring2_size_approx ( ) const
inlinenoexcept

Approximate number of signals waiting in Ring2.

Definition at line 147 of file stream_engine.hpp.

◆ running()

bool srfm::engine::StreamEngine::running ( ) const
inlinenoexcept

Whether the engine is currently running.

Definition at line 112 of file stream_engine.hpp.

◆ 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.

Member Data Documentation

◆ RING1_SIZE

constexpr std::size_t srfm::engine::StreamEngine::RING1_SIZE = 65536
staticconstexpr

Tick ring capacity.

Definition at line 66 of file stream_engine.hpp.

◆ RING2_SIZE

constexpr std::size_t srfm::engine::StreamEngine::RING2_SIZE = 65536
staticconstexpr

Signal ring capacity.

Definition at line 67 of file stream_engine.hpp.


The documentation for this class was generated from the following files: