37# define WIN32_LEAN_AND_MEAN
79 [[nodiscard]] virtual
bool is_open() const noexcept = 0;
86 virtual
void close() noexcept = 0;
106 QueueTickSource()
noexcept =
default;
114 queue_.push_back(tick);
120 template<
typename InputIt>
121 void push_range(InputIt first, InputIt last)
noexcept {
122 for (
auto it = first; it != last; ++it) {
123 queue_.push_back(*it);
127 [[nodiscard]] std::optional<OHLCVTick> read()
noexcept override {
128 if (queue_.empty())
return std::nullopt;
134 [[nodiscard]]
bool is_open()
const noexcept override {
return open_; }
136 void close()
noexcept override { open_ =
false; }
139 [[nodiscard]] std::size_t pending()
const noexcept {
return queue_.size(); }
142 std::deque<OHLCVTick> queue_;
172 unsigned int timeout_ms = 5000) noexcept {
173 open_pipe(name, timeout_ms);
178 [[nodiscard]] std::optional<OHLCVTick>
read() noexcept
override {
179 if (!is_open())
return std::nullopt;
182 const bool ok = read_exact(
reinterpret_cast<char*
>(&tick),
191 [[nodiscard]]
bool is_open() const noexcept
override {
193 return handle_ != INVALID_HANDLE_VALUE;
201 if (handle_ != INVALID_HANDLE_VALUE) {
202 CloseHandle(handle_);
203 handle_ = INVALID_HANDLE_VALUE;
214 void open_pipe(
const std::string& name,
unsigned int timeout_ms)
noexcept {
216 const std::string full =
"\\\\.\\pipe\\" + name;
217 const DWORD deadline = GetTickCount() + timeout_ms;
219 handle_ = CreateFileA(full.c_str(), GENERIC_READ, 0,
nullptr,
220 OPEN_EXISTING, 0,
nullptr);
221 if (handle_ != INVALID_HANDLE_VALUE)
break;
222 if (GetLastError() != ERROR_PIPE_BUSY)
break;
223 if (GetTickCount() >= deadline)
break;
224 WaitNamedPipeA(full.c_str(), 100);
227 fd_ = ::open(name.c_str(), O_RDONLY | O_NONBLOCK);
232 bool read_exact(
char* buf, std::size_t n)
noexcept {
235 while (total <
static_cast<DWORD
>(n)) {
237 if (!ReadFile(handle_, buf + total,
238 static_cast<DWORD
>(n) - total, &got,
nullptr)) {
245 std::size_t total = 0;
247 ssize_t got = ::read(fd_, buf + total, n - total);
248 if (got <= 0)
return false;
249 total +=
static_cast<std::size_t
>(got);
256 HANDLE handle_{INVALID_HANDLE_VALUE};
294template<
typename Source>
298 std::chrono::milliseconds initial_delay = std::chrono::milliseconds{100},
299 std::chrono::milliseconds max_delay = std::chrono::milliseconds{30'000})
300 : source_(std::move(source))
301 , initial_delay_(initial_delay)
302 , max_delay_(max_delay)
303 , current_delay_(initial_delay)
308 std::optional<OHLCVTick>
read() {
309 auto result = source_.read();
311 current_delay_ = initial_delay_;
315 std::this_thread::sleep_for(current_delay_);
316 current_delay_ = std::min(current_delay_ * 2, max_delay_);
324 Source&
source() noexcept {
return source_; }
325 const Source&
source() const noexcept {
return source_; }
329 std::chrono::milliseconds initial_delay_;
330 std::chrono::milliseconds max_delay_;
331 std::chrono::milliseconds current_delay_;
Tick source that reads raw OHLCVTick bytes from a named pipe / FIFO.
bool is_open() const noexcept override
Whether the source is still usable.
PipeTickSource(const std::string &name, unsigned int timeout_ms=5000) noexcept
Open the named pipe / FIFO for reading.
void close() noexcept override
Close the source and release any OS resources.
~PipeTickSource() noexcept override
std::optional< OHLCVTick > read() noexcept override
Read the next OHLCVTick from the source.
Wraps any TickSource with automatic reconnection on read failure.
void reset_backoff() noexcept
Reset the backoff state without reconnecting the underlying source.
Source & source() noexcept
Access the underlying source (e.g. to re-open a pipe).
ReconnectingTickSource(Source source, std::chrono::milliseconds initial_delay=std::chrono::milliseconds{100}, std::chrono::milliseconds max_delay=std::chrono::milliseconds{30 '000})
std::optional< OHLCVTick > read()
Read the next tick. Returns std::nullopt on transient failure (after sleeping current_delay_)....
const Source & source() const noexcept
Abstract source of raw OHLCVTick values.
virtual bool is_open() const noexcept=0
Whether the source is still usable.
virtual ~TickSource() noexcept=default
virtual void close() noexcept=0
Close the source and release any OS resources.
virtual std::optional< OHLCVTick > read() noexcept=0
Read the next OHLCVTick from the source.
Single OHLCV bar tick from the market data feed.
OHLCVTick — atomic market data unit for the lock-free streaming pipeline.