45#include "../stream/signal_consumer.hpp"
46#include "../stream/signal_processor.hpp"
47#include "../stream/spsc_ring.hpp"
48#include "../stream/stream_signal.hpp"
49#include "../stream/tick.hpp"
50#include "../stream/tick_ingester.hpp"
51#include "../stream/tick_source.hpp"
81 FILE* sink = stdout) noexcept
82 : ingester_{source, ring1_}
83 , processor_{ring1_, ring2_}
84 , consumer_{ring2_, sink}
101 void start() noexcept;
109 void stop() noexcept;
113 return running_.load(std::memory_order_acquire);
164 std::int64_t bar_index)
noexcept;
177 stream::TickIngester ingester_;
178 stream::SignalProcessor processor_;
179 stream::SignalConsumer consumer_;
181 std::atomic<
bool> running_{
false};
184 std::uint64_t inject_dropped_invalid_{0};
185 std::uint64_t inject_dropped_ring_full_{0};
186 std::uint64_t inject_pushed_{0};
Top-level streaming pipeline owner. Non-copyable, non-movable.
stream::SPSCRing< stream::OHLCVTick, RING1_SIZE > Ring1
stream::SignalConsumerCounters consumer_counters() const noexcept
StreamEngine(const StreamEngine &)=delete
StreamEngine(StreamEngine &&)=delete
static constexpr std::size_t RING2_SIZE
Signal ring capacity.
bool inject(stream::OHLCVTick tick) noexcept
Inject a tick directly into Ring1, bypassing the TickSource.
bool running() const noexcept
Whether the engine is currently running.
stream::TickIngesterCounters ingester_counters() const noexcept
StreamEngine & operator=(StreamEngine &&)=delete
StreamEngine(stream::TickSource &source, FILE *sink=stdout) noexcept
Construct the engine with a tick source and optional output sink.
stream::SPSCRing< stream::StreamRelativisticSignal, RING2_SIZE > Ring2
stream::SignalProcessorCounters processor_counters() const noexcept
void stop() noexcept
Stop all three pipeline threads.
std::size_t ring2_size_approx() const noexcept
Approximate number of signals waiting in Ring2.
void start() noexcept
Start all three pipeline threads.
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).
std::size_t ring1_size_approx() const noexcept
Approximate number of ticks waiting in Ring1.
static constexpr std::size_t RING1_SIZE
Tick ring capacity.
StreamEngine & operator=(const StreamEngine &)=delete
std::size_t size_approx() const noexcept
Approximate number of elements currently in the ring.
Abstract source of raw OHLCVTick values.
Single OHLCV bar tick from the market data feed.
Fully-characterised relativistic signal for one processed tick.
Diagnostic counters for the ingestion thread.