From 0a296aa7985a40e5f9be2df5ac0a4dd79130b0c5 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Mon, 7 Sep 2026 14:46:43 +0300 Subject: [PATCH 1/5] feat(market-data): add router tick continuity --- CMakeLists.txt | 1 + guides/api-and-header-contracts.md | 22 +- guides/market-data-router.md | 51 +- guides/market-data-router.ru.md | 53 +- guides/platform-api-guide.md | 14 +- guides/refactor-backlog.md | 8 +- .../market_data/MarketDataContinuity.hpp | 9 +- .../MarketDataContinuityOptions.hpp | 45 +- .../market_data/MarketDataSubscription.hpp | 1 + .../market_data/detail/MarketDataRouter.ipp | 1601 ++++++++++++++++- tests/market_data_tick_continuity_test.cpp | 527 ++++++ 11 files changed, 2227 insertions(+), 105 deletions(-) create mode 100644 tests/market_data_tick_continuity_test.cpp diff --git a/CMakeLists.txt b/CMakeLists.txt index 1e76dd0..0a3b54f 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -1110,6 +1110,7 @@ if(OPTIONX_BUILD_TESTS) telegram_worker_source_test trading_view_bridge_test market_data_tick_history_contract_test + market_data_tick_continuity_test intrade_observed_tick_history_test ) if(OPTIONX_LIGHTWEIGHT_BRIDGE_SMOKE_TESTS) diff --git a/guides/api-and-header-contracts.md b/guides/api-and-header-contracts.md index dc42f79..abafe2d 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -376,13 +376,21 @@ broker: plain `PREFILL` remains startup-only. A cached invalidating status applies the same transition before a newly accepted route may start prefill. Completed `PREFILL` routes do not acquire reconnect or timestamp-gap recovery implicitly. -- Bar continuity is currently implemented by Router. The provider API also - defines a separate `fetch_tick_history()` contract with inclusive - millisecond ranges, explicit `range_complete`, and non-decreasing timestamp - order validated by the adapter. A non-empty result symbol must match the - request. No - current provider implements authoritative tick history yet; Router does not - apply tick continuity or a universal timestamp deduplication policy. +- Router continuity supports both bar and tick routes. The provider API defines + a separate `fetch_tick_history()` contract with inclusive millisecond ranges, + explicit `range_complete`, and non-decreasing timestamp order validated by the + adapter. A non-empty result symbol must match the request. Tick routes opt in + through `TickSubscriptionRequest::continuity`: `PREFILL` holds live ticks + until the requested lookback completes, while `PREFILL_AND_RECOVER` also + detects suspicious timestamp gaps and repairs them in bounded ranges. +- Tick `expected_interval_ms` is only a gap-detection hint because tick streams + are event-oriented and may contain sub-second or equal-timestamp events. + `range_complete` remains the provider's continuity authority. Router reports + incomplete history as operation-level `FAILED` and sticky `DEGRADED`, while + still allowing returned observations to be delivered. Reconnect recovery waits + for `READY`; exact overlap identity is `(time_ms, ask, bid, last, volume)`, so + `received_ms`/flags do not distinguish duplicates and different same-second + observations remain distinct. `MarketDataRouter` is the subscription-scoped alternative to `MarketDataHub`: diff --git a/guides/market-data-router.md b/guides/market-data-router.md index 66d8be5..88cdac3 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -634,10 +634,11 @@ service.request_tick_history_batch( true); // require_complete_range ``` -`MarketDataRouter` still has its mature continuity state machine on bars. The -Intrade Bar provider now also implements `fetch_tick_history()` as a bounded, +`MarketDataRouter` also applies the history contract to tick routes. The +Intrade Bar provider implements `fetch_tick_history()` as a bounded, session-scoped archive of observed `/price_now` snapshots. This is useful for -short reconnect windows, but it is not an authoritative broker tick archive: +short prefill and reconnect windows, but it is not an authoritative broker tick +archive: - broker timestamps have one-second granularity; - the archive starts empty for a new authenticated session and evicts old data; @@ -648,9 +649,47 @@ short reconnect windows, but it is not an authoritative broker tick archive: - `trade_check2.php` remains a settlement/trade-result API and is not used for range history. -Router tick continuity is the next layer. Until that integration is enabled, -callers can use the provider operation directly and must treat -`range_complete=false` as an observation result rather than continuity proof. +### Tick continuity in Router + +Set `TickSubscriptionRequest::continuity` to enable history-first delivery for +one tick route: + +```cpp +md::TickSubscriptionRequest request("EURUSD"); +request.continuity.mode = md::MarketDataContinuityMode::PREFILL_AND_RECOVER; +request.continuity.prefill_lookback_ms = 60'000; +request.continuity.expected_interval_ms = 1'000; +request.continuity.max_backfill_ms = 60'000; + +auto route = router.subscribe_ticks(provider, bot, request); +``` + +`PREFILL` requests the configured lookback before releasing live ticks. +`PREFILL_AND_RECOVER` also holds the live tail when the difference between +successive observed timestamps exceeds `expected_interval_ms`. That value is a +gap-detection hint, not a claim that every tick must arrive on a fixed grid; +equal timestamps and sub-second live ticks are valid. The provider's +`range_complete` remains the authority for whether a requested history range +proves continuity. + +The Router sends historical ticks first, marks them `HISTORICAL`, and then +replays held live ticks as `LIVE_SOURCE | CATCHUP`. A complete result is required +before the route can report `LIVE`. An incomplete result may still be delivered +as observations, but it reports `FAILED` followed by sticky `DEGRADED`; a later +unrelated successful range cannot hide the earlier unresolved watermark. +History requests are bounded by `max_backfill_ms` and are scheduled by +`process()`, so tick continuity does not create a timer thread. + +On reconnect, tick continuity reports `STALE`, waits for `READY`, and requests +the unresolved range through the latest observed time. Exact overlap is removed +by `(time_ms, ask, bid, last, volume)` identity. `received_ms` and flags do not +make an otherwise identical observation distinct, while different observations +with the same second remain separate events. If the continuity buffer exceeds +its batch or item limit, Router releases the held live data, reports +`FAILED`/`DEGRADED`, disables continuity for that route, and resumes ordinary +live delivery. If transport is interrupted during the initial prefill, Router +restarts from the original lookback start after `READY` and extends the request +through the current time, so the interrupted interval is not silently skipped. ## Owner Loop And Bot Threads diff --git a/guides/market-data-router.ru.md b/guides/market-data-router.ru.md index fe03e14..a79a913 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -791,11 +791,11 @@ service.request_tick_history_batch( true); // require_complete_range ``` -`MarketDataRouter` уже имеет полноценную state machine continuity для bars. -Кроме того, Intrade Bar теперь реализует `fetch_tick_history(...)` через -ограниченный архив наблюдаемых snapshots из `/price_now`. Это полезно для -короткого восстановления после reconnect, но не является authoritative -broker tick archive: +`MarketDataRouter` теперь также применяет history contract к tick routes. +Intrade Bar реализует `fetch_tick_history(...)` через ограниченный +session-scoped архив наблюдаемых snapshots из `/price_now`. Это полезно для +короткого prefill и reconnect recovery, но не является authoritative broker +tick archive: - broker timestamps имеют гранулярность в одну секунду; - архив пуст для новой authenticated session и вытесняет старые данные; @@ -806,6 +806,43 @@ broker tick archive: - `trade_check2.php` остаётся settlement/trade-result API и не используется для range history. -Интеграция tick continuity в Router выполняется следующим слоем. До неё -вызывающий код может использовать provider operation напрямую и обязан -считать `range_complete=false` observations, а не доказательством continuity. +### Tick continuity в Router + +Чтобы включить history-first delivery для одного tick route, настройте +`TickSubscriptionRequest::continuity`: + +```cpp +md::TickSubscriptionRequest request("EURUSD"); +request.continuity.mode = md::MarketDataContinuityMode::PREFILL_AND_RECOVER; +request.continuity.prefill_lookback_ms = 60'000; +request.continuity.expected_interval_ms = 1'000; +request.continuity.max_backfill_ms = 60'000; + +auto route = router.subscribe_ticks(provider, bot, request); +``` + +`PREFILL` запрашивает заданный lookback до освобождения live ticks. +`PREFILL_AND_RECOVER` дополнительно удерживает live tail, когда дельта между +последовательными timestamps больше `expected_interval_ms`. Это только +gap-detection hint, а не требование плотной сетки: equal timestamps и live +ticks чаще одной секунды допустимы. Для доказательства continuity Router +доверяет только значению `range_complete` в provider result. + +Router сначала отправляет historical ticks с флагом `HISTORICAL`, затем +воспроизводит удержанные live ticks с флагами `LIVE_SOURCE | CATCHUP`. До +`LIVE` нужен complete result. Неполный result можно доставить как observations, +но Router отправит `FAILED`, затем sticky `DEGRADED`; поздний независимый +успешный запрос не скроет раннюю unresolved boundary. `max_backfill_ms` +ограничивает каждый history request, а `process()` обслуживает retries без +создания отдельного timer thread. + +После reconnect tick continuity публикует `STALE`, ждёт `READY` и запрашивает +unresolved range до последнего observed time. Exact overlap удаляется по +identity `(time_ms, ask, bid, last, volume)`; `received_ms` и flags не делают +полностью одинаковый snapshot новым, но разные observations той же секунды +сохраняются. При переполнении buffer Router освобождает live data, публикует +`FAILED`/`DEGRADED`, отключает continuity для этого route и возобновляет +обычную live delivery. +Если transport прервался во время initial prefill, после `READY` Router +начинает повторный запрос с исходного начала lookback и расширяет его до +текущего времени, поэтому прерванный интервал не пропускается молча. diff --git a/guides/platform-api-guide.md b/guides/platform-api-guide.md index e118a89..e4bef7f 100644 --- a/guides/platform-api-guide.md +++ b/guides/platform-api-guide.md @@ -148,15 +148,19 @@ Subscription rules: slots before emitting `LIVE`. Cached invalidating status replay blocks the same work until a later live `READY`, while a plain or completed `PREFILL` route does not gain outage - recovery. Tick routes do not have this guarantee because the provider contract - still lacks generic tick-history. -- Router continuity is currently bar-first. Intrade additionally exposes a - bounded, session-scoped observed-tick archive populated by `/price_now`. + recovery. Tick routes can use `TickSubscriptionRequest::continuity` with the + same history-first lifecycle when the provider implements + `fetch_tick_history()`. +- Router continuity supports bars and tick routes. Intrade additionally exposes + a bounded, session-scoped observed-tick archive populated by `/price_now`. Its timestamps have one-second broker granularity, and `range_complete=true` means every expected observed second is present in the retained archive, not that every broker micro-event was captured. The archive is non-persistent and starts empty for a new authenticated session. `trade_check2.php` remains a - settlement/trade-result endpoint and is not used for tick history. + settlement/trade-result endpoint and is not used for tick history. Router + tick continuity uses the archive's explicit `range_complete` assertion, + keeps incomplete ranges `DEGRADED`, and preserves distinct same-second + observations while removing only exact overlaps. - `BaseMarketDataProvider` is non-copyable and non-movable so provider identity cannot be duplicated after handles were issued. - Public subscriptions describe consumer routing. Internal platform polling or diff --git a/guides/refactor-backlog.md b/guides/refactor-backlog.md index b664f15..499bcb2 100644 --- a/guides/refactor-backlog.md +++ b/guides/refactor-backlog.md @@ -19,11 +19,15 @@ series. Keep it short and remove items once they are handled. `/price_now`. It preserves distinct same-second snapshots and reports `range_complete` only for proven one-second coverage; it is not a persistent or authoritative broker tick archive. +- Router tick continuity now uses that provider history contract for opt-in + `PREFILL` and `PREFILL_AND_RECOVER` routes. It buffers live ticks, treats + `range_complete` as the authority, retries transport failures through + `process()`, performs bounded recovery, preserves distinct same-second + events, removes only exact overlaps, and keeps incomplete history + `DEGRADED`. ## Next PR Candidates -- Integrate Intrade's observed-tick archive with Router continuity while keeping - incomplete ranges fail-closed and documenting the session/retention limit. - Add an authoritative provider tick-history implementation only if a broker later exposes one; do not treat `trade_check2.php` as a range-history API. - Replace the dense-bar assumption with an explicit provider completeness diff --git a/include/optionx_cpp/market_data/MarketDataContinuity.hpp b/include/optionx_cpp/market_data/MarketDataContinuity.hpp index 079ca4e..bad9964 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuity.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuity.hpp @@ -19,7 +19,7 @@ namespace optionx::market_data { UNKNOWN = 0, 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. + BACKFILLING, ///< Historical market data is being loaded for a gap. RETRYING, ///< A failed history request will be attempted again. LIVE, ///< No known unresolved history range remains for the route. FAILED = 6, ///< A specific history operation failed; the route may continue. @@ -119,7 +119,7 @@ namespace optionx::market_data { MarketDataContinuityStatus status = MarketDataContinuityStatus::UNKNOWN; std::uint64_t from_time_ms = 0; ///< Start of the requested history range, if known. std::uint64_t to_time_ms = 0; ///< End of the requested history range, if known. - std::size_t requested_items = 0; ///< Number of requested bars, when count-based. + std::size_t requested_items = 0; ///< Requested item count when the route uses count-based history; zero for timestamp ranges. std::size_t delivered_items = 0; ///< Number of history items delivered by the operation. std::string message; ///< Optional diagnostic text. }; @@ -135,12 +135,13 @@ namespace optionx::market_data { MarketDataType type = MarketDataType::UNKNOWN; ///< Routed payload type. std::string symbol; ///< Provider symbol. BarTimeframe timeframe = 0; ///< Bar timeframe, or zero for ticks. - bool enabled = false; ///< Whether bar continuity is enabled for this route. + bool enabled = false; ///< Whether configured history continuity is enabled for this route. MarketDataContinuityStatus last_status = MarketDataContinuityStatus::UNKNOWN; MarketDataContinuityPhase phase = MarketDataContinuityPhase::UNKNOWN; MarketDataContinuityOperation last_operation = MarketDataContinuityOperation::NONE; bool request_in_flight = false; ///< True while a history request is outstanding. - std::uint64_t last_observed_time_ms = 0; ///< Latest delivered or buffered bar timestamp; zero for tick routes. + std::uint64_t last_observed_time_ms = 0; ///< Latest delivered or buffered payload timestamp. + std::uint64_t expected_interval_ms = 0; ///< Tick gap-detection interval; zero for bar routes. std::uint64_t requested_from_time_ms = 0; ///< Last history range start. std::uint64_t requested_to_time_ms = 0; ///< Last history range end. std::size_t requested_items = 0; ///< Last history operation item count. diff --git a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp index 1ef35b0..53cf3fb 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp @@ -3,7 +3,7 @@ #define OPTIONX_HEADER_MARKET_DATA_MARKET_DATA_CONTINUITY_OPTIONS_HPP_INCLUDED /// \file MarketDataContinuityOptions.hpp -/// \brief Defines history prefill and gap-recovery options for bar routes. +/// \brief Defines history prefill and gap-recovery options for market-data routes. #include #include @@ -11,7 +11,7 @@ namespace optionx::market_data { /// \enum MarketDataContinuityMode - /// \brief Selects how a routed bar stream is initialized and recovered. + /// \brief Selects how a routed market-data stream is initialized and recovered. enum class MarketDataContinuityMode { LIVE_ONLY = 0, ///< Deliver live provider payloads immediately. PREFILL, ///< Deliver startup history without later gap recovery. @@ -68,6 +68,47 @@ namespace optionx::market_data { } }; + /// \struct MarketDataTickContinuityOptions + /// \brief Configures history prefill and recovery for a tick route. + /// + /// Tick streams do not have a dense timeframe grid. `expected_interval_ms` + /// is therefore only a hint for detecting a suspicious live gap; a history + /// result's `range_complete` value remains the authority for recovery. + struct MarketDataTickContinuityOptions { + MarketDataContinuityMode mode = MarketDataContinuityMode::LIVE_ONLY; + std::uint64_t prefill_lookback_ms = 0; + std::uint64_t expected_interval_ms = 1000; + std::uint64_t max_backfill_ms = 60000; ///< Maximum time span per history request; zero is unbounded. + MarketDataContinuityRetryPolicy retry; + std::size_t max_buffered_batches = 1024; + std::size_t max_buffered_items = 100000; + + /// \brief Returns true when the option combination is usable. + [[nodiscard]] bool valid() const noexcept { + if (!retry.valid()) return false; + if (mode == MarketDataContinuityMode::LIVE_ONLY) { + return prefill_lookback_ms == 0; + } + if (mode == MarketDataContinuityMode::PREFILL) { + return prefill_lookback_ms > 0; + } + if (mode == MarketDataContinuityMode::PREFILL_AND_RECOVER) { + return expected_interval_ms > 0; + } + return false; + } + + /// \brief Returns true when history work is enabled. + [[nodiscard]] bool enabled() const noexcept { + return mode != MarketDataContinuityMode::LIVE_ONLY; + } + + /// \brief Returns true when live gaps and reconnects are recovered. + [[nodiscard]] bool recovers_gaps() const noexcept { + return mode == MarketDataContinuityMode::PREFILL_AND_RECOVER; + } + }; + } // namespace optionx::market_data #endif // OPTIONX_HEADER_MARKET_DATA_MARKET_DATA_CONTINUITY_OPTIONS_HPP_INCLUDED diff --git a/include/optionx_cpp/market_data/MarketDataSubscription.hpp b/include/optionx_cpp/market_data/MarketDataSubscription.hpp index 4325606..13a5745 100644 --- a/include/optionx_cpp/market_data/MarketDataSubscription.hpp +++ b/include/optionx_cpp/market_data/MarketDataSubscription.hpp @@ -30,6 +30,7 @@ namespace optionx::market_data { struct TickSubscriptionRequest { std::string symbol; ///< Broker/provider symbol. MarketDataTransport transport = MarketDataTransport::AUTO; ///< Preferred transport. + MarketDataTickContinuityOptions continuity; ///< Optional tick history prefill and recovery. /// \brief Default constructor. TickSubscriptionRequest() = default; diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index 8b62592..2838d5d 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -64,6 +64,18 @@ namespace optionx::market_data { std::size_t attempt = 1; }; + struct PendingTickContinuityRequest { + RoutedSubscriptionId router_id; + MarketDataSubscriptionHandle subscription; + TickHistoryRequest 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 ContinuityState { MarketDataContinuityPhase phase = MarketDataContinuityPhase::LIVE; MarketDataContinuityStatus last_status = @@ -105,6 +117,47 @@ namespace optionx::market_data { } }; + struct TickContinuityState { + MarketDataContinuityPhase phase = MarketDataContinuityPhase::LIVE; + MarketDataContinuityStatus last_status = + MarketDataContinuityStatus::UNKNOWN; + MarketDataContinuityOperation last_operation = + MarketDataContinuityOperation::NONE; + bool initial_prefill_pending = false; + bool request_in_flight = false; + std::deque buffer; + std::size_t buffered_items = 0; + std::uint64_t last_observed_time_ms = 0; + std::uint64_t requested_from_time_ms = 0; + std::uint64_t requested_to_time_ms = 0; + std::size_t requested_items = 0; + std::uint64_t confirmed_from_time_ms = 0; + std::uint64_t confirmed_through_time_ms = 0; + std::size_t confirmed_items = 0; + std::uint64_t verified_through_time_ms = 0; + std::uint64_t unverified_from_time_ms = 0; + std::uint64_t initial_prefill_boundary_time_ms = 0; + std::uint64_t generation = 0; + std::uint64_t reconnect_target_time_ms = 0; + std::size_t history_request_count = 0; + std::size_t retry_count = 0; + std::size_t failure_count = 0; + std::string last_failure; + std::chrono::steady_clock::duration stale_duration{}; + std::chrono::steady_clock::duration degraded_duration{}; + std::chrono::steady_clock::time_point stale_since{}; + std::chrono::steady_clock::time_point degraded_since{}; + std::optional retry_request; + std::chrono::steady_clock::time_point retry_at; + + [[nodiscard]] bool buffers_live_data() const noexcept { + return phase == MarketDataContinuityPhase::PREFILLING || + phase == MarketDataContinuityPhase::WAITING_FOR_READY || + phase == MarketDataContinuityPhase::RECOVERING || + phase == MarketDataContinuityPhase::FLUSHING; + } + }; + struct Entry { RoutedSubscriptionId router_id; ProviderInstanceId provider_id = kInvalidProviderInstanceId; @@ -114,6 +167,8 @@ namespace optionx::market_data { StreamDescriptor stream; MarketDataContinuityOptions continuity; ContinuityState continuity_state; + MarketDataTickContinuityOptions tick_continuity; + TickContinuityState tick_continuity_state; MarketDataSubscriptionHandle retained_cleanup_subscription; MarketDataSubscriptionResult unsubscribe_completion; bool subscribe_completion_posted = false; @@ -248,13 +303,6 @@ namespace optionx::market_data { ProviderInstanceId provider_id, MarketDataStatusUpdate update); - template - void route_batch( - ProviderInstanceId provider_id, - std::unique_ptr batch, - MatchesStream matches_stream, - Deliver deliver); - private: mutable std::mutex m_mutex; std::unordered_map< @@ -303,6 +351,7 @@ namespace optionx::market_data { std::weak_ptr subscriber, StreamDescriptor stream, MarketDataContinuityOptions continuity, + MarketDataTickContinuityOptions tick_continuity, MarketDataProviderId registered_provider_id, std::string& error_message); @@ -322,6 +371,16 @@ namespace optionx::market_data { return {}; } + static MarketDataTickContinuityOptions tick_continuity_from( + const TickSubscriptionRequest& request) noexcept { + return request.continuity; + } + + static MarketDataTickContinuityOptions tick_continuity_from( + const BarSubscriptionRequest&) noexcept { + return {}; + } + static MarketDataContinuityOptions continuity_from( const BarSubscriptionRequest& request) noexcept { return request.continuity; @@ -330,7 +389,11 @@ namespace optionx::market_data { void start_continuity( RoutedSubscriptionId router_id, MarketDataSubscriptionHandle subscription); + void start_tick_continuity( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription); void start_reconnect_recovery(RoutedSubscriptionId router_id); + void start_tick_reconnect_recovery(RoutedSubscriptionId router_id); void request_continuity_history( RoutedSubscriptionId router_id, MarketDataSubscriptionHandle subscription, @@ -353,15 +416,43 @@ namespace optionx::market_data { std::size_t attempt, std::uint64_t generation, BarHistoryResult result); + void request_tick_continuity_history( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription, + TickHistoryRequest request, + ContinuityRequestKind kind, + bool announce_gap, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::size_t requested_items, + std::size_t attempt, + std::uint64_t generation = 0); + void complete_tick_continuity( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription, + TickHistoryRequest request, + ContinuityRequestKind kind, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::size_t requested_items, + std::size_t attempt, + std::uint64_t generation, + TickHistoryResult result); void notify_continuity( RoutedSubscriptionId router_id, MarketDataContinuityUpdate update); static void record_continuity_update_no_lock( const std::shared_ptr& entry, const MarketDataContinuityUpdate& update); + static void record_tick_continuity_update_no_lock( + const std::shared_ptr& entry, + const MarketDataContinuityUpdate& update); void record_continuity_update( RoutedSubscriptionId router_id, const MarketDataContinuityUpdate& update); + void notify_tick_continuity( + RoutedSubscriptionId router_id, + MarketDataContinuityUpdate update); static MarketDataContinuityUpdate make_continuity_update( const MarketDataSubscriptionHandle& subscription, MarketDataContinuityStatus status, @@ -383,6 +474,19 @@ namespace optionx::market_data { bool process_buffered = false, bool allow_gap_recovery = true, std::uint64_t confirmed_through_time_ms = 0); + bool route_tick_to_entry_no_lock( + const std::shared_ptr& entry, + const TickDataBatch& batch, + std::vector, + TickDataBatch>>& deliveries, + std::vector& continuity_requests, + std::vector, + MarketDataContinuityUpdate>>& continuity_deliveries, + bool process_buffered = false, + bool allow_gap_recovery = true, + std::uint64_t confirmed_through_time_ms = 0); static bool buffer_continuity_batch_no_lock( const std::shared_ptr& entry, std::shared_ptr subscriber, @@ -390,6 +494,17 @@ namespace optionx::market_data { std::vector, BarDataBatch>>& deliveries, + std::vector, + MarketDataContinuityUpdate>>& continuity_deliveries, + bool push_front); + static bool buffer_tick_continuity_batch_no_lock( + const std::shared_ptr& entry, + std::shared_ptr subscriber, + TickDataBatch batch, + std::vector, + TickDataBatch>>& deliveries, std::vector, MarketDataContinuityUpdate>>& continuity_deliveries, @@ -399,6 +514,16 @@ namespace optionx::market_data { std::uint64_t from_time_ms, std::uint64_t to_time_ms, std::uint64_t timeframe_ms) noexcept; + static bool tick_history_covers_range( + const TickDataBatch& batch, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::uint64_t expected_interval_ms) noexcept; + static bool same_tick_observation( + const Tick& lhs, + const Tick& rhs) noexcept; + static void deduplicate_tick_items( + std::vector& ticks); static void clip_history_to_range( BarDataBatch& batch, std::uint64_t from_time_ms, @@ -406,6 +531,9 @@ namespace optionx::market_data { static void mark_unverified_no_lock( const std::shared_ptr& entry, std::uint64_t from_time_ms) noexcept; + static void mark_tick_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, @@ -415,6 +543,13 @@ namespace optionx::market_data { static void record_bar_progress_no_lock( const std::shared_ptr& entry, const std::vector& bars) noexcept; + static void record_tick_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_tick_progress_no_lock( + const std::shared_ptr& entry, + const std::vector& ticks) noexcept; static std::uint64_t duration_to_milliseconds( std::chrono::steady_clock::duration duration) noexcept; static MarketDataContinuitySnapshot make_continuity_snapshot_no_lock( @@ -430,7 +565,9 @@ namespace optionx::market_data { std::shared_ptr, MarketDataContinuityUpdate>>& continuity_deliveries, std::vector& prefill_routes, - std::vector& reconnect_routes); + std::vector& reconnect_routes, + std::vector& tick_prefill_routes, + std::vector& tick_reconnect_routes); void fail_pending_subscribe( RoutedSubscriptionId router_id, @@ -936,6 +1073,7 @@ namespace optionx::market_data { std::weak_ptr subscriber, StreamDescriptor stream, MarketDataContinuityOptions continuity, + MarketDataTickContinuityOptions tick_continuity, MarketDataProviderId registered_provider_id, std::string& error_message) { if (subscriber.expired()) { @@ -991,12 +1129,20 @@ namespace optionx::market_data { entry->control = control; entry->stream = std::move(stream); entry->continuity = std::move(continuity); + entry->tick_continuity = std::move(tick_continuity); entry->continuity_state.initial_prefill_pending = entry->continuity.enabled() && entry->continuity.prefill_bars > 0; entry->continuity_state.phase = entry->continuity_state.initial_prefill_pending ? MarketDataContinuityPhase::PREFILLING : MarketDataContinuityPhase::LIVE; + entry->tick_continuity_state.initial_prefill_pending = + entry->tick_continuity.enabled() && + entry->tick_continuity.prefill_lookback_ms > 0; + entry->tick_continuity_state.phase = + entry->tick_continuity_state.initial_prefill_pending + ? MarketDataContinuityPhase::PREFILLING + : MarketDataContinuityPhase::LIVE; m_entries.emplace(router_id, entry); ++provider_it->second.route_count; } @@ -1015,7 +1161,9 @@ namespace optionx::market_data { const char* invalid_request_message, const char* operation_name, const char* not_accepted_message) { - if (!request.valid()) { + if (!request.valid() || + !continuity_from(request).valid() || + !tick_continuity_from(request).valid()) { dispatch_result( std::move(callback), MarketDataSubscriptionResult::failed( @@ -1032,6 +1180,7 @@ namespace optionx::market_data { std::move(subscriber), stream_from(request), continuity_from(request), + tick_continuity_from(request), registered_provider_id, error_message); if (!control) { @@ -1263,6 +1412,7 @@ namespace optionx::market_data { bool has_replay = false; bool release_requested = false; bool needs_continuity_prefill = false; + bool needs_tick_continuity_prefill = false; subscription_callback_t release_callback; BaseMarketDataProvider* unbind = nullptr; std::vector> replay_continuity_deliveries; std::vector replay_prefill_routes; std::vector replay_reconnect_routes; + std::vector replay_tick_prefill_routes; + std::vector replay_tick_reconnect_routes; if (result.success() && (!result.subscription.valid() || @@ -1313,6 +1465,9 @@ namespace optionx::market_data { needs_continuity_prefill = entry->stream.type == MarketDataType::BARS && entry->continuity.prefill_bars > 0; + needs_tick_continuity_prefill = + entry->stream.type == MarketDataType::TICKS && + entry->tick_continuity.prefill_lookback_ms > 0; release_requested = entry->release_requested; release_callback = std::move(entry->release_callback); @@ -1330,7 +1485,9 @@ namespace optionx::market_data { replay, replay_continuity_deliveries, replay_prefill_routes, - replay_reconnect_routes); + replay_reconnect_routes, + replay_tick_prefill_routes, + replay_tick_reconnect_routes); } } } @@ -1349,6 +1506,9 @@ namespace optionx::market_data { if (needs_continuity_prefill && result.success()) { start_continuity(router_id, result.subscription); } + if (needs_tick_continuity_prefill && result.success()) { + start_tick_continuity(router_id, result.subscription); + } if (has_replay && subscriber && result.success()) { bool still_active = false; { @@ -1373,6 +1533,12 @@ namespace optionx::market_data { for (const auto replay_router_id : replay_reconnect_routes) { start_reconnect_recovery(replay_router_id); } + for (const auto replay_router_id : replay_tick_prefill_routes) { + start_tick_continuity(replay_router_id, result.subscription); + } + for (const auto replay_router_id : replay_tick_reconnect_routes) { + start_tick_reconnect_recovery(replay_router_id); + } } inline void MarketDataRouterState::start_continuity( @@ -1445,7 +1611,79 @@ namespace optionx::market_data { pending.from_time_ms, pending.to_time_ms, pending.requested_items, - 1); + 1); + } + + inline void MarketDataRouterState::start_tick_continuity( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription) { + PendingTickContinuityRequest pending; + { + 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 || + entry_it->second->stream.type != MarketDataType::TICKS || + !entry_it->second->tick_continuity_state.initial_prefill_pending || + entry_it->second->tick_continuity_state.phase != + MarketDataContinuityPhase::PREFILLING || + entry_it->second->tick_continuity.prefill_lookback_ms == 0 || + entry_it->second->tick_continuity_state.request_in_flight) { + return; + } + + const auto& entry = entry_it->second; + const auto now_ms = static_cast(OPTIONX_TIMESTAMP_MS); + const auto lookback = entry->tick_continuity.prefill_lookback_ms; + pending.router_id = router_id; + pending.subscription = subscription; + auto& continuity = entry->tick_continuity_state; + if (continuity.initial_prefill_boundary_time_ms == 0) { + continuity.initial_prefill_boundary_time_ms = + now_ms > 0 ? now_ms : 1U; + } + const auto original_boundary_time_ms = + continuity.initial_prefill_boundary_time_ms; + const auto current_time_ms = now_ms > 0 ? now_ms : 1U; + const auto current_boundary_time_ms = std::max( + original_boundary_time_ms, + current_time_ms); + pending.from_time_ms = original_boundary_time_ms > lookback + ? original_boundary_time_ms - lookback + : 1U; + pending.to_time_ms = current_boundary_time_ms; + pending.request = TickHistoryRequest( + entry->stream.symbol, + pending.from_time_ms, + pending.to_time_ms); + pending.kind = ContinuityRequestKind::PREFILL; + pending.requested_items = 0; + pending.attempt = 1; + continuity.phase = + MarketDataContinuityPhase::PREFILLING; + continuity.request_in_flight = true; + } + + notify_tick_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::PREFILLING, + pending.from_time_ms, + pending.to_time_ms, + 0, + 0, + "Requesting observed tick history prefill.")); + request_tick_continuity_history( + pending.router_id, + std::move(pending.subscription), + std::move(pending.request), + pending.kind, + pending.announce_gap, + pending.from_time_ms, + pending.to_time_ms, + pending.requested_items, + pending.attempt); } inline void MarketDataRouterState::start_reconnect_recovery( @@ -1594,6 +1832,137 @@ namespace optionx::market_data { } request_continuity_history( + pending.router_id, + std::move(pending.subscription), + std::move(pending.request), + pending.kind, + pending.announce_gap, + pending.from_time_ms, + pending.to_time_ms, + pending.requested_items, + pending.attempt); + } + + inline void MarketDataRouterState::start_tick_reconnect_recovery( + RoutedSubscriptionId router_id) { + PendingTickContinuityRequest pending; + MarketDataSubscriptionHandle live_subscription; + bool notify_live = false; + bool notify_degraded = false; + + { + 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 || + entry_it->second->stream.type != MarketDataType::TICKS) { + return; + } + + const auto& entry = entry_it->second; + auto& continuity = entry->tick_continuity_state; + if (continuity.phase != MarketDataContinuityPhase::RECOVERING || + !entry->tick_continuity.recovers_gaps() || + continuity.initial_prefill_pending || + continuity.request_in_flight) { + return; + } + + std::uint64_t earliest_buffered_time_ms = 0; + std::uint64_t latest_buffered_time_ms = 0; + for (const auto& batch : continuity.buffer) { + for (const auto& tick : batch.items) { + if (tick.time_ms == 0) continue; + if (earliest_buffered_time_ms == 0 || + tick.time_ms < earliest_buffered_time_ms) { + earliest_buffered_time_ms = tick.time_ms; + } + latest_buffered_time_ms = std::max( + latest_buffered_time_ms, + tick.time_ms); + } + } + + const auto observed_time_ms = std::max( + continuity.last_observed_time_ms, + latest_buffered_time_ms); + const auto now_ms = static_cast(OPTIONX_TIMESTAMP_MS); + const auto target_time_ms = std::max(observed_time_ms, now_ms); + const auto interval = entry->tick_continuity.expected_interval_ms; + + auto from_time_ms = continuity.unverified_from_time_ms; + if (from_time_ms == 0 && continuity.verified_through_time_ms > 0 && + interval > 0 && + continuity.verified_through_time_ms <= + std::numeric_limits::max() - interval) { + from_time_ms = continuity.verified_through_time_ms + interval; + } else if (from_time_ms == 0 && earliest_buffered_time_ms > 0) { + from_time_ms = earliest_buffered_time_ms; + } else if (from_time_ms == 0 && observed_time_ms > 0 && interval > 0 && + observed_time_ms <= + std::numeric_limits::max() - interval) { + from_time_ms = observed_time_ms + interval; + } + + if (from_time_ms == 0 || from_time_ms > target_time_ms) { + continuity.phase = continuity.unverified_from_time_ms == 0 + ? MarketDataContinuityPhase::LIVE + : MarketDataContinuityPhase::DEGRADED; + continuity.reconnect_target_time_ms = 0; + live_subscription = entry->control->provider_subscription; + notify_live = continuity.unverified_from_time_ms == 0; + notify_degraded = !notify_live; + } else { + auto request_to_time_ms = target_time_ms; + const auto max_backfill_ms = entry->tick_continuity.max_backfill_ms; + if (max_backfill_ms > 0) { + const auto span = max_backfill_ms > interval + ? max_backfill_ms - interval + : 0U; + const auto bounded = from_time_ms > + std::numeric_limits::max() - span + ? std::numeric_limits::max() + : from_time_ms + span; + request_to_time_ms = std::min(request_to_time_ms, bounded); + } + if (request_to_time_ms < from_time_ms) return; + + continuity.reconnect_target_time_ms = target_time_ms; + continuity.phase = MarketDataContinuityPhase::RECOVERING; + continuity.request_in_flight = true; + pending.router_id = router_id; + pending.subscription = entry->control->provider_subscription; + pending.request = TickHistoryRequest( + entry->stream.symbol, + from_time_ms, + request_to_time_ms); + pending.kind = ContinuityRequestKind::RECONNECT_BACKFILL; + pending.from_time_ms = from_time_ms; + pending.to_time_ms = request_to_time_ms; + pending.requested_items = 0; + pending.attempt = 1; + } + } + + if (notify_live || notify_degraded) { + notify_tick_continuity( + router_id, + make_continuity_update( + live_subscription, + notify_live + ? MarketDataContinuityStatus::LIVE + : MarketDataContinuityStatus::DEGRADED, + 0, + 0, + 0, + 0, + notify_live + ? "No observed ticks required reconnect recovery." + : "Reconnect recovery could not prove the unresolved tick range.")); + return; + } + + request_tick_continuity_history( pending.router_id, std::move(pending.subscription), std::move(pending.request), @@ -1752,27 +2121,171 @@ namespace optionx::market_data { } } - inline bool MarketDataRouterState::history_covers_range( - const BarDataBatch& batch, + inline void MarketDataRouterState::request_tick_continuity_history( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription, + TickHistoryRequest request, + ContinuityRequestKind kind, + bool announce_gap, 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; + std::size_t requested_items, + std::size_t attempt, + std::uint64_t generation) { + BaseMarketDataProvider* provider = nullptr; + auto operation = std::make_shared(); + { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (m_shutdown || + entry_it == m_entries.end() || + entry_it->second->phase != EntryPhase::ACTIVE || + entry_it->second->stream.type != MarketDataType::TICKS || + !entry_it->second->tick_continuity_state.request_in_flight) { + return; } - expected_time_ms += timeframe_ms; + if (generation != 0 && + generation != entry_it->second->tick_continuity_state.generation) { + return; + } + generation = entry_it->second->tick_continuity_state.generation; + provider = entry_it->second->provider; + auto& continuity = entry_it->second->tick_continuity_state; + switch (kind) { + case ContinuityRequestKind::PREFILL: + continuity.last_operation = MarketDataContinuityOperation::PREFILL; + break; + case ContinuityRequestKind::GAP_BACKFILL: + continuity.last_operation = MarketDataContinuityOperation::GAP_BACKFILL; + break; + case ContinuityRequestKind::RECONNECT_BACKFILL: + continuity.last_operation = MarketDataContinuityOperation::RECONNECT_BACKFILL; + break; + } + continuity.requested_from_time_ms = from_time_ms; + continuity.requested_to_time_ms = to_time_ms; + continuity.requested_items = requested_items; + ++continuity.history_request_count; + if (attempt > 1) ++continuity.retry_count; + if (provider) ++m_continuity_operations_in_flight; + } + if (!provider) return; + + if (kind == ContinuityRequestKind::GAP_BACKFILL) { + if (announce_gap) { + notify_tick_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 tick stream.")); + } + notify_tick_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::BACKFILLING, + from_time_ms, + to_time_ms, + requested_items, + 0, + "Requesting historical ticks for the detected gap.")); + } else if (kind == ContinuityRequestKind::RECONNECT_BACKFILL) { + notify_tick_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::BACKFILLING, + from_time_ms, + to_time_ms, + requested_items, + 0, + "Revalidating observed ticks after reconnect.")); + } + + const auto state = shared_from_this(); + auto complete_on_owner = [state, + operation, + router_id, + subscription, + request, + kind, + from_time_ms, + to_time_ms, + requested_items, + attempt, + generation](TickHistoryResult result) mutable { + if (!state->record_continuity_completion(operation)) return; + auto task = [state, + router_id, + subscription, + request, + kind, + from_time_ms, + to_time_ms, + requested_items, + attempt, + generation, + result = std::move(result)]() mutable { + state->complete_tick_continuity( + router_id, + subscription, + request, + kind, + from_time_ms, + to_time_ms, + requested_items, + attempt, + generation, + std::move(result)); + }; + state->dispatch_or_run(std::move(task)); + }; + try { + const bool accepted = provider->fetch_tick_history( + request, + [complete_on_owner](TickHistoryResult result) mutable { + complete_on_owner(std::move(result)); + }); + if (!accepted) { + complete_on_owner(TickHistoryResult::fail( + "Market-data provider did not accept the tick history request.")); + } + } catch (const std::exception& exception) { + complete_on_owner(TickHistoryResult::fail( + std::string("Market-data provider tick history request threw: ") + + exception.what())); + } catch (...) { + complete_on_owner(TickHistoryResult::fail( + "Market-data provider tick history request threw.")); + } + } + + 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; } @@ -1803,6 +2316,16 @@ namespace optionx::market_data { } } + inline void MarketDataRouterState::mark_tick_unverified_no_lock( + const std::shared_ptr& entry, + std::uint64_t from_time_ms) noexcept { + if (!entry || from_time_ms == 0) return; + auto& unverified = entry->tick_continuity_state.unverified_from_time_ms; + if (unverified == 0 || from_time_ms < unverified) { + unverified = from_time_ms; + } + } + inline void MarketDataRouterState::record_verified_range_no_lock( const std::shared_ptr& entry, std::uint64_t from_time_ms, @@ -1857,6 +2380,83 @@ namespace optionx::market_data { } } + inline bool MarketDataRouterState::same_tick_observation( + const Tick& lhs, + const Tick& rhs) noexcept { + return lhs.time_ms == rhs.time_ms && + lhs.ask == rhs.ask && + lhs.bid == rhs.bid && + lhs.last == rhs.last && + lhs.volume == rhs.volume; + } + + inline void MarketDataRouterState::deduplicate_tick_items( + std::vector& ticks) { + std::vector unique; + unique.reserve(ticks.size()); + for (auto& tick : ticks) { + const auto duplicate = std::any_of( + unique.begin(), + unique.end(), + [&tick](const Tick& existing) { + return same_tick_observation(existing, tick); + }); + if (!duplicate) unique.push_back(std::move(tick)); + } + ticks = std::move(unique); + } + + inline void MarketDataRouterState::record_tick_progress_no_lock( + const std::shared_ptr& entry, + const std::vector& ticks) noexcept { + if (!entry) return; + auto& continuity = entry->tick_continuity_state; + for (const auto& tick : ticks) { + if (tick.time_ms > continuity.last_observed_time_ms) { + continuity.last_observed_time_ms = tick.time_ms; + } + if (continuity.unverified_from_time_ms == 0 && + tick.time_ms > continuity.verified_through_time_ms) { + continuity.verified_through_time_ms = tick.time_ms; + } + } + } + + inline void MarketDataRouterState::record_tick_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; + auto& continuity = entry->tick_continuity_state; + if (continuity.unverified_from_time_ms != 0 && + (from_time_ms > continuity.unverified_from_time_ms || + to_time_ms < continuity.unverified_from_time_ms)) { + return; + } + continuity.verified_through_time_ms = std::max( + continuity.verified_through_time_ms, + to_time_ms); + continuity.unverified_from_time_ms = 0; + } + + inline bool MarketDataRouterState::tick_history_covers_range( + const TickDataBatch& batch, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::uint64_t expected_interval_ms) noexcept { + (void)expected_interval_ms; + if (from_time_ms == 0 || to_time_ms < from_time_ms) return false; + std::uint64_t previous_time_ms = 0; + for (const auto& tick : batch.items) { + if (tick.time_ms < from_time_ms || tick.time_ms > to_time_ms || + (previous_time_ms != 0 && tick.time_ms < previous_time_ms)) { + return false; + } + previous_time_ms = tick.time_ms; + } + return true; + } + inline std::chrono::steady_clock::duration MarketDataRouterState::continuity_retry_delay( const MarketDataContinuityRetryPolicy& policy, @@ -2329,6 +2929,401 @@ namespace optionx::market_data { } } + inline void MarketDataRouterState::complete_tick_continuity( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription, + TickHistoryRequest request, + ContinuityRequestKind kind, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::size_t requested_items, + std::size_t attempt, + std::uint64_t generation, + TickHistoryResult result) { + StreamDescriptor expected_stream; + { + 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 || + entry_it->second->stream.type != MarketDataType::TICKS || + !entry_it->second->tick_continuity_state.request_in_flight || + entry_it->second->tick_continuity_state.generation != generation) { + return; + } + expected_stream = entry_it->second->stream; + entry_it->second->tick_continuity_state.request_in_flight = false; + entry_it->second->tick_continuity_state.phase = + MarketDataContinuityPhase::FLUSHING; + } + + const bool history_success = static_cast(result); + bool history_stream_matches = false; + bool history_range_valid = false; + TickDataBatch history_batch; + std::size_t delivered_history_items = 0; + if (history_success) { + const auto& sequence = result.sequence; + history_stream_matches = sequence.symbol.empty() || + sequence.symbol == expected_stream.symbol; + if (history_stream_matches) { + history_range_valid = true; + std::uint64_t previous_time_ms = 0; + for (const auto& tick : sequence.ticks) { + if (tick.time_ms < from_time_ms || + tick.time_ms > to_time_ms || + (previous_time_ms != 0 && + tick.time_ms < previous_time_ms)) { + history_range_valid = false; + break; + } + previous_time_ms = tick.time_ms; + } + } + if (history_range_valid) { + history_batch = *MarketDataContinuityService::make_tick_batch( + std::move(result.sequence), + request, + subscription, + kind != ContinuityRequestKind::PREFILL); + deduplicate_tick_items(history_batch.items); + delivered_history_items = history_batch.items.size(); + } + } + + const bool history_covers_range = history_range_valid && + tick_history_covers_range( + history_batch, + from_time_ms, + to_time_ms, + 0); + const bool usable_history = history_success && + history_stream_matches && + history_covers_range && + result.range_complete; + + bool retry_scheduled = false; + if (!history_success) { + PendingTickContinuityRequest 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 || + entry_it->second->tick_continuity_state.generation != generation) { + return; + } + + const auto& entry = entry_it->second; + if (attempt < entry->tick_continuity.retry.max_attempts) { + retry_request.router_id = router_id; + retry_request.subscription = subscription; + retry_request.request = request; + retry_request.kind = kind; + 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->tick_continuity.retry, + attempt); + const auto remaining = + std::chrono::steady_clock::time_point::max() - now; + entry->tick_continuity_state.retry_request = + std::move(retry_request); + entry->tick_continuity_state.retry_at = delay >= remaining + ? std::chrono::steady_clock::time_point::max() + : now + delay; + entry->tick_continuity_state.phase = + kind == ContinuityRequestKind::PREFILL + ? MarketDataContinuityPhase::PREFILLING + : MarketDataContinuityPhase::RECOVERING; + retry_scheduled = true; + } + } + + if (retry_scheduled) { + notify_tick_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::RETRYING, + from_time_ms, + to_time_ms, + requested_items, + 0, + "Retrying historical tick continuity request.")); + return; + } + } + + { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (entry_it == m_entries.end() || + entry_it->second->phase != EntryPhase::ACTIVE || + entry_it->second->tick_continuity_state.generation != generation) { + return; + } + + const auto& entry = entry_it->second; + auto& continuity = entry->tick_continuity_state; + if (kind == ContinuityRequestKind::PREFILL) { + continuity.initial_prefill_pending = false; + } + if (!usable_history) { + mark_tick_unverified_no_lock(entry, from_time_ms); + ++continuity.failure_count; + } else { + if (delivered_history_items > 0) { + std::uint64_t confirmed_from = 0; + std::uint64_t confirmed_to = 0; + for (const auto& tick : history_batch.items) { + if (confirmed_from == 0 || tick.time_ms < confirmed_from) { + confirmed_from = tick.time_ms; + } + confirmed_to = std::max(confirmed_to, tick.time_ms); + } + continuity.confirmed_from_time_ms = confirmed_from; + continuity.confirmed_through_time_ms = confirmed_to; + continuity.confirmed_items = delivered_history_items; + } + record_tick_verified_range_no_lock( + entry, + from_time_ms, + to_time_ms); + } + } + + if (!usable_history) { + notify_tick_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::FAILED, + from_time_ms, + to_time_ms, + requested_items, + delivered_history_items, + result.error_desc.empty() + ? (!history_success + ? "Historical tick continuity request failed." + : !history_stream_matches + ? "Historical tick response does not match the subscribed symbol." + : !history_range_valid + ? "Historical tick response is outside the requested range or unordered." + : "Historical tick response does not prove the requested range is complete.") + : result.error_desc)); + } + + if (!history_batch.items.empty()) { + std::shared_ptr subscriber; + bool active = false; + { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + active = entry_it != m_entries.end() && + entry_it->second->phase == EntryPhase::ACTIVE && + entry_it->second->tick_continuity_state.generation == generation && + !entry_it->second->release_requested; + if (active) { + auto& continuity = entry_it->second->tick_continuity_state; + for (auto batch_it = continuity.buffer.begin(); + batch_it != continuity.buffer.end();) { + auto& items = batch_it->items; + items.erase( + std::remove_if( + items.begin(), + items.end(), + [&history_batch](const Tick& tick) { + return std::any_of( + history_batch.items.begin(), + history_batch.items.end(), + [&tick](const Tick& history_tick) { + return same_tick_observation( + history_tick, + tick); + }); + }), + items.end()); + if (items.empty()) { + batch_it = continuity.buffer.erase(batch_it); + } else { + ++batch_it; + } + } + continuity.buffered_items = 0; + for (const auto& batch : continuity.buffer) { + continuity.buffered_items += batch.items.size(); + } + record_tick_progress_no_lock( + entry_it->second, + history_batch.items); + subscriber = entry_it->second->subscriber.lock(); + } + } + if (active && subscriber) subscriber->on_tick_data(history_batch); + } + + bool schedule_next = false; + PendingTickContinuityRequest next_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 || + entry_it->second->tick_continuity_state.generation != generation) { + return; + } + const auto& entry = entry_it->second; + auto& continuity = entry->tick_continuity_state; + if (usable_history && + kind != ContinuityRequestKind::PREFILL && + continuity.reconnect_target_time_ms > to_time_ms) { + const auto interval = entry->tick_continuity.expected_interval_ms; + const auto next_from = interval > 0 && + to_time_ms <= + std::numeric_limits::max() - interval + ? to_time_ms + interval + : to_time_ms == std::numeric_limits::max() + ? to_time_ms + : to_time_ms + 1U; + if (next_from <= to_time_ms) return; + auto next_to = continuity.reconnect_target_time_ms; + const auto max_backfill_ms = entry->tick_continuity.max_backfill_ms; + if (max_backfill_ms > 0) { + const auto span = max_backfill_ms > interval + ? max_backfill_ms - interval + : 0U; + const auto bounded = next_from > + std::numeric_limits::max() - span + ? std::numeric_limits::max() + : next_from + span; + next_to = std::min(next_to, bounded); + } + if (next_to < next_from) return; + next_request.router_id = router_id; + next_request.subscription = entry->control->provider_subscription; + next_request.request = TickHistoryRequest( + entry->stream.symbol, + next_from, + next_to); + next_request.kind = kind; + next_request.from_time_ms = next_from; + next_request.to_time_ms = next_to; + next_request.requested_items = 0; + next_request.attempt = 1; + continuity.phase = MarketDataContinuityPhase::RECOVERING; + continuity.request_in_flight = true; + schedule_next = true; + } + } + if (schedule_next) { + request_tick_continuity_history( + next_request.router_id, + std::move(next_request.subscription), + std::move(next_request.request), + next_request.kind, + false, + next_request.from_time_ms, + next_request.to_time_ms, + next_request.requested_items, + next_request.attempt); + return; + } + + for (;;) { + std::vector, + TickDataBatch>> deliveries; + std::vector, + MarketDataContinuityUpdate>> continuity_deliveries; + std::vector continuity_requests; + bool finished = false; + bool continuity_still_enabled = false; + bool continuity_verified = false; + { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (entry_it == m_entries.end() || + entry_it->second->phase != EntryPhase::ACTIVE || + entry_it->second->tick_continuity_state.generation != generation) { + return; + } + auto& continuity = entry_it->second->tick_continuity_state; + if (continuity.buffer.empty()) { + continuity_still_enabled = entry_it->second->tick_continuity.enabled(); + continuity_verified = usable_history && + continuity.unverified_from_time_ms == 0; + continuity.phase = continuity_verified + ? MarketDataContinuityPhase::LIVE + : MarketDataContinuityPhase::DEGRADED; + continuity.reconnect_target_time_ms = 0; + finished = true; + } else { + auto batch = std::move(continuity.buffer.front()); + continuity.buffer.pop_front(); + continuity.buffered_items = continuity.buffered_items >= batch.items.size() + ? continuity.buffered_items - batch.items.size() + : 0; + route_tick_to_entry_no_lock( + entry_it->second, + batch, + deliveries, + continuity_requests, + continuity_deliveries, + true, + usable_history, + 0); + } + } + + for (auto& continuity_delivery : continuity_deliveries) { + continuity_delivery.first->on_market_data_continuity( + continuity_delivery.second); + } + for (auto& delivery : deliveries) { + delivery.first->on_tick_data(delivery.second); + } + if (!continuity_requests.empty()) { + auto next = std::move(continuity_requests.front()); + request_tick_continuity_history( + next.router_id, + std::move(next.subscription), + std::move(next.request), + next.kind, + next.announce_gap, + next.from_time_ms, + next.to_time_ms, + next.requested_items, + next.attempt); + return; + } + if (finished) { + if (continuity_still_enabled) { + notify_tick_continuity( + router_id, + make_continuity_update( + subscription, + continuity_verified + ? MarketDataContinuityStatus::LIVE + : MarketDataContinuityStatus::DEGRADED, + from_time_ms, + to_time_ms, + requested_items, + delivered_history_items, + continuity_verified + ? "Historical tick continuity is ready." + : "Live tick delivery continues without verified continuity.")); + } + return; + } + } + } + inline void MarketDataRouterState::fail_pending_subscribe( RoutedSubscriptionId router_id, BaseMarketDataProvider& provider, @@ -2619,8 +3614,100 @@ namespace optionx::market_data { std::shared_ptr, MarketDataContinuityUpdate>>& continuity_deliveries, std::vector& prefill_routes, - std::vector& reconnect_routes) { + std::vector& reconnect_routes, + std::vector& tick_prefill_routes, + std::vector& tick_reconnect_routes) { auto subscriber = entry->subscriber.lock(); + + if (entry->stream.type == MarketDataType::TICKS) { + auto& continuity = entry->tick_continuity_state; + const bool invalidates_continuity = + update.status == MarketDataStreamStatus::DISCONNECTED || + update.status == MarketDataStreamStatus::RECONNECTING || + update.status == MarketDataStreamStatus::FAILED || + update.status == MarketDataStreamStatus::STOPPED; + const bool transport_ready = + update.status == MarketDataStreamStatus::READY; + const bool needs_transport_recovery = + continuity.initial_prefill_pending || + entry->tick_continuity.recovers_gaps(); + + if (entry->tick_continuity.enabled() && + needs_transport_recovery && invalidates_continuity) { + const bool announce_stale = + continuity.phase != MarketDataContinuityPhase::WAITING_FOR_READY; + ++continuity.generation; + continuity.request_in_flight = false; + continuity.retry_request.reset(); + continuity.retry_at = {}; + continuity.reconnect_target_time_ms = 0; + + if (entry->tick_continuity.recovers_gaps()) { + const auto interval = entry->tick_continuity.expected_interval_ms; + if (continuity.initial_prefill_pending && + continuity.initial_prefill_boundary_time_ms > 0) { + continuity.unverified_from_time_ms = + continuity.unverified_from_time_ms == 0 + ? continuity.initial_prefill_boundary_time_ms + : std::min( + continuity.unverified_from_time_ms, + continuity.initial_prefill_boundary_time_ms); + } else if (continuity.verified_through_time_ms > 0 && + interval > 0 && + continuity.verified_through_time_ms <= + std::numeric_limits::max() - interval) { + continuity.unverified_from_time_ms = + continuity.unverified_from_time_ms == 0 + ? continuity.verified_through_time_ms + interval + : std::min( + continuity.unverified_from_time_ms, + continuity.verified_through_time_ms + interval); + } else if (continuity.last_observed_time_ms > 0) { + const auto next = interval > 0 && + continuity.last_observed_time_ms <= + std::numeric_limits::max() - interval + ? continuity.last_observed_time_ms + interval + : continuity.last_observed_time_ms; + continuity.unverified_from_time_ms = + continuity.unverified_from_time_ms == 0 + ? next + : std::min(continuity.unverified_from_time_ms, next); + } + } + continuity.phase = MarketDataContinuityPhase::WAITING_FOR_READY; + + if (announce_stale) { + const auto stale_update = make_continuity_update( + entry->control->provider_subscription, + MarketDataContinuityStatus::STALE, + 0, + 0, + 0, + 0, + "Transport loss invalidated tick continuity."); + record_continuity_update_no_lock(entry, stale_update); + if (subscriber) { + continuity_deliveries.emplace_back( + subscriber, + stale_update); + } + } + } + + if (entry->tick_continuity.enabled() && transport_ready && + continuity.phase == MarketDataContinuityPhase::WAITING_FOR_READY && + !continuity.request_in_flight) { + if (continuity.initial_prefill_pending) { + continuity.phase = MarketDataContinuityPhase::PREFILLING; + tick_prefill_routes.push_back(entry->router_id); + } else if (entry->tick_continuity.recovers_gaps()) { + continuity.phase = MarketDataContinuityPhase::RECOVERING; + tick_reconnect_routes.push_back(entry->router_id); + } + } + return; + } + auto& continuity = entry->continuity_state; const bool invalidates_continuity = update.status == MarketDataStreamStatus::DISCONNECTED || @@ -2700,29 +3787,31 @@ namespace optionx::market_data { } } - template - inline void MarketDataRouterState::route_batch( + inline void MarketDataRouterState::route_ticks( ProviderInstanceId provider_id, - std::unique_ptr batch, - MatchesStream matches_stream, - Deliver deliver) { + std::unique_ptr batch) { if (!batch) return; - std::vector, Batch>> deliveries; + + std::vector, + TickDataBatch>> deliveries; + std::vector, + MarketDataContinuityUpdate>> continuity_deliveries; + std::vector continuity_requests; { std::lock_guard lock(m_mutex); const auto provider_it = m_providers.find(provider_id); if (provider_it == m_providers.end()) return; - auto add_delivery = [&](const std::shared_ptr& entry) { - if (entry->phase != EntryPhase::ACTIVE || - !matches_stream(*batch, entry->stream)) { - return; - } - auto subscriber = entry->subscriber.lock(); - if (!subscriber) return; - auto routed = *batch; - routed.subscription = entry->control->provider_subscription; - deliveries.emplace_back(std::move(subscriber), std::move(routed)); + auto route_one = [&](const std::shared_ptr& entry) { + route_tick_to_entry_no_lock( + entry, + *batch, + deliveries, + continuity_requests, + continuity_deliveries, + false); }; if (batch->subscription.valid()) { @@ -2731,38 +3820,34 @@ namespace optionx::market_data { batch->subscription.id); if (route_it == provider_it->second.provider_routes.end()) return; const auto entry_it = m_entries.find(route_it->second); - if (entry_it == m_entries.end()) return; - add_delivery(entry_it->second); + if (entry_it != m_entries.end()) route_one(entry_it->second); } else { for (const auto& [id, entry] : m_entries) { (void)id; - if (entry->provider_id == provider_id) { - add_delivery(entry); - } + if (entry->provider_id == provider_id) route_one(entry); } } } + for (auto& continuity_delivery : continuity_deliveries) { + continuity_delivery.first->on_market_data_continuity( + continuity_delivery.second); + } for (auto& delivery : deliveries) { - deliver(*delivery.first, delivery.second); + delivery.first->on_tick_data(delivery.second); + } + for (auto& request : continuity_requests) { + request_tick_continuity_history( + request.router_id, + std::move(request.subscription), + std::move(request.request), + request.kind, + request.announce_gap, + request.from_time_ms, + request.to_time_ms, + request.requested_items, + request.attempt); } - } - - inline void MarketDataRouterState::route_ticks( - ProviderInstanceId provider_id, - std::unique_ptr batch) { - route_batch( - provider_id, - std::move(batch), - [this](const TickDataBatch& data, const StreamDescriptor& stream) { - return batch_matches_stream(data, stream); - }, - [](IMarketDataSubscriber& subscriber, TickDataBatch& data) { - for (auto& tick : data.items) { - mark_live_payload(tick.flags); - } - subscriber.on_tick_data(data); - }); } inline MarketDataContinuityUpdate @@ -2811,7 +3896,6 @@ namespace optionx::market_data { MarketDataContinuitySnapshot snapshot; if (!entry) return snapshot; - const auto& continuity = entry->continuity_state; snapshot.route = entry->router_id; if (entry->control) { std::lock_guard control_lock(entry->control->mutex); @@ -2820,6 +3904,53 @@ namespace optionx::market_data { snapshot.type = entry->stream.type; snapshot.symbol = entry->stream.symbol; snapshot.timeframe = entry->stream.timeframe; + + if (entry->stream.type == MarketDataType::TICKS) { + const auto& continuity = entry->tick_continuity_state; + snapshot.enabled = entry->tick_continuity.enabled(); + snapshot.expected_interval_ms = + entry->tick_continuity.expected_interval_ms; + snapshot.last_status = continuity.last_status; + snapshot.last_operation = continuity.last_operation; + snapshot.request_in_flight = continuity.request_in_flight; + auto observed_time_ms = continuity.last_observed_time_ms; + for (const auto& batch : continuity.buffer) { + for (const auto& tick : batch.items) { + observed_time_ms = std::max(observed_time_ms, tick.time_ms); + } + } + snapshot.last_observed_time_ms = observed_time_ms; + snapshot.requested_from_time_ms = continuity.requested_from_time_ms; + snapshot.requested_to_time_ms = continuity.requested_to_time_ms; + snapshot.requested_items = continuity.requested_items; + snapshot.last_confirmed_from_time_ms = continuity.confirmed_from_time_ms; + snapshot.last_confirmed_to_time_ms = continuity.confirmed_through_time_ms; + snapshot.last_confirmed_items = continuity.confirmed_items; + snapshot.verified_through_time_ms = continuity.verified_through_time_ms; + snapshot.unverified_from_time_ms = continuity.unverified_from_time_ms; + snapshot.buffered_batches = continuity.buffer.size(); + snapshot.buffered_items = continuity.buffered_items; + snapshot.history_request_count = continuity.history_request_count; + snapshot.retry_count = continuity.retry_count; + snapshot.failure_count = continuity.failure_count; + snapshot.last_failure = continuity.last_failure; + snapshot.phase = continuity.phase; + + auto stale_duration = continuity.stale_duration; + if (continuity.stale_since != std::chrono::steady_clock::time_point{}) { + stale_duration += now - continuity.stale_since; + } + snapshot.stale_duration_ms = duration_to_milliseconds(stale_duration); + + auto degraded_duration = continuity.degraded_duration; + if (continuity.degraded_since != std::chrono::steady_clock::time_point{}) { + degraded_duration += now - continuity.degraded_since; + } + snapshot.degraded_duration_ms = duration_to_milliseconds(degraded_duration); + return snapshot; + } + + const auto& continuity = entry->continuity_state; snapshot.enabled = entry->continuity.enabled(); snapshot.last_status = continuity.last_status; snapshot.last_operation = continuity.last_operation; @@ -2869,6 +4000,11 @@ namespace optionx::market_data { const MarketDataContinuityUpdate& update) { if (!entry) return; + if (entry->stream.type == MarketDataType::TICKS) { + record_tick_continuity_update_no_lock(entry, update); + return; + } + auto& continuity = entry->continuity_state; const auto now = std::chrono::steady_clock::now(); if (continuity.stale_since != @@ -2899,6 +4035,41 @@ namespace optionx::market_data { } } + inline void MarketDataRouterState::record_tick_continuity_update_no_lock( + const std::shared_ptr& entry, + const MarketDataContinuityUpdate& update) { + if (!entry) return; + + auto& continuity = entry->tick_continuity_state; + const auto now = std::chrono::steady_clock::now(); + if (continuity.stale_since != + std::chrono::steady_clock::time_point{} && + update.status != MarketDataContinuityStatus::STALE) { + continuity.stale_duration += now - continuity.stale_since; + continuity.stale_since = {}; + } + if (continuity.degraded_since != + std::chrono::steady_clock::time_point{} && + update.status != MarketDataContinuityStatus::DEGRADED) { + continuity.degraded_duration += now - continuity.degraded_since; + continuity.degraded_since = {}; + } + if (update.status == MarketDataContinuityStatus::STALE && + continuity.stale_since == std::chrono::steady_clock::time_point{}) { + continuity.stale_since = now; + } + if (update.status == MarketDataContinuityStatus::DEGRADED && + continuity.degraded_since == std::chrono::steady_clock::time_point{}) { + continuity.degraded_since = now; + } + + continuity.last_status = update.status; + if (update.status == MarketDataContinuityStatus::FAILED && + !update.message.empty()) { + continuity.last_failure = update.message; + } + } + inline void MarketDataRouterState::record_continuity_update( RoutedSubscriptionId router_id, const MarketDataContinuityUpdate& update) { @@ -2926,6 +4097,23 @@ namespace optionx::market_data { if (subscriber) subscriber->on_market_data_continuity(update); } + inline void MarketDataRouterState::notify_tick_continuity( + RoutedSubscriptionId router_id, + MarketDataContinuityUpdate update) { + std::shared_ptr subscriber; + { + 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; + } + record_tick_continuity_update_no_lock(entry_it->second, update); + subscriber = entry_it->second->subscriber.lock(); + } + if (subscriber) subscriber->on_market_data_continuity(update); + } + inline bool MarketDataRouterState::buffer_continuity_batch_no_lock( const std::shared_ptr& entry, std::shared_ptr subscriber, @@ -3002,6 +4190,208 @@ namespace optionx::market_data { return false; } + inline bool MarketDataRouterState::buffer_tick_continuity_batch_no_lock( + const std::shared_ptr& entry, + std::shared_ptr subscriber, + TickDataBatch batch, + std::vector, + TickDataBatch>>& deliveries, + std::vector, + MarketDataContinuityUpdate>>& continuity_deliveries, + bool push_front) { + for (auto& tick : batch.items) { + mark_live_payload(tick.flags, true); + } + const auto batch_items = batch.items.size(); + const auto& options = entry->tick_continuity; + auto& continuity = entry->tick_continuity_state; + const bool exceeds_batch_limit = + options.max_buffered_batches > 0 && + continuity.buffer.size() >= options.max_buffered_batches; + const bool exceeds_item_limit = options.max_buffered_items > 0 && + (batch_items > options.max_buffered_items || + continuity.buffered_items > + options.max_buffered_items - batch_items); + + if (!exceeds_batch_limit && !exceeds_item_limit) { + if (push_front) { + continuity.buffer.push_front(std::move(batch)); + } else { + continuity.buffer.push_back(std::move(batch)); + } + continuity.buffered_items += batch_items; + return true; + } + + while (!continuity.buffer.empty()) { + auto buffered = std::move(continuity.buffer.front()); + continuity.buffer.pop_front(); + record_tick_progress_no_lock(entry, buffered.items); + deliveries.emplace_back(subscriber, std::move(buffered)); + } + continuity.buffered_items = 0; + record_tick_progress_no_lock(entry, batch.items); + deliveries.emplace_back(subscriber, std::move(batch)); + + continuity.phase = MarketDataContinuityPhase::DEGRADED; + continuity.initial_prefill_pending = false; + continuity.request_in_flight = false; + continuity.reconnect_target_time_ms = 0; + continuity.retry_request.reset(); + continuity.retry_at = {}; + entry->tick_continuity.mode = MarketDataContinuityMode::LIVE_ONLY; + + const auto failed_update = make_continuity_update( + entry->control->provider_subscription, + MarketDataContinuityStatus::FAILED, + 0, + 0, + 0, + 0, + "Tick continuity buffer limit exceeded; live delivery continues."); + const auto degraded_update = make_continuity_update( + entry->control->provider_subscription, + MarketDataContinuityStatus::DEGRADED, + 0, + 0, + 0, + 0, + "Live tick delivery resumed without verified continuity after buffer overflow."); + record_tick_continuity_update_no_lock(entry, failed_update); + record_tick_continuity_update_no_lock(entry, degraded_update); + continuity_deliveries.emplace_back(subscriber, failed_update); + continuity_deliveries.emplace_back(subscriber, degraded_update); + return false; + } + + inline bool MarketDataRouterState::route_tick_to_entry_no_lock( + const std::shared_ptr& entry, + const TickDataBatch& batch, + std::vector, + TickDataBatch>>& deliveries, + std::vector& continuity_requests, + std::vector, + MarketDataContinuityUpdate>>& continuity_deliveries, + bool process_buffered, + bool allow_gap_recovery, + std::uint64_t confirmed_through_time_ms) { + (void)confirmed_through_time_ms; + if (!entry || entry->phase != EntryPhase::ACTIVE || + !batch_matches_stream(batch, entry->stream)) { + return false; + } + + auto subscriber = entry->subscriber.lock(); + if (!subscriber) return false; + + auto routed = batch; + routed.subscription = entry->control->provider_subscription; + if (routed.items.empty()) return false; + + for (auto& tick : routed.items) { + mark_live_payload(tick.flags, process_buffered); + } + + auto& continuity = entry->tick_continuity_state; + if (entry->tick_continuity.enabled() && !process_buffered && + continuity.buffers_live_data()) { + buffer_tick_continuity_batch_no_lock( + entry, + std::move(subscriber), + std::move(routed), + deliveries, + continuity_deliveries, + false); + return false; + } + + const auto interval = entry->tick_continuity.expected_interval_ms; + if (allow_gap_recovery && entry->tick_continuity.recovers_gaps() && + !continuity.request_in_flight && + interval > 0 && continuity.last_observed_time_ms > 0) { + auto previous_time_ms = continuity.last_observed_time_ms; + for (std::size_t index = 0; index < routed.items.size(); ++index) { + const auto& tick = routed.items[index]; + if (tick.time_ms == 0) continue; + + const auto expected_time_ms = previous_time_ms <= + std::numeric_limits::max() - interval + ? previous_time_ms + interval + : std::numeric_limits::max(); + if (tick.time_ms > expected_time_ms) { + const auto gap_from_time_ms = expected_time_ms; + const auto gap_to_time_ms = tick.time_ms - interval; + if (gap_to_time_ms < gap_from_time_ms) return false; + + auto request_to_time_ms = gap_to_time_ms; + const auto max_backfill_ms = + entry->tick_continuity.max_backfill_ms; + if (max_backfill_ms > 0) { + const auto span = max_backfill_ms > interval + ? max_backfill_ms - interval + : 0U; + const auto bounded = gap_from_time_ms > + std::numeric_limits::max() - span + ? std::numeric_limits::max() + : gap_from_time_ms + span; + request_to_time_ms = std::min(request_to_time_ms, bounded); + } + if (request_to_time_ms < gap_from_time_ms) return false; + + if (index > 0) { + TickDataBatch prefix = routed; + prefix.items.resize(index); + record_tick_progress_no_lock(entry, prefix.items); + deliveries.emplace_back(subscriber, std::move(prefix)); + } + + routed.items.erase( + routed.items.begin(), + routed.items.begin() + static_cast(index)); + if (!buffer_tick_continuity_batch_no_lock( + entry, + subscriber, + std::move(routed), + deliveries, + continuity_deliveries, + process_buffered)) { + return false; + } + + mark_tick_unverified_no_lock(entry, gap_from_time_ms); + continuity.phase = MarketDataContinuityPhase::RECOVERING; + continuity.request_in_flight = true; + continuity.reconnect_target_time_ms = gap_to_time_ms; + continuity_requests.push_back(PendingTickContinuityRequest{ + entry->router_id, + entry->control->provider_subscription, + TickHistoryRequest( + entry->stream.symbol, + gap_from_time_ms, + request_to_time_ms), + ContinuityRequestKind::GAP_BACKFILL, + true, + gap_from_time_ms, + request_to_time_ms, + 0, + 1}); + return false; + } + if (tick.time_ms > previous_time_ms) { + previous_time_ms = tick.time_ms; + } + } + } + + record_tick_progress_no_lock(entry, routed.items); + deliveries.emplace_back(std::move(subscriber), std::move(routed)); + return true; + } + inline bool MarketDataRouterState::route_bar_to_entry_no_lock( const std::shared_ptr& entry, const BarDataBatch& batch, @@ -3233,6 +4623,8 @@ namespace optionx::market_data { MarketDataContinuityUpdate>> continuity_deliveries; std::vector prefill_routes; std::vector reconnect_routes; + std::vector tick_prefill_routes; + std::vector tick_reconnect_routes; { std::lock_guard lock(m_mutex); const auto provider_it = m_providers.find(provider_id); @@ -3265,7 +4657,9 @@ namespace optionx::market_data { update, continuity_deliveries, prefill_routes, - reconnect_routes); + reconnect_routes, + tick_prefill_routes, + tick_reconnect_routes); auto subscriber = entry->subscriber.lock(); if (subscriber) { @@ -3318,6 +4712,20 @@ namespace optionx::market_data { for (const auto router_id : reconnect_routes) { start_reconnect_recovery(router_id); } + for (const auto router_id : tick_prefill_routes) { + MarketDataSubscriptionHandle subscription; + { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (entry_it != m_entries.end()) { + subscription = entry_it->second->control->provider_subscription; + } + } + start_tick_continuity(router_id, std::move(subscription)); + } + for (const auto router_id : tick_reconnect_routes) { + start_tick_reconnect_recovery(router_id); + } } inline void MarketDataRouterState::set_control_active( @@ -3417,6 +4825,8 @@ namespace optionx::market_data { for (;;) { std::vector retry_requests; std::vector reconnect_routes; + std::vector tick_retry_requests; + std::vector tick_reconnect_routes; std::vector completions; bool shutting_down = false; { @@ -3445,6 +4855,26 @@ namespace optionx::market_data { continuity.retry_at = {}; continuity.request_in_flight = true; } + for (const auto& [id, entry] : m_entries) { + (void)id; + auto& continuity = entry->tick_continuity_state; + if (entry->phase != EntryPhase::ACTIVE || + entry->stream.type != MarketDataType::TICKS || + continuity.request_in_flight || + !continuity.retry_request || + continuity.retry_at > now) { + continue; + } + tick_retry_requests.push_back(*continuity.retry_request); + continuity.phase = + continuity.retry_request->kind == + ContinuityRequestKind::PREFILL + ? MarketDataContinuityPhase::PREFILLING + : MarketDataContinuityPhase::RECOVERING; + continuity.retry_request.reset(); + continuity.retry_at = {}; + continuity.request_in_flight = true; + } for (const auto& [id, entry] : m_entries) { const auto& continuity = entry->continuity_state; if (entry->phase == EntryPhase::ACTIVE && @@ -3455,6 +4885,18 @@ namespace optionx::market_data { reconnect_routes.push_back(id); } } + for (const auto& [id, entry] : m_entries) { + const auto& continuity = entry->tick_continuity_state; + if (entry->phase == EntryPhase::ACTIVE && + entry->stream.type == MarketDataType::TICKS && + entry->tick_continuity.recovers_gaps() && + continuity.phase == MarketDataContinuityPhase::RECOVERING && + !continuity.initial_prefill_pending && + !continuity.request_in_flight && + !continuity.retry_request) { + tick_reconnect_routes.push_back(id); + } + } } if (shutting_down) { @@ -3485,10 +4927,27 @@ namespace optionx::market_data { retry.attempt); } + for (auto& retry : tick_retry_requests) { + request_tick_continuity_history( + retry.router_id, + std::move(retry.subscription), + std::move(retry.request), + retry.kind, + retry.announce_gap, + retry.from_time_ms, + retry.to_time_ms, + retry.requested_items, + retry.attempt); + } + for (const auto router_id : reconnect_routes) { start_reconnect_recovery(router_id); } + for (const auto router_id : tick_reconnect_routes) { + start_tick_reconnect_recovery(router_id); + } + if (!shutting_down) return; for (auto& completion : completions) { diff --git a/tests/market_data_tick_continuity_test.cpp b/tests/market_data_tick_continuity_test.cpp new file mode 100644 index 0000000..d0a66cc --- /dev/null +++ b/tests/market_data_tick_continuity_test.cpp @@ -0,0 +1,527 @@ +#include + +#include + +namespace market_data_tick_continuity_test_clock { + inline std::uint64_t now_ms = 3000ULL; +} + +#ifndef OPTIONX_TIMESTAMP_MS +#define OPTIONX_TIMESTAMP_MS market_data_tick_continuity_test_clock::now_ms +#endif + +#include +#include +#include +#include +#include +#include +#include + +#include + +using namespace optionx; +using namespace optionx::market_data; + +namespace { + +class ScopedTestClock { +public: + explicit ScopedTestClock(std::uint64_t now_ms) { + market_data_tick_continuity_test_clock::now_ms = now_ms; + } + + ~ScopedTestClock() { + market_data_tick_continuity_test_clock::now_ms = 3000ULL; + } +}; + +class FakeTickHistoryProvider final : public BaseMarketDataProvider { +public: + ticks_callback_t& on_tick_data() override { + return m_tick_callback; + } + + status_callback_t& on_market_data_status() override { + return m_status_callback; + } + + bool subscribe_ticks( + TickSubscriptionRequest request, + subscription_callback_t callback) override { + m_active_subscription = MarketDataSubscriptionHandle::from_tick_request( + provider_id(), + m_next_subscription_id++, + request); + if (callback) { + callback(MarketDataSubscriptionResult::subscribed( + m_active_subscription)); + } + return true; + } + + bool unsubscribe( + MarketDataSubscriptionHandle subscription, + subscription_callback_t callback) override { + if (callback) { + callback(MarketDataSubscriptionResult::unsubscribed( + std::move(subscription))); + } + return true; + } + + bool fetch_tick_history( + const TickHistoryRequest& request, + tick_history_callback_t callback) override { + history_requests.push_back(request); + m_history_callbacks.push_back(std::move(callback)); + return true; + } + + void complete_history( + TickSequence sequence, + bool range_complete = true) { + ASSERT_FALSE(m_history_callbacks.empty()); + auto callback = std::move(m_history_callbacks.front()); + m_history_callbacks.pop_front(); + callback(TickHistoryResult::ok( + std::move(sequence), + range_complete)); + } + + void fail_history(std::string message) { + ASSERT_FALSE(m_history_callbacks.empty()); + auto callback = std::move(m_history_callbacks.front()); + m_history_callbacks.pop_front(); + callback(TickHistoryResult::fail(std::move(message))); + } + + void emit_ticks(std::vector ticks) { + ASSERT_TRUE(static_cast(m_tick_callback)); + auto batch = std::make_unique(); + batch->subscription = m_active_subscription; + batch->type = MarketDataType::TICKS; + batch->symbol = m_active_subscription.symbol; + batch->price_digits = 5; + batch->items = std::move(ticks); + for (auto& tick : batch->items) { + mark_live_payload(tick.flags); + } + m_tick_callback(std::move(batch)); + } + + void emit_status(MarketDataStreamStatus status) { + ASSERT_TRUE(static_cast(m_status_callback)); + MarketDataStatusUpdate update; + update.subscription = m_active_subscription; + update.type = MarketDataType::TICKS; + update.symbol = m_active_subscription.symbol; + update.transport = m_active_subscription.transport; + update.status = status; + m_status_callback(std::move(update)); + } + + const MarketDataSubscriptionHandle& active_subscription() const noexcept { + return m_active_subscription; + } + + std::vector history_requests; + +private: + SubscriptionId m_next_subscription_id = 1; + MarketDataSubscriptionHandle m_active_subscription; + std::deque m_history_callbacks; + ticks_callback_t m_tick_callback; + status_callback_t m_status_callback; +}; + +class RecordingSubscriber final : public IMarketDataSubscriber { +public: + void on_tick_data(const TickDataBatch& batch) override { + ticks.push_back(batch); + } + + void on_market_data_continuity( + const MarketDataContinuityUpdate& update) override { + continuity.push_back(update); + } + + std::vector ticks; + std::vector continuity; +}; + +Tick make_tick( + std::uint64_t time_ms, + double price = 1.0, + std::uint64_t received_ms = 0) { + Tick tick; + tick.ask = price + 0.01; + tick.bid = price; + tick.last = price; + tick.volume = 1.0; + tick.time_ms = time_ms; + tick.received_ms = received_ms; + return tick; +} + +TickSequence make_history(std::initializer_list times) { + TickSequence sequence; + sequence.symbol = "EURUSD"; + sequence.provider = "INTRADE_BAR"; + sequence.price_digits = 5; + for (const auto time_ms : times) { + sequence.ticks.push_back(make_tick(time_ms)); + } + return sequence; +} + +TickSubscriptionRequest continuity_request( + MarketDataContinuityMode mode = + MarketDataContinuityMode::PREFILL_AND_RECOVER, + std::uint64_t prefill_lookback_ms = 2000, + std::uint64_t max_backfill_ms = 60000) { + TickSubscriptionRequest request("EURUSD"); + request.continuity.mode = mode; + request.continuity.prefill_lookback_ms = prefill_lookback_ms; + request.continuity.expected_interval_ms = 1000; + request.continuity.max_backfill_ms = max_backfill_ms; + request.continuity.retry.max_attempts = 2; + request.continuity.retry.initial_backoff_ms = 0; + return request; +} + +std::size_t count_status( + 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; + })); +} + +const Tick& only_tick(const TickDataBatch& batch) { + EXPECT_EQ(batch.items.size(), 1U); + return batch.items.front(); +} + +TEST(MarketDataTickContinuity, CompletesPrefillAndDeliversDirectRealtime) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + ASSERT_EQ(provider.history_requests.size(), 1U); + EXPECT_EQ(provider.history_requests.front().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.front().to_time_ms, 3000U); + + provider.complete_history(make_history({1000, 2000, 3000})); + ASSERT_EQ(subscriber->ticks.size(), 1U); + ASSERT_EQ(subscriber->ticks.front().items.size(), 3U); + for (const auto& tick : subscriber->ticks.front().items) { + EXPECT_TRUE(tick.has_flag(MarketDataFlags::HISTORICAL)); + EXPECT_FALSE(tick.has_flag(MarketDataFlags::LIVE_SOURCE)); + } + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 1U); + + provider.emit_ticks({make_tick(4000)}); + ASSERT_EQ(subscriber->ticks.size(), 2U); + const auto& live_tick = only_tick(subscriber->ticks.back()); + EXPECT_EQ(live_tick.time_ms, 4000U); + EXPECT_TRUE(live_tick.has_flag(MarketDataFlags::LIVE_SOURCE)); + EXPECT_TRUE(live_tick.has_flag(MarketDataFlags::REALTIME)); + EXPECT_FALSE(live_tick.has_flag(MarketDataFlags::CATCHUP)); + + const auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::LIVE); + EXPECT_EQ(snapshot->verified_through_time_ms, 4000U); + EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); +} + +TEST(MarketDataTickContinuity, IncompletePrefillStaysDegraded) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + provider.complete_history(make_history({1000, 3000}), false); + + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 0U); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::FAILED), 1U); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::DEGRADED), 1U); + + const auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::DEGRADED); + EXPECT_EQ(snapshot->unverified_from_time_ms, 1000U); + EXPECT_TRUE(snapshot->enabled); +} + +TEST(MarketDataTickContinuity, RestartsInterruptedPrefillFromOriginalRange) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + ASSERT_EQ(provider.history_requests.size(), 1U); + EXPECT_EQ(provider.history_requests.front().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.front().to_time_ms, 3000U); + + provider.emit_status(MarketDataStreamStatus::DISCONNECTED); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::STALE), 1U); + + market_data_tick_continuity_test_clock::now_ms = 6000; + provider.emit_status(MarketDataStreamStatus::READY); + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 6000U); + + // The completion from before the transport loss is obsolete and must not + // satisfy the restarted request. + provider.complete_history(make_history({1000, 2000, 3000})); + EXPECT_TRUE(subscriber->ticks.empty()); + auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_TRUE(snapshot->request_in_flight); + + provider.complete_history(make_history({1000, 2000, 3000, 4000, 5000, 6000})); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 1U); + snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::LIVE); + EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); +} + +TEST(MarketDataTickContinuity, ChecksGapsInsidePrefillBacklog) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + provider.emit_ticks({make_tick(4000), make_tick(6000)}); + ASSERT_EQ(provider.history_requests.size(), 1U); + + provider.complete_history(make_history({1000, 2000, 3000})); + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 5000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 5000U); + provider.complete_history(make_history({5000})); + + ASSERT_FALSE(subscriber->ticks.empty()); + const auto& prefix = subscriber->ticks[1]; + ASSERT_EQ(prefix.items.size(), 1U); + EXPECT_EQ(prefix.items.front().time_ms, 4000U); + EXPECT_TRUE(prefix.items.front().has_flag(MarketDataFlags::CATCHUP)); + const auto& tail = subscriber->ticks.back(); + ASSERT_EQ(tail.items.size(), 1U); + EXPECT_EQ(tail.items.front().time_ms, 6000U); + EXPECT_TRUE(tail.items.front().has_flag(MarketDataFlags::CATCHUP)); + + const auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::LIVE); + EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); +} + +TEST(MarketDataTickContinuity, RepairsGapAndPreservesDistinctSameSecondTicks) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + provider.complete_history(make_history({1000, 2000, 3000})); + + provider.emit_ticks({make_tick(4000)}); + Tick first = make_tick(6000, 1.0); + Tick second = make_tick(6000, 1.1); + provider.emit_ticks({first, second}); + + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 5000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 5000U); + const auto before_repair_ticks = subscriber->ticks.size(); + + provider.complete_history(make_history({5000})); + ASSERT_EQ(subscriber->ticks.size(), before_repair_ticks + 2U); + const auto& history_batch = subscriber->ticks[before_repair_ticks]; + EXPECT_TRUE(only_tick(history_batch).has_flag(MarketDataFlags::HISTORICAL)); + const auto& catchup_batch = subscriber->ticks[before_repair_ticks + 1U]; + ASSERT_EQ(catchup_batch.items.size(), 2U); + EXPECT_EQ(catchup_batch.items[0].time_ms, 6000U); + EXPECT_EQ(catchup_batch.items[1].time_ms, 6000U); + EXPECT_NE(catchup_batch.items[0].bid, catchup_batch.items[1].bid); + EXPECT_TRUE(catchup_batch.items[0].has_flag(MarketDataFlags::CATCHUP)); + EXPECT_TRUE(catchup_batch.items[1].has_flag(MarketDataFlags::CATCHUP)); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 2U); +} + +TEST(MarketDataTickContinuity, UsesBoundedGapRequests) { + ScopedTestClock clock(10000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto request = continuity_request( + MarketDataContinuityMode::PREFILL_AND_RECOVER, + 0, + 3000); + auto route = router.subscribe_ticks(provider, subscriber, request); + ASSERT_TRUE(route.valid()); + provider.emit_ticks({make_tick(1000)}); + provider.emit_ticks({make_tick(10000)}); + + ASSERT_EQ(provider.history_requests.size(), 1U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 2000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 4000U); + provider.complete_history(make_history({2000, 3000, 4000})); + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 5000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 7000U); + provider.complete_history(make_history({5000, 6000, 7000})); + ASSERT_EQ(provider.history_requests.size(), 3U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 8000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 9000U); + provider.complete_history(make_history({8000, 9000})); + + const auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::LIVE); + EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); + EXPECT_EQ(snapshot->history_request_count, 3U); +} + +TEST(MarketDataTickContinuity, RetriesFailedHistoryFromProcess) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + provider.fail_history("temporary history failure"); + EXPECT_EQ(provider.history_requests.size(), 1U); + + auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->last_status, MarketDataContinuityStatus::RETRYING); + EXPECT_EQ(snapshot->retry_count, 0U); + + router.process(); + ASSERT_EQ(provider.history_requests.size(), 2U); + provider.complete_history(make_history({1000, 2000, 3000})); + snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::LIVE); + EXPECT_EQ(snapshot->retry_count, 1U); +} + +TEST(MarketDataTickContinuity, ReconnectDeduplicatesOnlyExactOverlap) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + provider.complete_history(make_history({1000, 2000, 3000})); + provider.emit_ticks({make_tick(3000)}); + provider.emit_status(MarketDataStreamStatus::DISCONNECTED); + provider.emit_ticks({make_tick(5000, 1.0), make_tick(5000, 1.1)}); + + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::STALE), 1U); + auto stale_snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(stale_snapshot.has_value()); + EXPECT_EQ(stale_snapshot->phase, MarketDataContinuityPhase::WAITING_FOR_READY); + EXPECT_FALSE(stale_snapshot->request_in_flight); + EXPECT_EQ(provider.history_requests.size(), 1U); + + provider.emit_status(MarketDataStreamStatus::READY); + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 4000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 5000U); + + TickSequence history = make_history({4000}); + history.ticks.push_back(make_tick(5000, 1.0)); + provider.complete_history(std::move(history)); + + ASSERT_FALSE(subscriber->ticks.empty()); + const auto& catchup_batch = subscriber->ticks.back(); + ASSERT_EQ(catchup_batch.items.size(), 1U); + EXPECT_EQ(catchup_batch.items.front().time_ms, 5000U); + EXPECT_DOUBLE_EQ(catchup_batch.items.front().bid, 1.1); + EXPECT_TRUE(catchup_batch.items.front().has_flag(MarketDataFlags::CATCHUP)); + + const auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::LIVE); + EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); +} + +TEST(MarketDataTickContinuity, DisablesContinuityAfterBufferOverflow) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto request = continuity_request(); + request.continuity.max_buffered_batches = 1; + request.continuity.max_buffered_items = 2; + auto route = router.subscribe_ticks(provider, subscriber, request); + ASSERT_TRUE(route.valid()); + + provider.emit_ticks({make_tick(1000)}); + provider.emit_ticks({make_tick(2000)}); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::FAILED), 1U); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::DEGRADED), 1U); + + const auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_FALSE(snapshot->enabled); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::DEGRADED); + + provider.complete_history(make_history({1000, 2000, 3000})); + provider.emit_ticks({make_tick(4000)}); + ASSERT_FALSE(subscriber->ticks.empty()); + const auto& live_tick = only_tick(subscriber->ticks.back()); + EXPECT_TRUE(live_tick.has_flag(MarketDataFlags::REALTIME)); +} + +} // namespace + +int main(int argc, char** argv) { + ::testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +} From ee388dab987fae82663b82824eda9d7609b040d1 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Mon, 7 Sep 2026 20:40:57 +0300 Subject: [PATCH 2/5] fix(market-data): align provider tick history recovery --- guides/api-and-header-contracts.md | 5 + guides/market-data-router.md | 22 ++ guides/market-data-router.ru.md | 21 ++ .../market_data/BaseMarketDataProvider.hpp | 19 ++ .../market_data/detail/MarketDataRouter.ipp | 242 ++++++++++++------ .../platforms/IntradeBarPlatform.hpp | 11 + .../ObservedTickHistory.hpp | 99 +++++++ tests/intrade_observed_tick_history_test.cpp | 15 ++ tests/market_data_tick_continuity_test.cpp | 188 ++++++++++++-- 9 files changed, 525 insertions(+), 97 deletions(-) diff --git a/guides/api-and-header-contracts.md b/guides/api-and-header-contracts.md index abafe2d..d9acd51 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -295,6 +295,11 @@ Contract rules: assertion; successful observations with `false` must not be used as proof of continuity. Equal timestamps are valid, and exact duplicate removal remains provider/consumer policy. +- `BaseMarketDataProvider::provider_time_ms()` and + `BaseMarketDataProvider::tick_history_interval_ms()` are optional metadata + hooks for providers whose history backend has its own clock or sampling grid. + A zero value means that Router should use its local clock or event-oriented + range semantics. ### Intrade observed tick history diff --git a/guides/market-data-router.md b/guides/market-data-router.md index 88cdac3..46f7c93 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -649,6 +649,17 @@ archive: - `trade_check2.php` remains a settlement/trade-result API and is not used for range history. +Providers may expose optional history-clock metadata through +`BaseMarketDataProvider::provider_time_ms()` and +`BaseMarketDataProvider::tick_history_interval_ms()`. The first is a provider +time estimate in milliseconds; the second describes the history backend's +sampling grid and is not the live gap-detection threshold. Intrade estimates +its provider time from recent `(Tick::time_ms - Tick::received_ms)` samples, +uses their median to reject polling jitter, and rounds the estimate down to +the one-second history grid. With no estimate, Router falls back to the local +clock. Providers without a discrete history grid leave both optional hooks at +their defaults. + ### Tick continuity in Router Set `TickSubscriptionRequest::continuity` to enable history-first delivery for @@ -672,6 +683,17 @@ equal timestamps and sub-second live ticks are valid. The provider's `range_complete` remains the authority for whether a requested history range proves continuity. +Recovery requests use inclusive timestamp ranges. For an event-oriented +provider, the suspicious interval is requested without inventing missing tick +slots, and the live tick that triggered recovery remains buffered. When a +provider declares a history grid, Router aligns the start down and (for an +unbounded request) the end up so an off-grid observation is not excluded. +Bounded chunks keep their size limit and overlap at the previous end point +whenever that overlap can advance the range; if the limit is smaller than a +provider grid step, Router advances to the next provider boundary instead of +repeating the same request. The overlap is removed only by exact observation +identity. + The Router sends historical ticks first, marks them `HISTORICAL`, and then replays held live ticks as `LIVE_SOURCE | CATCHUP`. A complete result is required before the route can report `LIVE`. An incomplete result may still be delivered diff --git a/guides/market-data-router.ru.md b/guides/market-data-router.ru.md index a79a913..3db4f50 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -806,6 +806,17 @@ tick archive: - `trade_check2.php` остаётся settlement/trade-result API и не используется для range history. +Провайдер может дополнительно реализовать метаданные времени через +`BaseMarketDataProvider::provider_time_ms()` и +`BaseMarketDataProvider::tick_history_interval_ms()`. Первый метод возвращает +оценку времени провайдера в миллисекундах, второй описывает сетку backend +истории и не является порогом live gap detection. Intrade получает эту оценку +по недавним samples `(Tick::time_ms - Tick::received_ms)`, берёт их медиану, +чтобы отфильтровать polling jitter, и округляет результат вниз до секундной +сетки истории. Если оценка недоступна, Router использует локальные часы. +Провайдеры без дискретной сетки оставляют оба optional hook со значениями по +умолчанию. + ### Tick continuity в Router Чтобы включить history-first delivery для одного tick route, настройте @@ -828,6 +839,16 @@ gap-detection hint, а не требование плотной сетки: equa ticks чаще одной секунды допустимы. Для доказательства continuity Router доверяет только значению `range_complete` в provider result. +Recovery использует inclusive timestamp ranges. Для event-oriented provider +Router не синтезирует пропущенные tick slots, а сохраняет в buffer live tick, +который запустил recovery. Если provider объявляет history grid, Router +округляет начало вниз, а конец вверх для unbounded request, чтобы observation +с timestamp между grid points не выпала из диапазона. Bounded chunks сохраняют +лимит размера и перекрываются в предыдущей конечной точке, когда такой overlap +позволяет продвинуть диапазон. Если лимит меньше шага provider grid, Router +переходит к следующей provider boundary, а не повторяет тот же запрос. Overlap +удаляется только по exact observation identity. + Router сначала отправляет historical ticks с флагом `HISTORICAL`, затем воспроизводит удержанные live ticks с флагами `LIVE_SOURCE | CATCHUP`. До `LIVE` нужен complete result. Неполный result можно доставить как observations, diff --git a/include/optionx_cpp/market_data/BaseMarketDataProvider.hpp b/include/optionx_cpp/market_data/BaseMarketDataProvider.hpp index 67d787b..6b58056 100644 --- a/include/optionx_cpp/market_data/BaseMarketDataProvider.hpp +++ b/include/optionx_cpp/market_data/BaseMarketDataProvider.hpp @@ -87,6 +87,25 @@ namespace optionx::market_data { return m_status_callback; } + /// \brief Returns the provider's current market-data timestamp. + /// \details The value is optional. A provider may return zero when it + /// cannot estimate its own clock. Router continuity falls back + /// to the local clock in that case. Providers with a discrete + /// history grid should return a timestamp aligned to that grid. + /// \return Current provider timestamp in milliseconds, or zero when unavailable. + virtual std::uint64_t provider_time_ms() const noexcept { + return 0; + } + + /// \brief Returns the provider tick-history sampling grid. + /// \details This describes the history backend, not a generic live-tick + /// gap threshold. Zero means that history uses continuous time + /// or has no provider-specific alignment requirement. + /// \return Positive history grid size in milliseconds, or zero. + virtual std::uint64_t tick_history_interval_ms() const noexcept { + return 0; + } + /// \brief Requests a live tick stream subscription. /// \param request Tick subscription parameters. /// \param callback Callback receiving desired-subscription acceptance or failure. diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index 2838d5d..8276987 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -550,6 +550,21 @@ namespace optionx::market_data { static void record_tick_progress_no_lock( const std::shared_ptr& entry, const std::vector& ticks) noexcept; + static std::uint64_t provider_time_ms( + BaseMarketDataProvider& provider) noexcept; + static std::uint64_t tick_history_interval_ms( + BaseMarketDataProvider& provider) noexcept; + static std::uint64_t align_tick_history_start_time_ms( + std::uint64_t time_ms, + std::uint64_t interval_ms) noexcept; + static std::uint64_t align_tick_history_end_time_ms( + std::uint64_t time_ms, + std::uint64_t interval_ms) noexcept; + static std::uint64_t bounded_tick_history_to_time_ms( + std::uint64_t from_time_ms, + std::uint64_t target_time_ms, + std::uint64_t max_backfill_ms, + std::uint64_t interval_ms) noexcept; static std::uint64_t duration_to_milliseconds( std::chrono::steady_clock::duration duration) noexcept; static MarketDataContinuitySnapshot make_continuity_snapshot_no_lock( @@ -1633,24 +1648,33 @@ namespace optionx::market_data { } const auto& entry = entry_it->second; - const auto now_ms = static_cast(OPTIONX_TIMESTAMP_MS); + const auto now_ms = provider_time_ms(*entry->provider); const auto lookback = entry->tick_continuity.prefill_lookback_ms; + const auto history_interval_ms = tick_history_interval_ms(*entry->provider); pending.router_id = router_id; pending.subscription = subscription; auto& continuity = entry->tick_continuity_state; if (continuity.initial_prefill_boundary_time_ms == 0) { continuity.initial_prefill_boundary_time_ms = - now_ms > 0 ? now_ms : 1U; + align_tick_history_start_time_ms(now_ms, history_interval_ms); } const auto original_boundary_time_ms = continuity.initial_prefill_boundary_time_ms; - const auto current_time_ms = now_ms > 0 ? now_ms : 1U; + const auto current_time_ms = + align_tick_history_start_time_ms(now_ms, history_interval_ms); const auto current_boundary_time_ms = std::max( original_boundary_time_ms, current_time_ms); - pending.from_time_ms = original_boundary_time_ms > lookback + const auto requested_from_time_ms = original_boundary_time_ms > lookback ? original_boundary_time_ms - lookback : 1U; + pending.from_time_ms = align_tick_history_start_time_ms( + requested_from_time_ms, + history_interval_ms); + if (pending.from_time_ms == 0 && history_interval_ms > 0 && + current_boundary_time_ms >= history_interval_ms) { + pending.from_time_ms = history_interval_ms; + } pending.to_time_ms = current_boundary_time_ms; pending.request = TickHistoryRequest( entry->stream.symbol, @@ -1886,24 +1910,25 @@ namespace optionx::market_data { const auto observed_time_ms = std::max( continuity.last_observed_time_ms, latest_buffered_time_ms); - const auto now_ms = static_cast(OPTIONX_TIMESTAMP_MS); + const auto now_ms = provider_time_ms(*entry->provider); const auto target_time_ms = std::max(observed_time_ms, now_ms); - const auto interval = entry->tick_continuity.expected_interval_ms; + const auto history_interval_ms = tick_history_interval_ms(*entry->provider); auto from_time_ms = continuity.unverified_from_time_ms; - if (from_time_ms == 0 && continuity.verified_through_time_ms > 0 && - interval > 0 && - continuity.verified_through_time_ms <= - std::numeric_limits::max() - interval) { - from_time_ms = continuity.verified_through_time_ms + interval; + if (from_time_ms == 0 && continuity.verified_through_time_ms > 0) { + // History ranges are inclusive. Reuse the last verified + // observation as an overlap instead of inventing a tick grid. + from_time_ms = continuity.verified_through_time_ms; } else if (from_time_ms == 0 && earliest_buffered_time_ms > 0) { from_time_ms = earliest_buffered_time_ms; - } else if (from_time_ms == 0 && observed_time_ms > 0 && interval > 0 && - observed_time_ms <= - std::numeric_limits::max() - interval) { - from_time_ms = observed_time_ms + interval; + } else if (from_time_ms == 0 && observed_time_ms > 0) { + from_time_ms = observed_time_ms; } + from_time_ms = align_tick_history_start_time_ms( + from_time_ms, + history_interval_ms); + if (from_time_ms == 0 || from_time_ms > target_time_ms) { continuity.phase = continuity.unverified_from_time_ms == 0 ? MarketDataContinuityPhase::LIVE @@ -1913,17 +1938,14 @@ namespace optionx::market_data { notify_live = continuity.unverified_from_time_ms == 0; notify_degraded = !notify_live; } else { - auto request_to_time_ms = target_time_ms; - const auto max_backfill_ms = entry->tick_continuity.max_backfill_ms; - if (max_backfill_ms > 0) { - const auto span = max_backfill_ms > interval - ? max_backfill_ms - interval - : 0U; - const auto bounded = from_time_ms > - std::numeric_limits::max() - span - ? std::numeric_limits::max() - : from_time_ms + span; - request_to_time_ms = std::min(request_to_time_ms, bounded); + auto request_to_time_ms = bounded_tick_history_to_time_ms( + from_time_ms, + target_time_ms, + entry->tick_continuity.max_backfill_ms, + history_interval_ms); + if (request_to_time_ms < from_time_ms && + history_interval_ms > 0) { + request_to_time_ms = from_time_ms; } if (request_to_time_ms < from_time_ms) return; @@ -2439,6 +2461,62 @@ namespace optionx::market_data { continuity.unverified_from_time_ms = 0; } + inline std::uint64_t MarketDataRouterState::provider_time_ms( + BaseMarketDataProvider& provider) noexcept { + const auto provider_time = provider.provider_time_ms(); + if (provider_time > 0) return provider_time; + const auto local_time = static_cast(OPTIONX_TIMESTAMP_MS); + return local_time > 0 ? local_time : 1U; + } + + inline std::uint64_t MarketDataRouterState::tick_history_interval_ms( + BaseMarketDataProvider& provider) noexcept { + return provider.tick_history_interval_ms(); + } + + inline std::uint64_t MarketDataRouterState::align_tick_history_start_time_ms( + std::uint64_t time_ms, + std::uint64_t interval_ms) noexcept { + if (time_ms == 0 || interval_ms == 0) return time_ms; + return time_ms - time_ms % interval_ms; + } + + inline std::uint64_t MarketDataRouterState::align_tick_history_end_time_ms( + std::uint64_t time_ms, + std::uint64_t interval_ms) noexcept { + if (time_ms == 0 || interval_ms == 0) return time_ms; + const auto remainder = time_ms % interval_ms; + if (remainder == 0) return time_ms; + const auto increment = interval_ms - remainder; + return time_ms > std::numeric_limits::max() - increment + ? std::numeric_limits::max() + : time_ms + increment; + } + + inline std::uint64_t MarketDataRouterState::bounded_tick_history_to_time_ms( + std::uint64_t from_time_ms, + std::uint64_t target_time_ms, + std::uint64_t max_backfill_ms, + std::uint64_t interval_ms) noexcept { + if (max_backfill_ms == 0) { + return align_tick_history_end_time_ms(target_time_ms, interval_ms); + } + const auto span = max_backfill_ms - 1U; + const auto bounded = from_time_ms > + std::numeric_limits::max() - span + ? std::numeric_limits::max() + : from_time_ms + span; + const auto target_end = align_tick_history_end_time_ms( + target_time_ms, + interval_ms); + if (target_end <= bounded) return target_end; + if (interval_ms == 0) return bounded; + const auto aligned_bound = align_tick_history_start_time_ms( + bounded, + interval_ms); + return aligned_bound >= from_time_ms ? aligned_bound : from_time_ms; + } + inline bool MarketDataRouterState::tick_history_covers_range( const TickDataBatch& batch, std::uint64_t from_time_ms, @@ -3182,26 +3260,38 @@ namespace optionx::market_data { if (usable_history && kind != ContinuityRequestKind::PREFILL && continuity.reconnect_target_time_ms > to_time_ms) { - const auto interval = entry->tick_continuity.expected_interval_ms; - const auto next_from = interval > 0 && - to_time_ms <= - std::numeric_limits::max() - interval - ? to_time_ms + interval - : to_time_ms == std::numeric_limits::max() - ? to_time_ms - : to_time_ms + 1U; - if (next_from <= to_time_ms) return; - auto next_to = continuity.reconnect_target_time_ms; const auto max_backfill_ms = entry->tick_continuity.max_backfill_ms; - if (max_backfill_ms > 0) { - const auto span = max_backfill_ms > interval - ? max_backfill_ms - interval - : 0U; - const auto bounded = next_from > - std::numeric_limits::max() - span - ? std::numeric_limits::max() - : next_from + span; - next_to = std::min(next_to, bounded); + const auto history_interval_ms = + tick_history_interval_ms(*entry->provider); + // Inclusive ranges overlap at the boundary whenever the + // configured limit can contain that overlap. If a + // provider grid is wider than the limit, advance to the + // next representable boundary to avoid repeating a range. + const bool can_overlap = max_backfill_ms == 0 || + (history_interval_ms == 0 + ? max_backfill_ms > 1U + : max_backfill_ms > history_interval_ms); + std::uint64_t next_from = to_time_ms; + if (!can_overlap) { + const auto step = history_interval_ms > 0 + ? history_interval_ms + : 1U; + if (to_time_ms > + std::numeric_limits::max() - step) { + return; + } + next_from = to_time_ms + step; + } + if (next_from > continuity.reconnect_target_time_ms) { + return; + } + auto next_to = bounded_tick_history_to_time_ms( + next_from, + continuity.reconnect_target_time_ms, + max_backfill_ms, + history_interval_ms); + if (next_to < next_from && history_interval_ms > 0) { + next_to = next_from; } if (next_to < next_from) return; next_request.router_id = router_id; @@ -3643,7 +3733,6 @@ namespace optionx::market_data { continuity.reconnect_target_time_ms = 0; if (entry->tick_continuity.recovers_gaps()) { - const auto interval = entry->tick_continuity.expected_interval_ms; if (continuity.initial_prefill_pending && continuity.initial_prefill_boundary_time_ms > 0) { continuity.unverified_from_time_ms = @@ -3653,25 +3742,20 @@ namespace optionx::market_data { continuity.unverified_from_time_ms, continuity.initial_prefill_boundary_time_ms); } else if (continuity.verified_through_time_ms > 0 && - interval > 0 && - continuity.verified_through_time_ms <= - std::numeric_limits::max() - interval) { + continuity.verified_through_time_ms > 0) { continuity.unverified_from_time_ms = continuity.unverified_from_time_ms == 0 - ? continuity.verified_through_time_ms + interval + ? continuity.verified_through_time_ms : std::min( continuity.unverified_from_time_ms, - continuity.verified_through_time_ms + interval); + continuity.verified_through_time_ms); } else if (continuity.last_observed_time_ms > 0) { - const auto next = interval > 0 && - continuity.last_observed_time_ms <= - std::numeric_limits::max() - interval - ? continuity.last_observed_time_ms + interval - : continuity.last_observed_time_ms; continuity.unverified_from_time_ms = continuity.unverified_from_time_ms == 0 - ? next - : std::min(continuity.unverified_from_time_ms, next); + ? continuity.last_observed_time_ms + : std::min( + continuity.unverified_from_time_ms, + continuity.last_observed_time_ms); } } continuity.phase = MarketDataContinuityPhase::WAITING_FOR_READY; @@ -4310,6 +4394,7 @@ namespace optionx::market_data { } const auto interval = entry->tick_continuity.expected_interval_ms; + const auto history_interval_ms = tick_history_interval_ms(*entry->provider); if (allow_gap_recovery && entry->tick_continuity.recovers_gaps() && !continuity.request_in_flight && interval > 0 && continuity.last_observed_time_ms > 0) { @@ -4323,24 +4408,25 @@ namespace optionx::market_data { ? previous_time_ms + interval : std::numeric_limits::max(); if (tick.time_ms > expected_time_ms) { - const auto gap_from_time_ms = expected_time_ms; - const auto gap_to_time_ms = tick.time_ms - interval; - if (gap_to_time_ms < gap_from_time_ms) return false; - - auto request_to_time_ms = gap_to_time_ms; - const auto max_backfill_ms = - entry->tick_continuity.max_backfill_ms; - if (max_backfill_ms > 0) { - const auto span = max_backfill_ms > interval - ? max_backfill_ms - interval - : 0U; - const auto bounded = gap_from_time_ms > - std::numeric_limits::max() - span - ? std::numeric_limits::max() - : gap_from_time_ms + span; - request_to_time_ms = std::min(request_to_time_ms, bounded); + // `expected_interval_ms` only identifies a suspicious + // distance. Ticks are events, so recovery must cover + // the continuous inclusive interval and retain the + // triggering observation in the buffer. + const auto gap_from_time_ms = previous_time_ms; + const auto gap_to_time_ms = tick.time_ms; + const auto request_from_time_ms = align_tick_history_start_time_ms( + gap_from_time_ms, + history_interval_ms); + auto request_to_time_ms = bounded_tick_history_to_time_ms( + request_from_time_ms, + gap_to_time_ms, + entry->tick_continuity.max_backfill_ms, + history_interval_ms); + if (request_to_time_ms < request_from_time_ms && + history_interval_ms > 0) { + request_to_time_ms = request_from_time_ms; } - if (request_to_time_ms < gap_from_time_ms) return false; + if (request_to_time_ms < request_from_time_ms) return false; if (index > 0) { TickDataBatch prefix = routed; @@ -4362,7 +4448,7 @@ namespace optionx::market_data { return false; } - mark_tick_unverified_no_lock(entry, gap_from_time_ms); + mark_tick_unverified_no_lock(entry, request_from_time_ms); continuity.phase = MarketDataContinuityPhase::RECOVERING; continuity.request_in_flight = true; continuity.reconnect_target_time_ms = gap_to_time_ms; @@ -4371,11 +4457,11 @@ namespace optionx::market_data { entry->control->provider_subscription, TickHistoryRequest( entry->stream.symbol, - gap_from_time_ms, + request_from_time_ms, request_to_time_ms), ContinuityRequestKind::GAP_BACKFILL, true, - gap_from_time_ms, + request_from_time_ms, request_to_time_ms, 0, 1}); diff --git a/include/optionx_cpp/platforms/IntradeBarPlatform.hpp b/include/optionx_cpp/platforms/IntradeBarPlatform.hpp index c1a9a75..6648d5f 100644 --- a/include/optionx_cpp/platforms/IntradeBarPlatform.hpp +++ b/include/optionx_cpp/platforms/IntradeBarPlatform.hpp @@ -168,6 +168,17 @@ namespace optionx::platforms { return true; } + /// \brief Returns the broker-aligned Intrade market-data time estimate. + std::uint64_t provider_time_ms() const noexcept override { + return m_tick_history.provider_time_ms( + static_cast(OPTIONX_TIMESTAMP_MS)); + } + + /// \brief Returns the sampling grid used by observed Intrade history. + std::uint64_t tick_history_interval_ms() const noexcept override { + return m_tick_history.options().sampling_interval_ms; + } + /// \brief Returns the live bar data callback. market_data::BaseMarketDataProvider::bars_callback_t& on_bar_data() override { return m_bar_data_callback; diff --git a/include/optionx_cpp/platforms/IntradeBarPlatform/ObservedTickHistory.hpp b/include/optionx_cpp/platforms/IntradeBarPlatform/ObservedTickHistory.hpp index ea5326b..d566128 100644 --- a/include/optionx_cpp/platforms/IntradeBarPlatform/ObservedTickHistory.hpp +++ b/include/optionx_cpp/platforms/IntradeBarPlatform/ObservedTickHistory.hpp @@ -32,6 +32,10 @@ namespace optionx::platforms::intrade_bar { /// It is used only to make a conservative completeness assertion. std::uint64_t sampling_interval_ms = 1000; + /// Number of recent local/provider offset samples used by the median + /// broker-clock estimate. + std::size_t clock_offset_sample_count = 9; + /// \brief Returns true when the archive can be used safely. [[nodiscard]] bool valid() const noexcept { return max_items_per_symbol > 0 && sampling_interval_ms > 0; @@ -71,6 +75,7 @@ namespace optionx::platforms::intrade_bar { if (tick.time_ms == 0 || contains_observation(history, tick)) { continue; } + record_clock_offset(tick); history.items.push_back(StoredTick{tick}); } prune(history); @@ -133,6 +138,7 @@ namespace optionx::platforms::intrade_bar { void clear() noexcept { std::lock_guard lock(m_mutex); m_symbols.clear(); + m_clock_offsets_ms.clear(); } /// \brief Returns the number of retained observations for a symbol. @@ -147,6 +153,24 @@ namespace optionx::platforms::intrade_bar { return m_options; } + /// \brief Returns the current broker-aligned time estimate. + /// \param local_time_ms Current local wall-clock time in milliseconds. + /// \details The estimate uses the median of recent + /// `tick.time_ms - tick.received_ms` observations and is then + /// rounded down to the provider's sampling grid. With no + /// usable samples this still provides a conservative aligned + /// local-time fallback. + [[nodiscard]] std::uint64_t provider_time_ms( + std::uint64_t local_time_ms) const noexcept { + std::lock_guard lock(m_mutex); + const auto adjusted = add_offset( + local_time_ms, + median_clock_offset_no_lock()); + const auto interval = m_options.sampling_interval_ms; + if (interval == 0) return adjusted; + return adjusted - adjusted % interval; + } + private: struct StoredTick { Tick tick; @@ -180,6 +204,80 @@ namespace optionx::platforms::intrade_bar { }); } + static std::int64_t observation_offset( + const Tick& tick) noexcept { + const auto max_int64 = static_cast( + (std::numeric_limits::max)()); + if (tick.time_ms > max_int64 || tick.received_ms > max_int64) { + return 0; + } + const auto broker_time = static_cast(tick.time_ms); + const auto received_time = static_cast(tick.received_ms); + if (broker_time >= received_time) { + return broker_time - received_time; + } + const auto difference = received_time - broker_time; + return difference == (std::numeric_limits::max)() + ? (std::numeric_limits::min)() + : -difference; + } + + void record_clock_offset(const Tick& tick) { + if (tick.received_ms == 0 || m_options.clock_offset_sample_count == 0) { + return; + } + m_clock_offsets_ms.push_back(observation_offset(tick)); + while (m_clock_offsets_ms.size() > m_options.clock_offset_sample_count) { + m_clock_offsets_ms.pop_front(); + } + } + + [[nodiscard]] std::int64_t median_clock_offset_no_lock() const noexcept { + if (m_clock_offsets_ms.empty()) return 0; + const auto select = [this](std::size_t rank) noexcept { + std::int64_t selected = 0; + bool has_selected = false; + for (const auto candidate : m_clock_offsets_ms) { + std::size_t not_greater = 0; + for (const auto value : m_clock_offsets_ms) { + if (value <= candidate) ++not_greater; + } + if (not_greater > rank && + (!has_selected || candidate < selected)) { + selected = candidate; + has_selected = true; + } + } + return selected; + }; + + const auto middle = m_clock_offsets_ms.size() / 2U; + if (m_clock_offsets_ms.size() % 2U != 0) { + return select(middle); + } + const auto lower = select(middle - 1U); + const auto upper = select(middle); + return static_cast( + (static_cast(lower) + + static_cast(upper)) / 2.0L); + } + + static std::uint64_t add_offset( + std::uint64_t local_time_ms, + std::int64_t offset_ms) noexcept { + if (offset_ms >= 0) { + const auto positive = static_cast(offset_ms); + return local_time_ms > + (std::numeric_limits::max)() - positive + ? (std::numeric_limits::max)() + : local_time_ms + positive; + } + const auto negative = offset_ms == (std::numeric_limits::min)() + ? static_cast((std::numeric_limits::max)()) + 1U + : static_cast(-offset_ms); + return local_time_ms < negative ? 0U : local_time_ms - negative; + } + void prune(SymbolHistory& history) { const auto max_items = m_options.max_items_per_symbol; while (history.items.size() > max_items) { @@ -245,6 +343,7 @@ namespace optionx::platforms::intrade_bar { IntradeObservedTickHistoryOptions m_options; mutable std::mutex m_mutex; std::unordered_map m_symbols; + std::deque m_clock_offsets_ms; }; } // namespace optionx::platforms::intrade_bar diff --git a/tests/intrade_observed_tick_history_test.cpp b/tests/intrade_observed_tick_history_test.cpp index fc8170f..4bf1b02 100644 --- a/tests/intrade_observed_tick_history_test.cpp +++ b/tests/intrade_observed_tick_history_test.cpp @@ -139,6 +139,21 @@ TEST(IntradeObservedTickHistory, ClearDropsSessionObservations) { EXPECT_EQ(archive.size("EURUSD"), 0U); } +TEST(IntradeObservedTickHistory, EstimatesBrokerAlignedTimeFromMedianOffsets) { + IntradeObservedTickHistoryOptions options; + options.clock_offset_sample_count = 5; + IntradeObservedTickHistory archive(options); + archive.record({make_batch( + "EURUSD", + {make_tick(1.1000, 1.1002, 1000, 1123), + make_tick(1.1001, 1.1003, 2000, 2123), + make_tick(1.1002, 1.1004, 3000, 3123), + make_tick(1.1003, 1.1005, 4000, 1000), + make_tick(1.1004, 1.1006, 5000, 5123)})}); + + EXPECT_EQ(archive.provider_time_ms(4123), 4000U); +} + int main(int argc, char** argv) { ::testing::InitGoogleTest(&argc, argv); return RUN_ALL_TESTS(); diff --git a/tests/market_data_tick_continuity_test.cpp b/tests/market_data_tick_continuity_test.cpp index d0a66cc..73c654f 100644 --- a/tests/market_data_tick_continuity_test.cpp +++ b/tests/market_data_tick_continuity_test.cpp @@ -38,6 +38,9 @@ class ScopedTestClock { class FakeTickHistoryProvider final : public BaseMarketDataProvider { public: + std::uint64_t provider_now_ms = 0; + std::uint64_t history_interval_ms = 0; + ticks_callback_t& on_tick_data() override { return m_tick_callback; } @@ -46,6 +49,14 @@ class FakeTickHistoryProvider final : public BaseMarketDataProvider { return m_status_callback; } + std::uint64_t provider_time_ms() const noexcept override { + return provider_now_ms; + } + + std::uint64_t tick_history_interval_ms() const noexcept override { + return history_interval_ms; + } + bool subscribe_ticks( TickSubscriptionRequest request, subscription_callback_t callback) override { @@ -245,6 +256,27 @@ TEST(MarketDataTickContinuity, CompletesPrefillAndDeliversDirectRealtime) { EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); } +TEST(MarketDataTickContinuity, UsesProviderAlignedClockForPrefill) { + ScopedTestClock clock(3123); + FakeTickHistoryProvider provider; + provider.provider_now_ms = 3000; + provider.history_interval_ms = 1000; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + ASSERT_EQ(provider.history_requests.size(), 1U); + EXPECT_EQ(provider.history_requests.front().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.front().to_time_ms, 3000U); + + provider.complete_history(make_history({1000, 2000, 3000})); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 1U); +} + TEST(MarketDataTickContinuity, IncompletePrefillStaysDegraded) { ScopedTestClock clock(3000); FakeTickHistoryProvider provider; @@ -325,9 +357,9 @@ TEST(MarketDataTickContinuity, ChecksGapsInsidePrefillBacklog) { provider.complete_history(make_history({1000, 2000, 3000})); ASSERT_EQ(provider.history_requests.size(), 2U); - EXPECT_EQ(provider.history_requests.back().from_time_ms, 5000U); - EXPECT_EQ(provider.history_requests.back().to_time_ms, 5000U); - provider.complete_history(make_history({5000})); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 4000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 6000U); + provider.complete_history(make_history({4000, 5000})); ASSERT_FALSE(subscriber->ticks.empty()); const auto& prefix = subscriber->ticks[1]; @@ -364,14 +396,16 @@ TEST(MarketDataTickContinuity, RepairsGapAndPreservesDistinctSameSecondTicks) { provider.emit_ticks({first, second}); ASSERT_EQ(provider.history_requests.size(), 2U); - EXPECT_EQ(provider.history_requests.back().from_time_ms, 5000U); - EXPECT_EQ(provider.history_requests.back().to_time_ms, 5000U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 4000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 6000U); const auto before_repair_ticks = subscriber->ticks.size(); - provider.complete_history(make_history({5000})); + provider.complete_history(make_history({4000, 5000})); ASSERT_EQ(subscriber->ticks.size(), before_repair_ticks + 2U); const auto& history_batch = subscriber->ticks[before_repair_ticks]; - EXPECT_TRUE(only_tick(history_batch).has_flag(MarketDataFlags::HISTORICAL)); + ASSERT_EQ(history_batch.items.size(), 2U); + EXPECT_TRUE(history_batch.items[0].has_flag(MarketDataFlags::HISTORICAL)); + EXPECT_TRUE(history_batch.items[1].has_flag(MarketDataFlags::HISTORICAL)); const auto& catchup_batch = subscriber->ticks[before_repair_ticks + 1U]; ASSERT_EQ(catchup_batch.items.size(), 2U); EXPECT_EQ(catchup_batch.items[0].time_ms, 6000U); @@ -382,6 +416,112 @@ TEST(MarketDataTickContinuity, RepairsGapAndPreservesDistinctSameSecondTicks) { EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 2U); } +TEST(MarketDataTickContinuity, PreservesNonGridTriggeringTickDuringRecovery) { + ScopedTestClock clock(1000); + FakeTickHistoryProvider provider; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + provider.complete_history(make_history({1000})); + + provider.emit_ticks({make_tick(2501)}); + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 2501U); + + provider.complete_history(make_history({1000, 2000})); + + bool saw_triggering_tick = false; + for (const auto& batch : subscriber->ticks) { + for (const auto& tick : batch.items) { + if (tick.time_ms == 2501U) saw_triggering_tick = true; + } + } + EXPECT_TRUE(saw_triggering_tick); +} + +TEST(MarketDataTickContinuity, AlignsProviderGridRecoveryEndOutward) { + ScopedTestClock clock(1000); + FakeTickHistoryProvider provider; + provider.history_interval_ms = 1000; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request()); + ASSERT_TRUE(route.valid()); + ASSERT_EQ(provider.history_requests.size(), 1U); + provider.complete_history(make_history({1000})); + + provider.emit_ticks({make_tick(2501)}); + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 3000U); + + provider.complete_history(make_history({1000, 2000})); + bool saw_triggering_tick = false; + for (const auto& batch : subscriber->ticks) { + for (const auto& tick : batch.items) { + if (tick.time_ms == 2501U) saw_triggering_tick = true; + } + } + EXPECT_TRUE(saw_triggering_tick); +} + +TEST(MarketDataTickContinuity, KeepsBoundedProviderGridChunksOverlapped) { + ScopedTestClock clock(1000); + FakeTickHistoryProvider provider; + provider.history_interval_ms = 1000; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request( + MarketDataContinuityMode::PREFILL_AND_RECOVER, + 1, + 1500)); + ASSERT_TRUE(route.valid()); + provider.complete_history(make_history({1000})); + provider.emit_ticks({make_tick(4501)}); + + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 2000U); + provider.complete_history(make_history({1000, 2000})); + + ASSERT_EQ(provider.history_requests.size(), 3U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 2000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 3000U); + provider.complete_history(make_history({2000, 3000})); + + ASSERT_EQ(provider.history_requests.size(), 4U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 3000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 4000U); + provider.complete_history(make_history({3000, 4000})); + + ASSERT_EQ(provider.history_requests.size(), 5U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 4000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 5000U); + provider.complete_history(make_history({4000, 5000})); + + bool saw_triggering_tick = false; + for (const auto& batch : subscriber->ticks) { + for (const auto& tick : batch.items) { + if (tick.time_ms == 4501U) saw_triggering_tick = true; + } + } + EXPECT_TRUE(saw_triggering_tick); +} + TEST(MarketDataTickContinuity, UsesBoundedGapRequests) { ScopedTestClock clock(10000); FakeTickHistoryProvider provider; @@ -398,23 +538,33 @@ TEST(MarketDataTickContinuity, UsesBoundedGapRequests) { provider.emit_ticks({make_tick(10000)}); ASSERT_EQ(provider.history_requests.size(), 1U); - EXPECT_EQ(provider.history_requests.back().from_time_ms, 2000U); - EXPECT_EQ(provider.history_requests.back().to_time_ms, 4000U); - provider.complete_history(make_history({2000, 3000, 4000})); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 3999U); + provider.complete_history(make_history({1000, 2000, 3000, 3999})); ASSERT_EQ(provider.history_requests.size(), 2U); - EXPECT_EQ(provider.history_requests.back().from_time_ms, 5000U); - EXPECT_EQ(provider.history_requests.back().to_time_ms, 7000U); - provider.complete_history(make_history({5000, 6000, 7000})); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 3999U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 6998U); + provider.complete_history(make_history({3999, 5000, 6000, 6998})); ASSERT_EQ(provider.history_requests.size(), 3U); - EXPECT_EQ(provider.history_requests.back().from_time_ms, 8000U); - EXPECT_EQ(provider.history_requests.back().to_time_ms, 9000U); - provider.complete_history(make_history({8000, 9000})); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 6998U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 9997U); + provider.complete_history(make_history({6998, 8000, 9000, 9997})); + ASSERT_EQ(provider.history_requests.size(), 4U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 9997U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 10000U); + provider.complete_history(make_history({9997, 10000})); const auto snapshot = router.continuity_snapshot(route.router_id()); ASSERT_TRUE(snapshot.has_value()); EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::LIVE); EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); - EXPECT_EQ(snapshot->history_request_count, 3U); + EXPECT_EQ(snapshot->history_request_count, 4U); + + for (std::size_t index = 1; index < provider.history_requests.size(); ++index) { + EXPECT_EQ( + provider.history_requests[index].from_time_ms, + provider.history_requests[index - 1U].to_time_ms); + } } TEST(MarketDataTickContinuity, RetriesFailedHistoryFromProcess) { @@ -470,10 +620,10 @@ TEST(MarketDataTickContinuity, ReconnectDeduplicatesOnlyExactOverlap) { provider.emit_status(MarketDataStreamStatus::READY); ASSERT_EQ(provider.history_requests.size(), 2U); - EXPECT_EQ(provider.history_requests.back().from_time_ms, 4000U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 3000U); EXPECT_EQ(provider.history_requests.back().to_time_ms, 5000U); - TickSequence history = make_history({4000}); + TickSequence history = make_history({3000, 4000}); history.ticks.push_back(make_tick(5000, 1.0)); provider.complete_history(std::move(history)); From e48543f7e76d82c24c0357e621387a084561a2af Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Mon, 7 Sep 2026 20:58:59 +0300 Subject: [PATCH 3/5] test(ci): run tick continuity regression suite --- .github/workflows/ubuntu-smoke.yml | 10 ++++++++++ .github/workflows/windows-smoke.yml | 7 +++++++ 2 files changed, 17 insertions(+) diff --git a/.github/workflows/ubuntu-smoke.yml b/.github/workflows/ubuntu-smoke.yml index d9f6485..638e3cf 100644 --- a/.github/workflows/ubuntu-smoke.yml +++ b/.github/workflows/ubuntu-smoke.yml @@ -212,6 +212,16 @@ jobs: - name: Run market_data_tick_history_contract_test run: ./build-linux/market_data_tick_history_contract_test --gtest_brief=1 + - name: Build market data tick continuity and observed history tests + run: > + cmake --build build-linux --target + market_data_tick_continuity_test intrade_observed_tick_history_test -j + + - name: Run market data tick continuity and observed history tests + run: | + ./build-linux/market_data_tick_continuity_test --gtest_brief=1 + ./build-linux/intrade_observed_tick_history_test --gtest_brief=1 + - name: Build market_data_subscriber_base_test and example run: cmake --build build-linux --target market_data_subscriber_base_test market_data_subscriber_base_example -j diff --git a/.github/workflows/windows-smoke.yml b/.github/workflows/windows-smoke.yml index df66f89..f2f3546 100644 --- a/.github/workflows/windows-smoke.yml +++ b/.github/workflows/windows-smoke.yml @@ -103,6 +103,11 @@ jobs: cmake --build build-windows --config Debug --target market_data_tick_history_contract_test + - name: Build market data tick continuity and observed history tests + run: > + cmake --build build-windows --config Debug --target + market_data_tick_continuity_test intrade_observed_tick_history_test + - name: Build TradingView extension bridge smoke example run: cmake --build build-windows --config Debug --target tradingview_extension_bridge_smoke @@ -126,6 +131,8 @@ jobs: .\build-windows\Debug\market_data_continuity_test.exe --gtest_brief=1 .\build-windows\Debug\market_data_continuity_example.exe .\build-windows\Debug\market_data_tick_history_contract_test.exe --gtest_brief=1 + .\build-windows\Debug\market_data_tick_continuity_test.exe --gtest_brief=1 + .\build-windows\Debug\intrade_observed_tick_history_test.exe --gtest_brief=1 .\build-windows\Debug\metatrader_file_bridge_smoke.exe --self-test .\build-windows\Debug\metatrader_file_command_writer_smoke.exe --self-test .\build-windows\Debug\metatrader_file_end_to_end_smoke.exe --self-test From 043db80a0da83d82411a08e046dd7789e15cc93b Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Mon, 7 Sep 2026 21:22:40 +0300 Subject: [PATCH 4/5] test(intrade): keep observed history test lightweight --- tests/intrade_observed_tick_history_test.cpp | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/tests/intrade_observed_tick_history_test.cpp b/tests/intrade_observed_tick_history_test.cpp index 4bf1b02..d897774 100644 --- a/tests/intrade_observed_tick_history_test.cpp +++ b/tests/intrade_observed_tick_history_test.cpp @@ -6,7 +6,11 @@ #include #include -#include +#include +#include +#include +#include +#include using namespace optionx; using namespace optionx::events; From e66a0882bcfa845450506c5fc509b27cfded2328 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Tue, 8 Sep 2026 13:17:33 +0300 Subject: [PATCH 5/5] fix(market-data): stop grid recovery at completed boundary --- guides/market-data-router.md | 15 +++++------ guides/market-data-router.ru.md | 14 ++++++----- .../market_data/detail/MarketDataRouter.ipp | 25 +++++++++++-------- tests/market_data_tick_continuity_test.cpp | 12 ++++----- 4 files changed, 36 insertions(+), 30 deletions(-) diff --git a/guides/market-data-router.md b/guides/market-data-router.md index 46f7c93..1600c04 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -686,13 +686,14 @@ proves continuity. Recovery requests use inclusive timestamp ranges. For an event-oriented provider, the suspicious interval is requested without inventing missing tick slots, and the live tick that triggered recovery remains buffered. When a -provider declares a history grid, Router aligns the start down and (for an -unbounded request) the end up so an off-grid observation is not excluded. -Bounded chunks keep their size limit and overlap at the previous end point -whenever that overlap can advance the range; if the limit is smaller than a -provider grid step, Router advances to the next provider boundary instead of -repeating the same request. The overlap is removed only by exact observation -identity. +provider declares a history grid, Router aligns the start down and the history +end down to the last completed provider boundary. An off-grid live observation +is not requested as a future history sample: it stays in the continuity buffer +and is released after the completed range is verified. Bounded chunks keep +their size limit and overlap at the previous end point whenever that overlap +can advance the range; if the limit is smaller than a provider grid step, +Router advances to the next provider boundary instead of repeating the same +request. The overlap is removed only by exact observation identity. The Router sends historical ticks first, marks them `HISTORICAL`, and then replays held live ticks as `LIVE_SOURCE | CATCHUP`. A complete result is required diff --git a/guides/market-data-router.ru.md b/guides/market-data-router.ru.md index 3db4f50..b690c9e 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -842,12 +842,14 @@ ticks чаще одной секунды допустимы. Для доказа Recovery использует inclusive timestamp ranges. Для event-oriented provider Router не синтезирует пропущенные tick slots, а сохраняет в buffer live tick, который запустил recovery. Если provider объявляет history grid, Router -округляет начало вниз, а конец вверх для unbounded request, чтобы observation -с timestamp между grid points не выпала из диапазона. Bounded chunks сохраняют -лимит размера и перекрываются в предыдущей конечной точке, когда такой overlap -позволяет продвинуть диапазон. Если лимит меньше шага provider grid, Router -переходит к следующей provider boundary, а не повторяет тот же запрос. Overlap -удаляется только по exact observation identity. +округляет начало вниз, а конец history - вниз до последней завершённой +provider boundary. Off-grid live observation не запрашивается как будущий +history sample: она остаётся в continuity buffer и выпускается после проверки +завершённого диапазона. Bounded chunks сохраняют лимит размера и перекрываются +в предыдущей конечной точке, когда такой overlap позволяет продвинуть диапазон. +Если лимит меньше шага provider grid, Router переходит к следующей provider +boundary, а не повторяет тот же запрос. Overlap удаляется только по exact +observation identity. Router сначала отправляет historical ticks с флагом `HISTORICAL`, затем воспроизводит удержанные live ticks с флагами `LIVE_SOURCE | CATCHUP`. До diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index 8276987..96c2335 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -2485,12 +2485,11 @@ namespace optionx::market_data { std::uint64_t time_ms, std::uint64_t interval_ms) noexcept { if (time_ms == 0 || interval_ms == 0) return time_ms; - const auto remainder = time_ms % interval_ms; - if (remainder == 0) return time_ms; - const auto increment = interval_ms - remainder; - return time_ms > std::numeric_limits::max() - increment - ? std::numeric_limits::max() - : time_ms + increment; + // History completeness can only be established through the last + // provider boundary that has already occurred. An off-grid live + // observation remains in the continuity buffer until it is + // released after this completed range is verified. + return time_ms - time_ms % interval_ms; } inline std::uint64_t MarketDataRouterState::bounded_tick_history_to_time_ms( @@ -3257,12 +3256,16 @@ namespace optionx::market_data { } const auto& entry = entry_it->second; auto& continuity = entry->tick_continuity_state; + const auto history_interval_ms = + tick_history_interval_ms(*entry->provider); + const auto history_target_time_ms = + align_tick_history_end_time_ms( + continuity.reconnect_target_time_ms, + history_interval_ms); if (usable_history && kind != ContinuityRequestKind::PREFILL && - continuity.reconnect_target_time_ms > to_time_ms) { + history_target_time_ms > to_time_ms) { const auto max_backfill_ms = entry->tick_continuity.max_backfill_ms; - const auto history_interval_ms = - tick_history_interval_ms(*entry->provider); // Inclusive ranges overlap at the boundary whenever the // configured limit can contain that overlap. If a // provider grid is wider than the limit, advance to the @@ -3282,12 +3285,12 @@ namespace optionx::market_data { } next_from = to_time_ms + step; } - if (next_from > continuity.reconnect_target_time_ms) { + if (next_from > history_target_time_ms) { return; } auto next_to = bounded_tick_history_to_time_ms( next_from, - continuity.reconnect_target_time_ms, + history_target_time_ms, max_backfill_ms, history_interval_ms); if (next_to < next_from && history_interval_ms > 0) { diff --git a/tests/market_data_tick_continuity_test.cpp b/tests/market_data_tick_continuity_test.cpp index 73c654f..f1a4839 100644 --- a/tests/market_data_tick_continuity_test.cpp +++ b/tests/market_data_tick_continuity_test.cpp @@ -445,7 +445,7 @@ TEST(MarketDataTickContinuity, PreservesNonGridTriggeringTickDuringRecovery) { EXPECT_TRUE(saw_triggering_tick); } -TEST(MarketDataTickContinuity, AlignsProviderGridRecoveryEndOutward) { +TEST(MarketDataTickContinuity, StopsProviderGridRecoveryAtCompletedBoundary) { ScopedTestClock clock(1000); FakeTickHistoryProvider provider; provider.history_interval_ms = 1000; @@ -463,9 +463,11 @@ TEST(MarketDataTickContinuity, AlignsProviderGridRecoveryEndOutward) { provider.emit_ticks({make_tick(2501)}); ASSERT_EQ(provider.history_requests.size(), 2U); EXPECT_EQ(provider.history_requests.back().from_time_ms, 1000U); - EXPECT_EQ(provider.history_requests.back().to_time_ms, 3000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 2000U); provider.complete_history(make_history({1000, 2000})); + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 2U); bool saw_triggering_tick = false; for (const auto& batch : subscriber->ticks) { for (const auto& tick : batch.items) { @@ -508,10 +510,8 @@ TEST(MarketDataTickContinuity, KeepsBoundedProviderGridChunksOverlapped) { EXPECT_EQ(provider.history_requests.back().to_time_ms, 4000U); provider.complete_history(make_history({3000, 4000})); - ASSERT_EQ(provider.history_requests.size(), 5U); - EXPECT_EQ(provider.history_requests.back().from_time_ms, 4000U); - EXPECT_EQ(provider.history_requests.back().to_time_ms, 5000U); - provider.complete_history(make_history({4000, 5000})); + ASSERT_EQ(provider.history_requests.size(), 4U); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 2U); bool saw_triggering_tick = false; for (const auto& batch : subscriber->ticks) {