From b12cc38e3e4d66a81d8ec9262564549bb1f8b505 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Fri, 4 Sep 2026 00:12:32 +0300 Subject: [PATCH 1/4] feat(market-data): harden bar continuity recovery --- examples/market_data_continuity_example.cpp | 5 + guides/api-and-header-contracts.md | 17 +- guides/market-data-router.md | 24 ++- guides/market-data-router.ru.md | 28 ++- guides/refactor-backlog.md | 21 +- .../market_data/MarketDataContinuity.hpp | 3 + .../MarketDataContinuityOptions.hpp | 54 ++++- .../market_data/MarketDataRouter.hpp | 4 +- .../market_data/detail/MarketDataRouter.ipp | 187 ++++++++++++++++-- tests/market_data_continuity_test.cpp | 129 ++++++++++++ 10 files changed, 412 insertions(+), 60 deletions(-) diff --git a/examples/market_data_continuity_example.cpp b/examples/market_data_continuity_example.cpp index e0917e1..7c8b585 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.bar_policy = + md::MarketDataContinuityBarPolicy::DROP_NON_MONOTONIC; + 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..ab06fc9 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -257,13 +257,22 @@ 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()`. +- `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. +- `MarketDataContinuityOptions::bar_policy` defaults to `KEEP_ALL`, preserving + provider revisions and overlapping snapshots. `DROP_NON_MONOTONIC` instead + filters zero, repeated, and out-of-order timestamps per route for both + history and live batches. Gap detection runs after this filtering. - 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 returns to `LIVE` after the retry budget is exhausted. +- Generic history continuity is currently defined for bars only. Tick history + remains a separate provider contract. The strict bar policy is a route-local + timestamp filter, not a universal semantic deduplication rule for all + providers. `MarketDataRouter` is the subscription-scoped alternative to `MarketDataHub`: diff --git a/guides/market-data-router.md b/guides/market-data-router.md index 40a3a6c..3a0c264 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -326,6 +326,14 @@ The options have these meanings: application wants recovery only. - `max_backfill_bars` bounds one gap request. Zero means that the provider request is not count-bounded by Router. +- `bar_policy` defaults to `KEEP_ALL`, preserving provider revisions. Set it to + `DROP_NON_MONOTONIC` when the consumer requires strictly increasing bar + timestamps; repeated and out-of-order bars are then discarded per 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 `FAILED -> LIVE` behavior. Retries are advanced by periodic + `MarketDataRouter::process()` calls on the owner loop. For initial prefill, the delivery order is: @@ -372,14 +380,14 @@ 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. +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 `LIVE`. The route stays usable +and live delivery continues, but an unrecovered historical range is reported to +the consumer. A provider may return overlapping snapshots for an in-progress +bar; `KEEP_ALL` preserves those revisions, while `DROP_NON_MONOTONIC` provides +a strict timestamp filter for consumers that need one. `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..4f4f4ad 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -326,6 +326,16 @@ auto route = router.subscribe_bars("intrade", chart, request); использовать ноль, если приложению нужно только восстановление разрывов. - `max_backfill_bars` ограничивает один запрос для разрыва. Ноль означает, что Router не ограничивает provider request количеством баров. +- `bar_policy` по умолчанию равен `KEEP_ALL` и сохраняет revisions провайдера. + Установите `DROP_NON_MONOTONIC`, если consumer требует строго возрастающие + timestamps: повторные и пришедшие не по порядку бары будут отброшены для + каждого route. +- `retry.max_attempts` задаёт общее количество попыток вместе с первой. + `retry.initial_backoff_ms` включает ограниченный exponential backoff, а + `retry.max_backoff_ms` задаёт его предел. По умолчанию выполняется одна + попытка, поэтому сохраняется прежняя политика `FAILED -> LIVE`. Повторы + запускаются периодическими вызовами `MarketDataRouter::process()` в owner + loop. Для initial prefill порядок доставки такой: @@ -373,15 +383,15 @@ Router продолжит разбирать очередь live batches и мо 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`/незавершённых баров. +Если history request завершился ошибкой или provider его отклонил, но попытки +ещё остались, Router публикует `RETRYING`, сохраняет накопленные live batches и +ждёт следующего вызова `process()` в owner loop. После исчерпания retry budget +Router публикует `FAILED`, выпускает накопленные live batches, затем публикует +`LIVE`. Route остаётся пригодным для работы, live-доставка продолжается, а +невосстановленный исторический диапазон явно сообщается consumer. Provider может +возвращать пересекающиеся snapshots для незавершённого бара; `KEEP_ALL` +сохраняет такие revisions, а `DROP_NON_MONOTONIC` даёт строгий timestamp-фильтр +для consumer, которому он нужен. `MarketDataContinuityService` остаётся низкоуровневым helper для приложений, которые хотят запрашивать историю напрямую. Он превращает `BarHistoryResult` в diff --git a/guides/refactor-backlog.md b/guides/refactor-backlog.md index c78295a..4886b20 100644 --- a/guides/refactor-backlog.md +++ b/guides/refactor-backlog.md @@ -5,24 +5,19 @@ series. Keep it short and remove items once they are handled. ## Next PR Candidates -- Extend market-data continuity beyond the first bar-only route implementation: - define provider support for tick history, retries, and a documented - history-to-live boundary for each provider. -- Add robust gap recovery policy with provider-aware retry/backoff, sequence or - timestamp validation, and an explicit deduplication policy for overlapping - historical, backfill, and live bar snapshots. - Add route-scoped continuity metrics and failure visibility for applications that need to prove that a chart or strategy has a complete time series. +- Add a generic tick-history provider contract and extend continuity from bars + to ticks where a provider can supply historical tick data. - Add a fuller CMake package/export story for consumers that do not use the project as a direct submodule. The current `optionx_cpp::optionx_cpp` interface target covers build-tree/submodule consumption. ## Explicitly Deferred -- Add a generic tick-history provider contract. Current continuity support is - intentionally bar-first because providers expose bar history only. -- Continue generation-safe lifecycle hardening for legacy bridge transports - when their behavior is changed; do not mix that work into market-data API - PRs. -- `TradeUpPlatform` remains a partial implementation. Do not refactor it as - part of generic cleanup PRs unless the task is specifically about TradeUp. +- `TradeUpPlatform` sources remain historical examples. The broker is no + longer available, so do not add new production work or generic cleanup for + this platform. +- Legacy bridge lifecycle hardening is complete. Preserve the existing + generation, callback, and shutdown patterns when changing those bridges, but + do not track the already completed audit as active work here. diff --git a/include/optionx_cpp/market_data/MarketDataContinuity.hpp b/include/optionx_cpp/market_data/MarketDataContinuity.hpp index 49c2f81..18c6470 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuity.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuity.hpp @@ -18,6 +18,7 @@ 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. + RETRYING, ///< A failed history request will be attempted again. LIVE, ///< Live delivery is current, with no pending history work. FAILED ///< History work failed; live delivery continues without it. }; @@ -31,6 +32,8 @@ 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: diff --git a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp index 1605971..f3040c1 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp @@ -6,6 +6,8 @@ /// \brief Defines history prefill and gap-recovery options for bar routes. #include +#include +#include namespace optionx::market_data { @@ -17,12 +19,58 @@ namespace optionx::market_data { PREFILL_AND_RECOVER ///< Prefill and repair timestamp gaps in live bars. }; + /// \enum MarketDataContinuityBarPolicy + /// \brief Selects how repeated or out-of-order bar timestamps are handled. + enum class MarketDataContinuityBarPolicy { + KEEP_ALL = 0, ///< Preserve every provider bar, including revisions. + DROP_NON_MONOTONIC ///< Keep only strictly increasing timestamps per route. + }; + + /// \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); + } + + /// \brief Calculates the delay after a failed one-based attempt. + [[nodiscard]] std::uint64_t delay_after_attempt( + std::size_t attempt) const noexcept { + if (attempt == 0 || initial_backoff_ms == 0) return 0; + + auto delay = initial_backoff_ms; + for (std::size_t index = 1; index < attempt; ++index) { + if (delay > std::numeric_limits::max() / 2U) { + delay = std::numeric_limits::max(); + break; + } + delay *= 2U; + if (max_backoff_ms > 0 && delay >= max_backoff_ms) { + delay = max_backoff_ms; + break; + } + } + return max_backoff_ms > 0 + ? (delay < max_backoff_ms ? delay : max_backoff_ms) + : delay; + } + }; + /// \struct MarketDataContinuityOptions - /// \brief Configures history prefill and timestamp-gap recovery for a bar route. + /// \brief Configures history prefill, recovery, ordering, and retries. 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. + MarketDataContinuityBarPolicy bar_policy = + MarketDataContinuityBarPolicy::KEEP_ALL; + MarketDataContinuityRetryPolicy retry; /// \brief Returns true when the option combination is usable. [[nodiscard]] bool valid() const noexcept { @@ -30,9 +78,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/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..9d8a1af 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -43,6 +43,8 @@ namespace optionx::market_data { CLEANUP_FAILED }; + struct PendingContinuityRequest; + struct Entry { RoutedSubscriptionId router_id; ProviderInstanceId provider_id = kInvalidProviderInstanceId; @@ -56,6 +58,8 @@ namespace optionx::market_data { bool continuity_flushing = false; std::deque continuity_buffer; std::uint64_t last_bar_time_ms = 0; + std::shared_ptr continuity_retry_request; + std::uint64_t continuity_retry_at_ms = 0; MarketDataSubscriptionHandle retained_cleanup_subscription; MarketDataSubscriptionResult unsubscribe_completion; bool subscribe_completion_posted = false; @@ -74,6 +78,7 @@ namespace optionx::market_data { 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 ContinuityOperation { @@ -276,7 +281,8 @@ namespace optionx::market_data { MarketDataContinuityStatus status, 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, @@ -285,6 +291,7 @@ namespace optionx::market_data { 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, @@ -306,6 +313,9 @@ namespace optionx::market_data { std::vector& continuity_requests, bool process_buffered = false, bool allow_gap_recovery = true); + static void apply_bar_policy( + const Entry& entry, + BarDataBatch& batch); void fail_pending_subscribe( RoutedSubscriptionId router_id, @@ -1249,7 +1259,8 @@ namespace optionx::market_data { pending.status, pending.from_time_ms, pending.to_time_ms, - pending.requested_items); + pending.requested_items, + 1); } inline void MarketDataRouterState::request_continuity_history( @@ -1259,13 +1270,15 @@ namespace optionx::market_data { MarketDataContinuityStatus status, 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; @@ -1307,7 +1320,8 @@ namespace optionx::market_data { status, 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, @@ -1318,6 +1332,7 @@ namespace optionx::market_data { from_time_ms, to_time_ms, requested_items, + attempt, result = std::move(result)]() mutable { state->complete_continuity( router_id, @@ -1327,6 +1342,7 @@ namespace optionx::market_data { from_time_ms, to_time_ms, requested_items, + attempt, std::move(result)); }; state->dispatch_or_run(std::move(task)); @@ -1359,6 +1375,7 @@ namespace optionx::market_data { 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; { @@ -1389,6 +1406,13 @@ namespace optionx::market_data { history_batch, expected_stream); if (history_stream_matches) { + 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; + } + apply_bar_policy(*entry_it->second, history_batch); delivered_history_items = history_batch.items.size(); has_history_batch = !history_batch.items.empty(); } @@ -1398,6 +1422,62 @@ namespace optionx::market_data { history_stream_matches && (operation_status != MarketDataContinuityStatus::GAP_DETECTED || delivered_history_items > 0); + + bool retry_scheduled = false; + if (!usable_history) { + std::shared_ptr 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 = std::make_shared(); + retry_request->router_id = router_id; + retry_request->subscription = subscription; + retry_request->request = request; + retry_request->status = operation_status; + 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 delay_ms = entry->continuity.retry + .delay_after_attempt(attempt); + auto retry_at_ms = static_cast( + OPTIONX_TIMESTAMP_MS); + if (retry_at_ms > + std::numeric_limits::max() - delay_ms) { + retry_at_ms = std::numeric_limits::max(); + } else { + retry_at_ms += delay_ms; + } + entry->continuity_retry_request = std::move(retry_request); + entry->continuity_retry_at_ms = retry_at_ms; + 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; + } + } + if (!usable_history) { notify_continuity( router_id, @@ -1484,7 +1564,8 @@ namespace optionx::market_data { pending.status, pending.from_time_ms, pending.to_time_ms, - pending.requested_items); + pending.requested_items, + pending.attempt); return; } @@ -1890,6 +1971,26 @@ namespace optionx::market_data { if (subscriber) subscriber->on_market_data_continuity(update); } + inline void MarketDataRouterState::apply_bar_policy( + const Entry& entry, + BarDataBatch& batch) { + if (entry.continuity.bar_policy != + MarketDataContinuityBarPolicy::DROP_NON_MONOTONIC) { + return; + } + + auto previous_time_ms = entry.last_bar_time_ms; + auto output = batch.items.begin(); + for (const auto& bar : batch.items) { + if (bar.time_ms == 0 || bar.time_ms <= previous_time_ms) { + continue; + } + previous_time_ms = bar.time_ms; + *output++ = bar; + } + batch.items.erase(output, batch.items.end()); + } + inline bool MarketDataRouterState::route_bar_to_entry_no_lock( const std::shared_ptr& entry, const BarDataBatch& batch, @@ -1909,6 +2010,8 @@ namespace optionx::market_data { auto routed = batch; routed.subscription = entry->control->provider_subscription; + apply_bar_policy(*entry, routed); + if (routed.items.empty()) return false; if (entry->continuity.enabled() && !process_buffered && (!entry->continuity_ready || entry->continuity_flushing)) { @@ -1923,8 +2026,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 @@ -1964,7 +2067,7 @@ namespace optionx::market_data { ++prefix_index) { entry->last_bar_time_ms = std::max( entry->last_bar_time_ms, - batch.items[prefix_index].time_ms); + routed.items[prefix_index].time_ms); } } @@ -1989,7 +2092,8 @@ namespace optionx::market_data { MarketDataContinuityStatus::GAP_DETECTED, gap_from_ms, gap_to_ms, - requested_items}); + requested_items, + 1}); return false; } } @@ -2057,7 +2161,8 @@ namespace optionx::market_data { request.status, request.from_time_ms, request.to_time_ms, - request.requested_items); + request.requested_items, + request.attempt); } } @@ -2195,24 +2300,64 @@ 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_ms = static_cast( + OPTIONX_TIMESTAMP_MS); + 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_ms > now_ms) { + continue; + } + retry_requests.push_back( + *entry->continuity_retry_request); + entry->continuity_retry_request.reset(); + 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.status, + retry.from_time_ms, + retry.to_time_ms, + retry.requested_items, + retry.attempt); + } + + if (!shutting_down) { + if (retry_requests.empty()) return; + continue; + } + 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..ed3ef8f 100644 --- a/tests/market_data_continuity_test.cpp +++ b/tests/market_data_continuity_test.cpp @@ -425,6 +425,135 @@ TEST(MarketDataContinuity, PrefillRequestUsesInclusiveBarCountRange) { 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, 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::LIVE); +} + +TEST(MarketDataContinuity, RetryPolicyUsesCappedExponentialBackoff) { + MarketDataContinuityRetryPolicy retry; + retry.max_attempts = 4; + retry.initial_backoff_ms = 100; + retry.max_backoff_ms = 250; + + EXPECT_EQ(retry.delay_after_attempt(1), 100U); + EXPECT_EQ(retry.delay_after_attempt(2), 200U); + EXPECT_EQ(retry.delay_after_attempt(3), 250U); + EXPECT_EQ(retry.delay_after_attempt(4), 250U); +} + +TEST(MarketDataContinuity, StrictBarPolicyDropsRepeatedAndOutOfOrderBars) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + auto request = continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER); + request.continuity.bar_policy = + MarketDataContinuityBarPolicy::DROP_NON_MONOTONIC; + + auto route = router.subscribe_bars(provider, subscriber, request); + ASSERT_TRUE(route.active()); + + provider.complete_history(make_history({100000, 100000, 160000})); + ASSERT_EQ(subscriber->bars.size(), 1U); + ASSERT_EQ(subscriber->bars.front().items.size(), 2U); + EXPECT_EQ(subscriber->bars.front().items[0].time_ms, 100000U); + EXPECT_EQ(subscriber->bars.front().items[1].time_ms, 160000U); + + provider.emit_live_bars({160000, 140000, 220000}); + + ASSERT_EQ(subscriber->bars.size(), 2U); + EXPECT_EQ(subscriber->bars.back().items.front().time_ms, 220000U); + EXPECT_EQ(provider.history_requests.size(), 1U); +} + TEST(MarketDataContinuity, UnsubscribeDuringHistoryDropsLateDelivery) { FakeHistoryProvider provider; MarketDataRouter router; From 82fa103db5f38118b18299dd0e0b50b9467415ae Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Fri, 4 Sep 2026 12:13:53 +0300 Subject: [PATCH 2/4] fix(market-data): bound continuity recovery state --- examples/market_data_continuity_example.cpp | 4 +- guides/api-and-header-contracts.md | 15 +- guides/market-data-router.md | 34 ++- guides/market-data-router.ru.md | 37 ++- .../MarketDataContinuityOptions.hpp | 13 +- .../MarketDataContinuityService.hpp | 15 +- .../market_data/detail/MarketDataRouter.ipp | 260 +++++++++++++----- tests/market_data_continuity_test.cpp | 106 ++++++- 8 files changed, 348 insertions(+), 136 deletions(-) diff --git a/examples/market_data_continuity_example.cpp b/examples/market_data_continuity_example.cpp index 7c8b585..981f9ba 100644 --- a/examples/market_data_continuity_example.cpp +++ b/examples/market_data_continuity_example.cpp @@ -139,8 +139,8 @@ int main() { request.continuity.mode = md::MarketDataContinuityMode::PREFILL_AND_RECOVER; request.continuity.prefill_bars = 2; request.continuity.max_backfill_bars = 10; - request.continuity.bar_policy = - md::MarketDataContinuityBarPolicy::DROP_NON_MONOTONIC; + 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; diff --git a/guides/api-and-header-contracts.md b/guides/api-and-header-contracts.md index ab06fc9..5c92fec 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -257,14 +257,21 @@ 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. - `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. -- `MarketDataContinuityOptions::bar_policy` defaults to `KEEP_ALL`, preserving - provider revisions and overlapping snapshots. `DROP_NON_MONOTONIC` instead - filters zero, repeated, and out-of-order timestamps per route for both - history and live batches. Gap detection runs after this filtering. +- 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 diff --git a/guides/market-data-router.md b/guides/market-data-router.md index 3a0c264..b6cd421 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -319,16 +319,21 @@ 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. -- `bar_policy` defaults to `KEEP_ALL`, preserving provider revisions. Set it to - `DROP_NON_MONOTONIC` when the consumer requires strictly increasing bar - timestamps; repeated and out-of-order bars are then discarded per route. +- `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 `LIVE`. This is an explicit loss-of-recovery fallback: live data + continues, but the route no longer promises historical ordering. - `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 @@ -377,17 +382,20 @@ stream-level `MarketDataStatusUpdate` events such as `READY` or `DISCONNECTED`. 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. +gap is treated as a terminal failed recovery, 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. 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 `LIVE`. The route stays usable and live delivery continues, but an unrecovered historical range is reported to -the consumer. A provider may return overlapping snapshots for an in-progress -bar; `KEEP_ALL` preserves those revisions, while `DROP_NON_MONOTONIC` provides -a strict timestamp filter for consumers that need one. +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 4f4f4ad..17fe50f 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -319,17 +319,21 @@ 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 количеством баров. -- `bar_policy` по умолчанию равен `KEEP_ALL` и сохраняет revisions провайдера. - Установите `DROP_NON_MONOTONIC`, если consumer требует строго возрастающие - timestamps: повторные и пришедшие не по порядку бары будут отброшены для - каждого route. +- `max_buffered_batches` и `max_buffered_items` ограничивают live data, + удерживаемые во время history request. Ноль отключает соответствующее + ограничение. При превышении любого лимита Router публикует continuity + `FAILED`, выпускает накопленные и текущий live batch, отключает continuity + для этого route и публикует `LIVE`. Это явный fallback с потерей обещания + historical ordering: live-доставка продолжается. - `retry.max_attempts` задаёт общее количество попыток вместе с первой. `retry.initial_backoff_ms` включает ограниченный exponential backoff, а `retry.max_backoff_ms` задаёт его предел. По умолчанию выполняется одна @@ -380,18 +384,21 @@ stream-level `MarketDataStatusUpdate`, например `READY` или `DISCONNE отсутствующий интервал. Если provider вернул неполный, но непустой backfill, Router продолжит разбирать очередь live batches и может выполнить следующий ограниченный запрос для оставшегося gap. Успешный пустой ответ для найденного -gap считается failed recovery: Router один раз выпускает накопленные live data и -не зацикливает запрос того же диапазона. +gap считается terminal failed recovery: Router один раз выпускает накопленные +live data и не зацикливает запрос того же диапазона. Успешный пустой prefill +тоже является terminal result: исторический batch не доставляется, а Router +переходит к накопленному live-потоку без retry. Если history request завершился ошибкой или provider его отклонил, но попытки ещё остались, Router публикует `RETRYING`, сохраняет накопленные live batches и ждёт следующего вызова `process()` в owner loop. После исчерпания retry budget Router публикует `FAILED`, выпускает накопленные live batches, затем публикует `LIVE`. Route остаётся пригодным для работы, live-доставка продолжается, а -невосстановленный исторический диапазон явно сообщается consumer. Provider может -возвращать пересекающиеся snapshots для незавершённого бара; `KEEP_ALL` -сохраняет такие revisions, а `DROP_NON_MONOTONIC` даёт строгий timestamp-фильтр -для consumer, которому он нужен. +невосстановленный исторический диапазон явно сообщается 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/MarketDataContinuityOptions.hpp b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp index f3040c1..83804ac 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp @@ -19,13 +19,6 @@ namespace optionx::market_data { PREFILL_AND_RECOVER ///< Prefill and repair timestamp gaps in live bars. }; - /// \enum MarketDataContinuityBarPolicy - /// \brief Selects how repeated or out-of-order bar timestamps are handled. - enum class MarketDataContinuityBarPolicy { - KEEP_ALL = 0, ///< Preserve every provider bar, including revisions. - DROP_NON_MONOTONIC ///< Keep only strictly increasing timestamps per route. - }; - /// \struct MarketDataContinuityRetryPolicy /// \brief Configures bounded history retry attempts and exponential backoff. struct MarketDataContinuityRetryPolicy { @@ -63,14 +56,14 @@ namespace optionx::market_data { }; /// \struct MarketDataContinuityOptions - /// \brief Configures history prefill, recovery, ordering, and retries. + /// \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. - MarketDataContinuityBarPolicy bar_policy = - MarketDataContinuityBarPolicy::KEEP_ALL; 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 { 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/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index 9d8a1af..cda6319 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -43,6 +43,11 @@ namespace optionx::market_data { CLEANUP_FAILED }; + enum class ContinuityRequestKind { + PREFILL, + GAP_BACKFILL + }; + struct PendingContinuityRequest; struct Entry { @@ -57,6 +62,7 @@ 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::shared_ptr continuity_retry_request; std::uint64_t continuity_retry_at_ms = 0; @@ -74,7 +80,8 @@ namespace optionx::market_data { RoutedSubscriptionId router_id; MarketDataSubscriptionHandle subscription; BarHistoryRequest request; - MarketDataContinuityStatus status = MarketDataContinuityStatus::UNKNOWN; + 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; @@ -278,7 +285,8 @@ 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, @@ -287,7 +295,7 @@ namespace optionx::market_data { 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, @@ -311,11 +319,22 @@ 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 void apply_bar_policy( - const Entry& entry, - BarDataBatch& batch); + 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); void fail_pending_subscribe( RoutedSubscriptionId router_id, @@ -1229,7 +1248,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); @@ -1256,7 +1276,8 @@ 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, @@ -1267,7 +1288,8 @@ 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, @@ -1288,17 +1310,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( @@ -1317,7 +1341,7 @@ namespace optionx::market_data { router_id, subscription, request, - status, + kind, from_time_ms, to_time_ms, requested_items, @@ -1328,7 +1352,7 @@ namespace optionx::market_data { router_id, subscription, request, - status, + kind, from_time_ms, to_time_ms, requested_items, @@ -1338,7 +1362,7 @@ namespace optionx::market_data { router_id, subscription, request, - status, + kind, from_time_ms, to_time_ms, requested_items, @@ -1371,7 +1395,7 @@ namespace optionx::market_data { 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, @@ -1401,7 +1425,7 @@ 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); @@ -1412,7 +1436,6 @@ namespace optionx::market_data { entry_it->second->phase != EntryPhase::ACTIVE) { return; } - apply_bar_policy(*entry_it->second, history_batch); delivered_history_items = history_batch.items.size(); has_history_batch = !history_batch.items.empty(); } @@ -1420,11 +1443,11 @@ namespace optionx::market_data { const bool usable_history = history_success && history_stream_matches && - (operation_status != MarketDataContinuityStatus::GAP_DETECTED || + (kind != ContinuityRequestKind::GAP_BACKFILL || delivered_history_items > 0); bool retry_scheduled = false; - if (!usable_history) { + if (!history_success) { std::shared_ptr retry_request; { std::lock_guard lock(m_mutex); @@ -1440,7 +1463,8 @@ namespace optionx::market_data { retry_request->router_id = router_id; retry_request->subscription = subscription; retry_request->request = request; - retry_request->status = operation_status; + 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; @@ -1493,7 +1517,9 @@ 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 + ? "No bars were returned for the detected gap." + : "No historical bars were returned for prefill.") : result.error_desc)); } @@ -1524,8 +1550,12 @@ 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; { std::lock_guard lock(m_mutex); const auto entry_it = m_entries.find(router_id); @@ -1536,21 +1566,35 @@ 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(); 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); } @@ -1561,7 +1605,8 @@ 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, @@ -1570,18 +1615,20 @@ namespace optionx::market_data { } 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, + 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.")); + } return; } } @@ -1971,24 +2018,73 @@ namespace optionx::market_data { if (subscriber) subscriber->on_market_data_continuity(update); } - inline void MarketDataRouterState::apply_bar_policy( - const Entry& entry, - BarDataBatch& batch) { - if (entry.continuity.bar_policy != - MarketDataContinuityBarPolicy::DROP_NON_MONOTONIC) { - return; + 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; } - auto previous_time_ms = entry.last_bar_time_ms; - auto output = batch.items.begin(); - for (const auto& bar : batch.items) { - if (bar.time_ms == 0 || bar.time_ms <= previous_time_ms) { - continue; - } - previous_time_ms = bar.time_ms; - *output++ = bar; + 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)); } - batch.items.erase(output, batch.items.end()); + 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_ms = 0; + 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::LIVE, + 0, + 0, + 0, + 0, + "Live delivery resumed after continuity buffer overflow.")); + return false; } inline bool MarketDataRouterState::route_bar_to_entry_no_lock( @@ -1998,6 +2094,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 || @@ -2010,12 +2109,17 @@ namespace optionx::market_data { auto routed = batch; routed.subscription = entry->control->provider_subscription; - apply_bar_policy(*entry, routed); 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; } @@ -2074,13 +2178,17 @@ namespace optionx::market_data { 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, @@ -2089,7 +2197,8 @@ namespace optionx::market_data { gap_from_ms, gap_to_ms, entry->continuity.max_backfill_bars), - MarketDataContinuityStatus::GAP_DETECTED, + ContinuityRequestKind::GAP_BACKFILL, + true, gap_from_ms, gap_to_ms, requested_items, @@ -2120,6 +2229,9 @@ namespace optionx::market_data { std::vector, BarDataBatch>> deliveries; + std::vector, + MarketDataContinuityUpdate>> continuity_deliveries; std::vector continuity_requests; { std::lock_guard lock(m_mutex); @@ -2132,6 +2244,7 @@ namespace optionx::market_data { *batch, deliveries, continuity_requests, + continuity_deliveries, false); }; @@ -2150,6 +2263,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); } @@ -2158,7 +2275,8 @@ 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, @@ -2346,17 +2464,15 @@ namespace optionx::market_data { retry.router_id, std::move(retry.subscription), std::move(retry.request), - retry.status, + retry.kind, + retry.announce_gap, retry.from_time_ms, retry.to_time_ms, retry.requested_items, retry.attempt); } - if (!shutting_down) { - if (retry_requests.empty()) return; - continue; - } + if (!shutting_down) return; for (auto& completion : completions) { complete_unsubscribe( diff --git a/tests/market_data_continuity_test.cpp b/tests/market_data_continuity_test.cpp index ed3ef8f..02284ae 100644 --- a/tests/market_data_continuity_test.cpp +++ b/tests/market_data_continuity_test.cpp @@ -1,5 +1,6 @@ #include +#include #include #include #include @@ -278,11 +279,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})); @@ -418,7 +421,7 @@ TEST(MarketDataContinuity, PrefillRequestUsesInclusiveBarCountRange) { const auto history = MarketDataContinuityService::make_prefill_request( request, - 600000, + 600013, 3); EXPECT_EQ(history.from_ts, 480); @@ -491,6 +494,69 @@ TEST(MarketDataContinuity, RetriesRejectedHistoryBeforeReleasingLive) { 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; @@ -530,28 +596,38 @@ TEST(MarketDataContinuity, RetryPolicyUsesCappedExponentialBackoff) { EXPECT_EQ(retry.delay_after_attempt(4), 250U); } -TEST(MarketDataContinuity, StrictBarPolicyDropsRepeatedAndOutOfOrderBars) { +TEST(MarketDataContinuity, BufferLimitReleasesLiveDataAndDisablesContinuity) { FakeHistoryProvider provider; MarketDataRouter router; auto subscriber = std::make_shared(); - auto request = continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER); - request.continuity.bar_policy = - MarketDataContinuityBarPolicy::DROP_NON_MONOTONIC; + 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.complete_history(make_history({100000, 100000, 160000})); - ASSERT_EQ(subscriber->bars.size(), 1U); - ASSERT_EQ(subscriber->bars.front().items.size(), 2U); - EXPECT_EQ(subscriber->bars.front().items[0].time_ms, 100000U); - EXPECT_EQ(subscriber->bars.front().items[1].time_ms, 160000U); - - provider.emit_live_bars({160000, 140000, 220000}); + 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.back().items.front().time_ms, 220000U); - EXPECT_EQ(provider.history_requests.size(), 1U); + 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::LIVE); + + 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) { From a7d7d3568207b3d253549e6f340b57566813803d Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Sat, 5 Sep 2026 01:01:45 +0300 Subject: [PATCH 3/4] fix(market-data): validate bounded continuity history --- guides/api-and-header-contracts.md | 16 +- guides/market-data-router.md | 34 ++- guides/market-data-router.ru.md | 18 +- guides/refactor-backlog.md | 21 +- .../market_data/MarketDataContinuity.hpp | 7 +- .../MarketDataContinuityOptions.hpp | 22 -- .../market_data/detail/MarketDataRouter.ipp | 257 +++++++++++++----- tests/market_data_continuity_test.cpp | 168 +++++++++++- 8 files changed, 406 insertions(+), 137 deletions(-) diff --git a/guides/api-and-header-contracts.md b/guides/api-and-header-contracts.md index 5c92fec..7a68cc1 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -266,6 +266,15 @@ Contract rules: 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. - 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. @@ -275,11 +284,10 @@ Contract rules: - 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` after the retry budget is exhausted. + 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. The strict bar policy is a route-local - timestamp filter, not a universal semantic deduplication rule for all - providers. + 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 b6cd421..37021ed 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -337,7 +337,8 @@ The options have these meanings: - `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 `FAILED -> LIVE` behavior. Retries are advanced by periodic + 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: @@ -380,22 +381,29 @@ 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 terminal failed recovery, 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. +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. 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 `LIVE`. 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. +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 17fe50f..c4d31d9 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -337,7 +337,8 @@ auto route = router.subscribe_bars("intrade", chart, request); - `retry.max_attempts` задаёт общее количество попыток вместе с первой. `retry.initial_backoff_ms` включает ограниченный exponential backoff, а `retry.max_backoff_ms` задаёт его предел. По умолчанию выполняется одна - попытка, поэтому сохраняется прежняя политика `FAILED -> LIVE`. Повторы + попытка, поэтому live-доставка продолжается, а итоговый статус становится + `DEGRADED`. Повторы запускаются периодическими вызовами `MarketDataRouter::process()` в owner loop. @@ -382,18 +383,25 @@ stream-level `MarketDataStatusUpdate`, например `READY` или `DISCONNE `max_backfill_bars` ограничивает каждый отдельный provider request, а не весь отсутствующий интервал. Если provider вернул неполный, но непустой backfill, -Router продолжит разбирать очередь live batches и может выполнить следующий -ограниченный запрос для оставшегося gap. Успешный пустой ответ для найденного -gap считается terminal failed recovery: Router один раз выпускает накопленные +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. + Если history request завершился ошибкой или provider его отклонил, но попытки ещё остались, Router публикует `RETRYING`, сохраняет накопленные live batches и ждёт следующего вызова `process()` в owner loop. После исчерпания retry budget Router публикует `FAILED`, выпускает накопленные live batches, затем публикует -`LIVE`. Route остаётся пригодным для работы, live-доставка продолжается, а +`DEGRADED`. Route остаётся пригодным для работы, live-доставка продолжается, а невосстановленный исторический диапазон явно сообщается consumer. Router сохраняет revisions баров от provider и не применяет универсальную timestamp-дедупликацию. Если графику или storage нужно одно текущее значение diff --git a/guides/refactor-backlog.md b/guides/refactor-backlog.md index 4886b20..c78295a 100644 --- a/guides/refactor-backlog.md +++ b/guides/refactor-backlog.md @@ -5,19 +5,24 @@ series. Keep it short and remove items once they are handled. ## Next PR Candidates +- Extend market-data continuity beyond the first bar-only route implementation: + define provider support for tick history, retries, and a documented + history-to-live boundary for each provider. +- Add robust gap recovery policy with provider-aware retry/backoff, sequence or + timestamp validation, and an explicit deduplication policy for overlapping + historical, backfill, and live bar snapshots. - Add route-scoped continuity metrics and failure visibility for applications that need to prove that a chart or strategy has a complete time series. -- Add a generic tick-history provider contract and extend continuity from bars - to ticks where a provider can supply historical tick data. - Add a fuller CMake package/export story for consumers that do not use the project as a direct submodule. The current `optionx_cpp::optionx_cpp` interface target covers build-tree/submodule consumption. ## Explicitly Deferred -- `TradeUpPlatform` sources remain historical examples. The broker is no - longer available, so do not add new production work or generic cleanup for - this platform. -- Legacy bridge lifecycle hardening is complete. Preserve the existing - generation, callback, and shutdown patterns when changing those bridges, but - do not track the already completed audit as active work here. +- Add a generic tick-history provider contract. Current continuity support is + intentionally bar-first because providers expose bar history only. +- Continue generation-safe lifecycle hardening for legacy bridge transports + when their behavior is changed; do not mix that work into market-data API + PRs. +- `TradeUpPlatform` remains a partial implementation. Do not refactor it as + part of generic cleanup PRs unless the task is specifically about TradeUp. diff --git a/include/optionx_cpp/market_data/MarketDataContinuity.hpp b/include/optionx_cpp/market_data/MarketDataContinuity.hpp index 18c6470..de27584 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuity.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuity.hpp @@ -19,8 +19,9 @@ namespace optionx::market_data { GAP_DETECTED, ///< A timestamp gap was found in the live stream. BACKFILLING, ///< Historical bars are being loaded for a gap. RETRYING, ///< A failed history request will be attempted again. - LIVE, ///< Live delivery is current, with no pending history work. - FAILED ///< History work failed; live delivery continues without it. + 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. }; /// \brief Converts a continuity status to stable text. @@ -38,6 +39,8 @@ namespace optionx::market_data { 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 83804ac..2baa132 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp @@ -7,7 +7,6 @@ #include #include -#include namespace optionx::market_data { @@ -32,27 +31,6 @@ namespace optionx::market_data { (max_backoff_ms == 0 || max_backoff_ms >= initial_backoff_ms); } - /// \brief Calculates the delay after a failed one-based attempt. - [[nodiscard]] std::uint64_t delay_after_attempt( - std::size_t attempt) const noexcept { - if (attempt == 0 || initial_backoff_ms == 0) return 0; - - auto delay = initial_backoff_ms; - for (std::size_t index = 1; index < attempt; ++index) { - if (delay > std::numeric_limits::max() / 2U) { - delay = std::numeric_limits::max(); - break; - } - delay *= 2U; - if (max_backoff_ms > 0 && delay >= max_backoff_ms) { - delay = max_backoff_ms; - break; - } - } - return max_backoff_ms > 0 - ? (delay < max_backoff_ms ? delay : max_backoff_ms) - : delay; - } }; /// \struct MarketDataContinuityOptions diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index cda6319..86a5a30 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 { @@ -48,7 +51,17 @@ namespace optionx::market_data { GAP_BACKFILL }; - struct PendingContinuityRequest; + 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; @@ -64,8 +77,8 @@ namespace optionx::market_data { std::deque continuity_buffer; std::size_t continuity_buffered_items = 0; std::uint64_t last_bar_time_ms = 0; - std::shared_ptr continuity_retry_request; - std::uint64_t continuity_retry_at_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; @@ -76,18 +89,6 @@ namespace optionx::market_data { subscription_callback_t release_callback; }; - 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 ContinuityOperation { bool completed = false; }; @@ -333,8 +334,20 @@ namespace optionx::market_data { BarDataBatch>>& deliveries, std::vector, - MarketDataContinuityUpdate>>& continuity_deliveries, + 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 std::chrono::steady_clock::duration continuity_retry_delay( + const MarketDataContinuityRetryPolicy& policy, + std::size_t attempt) noexcept; void fail_pending_subscribe( RoutedSubscriptionId router_id, @@ -1391,6 +1404,82 @@ 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 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, @@ -1416,6 +1505,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; @@ -1430,6 +1522,12 @@ namespace optionx::market_data { 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() || @@ -1441,14 +1539,27 @@ namespace optionx::market_data { } } + 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 && (kind != ContinuityRequestKind::GAP_BACKFILL || - delivered_history_items > 0); + delivered_history_items > 0) && + gap_history_covers_range; bool retry_scheduled = false; if (!history_success) { - std::shared_ptr retry_request; + PendingContinuityRequest retry_request; { std::lock_guard lock(m_mutex); const auto entry_it = m_entries.find(router_id); @@ -1459,29 +1570,26 @@ namespace optionx::market_data { const auto& entry = entry_it->second; if (attempt < entry->continuity.retry.max_attempts) { - retry_request = std::make_shared(); - 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 delay_ms = entry->continuity.retry - .delay_after_attempt(attempt); - auto retry_at_ms = static_cast( - OPTIONX_TIMESTAMP_MS); - if (retry_at_ms > - std::numeric_limits::max() - delay_ms) { - retry_at_ms = std::numeric_limits::max(); - } else { - retry_at_ms += delay_ms; - } + 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_ms = retry_at_ms; + entry->continuity_retry_at = delay >= remaining + ? std::chrono::steady_clock::time_point::max() + : now + delay; entry->continuity_flushing = false; retry_scheduled = true; } @@ -1517,6 +1625,9 @@ namespace optionx::market_data { ? "Historical market-data continuity request failed." : !history_stream_matches ? "Historical bar response does not match the subscribed stream." + : 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.") @@ -1620,14 +1731,16 @@ namespace optionx::market_data { router_id, make_continuity_update( subscription, - MarketDataContinuityStatus::LIVE, + usable_history + ? MarketDataContinuityStatus::LIVE + : MarketDataContinuityStatus::DEGRADED, 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.")); + : "Live delivery continues without verified continuity.")); } return; } @@ -2061,7 +2174,7 @@ namespace optionx::market_data { entry->continuity_request_in_flight = false; entry->continuity_flushing = false; entry->continuity_retry_request.reset(); - entry->continuity_retry_at_ms = 0; + entry->continuity_retry_at = {}; entry->continuity.mode = MarketDataContinuityMode::LIVE_ONLY; continuity_deliveries.emplace_back( @@ -2075,15 +2188,15 @@ namespace optionx::market_data { 0, "Continuity buffer limit exceeded; live delivery continues.")); continuity_deliveries.emplace_back( - subscriber, - make_continuity_update( - entry->control->provider_subscription, - MarketDataContinuityStatus::LIVE, - 0, + subscriber, + make_continuity_update( + entry->control->provider_subscription, + MarketDataContinuityStatus::DEGRADED, + 0, 0, 0, 0, - "Live delivery resumed after continuity buffer overflow.")); + "Live delivery resumed without verified continuity after buffer overflow.")); return false; } @@ -2148,19 +2261,29 @@ 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; @@ -2192,15 +2315,11 @@ namespace optionx::market_data { 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), + history_request, ContinuityRequestKind::GAP_BACKFILL, true, - gap_from_ms, - gap_to_ms, + request_from_time_ms, + request_to_time_ms, requested_items, 1}); return false; @@ -2427,19 +2546,19 @@ namespace optionx::market_data { shutting_down = m_shutdown; if (!shutting_down) { - const auto now_ms = static_cast( - OPTIONX_TIMESTAMP_MS); + 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_ms > now_ms) { + 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; } } diff --git a/tests/market_data_continuity_test.cpp b/tests/market_data_continuity_test.cpp index 02284ae..143b32f 100644 --- a/tests/market_data_continuity_test.cpp +++ b/tests/market_data_continuity_test.cpp @@ -171,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; @@ -238,6 +249,125 @@ 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, 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; @@ -298,7 +428,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) { @@ -321,7 +451,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); @@ -404,7 +534,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); @@ -581,19 +711,29 @@ TEST(MarketDataContinuity, ExhaustedRetriesReleaseBufferedLiveOnce) { MarketDataContinuityStatus::FAILED); EXPECT_EQ( subscriber->continuity[3].status, - MarketDataContinuityStatus::LIVE); + MarketDataContinuityStatus::DEGRADED); } -TEST(MarketDataContinuity, RetryPolicyUsesCappedExponentialBackoff) { - MarketDataContinuityRetryPolicy retry; - retry.max_attempts = 4; - retry.initial_backoff_ms = 100; - retry.max_backoff_ms = 250; +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(retry.delay_after_attempt(1), 100U); - EXPECT_EQ(retry.delay_after_attempt(2), 200U); - EXPECT_EQ(retry.delay_after_attempt(3), 250U); - EXPECT_EQ(retry.delay_after_attempt(4), 250U); + EXPECT_EQ(provider.history_requests.size(), 1U); + ASSERT_FALSE(subscriber->continuity.empty()); + EXPECT_EQ( + subscriber->continuity.back().status, + MarketDataContinuityStatus::RETRYING); } TEST(MarketDataContinuity, BufferLimitReleasesLiveDataAndDisablesContinuity) { @@ -620,7 +760,7 @@ TEST(MarketDataContinuity, BufferLimitReleasesLiveDataAndDisablesContinuity) { MarketDataContinuityStatus::FAILED); EXPECT_EQ( subscriber->continuity[2].status, - MarketDataContinuityStatus::LIVE); + MarketDataContinuityStatus::DEGRADED); provider.complete_history(make_history({100000})); EXPECT_EQ(subscriber->bars.size(), 2U); From bd98d8473beb56ff6196df8d676aab5fe86046b3 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Sat, 5 Sep 2026 02:30:47 +0300 Subject: [PATCH 4/4] fix(market-data): retain continuity trust loss --- guides/api-and-header-contracts.md | 8 +- guides/market-data-router.md | 11 +- guides/market-data-router.ru.md | 11 +- .../market_data/MarketDataContinuity.hpp | 2 +- .../market_data/detail/MarketDataRouter.ipp | 115 +++++++++++++++--- tests/market_data_continuity_test.cpp | 75 ++++++++++++ 6 files changed, 197 insertions(+), 25 deletions(-) diff --git a/guides/api-and-header-contracts.md b/guides/api-and-header-contracts.md index 7a68cc1..bf274a8 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -261,7 +261,8 @@ Contract rules: 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. + 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 @@ -275,6 +276,11 @@ Contract rules: - `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. diff --git a/guides/market-data-router.md b/guides/market-data-router.md index 37021ed..199bfd6 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -332,8 +332,9 @@ The options have these meanings: 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 `LIVE`. This is an explicit loss-of-recovery fallback: live data - continues, but the route no longer promises historical ordering. + 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 @@ -395,6 +396,12 @@ live stream without retry. 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`, diff --git a/guides/market-data-router.ru.md b/guides/market-data-router.ru.md index c4d31d9..54fba91 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -332,8 +332,9 @@ auto route = router.subscribe_bars("intrade", chart, request); удерживаемые во время history request. Ноль отключает соответствующее ограничение. При превышении любого лимита Router публикует continuity `FAILED`, выпускает накопленные и текущий live batch, отключает continuity - для этого route и публикует `LIVE`. Это явный fallback с потерей обещания - historical ordering: live-доставка продолжается. + для этого route и публикует `DEGRADED`, продолжая live-доставку. Это явный + fallback с потерей обещания historical ordering: continuity автоматически + для этого route не включается. - `retry.max_attempts` задаёт общее количество попыток вместе с первой. `retry.initial_backoff_ms` включает ограниченный exponential backoff, а `retry.max_backoff_ms` задаёт его предел. По умолчанию выполняется одна @@ -397,6 +398,12 @@ live data и не зацикливает запрос того же диапаз 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 diff --git a/include/optionx_cpp/market_data/MarketDataContinuity.hpp b/include/optionx_cpp/market_data/MarketDataContinuity.hpp index de27584..33f2c13 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuity.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuity.hpp @@ -21,7 +21,7 @@ namespace optionx::market_data { 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. + 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. diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index 86a5a30..f3a319d 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -77,6 +77,8 @@ namespace optionx::market_data { 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; @@ -345,6 +347,16 @@ namespace optionx::market_data { 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; @@ -1444,6 +1456,61 @@ namespace optionx::market_data { 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, @@ -1610,6 +1677,23 @@ namespace optionx::market_data { } } + { + 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, @@ -1645,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) { @@ -1667,6 +1749,7 @@ namespace optionx::market_data { 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); @@ -1679,6 +1762,8 @@ namespace optionx::market_data { 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( @@ -1731,14 +1816,14 @@ namespace optionx::market_data { router_id, make_continuity_update( subscription, - usable_history + continuity_verified ? MarketDataContinuityStatus::LIVE : MarketDataContinuityStatus::DEGRADED, from_time_ms, to_time_ms, requested_items, delivered_history_items, - usable_history + continuity_verified ? "Historical market-data continuity is ready." : "Live delivery continues without verified continuity.")); } @@ -2288,14 +2373,10 @@ namespace optionx::market_data { 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, - routed.items[prefix_index].time_ms); - } } routed.items.erase( @@ -2331,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; } diff --git a/tests/market_data_continuity_test.cpp b/tests/market_data_continuity_test.cpp index 143b32f..697533c 100644 --- a/tests/market_data_continuity_test.cpp +++ b/tests/market_data_continuity_test.cpp @@ -311,6 +311,81 @@ TEST(MarketDataContinuity, RejectsGapHistoryMissingRangeEnd) { 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;