diff --git a/examples/market_data_continuity_example.cpp b/examples/market_data_continuity_example.cpp index e0917e1..981f9ba 100644 --- a/examples/market_data_continuity_example.cpp +++ b/examples/market_data_continuity_example.cpp @@ -139,6 +139,11 @@ int main() { request.continuity.mode = md::MarketDataContinuityMode::PREFILL_AND_RECOVER; request.continuity.prefill_bars = 2; request.continuity.max_backfill_bars = 10; + request.continuity.max_buffered_batches = 32; + request.continuity.max_buffered_items = 256; + request.continuity.retry.max_attempts = 3; + request.continuity.retry.initial_backoff_ms = 100; + request.continuity.retry.max_backoff_ms = 1000; auto route = router.subscribe_bars(provider, chart, request); if (!route.active()) { diff --git a/guides/api-and-header-contracts.md b/guides/api-and-header-contracts.md index 54c22f1..bf274a8 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -257,13 +257,43 @@ Contract rules: optional timestamp-gap recovery. Router buffers live batches until the corresponding history operation completes and reports route-scoped progress through `IMarketDataSubscriber::on_market_data_continuity()`. +- `prefill_bars` requests inclusive timeframe slots ending at the start of the + current timeframe bucket. `max_buffered_batches` and `max_buffered_items` + bound live batches held during history; zero disables each limit. On buffer + overflow Router reports continuity `FAILED`, releases the held live data, + disables continuity for that route, and resumes with live delivery in + `DEGRADED`; continuity is not automatically re-enabled for that route. +- `MarketDataContinuityOptions::retry` bounds failed history requests with a + capped exponential backoff. Retries are scheduled by the periodic Router + `process()` call, so no continuity timer thread is created. Buffered live + batches remain withheld while retrying. +- Each bounded gap-history response must cover every expected timeframe slot in + the actual requested range. Router clips provider payloads to that range + before validating it; a non-empty partial response is not delivered and does + not advance the route's continuity watermark. Router reports `FAILED`, + releases buffered live data, and ends in `DEGRADED` while live delivery + continues. +- `max_backfill_bars` applies to each individual gap request. Continuity + telemetry reports the actual bounded `from_time_ms`, `to_time_ms`, and + `requested_items`, rather than the original unbounded gap. +- `last_bar_time_ms` tracks the latest delivered bar, not continuity trust. + After any terminal unusable history operation, Router retains the earliest + unresolved boundary. A later successful history range cannot clear that + trust loss unless it covers the unresolved boundary; `LIVE` is emitted only + when no unresolved boundary remains known. +- Successful empty history is terminal: an empty prefill proceeds to live + delivery, while an empty gap backfill reports failed recovery and releases + buffered live data without retrying the same range. +- Router preserves provider bar revisions and leaves timestamp upsert or + deduplication to the consumer. A chart or storage component should upsert by + stream and `time_ms` when it needs one current candle value. - Continuity updates carry the concrete provider subscription handle and must not be confused with stream-level `MarketDataStatusUpdate`. History failure does not terminate the live route: Router reports `FAILED`, releases buffered - live batches, and returns to `LIVE`. -- Generic history continuity is currently defined for bars only. Tick history, - provider-aware retries, and a universal overlap-deduplication policy remain - separate contracts. + live batches, and ends in `DEGRADED` after the retry budget is exhausted. +- Generic history continuity is currently defined for bars only. Tick history + remains a separate provider contract. Router does not apply a universal + timestamp deduplication policy; consumers decide how to upsert revisions. `MarketDataRouter` is the subscription-scoped alternative to `MarketDataHub`: diff --git a/guides/market-data-router.md b/guides/market-data-router.md index 40a3a6c..199bfd6 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -319,13 +319,28 @@ The options have these meanings: - `LIVE_ONLY` leaves live delivery unchanged. `PREFILL` requests `prefill_bars` before the first live delivery. `PREFILL_AND_RECOVER` does both and also repairs timestamp gaps. -- `prefill_bars` is the count-based initial history depth. Router builds an - inclusive timeframe range for that many slots; a provider may still return - fewer bars when part of the range has no data. A `PREFILL` request must - specify a positive count. `PREFILL_AND_RECOVER` may use zero when the - application wants recovery only. +- `prefill_bars` is the count-based initial history depth. Router aligns the + end of the inclusive range to the start of the current timeframe bucket and + requests `prefill_bars` slots, so the range is `boundary - (N - 1) * + timeframe` through `boundary`. A provider may still return fewer bars when + part of the range has no data. A `PREFILL` request must specify a positive + count. `PREFILL_AND_RECOVER` may use zero when the application wants + recovery only. - `max_backfill_bars` bounds one gap request. Zero means that the provider request is not count-bounded by Router. +- `max_buffered_batches` and `max_buffered_items` bound live data retained + while a history request is in flight. Zero disables the corresponding limit. + When either limit is exceeded, Router reports continuity `FAILED`, releases + the queued and current live batches, disables continuity for that route, and + reports `DEGRADED` while live data continues. This is an explicit + loss-of-recovery fallback: continuity is not automatically re-enabled for + that route. +- `retry.max_attempts` is the total number of attempts, including the initial + request. `retry.initial_backoff_ms` starts a capped exponential backoff, and + `retry.max_backoff_ms` limits it. The default is one attempt, preserving the + existing live-delivery fallback with a final `DEGRADED` status. Retries are + advanced by periodic + `MarketDataRouter::process()` calls on the owner loop. For initial prefill, the delivery order is: @@ -367,19 +382,35 @@ stream-level `MarketDataStatusUpdate` events such as `READY` or `DISCONNECTED`. `max_backfill_bars` limits each individual provider request, not the complete missing interval. When a provider returns a partial but non-empty backfill, -Router keeps draining the queued live batches and can issue another bounded -request for the next remaining gap. An empty successful response for a detected -gap is treated as a failed recovery, so Router releases the queued live data -once instead of retrying the same range forever. - -If a history request fails or is rejected, Router emits `FAILED`, releases any -buffered live batches, and then emits `LIVE`. The route stays usable and live -delivery continues, but the missing historical range is not reconstructed. -Applications that require a complete time series should record the failure and -apply their own retry policy. A provider may return overlapping snapshots for -an in-progress bar; this first continuity layer does not impose a universal -payload deduplication policy, so consumers should correlate bars by stream and -`time_ms` according to their finalized/incomplete-bar policy. +Router treats that response as an unsuccessful recovery: it does not deliver +the partial batch or advance the continuity watermark, reports `FAILED`, +releases the queued live data, and finishes in `DEGRADED`. A bounded gap +response is accepted only when, after clipping to the actual requested range, +it contains every expected timeframe slot. An empty successful response for a +detected gap is treated the same way, so Router releases the queued live data +once instead of retrying the same range forever. A successful empty prefill is +terminal as well: it carries no historical batch and proceeds to the buffered +live stream without retry. + +`max_backfill_bars` is applied before the provider request is sent. Therefore +continuity updates report the actual bounded `from_time_ms`, `to_time_ms`, and +`requested_items` for the current chunk, not the original full gap. + +`last_bar_time_ms` is the latest delivered bar, not a continuity trust +watermark. After any terminal unusable history operation, Router retains the +earliest unresolved boundary. A later successful history range cannot clear +that trust loss unless it covers the unresolved boundary; Router emits `LIVE` +only when no unresolved boundary remains known. + +If a history request fails or is rejected and attempts remain, Router emits +`RETRYING`, keeps buffered live batches, and waits for a later owner-loop +`process()` call. After the retry budget is exhausted, Router emits `FAILED`, +releases any buffered live batches, and then emits `DEGRADED`. The route stays +usable and live delivery continues, but an unrecovered historical range is +reported to the consumer. Router preserves provider bar revisions and does not +apply a generic timestamp deduplication policy; consumers such as charts or +storage should upsert by stream and `time_ms` when they need one current value +per candle while still accepting later finalized revisions. `MarketDataContinuityService` is the lower-level helper for applications that want to request history directly. It converts a `BarHistoryResult` into a diff --git a/guides/market-data-router.ru.md b/guides/market-data-router.ru.md index 13045a1..54fba91 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -319,13 +319,29 @@ auto route = router.subscribe_bars("intrade", chart, request); - `LIVE_ONLY` оставляет live-доставку без изменений. `PREFILL` запрашивает `prefill_bars` перед первой live-доставкой. `PREFILL_AND_RECOVER` делает оба действия и дополнительно восстанавливает разрывы по timestamp. -- `prefill_bars` задаёт глубину начальной истории в барах. Router строит - inclusive-диапазон на это количество timeframe slots; провайдер всё равно - может вернуть меньше баров, если часть диапазона не содержит данных. Для - `PREFILL` требуется положительное значение. `PREFILL_AND_RECOVER` может - использовать ноль, если приложению нужно только восстановление разрывов. +- `prefill_bars` задаёт глубину начальной истории в барах. Router выравнивает + конец inclusive-диапазона к началу текущего timeframe bucket и запрашивает + `prefill_bars` slots: от `boundary - (N - 1) * timeframe` до `boundary`. + Провайдер всё равно может вернуть меньше баров, если часть диапазона не + содержит данных. Для `PREFILL` требуется положительное значение. + `PREFILL_AND_RECOVER` может использовать ноль, если приложению нужно только + восстановление разрывов. - `max_backfill_bars` ограничивает один запрос для разрыва. Ноль означает, что Router не ограничивает provider request количеством баров. +- `max_buffered_batches` и `max_buffered_items` ограничивают live data, + удерживаемые во время history request. Ноль отключает соответствующее + ограничение. При превышении любого лимита Router публикует continuity + `FAILED`, выпускает накопленные и текущий live batch, отключает continuity + для этого route и публикует `DEGRADED`, продолжая live-доставку. Это явный + fallback с потерей обещания historical ordering: continuity автоматически + для этого route не включается. +- `retry.max_attempts` задаёт общее количество попыток вместе с первой. + `retry.initial_backoff_ms` включает ограниченный exponential backoff, а + `retry.max_backoff_ms` задаёт его предел. По умолчанию выполняется одна + попытка, поэтому live-доставка продолжается, а итоговый статус становится + `DEGRADED`. Повторы + запускаются периодическими вызовами `MarketDataRouter::process()` в owner + loop. Для initial prefill порядок доставки такой: @@ -368,20 +384,36 @@ stream-level `MarketDataStatusUpdate`, например `READY` или `DISCONNE `max_backfill_bars` ограничивает каждый отдельный provider request, а не весь отсутствующий интервал. Если provider вернул неполный, но непустой backfill, -Router продолжит разбирать очередь live batches и может выполнить следующий -ограниченный запрос для оставшегося gap. Успешный пустой ответ для найденного -gap считается failed recovery: Router один раз выпускает накопленные live data и -не зацикливает запрос того же диапазона. - -Если history request завершился ошибкой или provider его отклонил, Router -публикует `FAILED`, выпускает накопленные live batches, затем публикует `LIVE`. -Route остаётся пригодным для работы и live-доставка продолжается, но -отсутствующий исторический диапазон не восстанавливается. Приложение, которому -нужен полный временной ряд, должно записать ошибку и применить собственную -политику повторной попытки. Provider может возвращать пересекающиеся snapshots -для незавершённого бара; этот первый слой continuity не вводит универсальную -политику deduplication, поэтому consumer должен сам сопоставлять бары по -stream и `time_ms` с учётом политики `FINALIZED`/незавершённых баров. +Router считает такое восстановление неуспешным: неполный batch не доставляется +и не продвигает continuity watermark, Router публикует `FAILED`, выпускает +накопленные live data и завершает операцию в `DEGRADED`. Bounded gap response +принимается только если после обрезки до фактического запрошенного диапазона +он содержит каждый ожидаемый timeframe slot. Успешный пустой ответ для +найденного gap обрабатывается так же: Router один раз выпускает накопленные +live data и не зацикливает запрос того же диапазона. Успешный пустой prefill +тоже является terminal result: исторический batch не доставляется, а Router +переходит к накопленному live-потоку без retry. + +`max_backfill_bars` применяется до отправки provider request. Поэтому +continuity updates содержат фактические bounded `from_time_ms`, +`to_time_ms` и `requested_items` текущего chunk, а не исходный полный gap. + +`last_bar_time_ms` является последним доставленным bar, а не watermark доверия +continuity. После любой terminal unusable history operation Router сохраняет +самую раннюю неподтверждённую границу. Более поздний успешный history range не +может убрать эту потерю доверия, пока не покрывает unresolved boundary; Router +публикует `LIVE` только когда известной неподтверждённой границы не осталось. + +Если history request завершился ошибкой или provider его отклонил, но попытки +ещё остались, Router публикует `RETRYING`, сохраняет накопленные live batches и +ждёт следующего вызова `process()` в owner loop. После исчерпания retry budget +Router публикует `FAILED`, выпускает накопленные live batches, затем публикует +`DEGRADED`. Route остаётся пригодным для работы, live-доставка продолжается, а +невосстановленный исторический диапазон явно сообщается consumer. Router +сохраняет revisions баров от provider и не применяет универсальную +timestamp-дедупликацию. Если графику или storage нужно одно текущее значение +на candle, consumer должен делать upsert по stream и `time_ms`, сохраняя +возможность принять более позднюю finalized revision. `MarketDataContinuityService` остаётся низкоуровневым helper для приложений, которые хотят запрашивать историю напрямую. Он превращает `BarHistoryResult` в diff --git a/include/optionx_cpp/market_data/MarketDataContinuity.hpp b/include/optionx_cpp/market_data/MarketDataContinuity.hpp index 49c2f81..33f2c13 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuity.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuity.hpp @@ -18,8 +18,10 @@ namespace optionx::market_data { PREFILLING, ///< Historical initialization is being requested. GAP_DETECTED, ///< A timestamp gap was found in the live stream. BACKFILLING, ///< Historical bars are being loaded for a gap. - LIVE, ///< Live delivery is current, with no pending history work. - FAILED ///< History work failed; live delivery continues without it. + RETRYING, ///< A failed history request will be attempted again. + LIVE, ///< No known unresolved history range remains for the route. + FAILED, ///< A specific history operation failed; the route may continue. + DEGRADED ///< Live delivery continues while continuity remains unverified; the status is sticky until the unresolved range is verified. }; /// \brief Converts a continuity status to stable text. @@ -31,10 +33,14 @@ namespace optionx::market_data { return "GAP_DETECTED"; case MarketDataContinuityStatus::BACKFILLING: return "BACKFILLING"; + case MarketDataContinuityStatus::RETRYING: + return "RETRYING"; case MarketDataContinuityStatus::LIVE: return "LIVE"; case MarketDataContinuityStatus::FAILED: return "FAILED"; + case MarketDataContinuityStatus::DEGRADED: + return "DEGRADED"; case MarketDataContinuityStatus::UNKNOWN: default: return "UNKNOWN"; diff --git a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp index 1605971..2baa132 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp @@ -6,6 +6,7 @@ /// \brief Defines history prefill and gap-recovery options for bar routes. #include +#include namespace optionx::market_data { @@ -17,12 +18,30 @@ namespace optionx::market_data { PREFILL_AND_RECOVER ///< Prefill and repair timestamp gaps in live bars. }; + /// \struct MarketDataContinuityRetryPolicy + /// \brief Configures bounded history retry attempts and exponential backoff. + struct MarketDataContinuityRetryPolicy { + std::size_t max_attempts = 1; ///< Total attempts, including the first request. + std::uint64_t initial_backoff_ms = 0; ///< Delay before the second attempt. + std::uint64_t max_backoff_ms = 30000; ///< Backoff cap; zero means no cap. + + /// \brief Returns true when retry settings can be applied safely. + [[nodiscard]] bool valid() const noexcept { + return max_attempts > 0 && + (max_backoff_ms == 0 || max_backoff_ms >= initial_backoff_ms); + } + + }; + /// \struct MarketDataContinuityOptions - /// \brief Configures history prefill and timestamp-gap recovery for a bar route. + /// \brief Configures history prefill, recovery, retries, and buffering. struct MarketDataContinuityOptions { MarketDataContinuityMode mode = MarketDataContinuityMode::LIVE_ONLY; std::size_t prefill_bars = 0; ///< Number of historical bars requested before live delivery. std::size_t max_backfill_bars = 1000; ///< Maximum bars per detected gap; zero is unbounded. + MarketDataContinuityRetryPolicy retry; + std::size_t max_buffered_batches = 1024; ///< Maximum live batches held during history; zero is unbounded. + std::size_t max_buffered_items = 100000; ///< Maximum live items held during history; zero is unbounded. /// \brief Returns true when the option combination is usable. [[nodiscard]] bool valid() const noexcept { @@ -30,9 +49,9 @@ namespace optionx::market_data { case MarketDataContinuityMode::LIVE_ONLY: return prefill_bars == 0; case MarketDataContinuityMode::PREFILL: - return prefill_bars > 0; + return prefill_bars > 0 && retry.valid(); case MarketDataContinuityMode::PREFILL_AND_RECOVER: - return true; + return retry.valid(); default: return false; } diff --git a/include/optionx_cpp/market_data/MarketDataContinuityService.hpp b/include/optionx_cpp/market_data/MarketDataContinuityService.hpp index 1e1f566..c9e163e 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuityService.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuityService.hpp @@ -19,9 +19,10 @@ namespace optionx::market_data { /// \class MarketDataContinuityService /// \brief Bridges historical bar requests into the live market-data batch pipeline. /// - /// The service intentionally stays thin: providers still own transport, - /// retry, and stream lifecycle. This helper only tags recovered payloads - /// as historical/backfill data and packages them into BarDataBatch objects. + /// The service intentionally stays thin: providers still own transport and + /// stream lifecycle, while Router owns route-level retry policy. This helper + /// only tags recovered payloads as historical/backfill data and packages + /// them into BarDataBatch objects. class MarketDataContinuityService { public: /// \brief Callback that receives a failed history request. @@ -61,9 +62,13 @@ namespace optionx::market_data { const auto depth_ms = safe_multiply( timeframe_ms, bars > 0 ? bars - 1 : 0); - const auto from_ms = now_ms > depth_ms ? now_ms - depth_ms : 1U; + const auto aligned_now_ms = timeframe_ms > 0 + ? now_ms - (now_ms % timeframe_ms) + : now_ms; + const auto boundary_ms = aligned_now_ms > 0 ? aligned_now_ms : now_ms; + const auto from_ms = boundary_ms > depth_ms ? boundary_ms - depth_ms : 1U; - return make_history_request(request, from_ms, now_ms); + return make_history_request(request, from_ms, boundary_ms); } /// \brief Builds a bounded request for a missing bar range. diff --git a/include/optionx_cpp/market_data/MarketDataRouter.hpp b/include/optionx_cpp/market_data/MarketDataRouter.hpp index d7048f8..5291520 100644 --- a/include/optionx_cpp/market_data/MarketDataRouter.hpp +++ b/include/optionx_cpp/market_data/MarketDataRouter.hpp @@ -218,8 +218,8 @@ namespace optionx::market_data { std::size_t retry_failed_unsubscribes(); /// \brief Advances deferred Router lifecycle work on the owner loop. - /// \details Processes provider completions retained during shutdown and - /// starts or completes physical subscription cleanup. This method + /// \details Processes due continuity retries, provider completions retained + /// during shutdown, and physical subscription cleanup. This method /// does not poll providers or transport data. void process() override; diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index 5d2a00f..f3a319d 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -5,6 +5,9 @@ /// \file MarketDataRouter.ipp /// \brief Implements subscription-scoped market-data routing utilities. +#include +#include + namespace optionx::market_data { namespace detail { @@ -43,6 +46,23 @@ namespace optionx::market_data { CLEANUP_FAILED }; + enum class ContinuityRequestKind { + PREFILL, + GAP_BACKFILL + }; + + struct PendingContinuityRequest { + RoutedSubscriptionId router_id; + MarketDataSubscriptionHandle subscription; + BarHistoryRequest request; + ContinuityRequestKind kind = ContinuityRequestKind::PREFILL; + bool announce_gap = false; + std::uint64_t from_time_ms = 0; + std::uint64_t to_time_ms = 0; + std::size_t requested_items = 0; + std::size_t attempt = 1; + }; + struct Entry { RoutedSubscriptionId router_id; ProviderInstanceId provider_id = kInvalidProviderInstanceId; @@ -55,7 +75,12 @@ namespace optionx::market_data { bool continuity_request_in_flight = false; bool continuity_flushing = false; std::deque continuity_buffer; + std::size_t continuity_buffered_items = 0; std::uint64_t last_bar_time_ms = 0; + std::uint64_t verified_through_time_ms = 0; + std::uint64_t unverified_from_time_ms = 0; + std::optional continuity_retry_request; + std::chrono::steady_clock::time_point continuity_retry_at; MarketDataSubscriptionHandle retained_cleanup_subscription; MarketDataSubscriptionResult unsubscribe_completion; bool subscribe_completion_posted = false; @@ -66,16 +91,6 @@ namespace optionx::market_data { subscription_callback_t release_callback; }; - struct PendingContinuityRequest { - RoutedSubscriptionId router_id; - MarketDataSubscriptionHandle subscription; - BarHistoryRequest request; - MarketDataContinuityStatus status = MarketDataContinuityStatus::UNKNOWN; - std::uint64_t from_time_ms = 0; - std::uint64_t to_time_ms = 0; - std::size_t requested_items = 0; - }; - struct ContinuityOperation { bool completed = false; }; @@ -273,18 +288,21 @@ namespace optionx::market_data { RoutedSubscriptionId router_id, MarketDataSubscriptionHandle subscription, BarHistoryRequest request, - MarketDataContinuityStatus status, + ContinuityRequestKind kind, + bool announce_gap, std::uint64_t from_time_ms, std::uint64_t to_time_ms, - std::size_t requested_items); + std::size_t requested_items, + std::size_t attempt); void complete_continuity( RoutedSubscriptionId router_id, MarketDataSubscriptionHandle subscription, BarHistoryRequest request, - MarketDataContinuityStatus operation_status, + ContinuityRequestKind kind, std::uint64_t from_time_ms, std::uint64_t to_time_ms, std::size_t requested_items, + std::size_t attempt, BarHistoryResult result); void notify_continuity( RoutedSubscriptionId router_id, @@ -304,8 +322,44 @@ namespace optionx::market_data { std::shared_ptr, BarDataBatch>>& deliveries, std::vector& continuity_requests, + std::vector, + MarketDataContinuityUpdate>>& continuity_deliveries, bool process_buffered = false, bool allow_gap_recovery = true); + static bool buffer_continuity_batch_no_lock( + const std::shared_ptr& entry, + std::shared_ptr subscriber, + BarDataBatch batch, + std::vector, + BarDataBatch>>& deliveries, + std::vector, + MarketDataContinuityUpdate>>& continuity_deliveries, + bool push_front); + static bool history_covers_range( + const BarDataBatch& batch, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::uint64_t timeframe_ms) noexcept; + static void clip_history_to_range( + BarDataBatch& batch, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms); + static void mark_unverified_no_lock( + const std::shared_ptr& entry, + std::uint64_t from_time_ms) noexcept; + static void record_verified_range_no_lock( + const std::shared_ptr& entry, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms) noexcept; + static void record_bar_progress_no_lock( + const std::shared_ptr& entry, + const std::vector& bars) noexcept; + static std::chrono::steady_clock::duration continuity_retry_delay( + const MarketDataContinuityRetryPolicy& policy, + std::size_t attempt) noexcept; void fail_pending_subscribe( RoutedSubscriptionId router_id, @@ -1219,7 +1273,8 @@ namespace optionx::market_data { bar_request, now_ms, entry->continuity.prefill_bars); - pending.status = MarketDataContinuityStatus::PREFILLING; + pending.kind = ContinuityRequestKind::PREFILL; + pending.announce_gap = false; pending.from_time_ms = MarketDataContinuityService::seconds_to_milliseconds( pending.request.from_ts); @@ -1246,26 +1301,31 @@ namespace optionx::market_data { pending.router_id, std::move(pending.subscription), std::move(pending.request), - pending.status, + pending.kind, + pending.announce_gap, pending.from_time_ms, pending.to_time_ms, - pending.requested_items); + pending.requested_items, + 1); } inline void MarketDataRouterState::request_continuity_history( RoutedSubscriptionId router_id, MarketDataSubscriptionHandle subscription, BarHistoryRequest request, - MarketDataContinuityStatus status, + ContinuityRequestKind kind, + bool announce_gap, std::uint64_t from_time_ms, std::uint64_t to_time_ms, - std::size_t requested_items) { + std::size_t requested_items, + std::size_t attempt) { BaseMarketDataProvider* provider = nullptr; auto operation = std::make_shared(); { std::lock_guard lock(m_mutex); const auto entry_it = m_entries.find(router_id); - if (entry_it == m_entries.end() || + if (m_shutdown || + entry_it == m_entries.end() || entry_it->second->phase != EntryPhase::ACTIVE || !entry_it->second->continuity_request_in_flight) { return; @@ -1275,17 +1335,19 @@ namespace optionx::market_data { } if (!provider) return; - if (status == MarketDataContinuityStatus::GAP_DETECTED) { - notify_continuity( - router_id, - make_continuity_update( - subscription, - MarketDataContinuityStatus::GAP_DETECTED, - from_time_ms, - to_time_ms, - requested_items, - 0, - "A gap was detected in the live bar stream.")); + if (kind == ContinuityRequestKind::GAP_BACKFILL) { + if (announce_gap) { + notify_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::GAP_DETECTED, + from_time_ms, + to_time_ms, + requested_items, + 0, + "A gap was detected in the live bar stream.")); + } notify_continuity( router_id, make_continuity_update( @@ -1304,29 +1366,32 @@ namespace optionx::market_data { router_id, subscription, request, - status, + kind, from_time_ms, to_time_ms, - requested_items](BarHistoryResult result) mutable { + requested_items, + attempt](BarHistoryResult result) mutable { if (!state->record_continuity_completion(operation)) return; auto task = [state, router_id, subscription, request, - status, + kind, from_time_ms, to_time_ms, requested_items, + attempt, result = std::move(result)]() mutable { state->complete_continuity( router_id, subscription, request, - status, + kind, from_time_ms, to_time_ms, requested_items, + attempt, std::move(result)); }; state->dispatch_or_run(std::move(task)); @@ -1351,14 +1416,146 @@ namespace optionx::market_data { } } + inline bool MarketDataRouterState::history_covers_range( + const BarDataBatch& batch, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::uint64_t timeframe_ms) noexcept { + if (from_time_ms == 0 || + to_time_ms < from_time_ms || + timeframe_ms == 0) { + return false; + } + + auto expected_time_ms = from_time_ms; + for (const auto& bar : batch.items) { + if (bar.time_ms < expected_time_ms) continue; + if (bar.time_ms != expected_time_ms) return false; + if (expected_time_ms == to_time_ms) return true; + if (expected_time_ms > + std::numeric_limits::max() - timeframe_ms) { + return false; + } + expected_time_ms += timeframe_ms; + } + return false; + } + + inline void MarketDataRouterState::clip_history_to_range( + BarDataBatch& batch, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms) { + batch.items.erase( + std::remove_if( + batch.items.begin(), + batch.items.end(), + [from_time_ms, to_time_ms](const Bar& bar) { + return bar.time_ms < from_time_ms || + bar.time_ms > to_time_ms; + }), + batch.items.end()); + } + + inline void MarketDataRouterState::mark_unverified_no_lock( + const std::shared_ptr& entry, + std::uint64_t from_time_ms) noexcept { + if (!entry || from_time_ms == 0) return; + auto& unverified_from = entry->unverified_from_time_ms; + if (unverified_from == 0 || from_time_ms < unverified_from) { + unverified_from = from_time_ms; + } + } + + inline void MarketDataRouterState::record_verified_range_no_lock( + const std::shared_ptr& entry, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms) noexcept { + if (!entry || from_time_ms == 0 || to_time_ms < from_time_ms) { + return; + } + + if (entry->unverified_from_time_ms != 0 && + (from_time_ms > entry->unverified_from_time_ms || + to_time_ms < entry->unverified_from_time_ms)) { + return; + } + + entry->verified_through_time_ms = std::max( + entry->verified_through_time_ms, + to_time_ms); + if (entry->unverified_from_time_ms != 0) { + entry->unverified_from_time_ms = 0; + } + } + + inline void MarketDataRouterState::record_bar_progress_no_lock( + const std::shared_ptr& entry, + const std::vector& bars) noexcept { + if (!entry) return; + + for (const auto& bar : bars) { + if (bar.time_ms > entry->last_bar_time_ms) { + entry->last_bar_time_ms = bar.time_ms; + } + + const bool trusted_finalized = + bar.has_flag(MarketDataFlags::FINALIZED) || + (bar.has_flag(MarketDataFlags::HISTORICAL) && + !bar.has_flag(MarketDataFlags::INCOMPLETE)); + if (trusted_finalized && + (entry->unverified_from_time_ms == 0 || + bar.time_ms < entry->unverified_from_time_ms) && + bar.time_ms > entry->verified_through_time_ms) { + entry->verified_through_time_ms = bar.time_ms; + } + } + } + + inline std::chrono::steady_clock::duration + MarketDataRouterState::continuity_retry_delay( + const MarketDataContinuityRetryPolicy& policy, + std::size_t attempt) noexcept { + if (attempt == 0 || policy.initial_backoff_ms == 0) { + return std::chrono::steady_clock::duration::zero(); + } + + auto delay_ms = policy.initial_backoff_ms; + for (std::size_t index = 1; index < attempt; ++index) { + if (delay_ms > std::numeric_limits::max() / 2U) { + delay_ms = std::numeric_limits::max(); + break; + } + delay_ms *= 2U; + if (policy.max_backoff_ms > 0 && + delay_ms >= policy.max_backoff_ms) { + delay_ms = policy.max_backoff_ms; + break; + } + } + if (policy.max_backoff_ms > 0) { + delay_ms = std::min(delay_ms, policy.max_backoff_ms); + } + + using milliseconds = std::chrono::milliseconds; + const auto max_delay_count = std::chrono::duration_cast( + std::chrono::steady_clock::duration::max()).count(); + const auto max_delay_ms = max_delay_count > 0 + ? static_cast(max_delay_count) + : 0U; + return std::chrono::duration_cast( + milliseconds(static_cast( + std::min(delay_ms, max_delay_ms)))); + } + inline void MarketDataRouterState::complete_continuity( RoutedSubscriptionId router_id, MarketDataSubscriptionHandle subscription, BarHistoryRequest request, - MarketDataContinuityStatus operation_status, + ContinuityRequestKind kind, std::uint64_t from_time_ms, std::uint64_t to_time_ms, std::size_t requested_items, + std::size_t attempt, BarHistoryResult result) { StreamDescriptor expected_stream; { @@ -1375,6 +1572,9 @@ namespace optionx::market_data { } const bool history_success = static_cast(result); + const auto timeframe_ms = expected_stream.timeframe > 0 + ? static_cast(expected_stream.timeframe) * 1000U + : 0U; std::size_t delivered_history_items = 0; BarDataBatch history_batch; bool has_history_batch = false; @@ -1384,20 +1584,116 @@ namespace optionx::market_data { std::move(result.sequence), request, subscription, - operation_status == MarketDataContinuityStatus::GAP_DETECTED); + kind == ContinuityRequestKind::GAP_BACKFILL); history_stream_matches = batch_matches_stream( history_batch, expected_stream); if (history_stream_matches) { + if (kind == ContinuityRequestKind::GAP_BACKFILL) { + clip_history_to_range( + history_batch, + from_time_ms, + to_time_ms); + } + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (entry_it == m_entries.end() || + entry_it->second->phase != EntryPhase::ACTIVE) { + return; + } delivered_history_items = history_batch.items.size(); has_history_batch = !history_batch.items.empty(); } } + const bool gap_history_covers_range = + kind != ContinuityRequestKind::GAP_BACKFILL || + history_covers_range( + history_batch, + from_time_ms, + to_time_ms, + timeframe_ms); + if (kind == ContinuityRequestKind::GAP_BACKFILL && + !gap_history_covers_range) { + has_history_batch = false; + delivered_history_items = 0; + } const bool usable_history = history_success && history_stream_matches && - (operation_status != MarketDataContinuityStatus::GAP_DETECTED || - delivered_history_items > 0); + (kind != ContinuityRequestKind::GAP_BACKFILL || + delivered_history_items > 0) && + gap_history_covers_range; + + bool retry_scheduled = false; + if (!history_success) { + PendingContinuityRequest retry_request; + { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (entry_it == m_entries.end() || + entry_it->second->phase != EntryPhase::ACTIVE) { + return; + } + + const auto& entry = entry_it->second; + if (attempt < entry->continuity.retry.max_attempts) { + retry_request.router_id = router_id; + retry_request.subscription = subscription; + retry_request.request = request; + retry_request.kind = kind; + retry_request.announce_gap = false; + retry_request.from_time_ms = from_time_ms; + retry_request.to_time_ms = to_time_ms; + retry_request.requested_items = requested_items; + retry_request.attempt = attempt + 1; + + const auto now = std::chrono::steady_clock::now(); + const auto delay = continuity_retry_delay( + entry->continuity.retry, + attempt); + const auto remaining = + std::chrono::steady_clock::time_point::max() - now; + entry->continuity_retry_request = std::move(retry_request); + entry->continuity_retry_at = delay >= remaining + ? std::chrono::steady_clock::time_point::max() + : now + delay; + entry->continuity_flushing = false; + retry_scheduled = true; + } + } + + if (retry_scheduled) { + notify_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::RETRYING, + from_time_ms, + to_time_ms, + requested_items, + 0, + "Retrying historical market-data continuity request.")); + return; + } + } + + { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (entry_it == m_entries.end() || + entry_it->second->phase != EntryPhase::ACTIVE) { + return; + } + if (!usable_history) { + mark_unverified_no_lock(entry_it->second, from_time_ms); + } else if (kind == ContinuityRequestKind::GAP_BACKFILL) { + record_verified_range_no_lock( + entry_it->second, + from_time_ms, + to_time_ms); + } + } + if (!usable_history) { notify_continuity( router_id, @@ -1413,7 +1709,12 @@ namespace optionx::market_data { ? "Historical market-data continuity request failed." : !history_stream_matches ? "Historical bar response does not match the subscribed stream." - : "No bars were returned for the detected gap.") + : kind == ContinuityRequestKind::GAP_BACKFILL && + !gap_history_covers_range + ? "Historical bars do not cover the requested gap range." + : kind == ContinuityRequestKind::GAP_BACKFILL + ? "No bars were returned for the detected gap." + : "No historical bars were returned for prefill.") : result.error_desc)); } @@ -1428,11 +1729,9 @@ namespace optionx::market_data { !entry_it->second->release_requested; if (active) { subscriber = entry_it->second->subscriber.lock(); - for (const auto& bar : history_batch.items) { - if (bar.time_ms > entry_it->second->last_bar_time_ms) { - entry_it->second->last_bar_time_ms = bar.time_ms; - } - } + record_bar_progress_no_lock( + entry_it->second, + history_batch.items); } } if (active && subscriber) { @@ -1444,8 +1743,13 @@ namespace optionx::market_data { std::vector, BarDataBatch>> deliveries; + std::vector, + MarketDataContinuityUpdate>> continuity_deliveries; std::vector continuity_requests; bool finished = false; + bool continuity_still_enabled = false; + bool continuity_verified = false; { std::lock_guard lock(m_mutex); const auto entry_it = m_entries.find(router_id); @@ -1456,21 +1760,37 @@ namespace optionx::market_data { if (entry_it->second->continuity_buffer.empty()) { entry_it->second->continuity_ready = true; entry_it->second->continuity_flushing = false; + continuity_still_enabled = + entry_it->second->continuity.enabled(); + continuity_verified = usable_history && + entry_it->second->unverified_from_time_ms == 0; finished = true; } else { auto batch = std::move( entry_it->second->continuity_buffer.front()); entry_it->second->continuity_buffer.pop_front(); + if (entry_it->second->continuity_buffered_items >= + batch.items.size()) { + entry_it->second->continuity_buffered_items -= + batch.items.size(); + } else { + entry_it->second->continuity_buffered_items = 0; + } route_bar_to_entry_no_lock( entry_it->second, batch, deliveries, continuity_requests, + continuity_deliveries, true, usable_history); } } + for (auto& continuity_delivery : continuity_deliveries) { + continuity_delivery.first->on_market_data_continuity( + continuity_delivery.second); + } for (auto& delivery : deliveries) { delivery.first->on_bar_data(delivery.second); } @@ -1481,26 +1801,32 @@ namespace optionx::market_data { pending.router_id, std::move(pending.subscription), std::move(pending.request), - pending.status, + pending.kind, + pending.announce_gap, pending.from_time_ms, pending.to_time_ms, - pending.requested_items); + pending.requested_items, + pending.attempt); return; } if (finished) { - notify_continuity( - router_id, - make_continuity_update( - subscription, - MarketDataContinuityStatus::LIVE, - from_time_ms, - to_time_ms, - requested_items, - delivered_history_items, - usable_history - ? "Historical market-data continuity is ready." - : "Live delivery continues after history failure.")); + if (continuity_still_enabled) { + notify_continuity( + router_id, + make_continuity_update( + subscription, + continuity_verified + ? MarketDataContinuityStatus::LIVE + : MarketDataContinuityStatus::DEGRADED, + from_time_ms, + to_time_ms, + requested_items, + delivered_history_items, + continuity_verified + ? "Historical market-data continuity is ready." + : "Live delivery continues without verified continuity.")); + } return; } } @@ -1890,6 +2216,75 @@ namespace optionx::market_data { if (subscriber) subscriber->on_market_data_continuity(update); } + inline bool MarketDataRouterState::buffer_continuity_batch_no_lock( + const std::shared_ptr& entry, + std::shared_ptr subscriber, + BarDataBatch batch, + std::vector, + BarDataBatch>>& deliveries, + std::vector, + MarketDataContinuityUpdate>>& continuity_deliveries, + bool push_front) { + const auto batch_items = batch.items.size(); + const auto& continuity = entry->continuity; + const bool exceeds_batch_limit = + continuity.max_buffered_batches > 0 && + entry->continuity_buffer.size() >= continuity.max_buffered_batches; + const bool exceeds_item_limit = continuity.max_buffered_items > 0 && + (batch_items > continuity.max_buffered_items || + entry->continuity_buffered_items > + continuity.max_buffered_items - batch_items); + + if (!exceeds_batch_limit && !exceeds_item_limit) { + if (push_front) { + entry->continuity_buffer.push_front(std::move(batch)); + } else { + entry->continuity_buffer.push_back(std::move(batch)); + } + entry->continuity_buffered_items += batch_items; + return true; + } + + while (!entry->continuity_buffer.empty()) { + auto buffered = std::move(entry->continuity_buffer.front()); + entry->continuity_buffer.pop_front(); + deliveries.emplace_back(subscriber, std::move(buffered)); + } + entry->continuity_buffered_items = 0; + deliveries.emplace_back(subscriber, std::move(batch)); + + entry->continuity_ready = true; + entry->continuity_request_in_flight = false; + entry->continuity_flushing = false; + entry->continuity_retry_request.reset(); + entry->continuity_retry_at = {}; + entry->continuity.mode = MarketDataContinuityMode::LIVE_ONLY; + + continuity_deliveries.emplace_back( + subscriber, + make_continuity_update( + entry->control->provider_subscription, + MarketDataContinuityStatus::FAILED, + 0, + 0, + 0, + 0, + "Continuity buffer limit exceeded; live delivery continues.")); + continuity_deliveries.emplace_back( + subscriber, + make_continuity_update( + entry->control->provider_subscription, + MarketDataContinuityStatus::DEGRADED, + 0, + 0, + 0, + 0, + "Live delivery resumed without verified continuity after buffer overflow.")); + return false; + } + inline bool MarketDataRouterState::route_bar_to_entry_no_lock( const std::shared_ptr& entry, const BarDataBatch& batch, @@ -1897,6 +2292,9 @@ namespace optionx::market_data { std::shared_ptr, BarDataBatch>>& deliveries, std::vector& continuity_requests, + std::vector, + MarketDataContinuityUpdate>>& continuity_deliveries, bool process_buffered, bool allow_gap_recovery) { if (!entry || entry->phase != EntryPhase::ACTIVE || @@ -1909,10 +2307,17 @@ namespace optionx::market_data { auto routed = batch; routed.subscription = entry->control->provider_subscription; + if (routed.items.empty()) return false; if (entry->continuity.enabled() && !process_buffered && (!entry->continuity_ready || entry->continuity_flushing)) { - entry->continuity_buffer.push_back(std::move(routed)); + buffer_continuity_batch_no_lock( + entry, + std::move(subscriber), + std::move(routed), + deliveries, + continuity_deliveries, + false); return false; } @@ -1923,8 +2328,8 @@ namespace optionx::market_data { const auto timeframe_ms = static_cast(entry->stream.timeframe) * 1000U; std::uint64_t previous_time_ms = entry->last_bar_time_ms; - for (std::size_t index = 0; index < batch.items.size(); ++index) { - const auto& bar = batch.items[index]; + for (std::size_t index = 0; index < routed.items.size(); ++index) { + const auto& bar = routed.items[index]; if (bar.time_ms == 0) continue; const auto expected_time_ms = previous_time_ms > std::numeric_limits::max() - timeframe_ms @@ -1941,55 +2346,63 @@ namespace optionx::market_data { entry->stream.transport); request.continuity = entry->continuity; - auto requested_items = static_cast(0); - const auto gap_items = - ((gap_to_ms - gap_from_ms) / timeframe_ms) + 1U; - requested_items = gap_items > + const auto history_request = + MarketDataContinuityService::make_gap_request( + request, + gap_from_ms, + gap_to_ms, + entry->continuity.max_backfill_bars); + const auto request_from_time_ms = + MarketDataContinuityService::seconds_to_milliseconds( + history_request.from_ts); + const auto request_to_time_ms = + MarketDataContinuityService::seconds_to_milliseconds( + history_request.to_ts); + const auto request_items = + request_from_time_ms > 0 && + request_to_time_ms >= request_from_time_ms + ? ((request_to_time_ms - request_from_time_ms) / + timeframe_ms) + 1U + : 0U; + const auto requested_items = request_items > static_cast( std::numeric_limits::max()) ? std::numeric_limits::max() - : static_cast(gap_items); - if (entry->continuity.max_backfill_bars > 0) { - requested_items = std::min( - requested_items, - entry->continuity.max_backfill_bars); - } + : static_cast(request_items); if (index > 0) { BarDataBatch prefix = routed; prefix.items.resize(index); + record_bar_progress_no_lock( + entry, + prefix.items); deliveries.emplace_back(subscriber, std::move(prefix)); - for (std::size_t prefix_index = 0; - prefix_index < index; - ++prefix_index) { - entry->last_bar_time_ms = std::max( - entry->last_bar_time_ms, - batch.items[prefix_index].time_ms); - } } routed.items.erase( routed.items.begin(), routed.items.begin() + static_cast(index)); + if (!buffer_continuity_batch_no_lock( + entry, + subscriber, + std::move(routed), + deliveries, + continuity_deliveries, + process_buffered)) { + return false; + } entry->continuity_ready = false; entry->continuity_request_in_flight = true; - if (process_buffered) { - entry->continuity_buffer.push_front(std::move(routed)); - } else { - entry->continuity_buffer.push_back(std::move(routed)); - } continuity_requests.push_back(PendingContinuityRequest{ entry->router_id, entry->control->provider_subscription, - MarketDataContinuityService::make_gap_request( - request, - gap_from_ms, - gap_to_ms, - entry->continuity.max_backfill_bars), - MarketDataContinuityStatus::GAP_DETECTED, - gap_from_ms, - gap_to_ms, - requested_items}); + history_request, + ContinuityRequestKind::GAP_BACKFILL, + true, + request_from_time_ms, + request_to_time_ms, + requested_items, + 1}); return false; } } @@ -1999,11 +2412,7 @@ namespace optionx::market_data { } } - for (const auto& bar : routed.items) { - if (bar.time_ms > entry->last_bar_time_ms) { - entry->last_bar_time_ms = bar.time_ms; - } - } + record_bar_progress_no_lock(entry, routed.items); deliveries.emplace_back(std::move(subscriber), std::move(routed)); return true; } @@ -2016,6 +2425,9 @@ namespace optionx::market_data { std::vector, BarDataBatch>> deliveries; + std::vector, + MarketDataContinuityUpdate>> continuity_deliveries; std::vector continuity_requests; { std::lock_guard lock(m_mutex); @@ -2028,6 +2440,7 @@ namespace optionx::market_data { *batch, deliveries, continuity_requests, + continuity_deliveries, false); }; @@ -2046,6 +2459,10 @@ namespace optionx::market_data { } } + for (auto& continuity_delivery : continuity_deliveries) { + continuity_delivery.first->on_market_data_continuity( + continuity_delivery.second); + } for (auto& delivery : deliveries) { delivery.first->on_bar_data(delivery.second); } @@ -2054,10 +2471,12 @@ namespace optionx::market_data { request.router_id, std::move(request.subscription), std::move(request.request), - request.status, + request.kind, + request.announce_gap, request.from_time_ms, request.to_time_ms, - request.requested_items); + request.requested_items, + request.attempt); } } @@ -2195,24 +2614,62 @@ namespace optionx::market_data { }; for (;;) { + std::vector retry_requests; std::vector completions; + bool shutting_down = false; { std::lock_guard lock(m_mutex); - if (!m_shutdown || m_shutdown_complete) return; + if (m_shutdown_complete) return; + shutting_down = m_shutdown; + + if (!shutting_down) { + const auto now = std::chrono::steady_clock::now(); + for (const auto& [id, entry] : m_entries) { + (void)id; + if (entry->phase != EntryPhase::ACTIVE || + entry->continuity_request_in_flight || + !entry->continuity_retry_request || + entry->continuity_retry_at > now) { + continue; + } + retry_requests.push_back( + *entry->continuity_retry_request); + entry->continuity_retry_request.reset(); + entry->continuity_retry_at = {}; + entry->continuity_request_in_flight = true; + } + } - completions.reserve(m_entries.size()); - for (const auto& [id, entry] : m_entries) { - if (entry->phase != EntryPhase::UNSUBSCRIBING || - !entry->unsubscribe_completion_received) { - continue; + if (shutting_down) { + completions.reserve(m_entries.size()); + for (const auto& [id, entry] : m_entries) { + if (entry->phase != EntryPhase::UNSUBSCRIBING || + !entry->unsubscribe_completion_received) { + continue; + } + completions.push_back(PendingCompletion{ + id, + entry->unsubscribe_completion.subscription, + entry->unsubscribe_completion}); } - completions.push_back(PendingCompletion{ - id, - entry->unsubscribe_completion.subscription, - entry->unsubscribe_completion}); } } + for (auto& retry : retry_requests) { + request_continuity_history( + retry.router_id, + std::move(retry.subscription), + std::move(retry.request), + retry.kind, + retry.announce_gap, + retry.from_time_ms, + retry.to_time_ms, + retry.requested_items, + retry.attempt); + } + + if (!shutting_down) return; + for (auto& completion : completions) { complete_unsubscribe( completion.router_id, diff --git a/tests/market_data_continuity_test.cpp b/tests/market_data_continuity_test.cpp index 1c2cfb2..697533c 100644 --- a/tests/market_data_continuity_test.cpp +++ b/tests/market_data_continuity_test.cpp @@ -1,5 +1,6 @@ #include +#include #include #include #include @@ -170,6 +171,17 @@ BarSubscriptionRequest continuity_request(MarketDataContinuityMode mode) { return request; } +std::size_t continuity_status_count( + const RecordingSubscriber& subscriber, + MarketDataContinuityStatus status) { + return static_cast(std::count_if( + subscriber.continuity.begin(), + subscriber.continuity.end(), + [status](const MarketDataContinuityUpdate& update) { + return update.status == status; + })); +} + TEST(MarketDataContinuity, PrefillDeliversHistoryBeforeBufferedLiveBars) { FakeHistoryProvider provider; MarketDataRouter router; @@ -237,6 +249,200 @@ TEST(MarketDataContinuity, RecoversTimestampGapBeforeReleasingLiveBatch) { route.provider_subscription().id); } +TEST(MarketDataContinuity, RejectsGapHistoryMissingRangeStart) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + ASSERT_TRUE(route.active()); + provider.complete_history(make_history({100000})); + + provider.emit_live_bar(280000); + ASSERT_EQ(provider.history_requests.size(), 2U); + provider.complete_history(make_history({220000})); + + ASSERT_EQ(subscriber->bars.size(), 2U); + EXPECT_EQ(subscriber->bars[0].items.front().time_ms, 100000U); + EXPECT_EQ(subscriber->bars[1].items.front().time_ms, 280000U); + EXPECT_EQ( + continuity_status_count(*subscriber, MarketDataContinuityStatus::LIVE), + 1U); + ASSERT_GE(subscriber->continuity.size(), 2U); + EXPECT_EQ( + subscriber->continuity[subscriber->continuity.size() - 2].status, + MarketDataContinuityStatus::FAILED); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::DEGRADED); +} + +TEST(MarketDataContinuity, RejectsGapHistoryMissingRangeEnd) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + ASSERT_TRUE(route.active()); + provider.complete_history(make_history({100000})); + + provider.emit_live_bar(280000); + ASSERT_EQ(provider.history_requests.size(), 2U); + provider.complete_history(make_history({160000})); + + ASSERT_EQ(subscriber->bars.size(), 2U); + EXPECT_EQ(subscriber->bars[0].items.front().time_ms, 100000U); + EXPECT_EQ(subscriber->bars[1].items.front().time_ms, 280000U); + EXPECT_EQ( + continuity_status_count(*subscriber, MarketDataContinuityStatus::LIVE), + 1U); + ASSERT_GE(subscriber->continuity.size(), 2U); + EXPECT_EQ( + subscriber->continuity[subscriber->continuity.size() - 2].status, + MarketDataContinuityStatus::FAILED); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::DEGRADED); +} + +TEST(MarketDataContinuity, LaterGapRepairDoesNotHideEarlierMissingRange) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + ASSERT_TRUE(route.active()); + provider.complete_history(make_history({100000})); + + provider.emit_live_bar(280000); + ASSERT_EQ(provider.history_requests.size(), 2U); + provider.complete_history(make_history({220000})); + ASSERT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::DEGRADED); + + const auto live_status_count = continuity_status_count( + *subscriber, + MarketDataContinuityStatus::LIVE); + provider.emit_live_bar(400000); + ASSERT_EQ(provider.history_requests.size(), 3U); + EXPECT_EQ(provider.history_requests[2].from_ts, 340); + EXPECT_EQ(provider.history_requests[2].to_ts, 340); + + provider.complete_history(make_history({340000})); + + EXPECT_EQ( + continuity_status_count( + *subscriber, + MarketDataContinuityStatus::LIVE), + live_status_count); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::DEGRADED); +} + +TEST(MarketDataContinuity, FailedPrefillKeepsLaterGapRecoveryDegraded) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + ASSERT_TRUE(route.active()); + provider.fail_history("initial history unavailable"); + ASSERT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::DEGRADED); + + const auto live_status_count = continuity_status_count( + *subscriber, + MarketDataContinuityStatus::LIVE); + provider.emit_live_bar(100000); + provider.emit_live_bar(220000); + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests[1].from_ts, 160); + EXPECT_EQ(provider.history_requests[1].to_ts, 160); + + provider.complete_history(make_history({160000})); + + EXPECT_EQ( + continuity_status_count( + *subscriber, + MarketDataContinuityStatus::LIVE), + live_status_count); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::DEGRADED); +} + +TEST(MarketDataContinuity, ClipsGapHistoryToRequestedRange) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + ASSERT_TRUE(route.active()); + provider.complete_history(make_history({100000})); + + provider.emit_live_bar(280000); + provider.complete_history(make_history({100000, 160000, 220000, 280000})); + + ASSERT_EQ(subscriber->bars.size(), 3U); + ASSERT_EQ(subscriber->bars[1].items.size(), 2U); + EXPECT_EQ(subscriber->bars[1].items[0].time_ms, 160000U); + EXPECT_EQ(subscriber->bars[1].items[1].time_ms, 220000U); + EXPECT_EQ(subscriber->bars[2].items.front().time_ms, 280000U); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::LIVE); +} + +TEST(MarketDataContinuity, ReportsActualBoundedGapRange) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER); + request.continuity.max_backfill_bars = 2; + + auto route = router.subscribe_bars(provider, subscriber, request); + ASSERT_TRUE(route.active()); + provider.complete_history(make_history({100000})); + + provider.emit_live_bar(400000); + + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests[1].from_ts, 160); + EXPECT_EQ(provider.history_requests[1].to_ts, 220); + ASSERT_FALSE(subscriber->continuity.empty()); + EXPECT_EQ(subscriber->continuity.back().from_time_ms, 160000U); + EXPECT_EQ(subscriber->continuity.back().to_time_ms, 220000U); + EXPECT_EQ(subscriber->continuity.back().requested_items, 2U); + + provider.complete_history(make_history({160000, 220000})); + ASSERT_EQ(provider.history_requests.size(), 3U); + EXPECT_EQ(provider.history_requests[2].from_ts, 280); + EXPECT_EQ(provider.history_requests[2].to_ts, 340); + provider.complete_history(make_history({280000, 340000})); + + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::LIVE); +} + TEST(MarketDataContinuity, ChecksGapsInsideBufferedBatchesInOrder) { FakeHistoryProvider provider; MarketDataRouter router; @@ -278,11 +484,13 @@ TEST(MarketDataContinuity, EmptyBackfillFailsOnceAndReleasesLiveBatch) { FakeHistoryProvider provider; MarketDataRouter router; auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER); + request.continuity.retry.max_attempts = 3; auto route = router.subscribe_bars( provider, subscriber, - continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + request); ASSERT_TRUE(route.active()); provider.complete_history(make_history({100000})); @@ -295,7 +503,7 @@ TEST(MarketDataContinuity, EmptyBackfillFailsOnceAndReleasesLiveBatch) { EXPECT_EQ(subscriber->bars.back().items.front().time_ms, 280000U); ASSERT_EQ(subscriber->continuity.size(), 6U); EXPECT_EQ(subscriber->continuity[4].status, MarketDataContinuityStatus::FAILED); - EXPECT_EQ(subscriber->continuity[5].status, MarketDataContinuityStatus::LIVE); + EXPECT_EQ(subscriber->continuity[5].status, MarketDataContinuityStatus::DEGRADED); } TEST(MarketDataContinuity, HistoryFailureKeepsLiveRouteUsable) { @@ -318,7 +526,7 @@ TEST(MarketDataContinuity, HistoryFailureKeepsLiveRouteUsable) { ASSERT_EQ(subscriber->continuity.size(), 3U); EXPECT_EQ(subscriber->continuity[1].status, MarketDataContinuityStatus::FAILED); EXPECT_EQ(subscriber->continuity[1].message, "history endpoint unavailable"); - EXPECT_EQ(subscriber->continuity[2].status, MarketDataContinuityStatus::LIVE); + EXPECT_EQ(subscriber->continuity[2].status, MarketDataContinuityStatus::DEGRADED); provider.emit_live_bar(260000); EXPECT_EQ(subscriber->bars.size(), 2U); @@ -401,7 +609,7 @@ TEST(MarketDataContinuity, RejectsHistoryResponseFromAnotherStream) { "Historical bar response does not match the subscribed stream."); EXPECT_EQ( subscriber->continuity[2].status, - MarketDataContinuityStatus::LIVE); + MarketDataContinuityStatus::DEGRADED); provider.emit_live_bar(280000); EXPECT_EQ(provider.history_requests.size(), 1U); @@ -418,13 +626,225 @@ TEST(MarketDataContinuity, PrefillRequestUsesInclusiveBarCountRange) { const auto history = MarketDataContinuityService::make_prefill_request( request, - 600000, + 600013, 3); EXPECT_EQ(history.from_ts, 480); EXPECT_EQ(history.to_ts, 600); } +TEST(MarketDataContinuity, RetriesFailedHistoryBeforeReleasingLive) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL); + request.continuity.retry.max_attempts = 2; + + auto route = router.subscribe_bars(provider, subscriber, request); + ASSERT_TRUE(route.active()); + + provider.emit_live_bar(200000); + provider.fail_history("temporary history failure"); + + ASSERT_EQ(provider.history_requests.size(), 1U); + EXPECT_TRUE(subscriber->bars.empty()); + ASSERT_EQ(subscriber->continuity.size(), 2U); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::RETRYING); + + router.process(); + ASSERT_EQ(provider.history_requests.size(), 2U); + + provider.complete_history(make_history({100000})); + + ASSERT_EQ(subscriber->bars.size(), 2U); + EXPECT_TRUE(subscriber->bars[0].items[0].has_flag(MarketDataFlags::HISTORICAL)); + EXPECT_TRUE(subscriber->bars[1].items[0].has_flag(MarketDataFlags::REALTIME)); + ASSERT_EQ(subscriber->continuity.size(), 3U); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::LIVE); +} + +TEST(MarketDataContinuity, RetriesRejectedHistoryBeforeReleasingLive) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL); + request.continuity.retry.max_attempts = 2; + + provider.reject_history = true; + auto route = router.subscribe_bars(provider, subscriber, request); + ASSERT_TRUE(route.active()); + + provider.emit_live_bar(200000); + ASSERT_EQ(provider.history_requests.size(), 1U); + ASSERT_EQ(subscriber->continuity.size(), 2U); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::RETRYING); + EXPECT_TRUE(subscriber->bars.empty()); + + provider.reject_history = false; + router.process(); + ASSERT_EQ(provider.history_requests.size(), 2U); + provider.complete_history(make_history({100000})); + + ASSERT_EQ(subscriber->bars.size(), 2U); + EXPECT_TRUE(subscriber->bars[0].items[0].has_flag(MarketDataFlags::HISTORICAL)); + EXPECT_TRUE(subscriber->bars[1].items[0].has_flag(MarketDataFlags::REALTIME)); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::LIVE); +} + +TEST(MarketDataContinuity, ProcessRunsOnlyOneRetryPass) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL); + request.continuity.retry.max_attempts = 3; + + provider.reject_history = true; + auto route = router.subscribe_bars(provider, subscriber, request); + ASSERT_TRUE(route.active()); + ASSERT_EQ(provider.history_requests.size(), 1U); + + router.process(); + + EXPECT_EQ(provider.history_requests.size(), 2U); + ASSERT_FALSE(subscriber->continuity.empty()); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::RETRYING); +} + +TEST(MarketDataContinuity, GapRetryDoesNotRepeatGapDetectedStatus) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER); + request.continuity.retry.max_attempts = 2; + + auto route = router.subscribe_bars(provider, subscriber, request); + ASSERT_TRUE(route.active()); + provider.complete_history(make_history({100000})); + + provider.emit_live_bar(280000); + provider.fail_history("temporary backfill failure"); + ASSERT_EQ(subscriber->continuity.size(), 5U); + EXPECT_EQ( + subscriber->continuity[2].status, + MarketDataContinuityStatus::GAP_DETECTED); + EXPECT_EQ( + subscriber->continuity[4].status, + MarketDataContinuityStatus::RETRYING); + + router.process(); + ASSERT_EQ(provider.history_requests.size(), 3U); + ASSERT_EQ(subscriber->continuity.size(), 6U); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::BACKFILLING); + + provider.complete_history(make_history({160000, 220000})); + + const auto gap_detected_count = std::count_if( + subscriber->continuity.begin(), + subscriber->continuity.end(), + [](const MarketDataContinuityUpdate& update) { + return update.status == MarketDataContinuityStatus::GAP_DETECTED; + }); + EXPECT_EQ(gap_detected_count, 1); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::LIVE); +} + +TEST(MarketDataContinuity, ExhaustedRetriesReleaseBufferedLiveOnce) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL); + request.continuity.retry.max_attempts = 2; + + auto route = router.subscribe_bars(provider, subscriber, request); + ASSERT_TRUE(route.active()); + + provider.emit_live_bar(200000); + provider.fail_history("first failure"); + router.process(); + provider.fail_history("second failure"); + + ASSERT_EQ(subscriber->bars.size(), 1U); + EXPECT_TRUE(subscriber->bars.front().items.front().has_flag( + MarketDataFlags::REALTIME)); + ASSERT_EQ(subscriber->continuity.size(), 4U); + EXPECT_EQ( + subscriber->continuity[2].status, + MarketDataContinuityStatus::FAILED); + EXPECT_EQ( + subscriber->continuity[3].status, + MarketDataContinuityStatus::DEGRADED); +} + +TEST(MarketDataContinuity, RetryBackoffDefersProviderRequest) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL); + request.continuity.retry.max_attempts = 2; + request.continuity.retry.initial_backoff_ms = 60000; + request.continuity.retry.max_backoff_ms = 60000; + + auto route = router.subscribe_bars(provider, subscriber, request); + ASSERT_TRUE(route.active()); + provider.fail_history("temporary history failure"); + + router.process(); + + EXPECT_EQ(provider.history_requests.size(), 1U); + ASSERT_FALSE(subscriber->continuity.empty()); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::RETRYING); +} + +TEST(MarketDataContinuity, BufferLimitReleasesLiveDataAndDisablesContinuity) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL); + request.continuity.max_buffered_batches = 1; + request.continuity.max_buffered_items = 1; + + auto route = router.subscribe_bars(provider, subscriber, request); + ASSERT_TRUE(route.active()); + + provider.emit_live_bar(200000); + EXPECT_TRUE(subscriber->bars.empty()); + provider.emit_live_bar(260000); + + ASSERT_EQ(subscriber->bars.size(), 2U); + EXPECT_EQ(subscriber->bars[0].items.front().time_ms, 200000U); + EXPECT_EQ(subscriber->bars[1].items.front().time_ms, 260000U); + ASSERT_EQ(subscriber->continuity.size(), 3U); + EXPECT_EQ( + subscriber->continuity[1].status, + MarketDataContinuityStatus::FAILED); + EXPECT_EQ( + subscriber->continuity[2].status, + MarketDataContinuityStatus::DEGRADED); + + provider.complete_history(make_history({100000})); + EXPECT_EQ(subscriber->bars.size(), 2U); + + provider.emit_live_bar(320000); + ASSERT_EQ(subscriber->bars.size(), 3U); + EXPECT_EQ(subscriber->bars.back().items.front().time_ms, 320000U); +} + TEST(MarketDataContinuity, UnsubscribeDuringHistoryDropsLateDelivery) { FakeHistoryProvider provider; MarketDataRouter router;