28#include "../../include/srfm/stream/tick_ingester.hpp"
38 if (running_.load(std::memory_order_acquire))
return;
40 stop_requested_.store(
false, std::memory_order_release);
41 running_.store(
true, std::memory_order_release);
43 thread_ = std::thread([
this]()
noexcept { run_loop(); });
50 stop_requested_.store(
true, std::memory_order_release);
53 if (thread_.joinable()) {
57 running_.store(
false, std::memory_order_release);
62void TickIngester::run_loop() noexcept {
63 while (!stop_requested_.load(std::memory_order_acquire)) {
66 auto maybe_tick = source_.
read();
68 if (!maybe_tick.has_value()) {
74 std::this_thread::yield();
87 if (!ring_.
push(std::move(*maybe_tick))) {
95 running_.store(
false, std::memory_order_release);
bool push(T &&item) noexcept
Try to enqueue one element.
void stop() noexcept
Signal the ingestion thread to stop and wait for it to finish.
void start() noexcept
Launch the ingestion thread.
virtual bool is_open() const noexcept=0
Whether the source is still usable.
virtual std::optional< OHLCVTick > read() noexcept=0
Read the next OHLCVTick from the source.
bool tick_is_valid(const OHLCVTick &t) noexcept
Validate all fields of an OHLCVTick.
std::uint64_t ticks_dropped_ring_full
Dropped: ring was full.
std::uint64_t ticks_received
Raw ticks read from source.
std::uint64_t ticks_dropped_invalid
Dropped: failed tick_is_valid().
std::uint64_t ticks_pushed
Valid ticks pushed to ring.