From 1cc1ab704cd91c7cae252cb585141ea68bd1b20d Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Tue, 8 Sep 2026 22:51:31 +0300 Subject: [PATCH] fix(market-data): deduplicate tick history overlaps --- guides/api-and-header-contracts.md | 10 +- guides/market-data-router.md | 29 ++- guides/market-data-router.ru.md | 31 ++- .../market_data/BaseMarketDataProvider.hpp | 11 + .../MarketDataContinuityOptions.hpp | 22 ++ .../market_data/detail/MarketDataRouter.ipp | 214 +++++++++++++++--- .../platforms/IntradeBarPlatform.hpp | 6 + ...market_data_subscription_contract_test.cpp | 20 ++ tests/market_data_tick_continuity_test.cpp | 146 +++++++++++- 9 files changed, 445 insertions(+), 44 deletions(-) diff --git a/guides/api-and-header-contracts.md b/guides/api-and-header-contracts.md index d9acd51..e85400b 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -393,9 +393,13 @@ broker: `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. + for `READY`; history overlap identity is selected by + `TickSubscriptionRequest::continuity.deduplication_mode`. `TIMESTAMP` uses + only `time_ms`, `TIME_AND_PRICES` adds quote/trade prices, and + `EXACT_OBSERVATION` also adds volume. `received_ms`/flags never distinguish + duplicates. `PROVIDER_DEFAULT` uses the provider hook, which is timestamp + based for Intrade; choose `EXACT_OBSERVATION` when same-second observations + must 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 1600c04..e6b8595 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -675,6 +675,24 @@ request.continuity.max_backfill_ms = 60'000; auto route = router.subscribe_ticks(provider, bot, request); ``` +History-overlap identity is configurable per route through +`TickSubscriptionRequest::continuity.deduplication_mode`. The default +`PROVIDER_DEFAULT` uses `BaseMarketDataProvider::tick_deduplication_mode()`; +Intrade uses `TIMESTAMP` because its broker observations are one-second +snapshots. A route can override that choice when it needs richer identity: + +```cpp +request.continuity.deduplication_mode = + md::MarketDataTickDeduplicationMode::TIME_AND_PRICES; +``` + +`TIMESTAMP` matches only `time_ms`, `TIME_AND_PRICES` matches `time_ms` plus +`ask`, `bid`, and `last`, and `EXACT_OBSERVATION` also matches `volume`. +`received_ms` and flags never identify an observation. The policy is applied +when resolving history overlap; it does not collapse distinct live events +before they are buffered or delivered. Use `EXACT_OBSERVATION` when distinct +same-timestamp snapshots must remain separate. + `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 @@ -693,7 +711,7 @@ 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. +request. History overlap is resolved with the selected tick identity policy. The Router sends historical ticks first, marks them `HISTORICAL`, and then replays held live ticks as `LIVE_SOURCE | CATCHUP`. A complete result is required @@ -704,10 +722,11 @@ 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 +the unresolved range through the latest observed time. History overlap is +removed according to the selected tick identity policy. The default provider +policy for Intrade is timestamp-based; use `EXACT_OBSERVATION` to retain +different observations with the same second. `received_ms` and flags do not +make an otherwise identical observation distinct. 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 diff --git a/guides/market-data-router.ru.md b/guides/market-data-router.ru.md index b690c9e..7816f89 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -832,6 +832,24 @@ request.continuity.max_backfill_ms = 60'000; auto route = router.subscribe_ticks(provider, bot, request); ``` +Identity history-overlap настраивается для каждого route через +`TickSubscriptionRequest::continuity.deduplication_mode`. Значение по умолчанию +`PROVIDER_DEFAULT` использует `BaseMarketDataProvider::tick_deduplication_mode()`; +для Intrade это `TIMESTAMP`, потому что broker observations имеют секундную +метку. При необходимости route может выбрать более подробную identity: + +```cpp +request.continuity.deduplication_mode = + md::MarketDataTickDeduplicationMode::TIME_AND_PRICES; +``` + +`TIMESTAMP` сравнивает только `time_ms`, `TIME_AND_PRICES` добавляет `ask`, +`bid` и `last`, а `EXACT_OBSERVATION` также сравнивает `volume`. +`received_ms` и flags не являются частью identity. Политика используется при +разрешении history overlap и не схлопывает разные live events до их +buffering или delivery. Если разные snapshots одной секунды должны сохраниться, +выберите `EXACT_OBSERVATION`. + `PREFILL` запрашивает заданный lookback до освобождения live ticks. `PREFILL_AND_RECOVER` дополнительно удерживает live tail, когда дельта между последовательными timestamps больше `expected_interval_ms`. Это только @@ -848,8 +866,8 @@ history sample: она остаётся в continuity buffer и выпускае завершённого диапазона. Bounded chunks сохраняют лимит размера и перекрываются в предыдущей конечной точке, когда такой overlap позволяет продвинуть диапазон. Если лимит меньше шага provider grid, Router переходит к следующей provider -boundary, а не повторяет тот же запрос. Overlap удаляется только по exact -observation identity. +boundary, а не повторяет тот же запрос. History overlap разрешается по +выбранной tick identity policy. Router сначала отправляет historical ticks с флагом `HISTORICAL`, затем воспроизводит удержанные live ticks с флагами `LIVE_SOURCE | CATCHUP`. До @@ -860,10 +878,11 @@ Router сначала отправляет historical ticks с флагом `HIS создания отдельного 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, публикует +unresolved range до последнего observed time. History overlap удаляется по +выбранной tick identity policy. Provider default для Intrade использует только +timestamp; чтобы сохранить разные observations той же секунды, выберите +`EXACT_OBSERVATION`. `received_ms` и flags не делают полностью одинаковый +snapshot новым. При переполнении buffer Router освобождает live data, публикует `FAILED`/`DEGRADED`, отключает continuity для этого route и возобновляет обычную live delivery. Если transport прервался во время initial prefill, после `READY` Router diff --git a/include/optionx_cpp/market_data/BaseMarketDataProvider.hpp b/include/optionx_cpp/market_data/BaseMarketDataProvider.hpp index 6b58056..2819704 100644 --- a/include/optionx_cpp/market_data/BaseMarketDataProvider.hpp +++ b/include/optionx_cpp/market_data/BaseMarketDataProvider.hpp @@ -106,6 +106,17 @@ namespace optionx::market_data { return 0; } + /// \brief Returns the provider's default tick identity policy. + /// \details Router uses this policy when a tick request leaves + /// `deduplication_mode` at PROVIDER_DEFAULT. The policy is + /// used only to remove observations repeated by inclusive + /// history overlap; it never collapses distinct live events + /// before they are buffered or delivered. + /// \return Provider-specific tick identity policy. + virtual MarketDataTickDeduplicationMode tick_deduplication_mode() const noexcept { + return MarketDataTickDeduplicationMode::EXACT_OBSERVATION; + } + /// \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/MarketDataContinuityOptions.hpp b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp index 53cf3fb..be4f8db 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp @@ -18,6 +18,15 @@ namespace optionx::market_data { PREFILL_AND_RECOVER ///< Prefill and repair live/reconnect timestamp gaps. }; + /// \enum MarketDataTickDeduplicationMode + /// \brief Selects which tick fields identify an already delivered observation. + enum class MarketDataTickDeduplicationMode { + PROVIDER_DEFAULT = 0, ///< Use the provider's identity contract. + TIMESTAMP, ///< Treat one timestamp as one observation. + TIME_AND_PRICES, ///< Match timestamp, ask, bid, and last price. + EXACT_OBSERVATION ///< Match timestamp, prices, and volume. + }; + /// \struct MarketDataContinuityRetryPolicy /// \brief Configures bounded history retry attempts and exponential backoff. struct MarketDataContinuityRetryPolicy { @@ -82,10 +91,23 @@ namespace optionx::market_data { MarketDataContinuityRetryPolicy retry; std::size_t max_buffered_batches = 1024; std::size_t max_buffered_items = 100000; + /// Identity policy used to remove inclusive history overlap. + /// Provider default keeps the policy provider-specific. + MarketDataTickDeduplicationMode deduplication_mode = + MarketDataTickDeduplicationMode::PROVIDER_DEFAULT; /// \brief Returns true when the option combination is usable. [[nodiscard]] bool valid() const noexcept { if (!retry.valid()) return false; + switch (deduplication_mode) { + case MarketDataTickDeduplicationMode::PROVIDER_DEFAULT: + case MarketDataTickDeduplicationMode::TIMESTAMP: + case MarketDataTickDeduplicationMode::TIME_AND_PRICES: + case MarketDataTickDeduplicationMode::EXACT_OBSERVATION: + break; + default: + return false; + } if (mode == MarketDataContinuityMode::LIVE_ONLY) { return prefill_lookback_ms == 0; } diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index 96c2335..05457e1 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -6,12 +6,47 @@ /// \brief Implements subscription-scoped market-data routing utilities. #include +#include #include +#include namespace optionx::market_data { namespace detail { + struct TickObservationKey { + std::uint64_t time_ms = 0; + double ask = 0.0; + double bid = 0.0; + double last = 0.0; + double volume = 0.0; + + [[nodiscard]] bool operator==( + const TickObservationKey& other) const noexcept { + return time_ms == other.time_ms && + ask == other.ask && + bid == other.bid && + last == other.last && + volume == other.volume; + } + }; + + struct TickObservationKeyHash { + [[nodiscard]] std::size_t operator()( + const TickObservationKey& key) const noexcept { + const auto combine = [](std::size_t seed, std::size_t value) noexcept { + return seed ^ (value + static_cast(0x9e3779b9) + + (seed << 6U) + (seed >> 2U)); + }; + auto hash = std::hash{}(key.time_ms); + hash = combine(hash, std::hash{}(key.ask)); + hash = combine(hash, std::hash{}(key.bid)); + hash = combine(hash, std::hash{}(key.last)); + hash = combine(hash, std::hash{}(key.volume)); + return hash; + } + }; + struct MarketDataRouterSubscriptionControl { mutable std::mutex mutex; RoutedSubscriptionId router_id; @@ -147,6 +182,8 @@ namespace optionx::market_data { std::chrono::steady_clock::duration degraded_duration{}; std::chrono::steady_clock::time_point stale_since{}; std::chrono::steady_clock::time_point degraded_since{}; + std::uint64_t last_delivered_tick_time_ms = 0; + std::vector last_delivered_tick_observations; std::optional retry_request; std::chrono::steady_clock::time_point retry_at; @@ -169,6 +206,8 @@ namespace optionx::market_data { ContinuityState continuity_state; MarketDataTickContinuityOptions tick_continuity; TickContinuityState tick_continuity_state; + MarketDataTickDeduplicationMode tick_deduplication_mode = + MarketDataTickDeduplicationMode::EXACT_OBSERVATION; MarketDataSubscriptionHandle retained_cleanup_subscription; MarketDataSubscriptionResult unsubscribe_completion; bool subscribe_completion_posted = false; @@ -326,6 +365,9 @@ namespace optionx::market_data { static StreamDescriptor stream_from(const TickSubscriptionRequest& request); static StreamDescriptor stream_from(const BarSubscriptionRequest& request); static StreamDescriptor stream_from(const MarketDataSubscriptionHandle& subscription); + static MarketDataTickDeduplicationMode resolve_tick_deduplication_mode( + BaseMarketDataProvider& provider, + MarketDataTickDeduplicationMode requested) noexcept; static bool same_status_stream( const MarketDataStatusUpdate& lhs, @@ -521,8 +563,16 @@ namespace optionx::market_data { std::uint64_t expected_interval_ms) noexcept; static bool same_tick_observation( const Tick& lhs, - const Tick& rhs) noexcept; + const Tick& rhs, + MarketDataTickDeduplicationMode mode) noexcept; + static TickObservationKey tick_observation_key( + const Tick& tick, + MarketDataTickDeduplicationMode mode) noexcept; static void deduplicate_tick_items( + std::vector& ticks, + MarketDataTickDeduplicationMode mode); + static void deduplicate_tick_items_against_last_delivery_no_lock( + const std::shared_ptr& entry, std::vector& ticks); static void clip_history_to_range( BarDataBatch& batch, @@ -549,7 +599,7 @@ namespace optionx::market_data { std::uint64_t to_time_ms) noexcept; static void record_tick_progress_no_lock( const std::shared_ptr& entry, - const std::vector& ticks) noexcept; + const std::vector& ticks); static std::uint64_t provider_time_ms( BaseMarketDataProvider& provider) noexcept; static std::uint64_t tick_history_interval_ms( @@ -814,6 +864,24 @@ namespace optionx::market_data { return stream; } + inline MarketDataTickDeduplicationMode + MarketDataRouterState::resolve_tick_deduplication_mode( + BaseMarketDataProvider& provider, + MarketDataTickDeduplicationMode requested) noexcept { + auto mode = requested == MarketDataTickDeduplicationMode::PROVIDER_DEFAULT + ? provider.tick_deduplication_mode() + : requested; + switch (mode) { + case MarketDataTickDeduplicationMode::TIMESTAMP: + case MarketDataTickDeduplicationMode::TIME_AND_PRICES: + case MarketDataTickDeduplicationMode::EXACT_OBSERVATION: + return mode; + case MarketDataTickDeduplicationMode::PROVIDER_DEFAULT: + default: + return MarketDataTickDeduplicationMode::EXACT_OBSERVATION; + } + } + inline bool MarketDataRouterState::register_provider( MarketDataProviderId id, BaseMarketDataProvider& provider, @@ -1145,6 +1213,9 @@ namespace optionx::market_data { entry->stream = std::move(stream); entry->continuity = std::move(continuity); entry->tick_continuity = std::move(tick_continuity); + entry->tick_deduplication_mode = resolve_tick_deduplication_mode( + provider, + entry->tick_continuity.deduplication_mode); entry->continuity_state.initial_prefill_pending = entry->continuity.enabled() && entry->continuity.prefill_bars > 0; entry->continuity_state.phase = @@ -2402,46 +2473,122 @@ namespace optionx::market_data { } } + inline TickObservationKey MarketDataRouterState::tick_observation_key( + const Tick& tick, + MarketDataTickDeduplicationMode mode) noexcept { + TickObservationKey key; + key.time_ms = tick.time_ms; + switch (mode) { + case MarketDataTickDeduplicationMode::TIMESTAMP: + break; + case MarketDataTickDeduplicationMode::TIME_AND_PRICES: + key.ask = tick.ask; + key.bid = tick.bid; + key.last = tick.last; + break; + case MarketDataTickDeduplicationMode::EXACT_OBSERVATION: + case MarketDataTickDeduplicationMode::PROVIDER_DEFAULT: + default: + key.ask = tick.ask; + key.bid = tick.bid; + key.last = tick.last; + key.volume = tick.volume; + break; + } + return key; + } + 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; + const Tick& rhs, + MarketDataTickDeduplicationMode mode) noexcept { + return tick_observation_key(lhs, mode) == + tick_observation_key(rhs, mode); } inline void MarketDataRouterState::deduplicate_tick_items( + std::vector& ticks, + MarketDataTickDeduplicationMode mode) { + std::unordered_set seen; + seen.reserve(ticks.size()); + auto unique_end = std::remove_if( + ticks.begin(), + ticks.end(), + [&seen, mode](Tick& tick) { + return !seen.insert(tick_observation_key(tick, mode)).second; + }); + ticks.erase(unique_end, ticks.end()); + } + + inline void + MarketDataRouterState::deduplicate_tick_items_against_last_delivery_no_lock( + const std::shared_ptr& entry, 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)); + if (!entry || ticks.empty()) return; + const auto& continuity = entry->tick_continuity_state; + if (continuity.last_delivered_tick_time_ms == 0 || + continuity.last_delivered_tick_observations.empty()) { + return; + } + + const auto mode = entry->tick_deduplication_mode; + std::unordered_set delivered; + delivered.reserve( + continuity.last_delivered_tick_observations.size()); + for (const auto& tick : continuity.last_delivered_tick_observations) { + delivered.insert(tick_observation_key(tick, mode)); } - ticks = std::move(unique); + + const auto unique_end = std::remove_if( + ticks.begin(), + ticks.end(), + [&delivered, mode](const Tick& tick) { + return delivered.find(tick_observation_key(tick, mode)) != + delivered.end(); + }); + ticks.erase(unique_end, ticks.end()); } inline void MarketDataRouterState::record_tick_progress_no_lock( const std::shared_ptr& entry, - const std::vector& ticks) noexcept { + const std::vector& ticks) { if (!entry) return; auto& continuity = entry->tick_continuity_state; + std::uint64_t latest_delivery_time_ms = 0; for (const auto& tick : ticks) { if (tick.time_ms > continuity.last_observed_time_ms) { continuity.last_observed_time_ms = tick.time_ms; } + latest_delivery_time_ms = std::max( + latest_delivery_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; } } + + if (latest_delivery_time_ms == 0) return; + const auto mode = entry->tick_deduplication_mode; + if (latest_delivery_time_ms > continuity.last_delivered_tick_time_ms) { + continuity.last_delivered_tick_time_ms = latest_delivery_time_ms; + continuity.last_delivered_tick_observations.clear(); + } + if (latest_delivery_time_ms != continuity.last_delivered_tick_time_ms) { + return; + } + + std::unordered_set seen; + seen.reserve(continuity.last_delivered_tick_observations.size() + ticks.size()); + for (const auto& tick : continuity.last_delivered_tick_observations) { + seen.insert(tick_observation_key(tick, mode)); + } + for (const auto& tick : ticks) { + if (tick.time_ms == latest_delivery_time_ms && + seen.insert(tick_observation_key(tick, mode)).second) { + continuity.last_delivered_tick_observations.push_back(tick); + } + } } inline void MarketDataRouterState::record_tick_verified_range_no_lock( @@ -3018,6 +3165,8 @@ namespace optionx::market_data { std::uint64_t generation, TickHistoryResult result) { StreamDescriptor expected_stream; + MarketDataTickDeduplicationMode deduplication_mode = + MarketDataTickDeduplicationMode::EXACT_OBSERVATION; { std::lock_guard lock(m_mutex); const auto entry_it = m_entries.find(router_id); @@ -3029,6 +3178,7 @@ namespace optionx::market_data { return; } expected_stream = entry_it->second->stream; + deduplication_mode = entry_it->second->tick_deduplication_mode; entry_it->second->tick_continuity_state.request_in_flight = false; entry_it->second->tick_continuity_state.phase = MarketDataContinuityPhase::FLUSHING; @@ -3063,7 +3213,7 @@ namespace optionx::market_data { request, subscription, kind != ContinuityRequestKind::PREFILL); - deduplicate_tick_items(history_batch.items); + deduplicate_tick_items(history_batch.items, deduplication_mode); delivered_history_items = history_batch.items.size(); } } @@ -3207,6 +3357,8 @@ namespace optionx::market_data { !entry_it->second->release_requested; if (active) { auto& continuity = entry_it->second->tick_continuity_state; + const auto current_deduplication_mode = + entry_it->second->tick_deduplication_mode; for (auto batch_it = continuity.buffer.begin(); batch_it != continuity.buffer.end();) { auto& items = batch_it->items; @@ -3214,14 +3366,15 @@ namespace optionx::market_data { std::remove_if( items.begin(), items.end(), - [&history_batch](const Tick& tick) { + [&history_batch, current_deduplication_mode](const Tick& tick) { return std::any_of( history_batch.items.begin(), history_batch.items.end(), - [&tick](const Tick& history_tick) { + [&tick, current_deduplication_mode](const Tick& history_tick) { return same_tick_observation( history_tick, - tick); + tick, + current_deduplication_mode); }); }), items.end()); @@ -3235,13 +3388,20 @@ namespace optionx::market_data { for (const auto& batch : continuity.buffer) { continuity.buffered_items += batch.items.size(); } - record_tick_progress_no_lock( + deduplicate_tick_items_against_last_delivery_no_lock( entry_it->second, history_batch.items); - subscriber = entry_it->second->subscriber.lock(); + if (!history_batch.items.empty()) { + 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); + if (active && subscriber && !history_batch.items.empty()) { + subscriber->on_tick_data(history_batch); + } } bool schedule_next = false; diff --git a/include/optionx_cpp/platforms/IntradeBarPlatform.hpp b/include/optionx_cpp/platforms/IntradeBarPlatform.hpp index 6648d5f..4d4bb6b 100644 --- a/include/optionx_cpp/platforms/IntradeBarPlatform.hpp +++ b/include/optionx_cpp/platforms/IntradeBarPlatform.hpp @@ -179,6 +179,12 @@ namespace optionx::platforms { return m_tick_history.options().sampling_interval_ms; } + /// \brief Returns the identity policy for one-second Intrade snapshots. + market_data::MarketDataTickDeduplicationMode tick_deduplication_mode() + const noexcept override { + return market_data::MarketDataTickDeduplicationMode::TIMESTAMP; + } + /// \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/tests/market_data_subscription_contract_test.cpp b/tests/market_data_subscription_contract_test.cpp index 9f1e909..64b7d10 100644 --- a/tests/market_data_subscription_contract_test.cpp +++ b/tests/market_data_subscription_contract_test.cpp @@ -39,11 +39,31 @@ TEST(TickSubscriptionRequest, BuildsAndValidatesTickRequests) { EXPECT_TRUE(request.valid()); EXPECT_EQ(request.symbol, "EUR/USD"); EXPECT_EQ(request.transport, MarketDataTransport::WEBSOCKET); + EXPECT_EQ( + request.continuity.deduplication_mode, + MarketDataTickDeduplicationMode::PROVIDER_DEFAULT); EXPECT_EQ(to_str(request.transport), std::string("WEBSOCKET")); EXPECT_FALSE(TickSubscriptionRequest("").valid()); } +TEST(MarketDataTickContinuityOptions, ValidatesIdentityPolicies) { + MarketDataTickContinuityOptions options; + EXPECT_TRUE(options.valid()); + + options.deduplication_mode = MarketDataTickDeduplicationMode::TIMESTAMP; + EXPECT_TRUE(options.valid()); + options.deduplication_mode = + MarketDataTickDeduplicationMode::TIME_AND_PRICES; + EXPECT_TRUE(options.valid()); + options.deduplication_mode = + MarketDataTickDeduplicationMode::EXACT_OBSERVATION; + EXPECT_TRUE(options.valid()); + options.deduplication_mode = + static_cast(99); + EXPECT_FALSE(options.valid()); +} + TEST(BarSubscriptionRequest, BuildsAndValidatesBarRequests) { const BarSubscriptionRequest request( "BTCUSDT", diff --git a/tests/market_data_tick_continuity_test.cpp b/tests/market_data_tick_continuity_test.cpp index f1a4839..c963277 100644 --- a/tests/market_data_tick_continuity_test.cpp +++ b/tests/market_data_tick_continuity_test.cpp @@ -40,6 +40,8 @@ class FakeTickHistoryProvider final : public BaseMarketDataProvider { public: std::uint64_t provider_now_ms = 0; std::uint64_t history_interval_ms = 0; + MarketDataTickDeduplicationMode deduplication_mode = + MarketDataTickDeduplicationMode::EXACT_OBSERVATION; ticks_callback_t& on_tick_data() override { return m_tick_callback; @@ -57,6 +59,10 @@ class FakeTickHistoryProvider final : public BaseMarketDataProvider { return history_interval_ms; } + MarketDataTickDeduplicationMode tick_deduplication_mode() const noexcept override { + return deduplication_mode; + } + bool subscribe_ticks( TickSubscriptionRequest request, subscription_callback_t callback) override { @@ -186,6 +192,17 @@ TickSequence make_history(std::initializer_list times) { return sequence; } +TickSequence make_history(std::initializer_list ticks) { + TickSequence sequence; + sequence.symbol = "EURUSD"; + sequence.provider = "INTRADE_BAR"; + sequence.price_digits = 5; + for (const auto& tick : ticks) { + sequence.ticks.push_back(tick); + } + return sequence; +} + TickSubscriptionRequest continuity_request( MarketDataContinuityMode mode = MarketDataContinuityMode::PREFILL_AND_RECOVER, @@ -403,9 +420,9 @@ TEST(MarketDataTickContinuity, RepairsGapAndPreservesDistinctSameSecondTicks) { 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]; - 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)); + ASSERT_EQ(history_batch.items.size(), 1U); + EXPECT_EQ(history_batch.items.front().time_ms, 5000U); + EXPECT_TRUE(history_batch.items.front().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); @@ -522,6 +539,129 @@ TEST(MarketDataTickContinuity, KeepsBoundedProviderGridChunksOverlapped) { EXPECT_TRUE(saw_triggering_tick); } +TEST(MarketDataTickContinuity, DeduplicatesInclusiveHistoryAcrossDeliveredBatches) { + ScopedTestClock clock(3000); + 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, + 2000, + 1500)); + ASSERT_TRUE(route.valid()); + + provider.complete_history(make_history({make_tick(1000, 1.0)})); + provider.emit_ticks({make_tick(4501, 5.0)}); + + 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({ + make_tick(1000, 1.0), + make_tick(2000, 2.0)})); + + 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({ + make_tick(2000, 2.0), + make_tick(2000, 2.1), + make_tick(3000, 3.0)})); + + 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({ + make_tick(3000, 3.0), + make_tick(4000, 4.0)})); + + ASSERT_EQ(provider.history_requests.size(), 4U); + std::vector delivered; + for (const auto& batch : subscriber->ticks) { + delivered.insert(delivered.end(), batch.items.begin(), batch.items.end()); + } + + const auto count_observations = [&delivered]( + std::uint64_t time_ms, + double bid) { + return static_cast(std::count_if( + delivered.begin(), + delivered.end(), + [time_ms, bid](const Tick& tick) { + return tick.time_ms == time_ms && tick.bid == bid; + })); + }; + EXPECT_EQ(count_observations(1000, 1.0), 1U); + EXPECT_EQ(count_observations(2000, 2.0), 1U); + EXPECT_EQ(count_observations(2000, 2.1), 1U); + EXPECT_EQ(count_observations(3000, 3.0), 1U); + EXPECT_EQ(count_observations(4000, 4.0), 1U); + EXPECT_EQ(count_observations(4501, 5.0), 1U); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 2U); +} + +TEST(MarketDataTickContinuity, UsesProviderDefaultTimestampIdentityPolicy) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + provider.deduplication_mode = + MarketDataTickDeduplicationMode::TIMESTAMP; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto route = router.subscribe_ticks( + provider, + subscriber, + continuity_request( + MarketDataContinuityMode::PREFILL, + 2000)); + ASSERT_TRUE(route.valid()); + + provider.complete_history(make_history({ + make_tick(1000, 1.0), + make_tick(1000, 1.1), + make_tick(2000, 2.0)})); + + ASSERT_EQ(subscriber->ticks.size(), 1U); + ASSERT_EQ(subscriber->ticks.front().items.size(), 2U); + EXPECT_DOUBLE_EQ(subscriber->ticks.front().items[0].bid, 1.0); + EXPECT_DOUBLE_EQ(subscriber->ticks.front().items[1].bid, 2.0); +} + +TEST(MarketDataTickContinuity, SupportsRouteTimeAndPricesIdentityOverride) { + ScopedTestClock clock(3000); + FakeTickHistoryProvider provider; + provider.deduplication_mode = + MarketDataTickDeduplicationMode::TIMESTAMP; + auto subscriber = std::make_shared(); + MarketDataRouter router; + + auto request = continuity_request( + MarketDataContinuityMode::PREFILL, + 2000); + request.continuity.deduplication_mode = + MarketDataTickDeduplicationMode::TIME_AND_PRICES; + auto route = router.subscribe_ticks(provider, subscriber, request); + ASSERT_TRUE(route.valid()); + + Tick same_prices_new_volume = make_tick(1000, 1.0); + same_prices_new_volume.volume = 2.0; + Tick changed_price = make_tick(1000, 1.1); + provider.complete_history(make_history({ + make_tick(1000, 1.0), + same_prices_new_volume, + changed_price})); + + ASSERT_EQ(subscriber->ticks.size(), 1U); + ASSERT_EQ(subscriber->ticks.front().items.size(), 2U); + EXPECT_DOUBLE_EQ(subscriber->ticks.front().items[0].bid, 1.0); + EXPECT_DOUBLE_EQ(subscriber->ticks.front().items[1].bid, 1.1); +} + TEST(MarketDataTickContinuity, UsesBoundedGapRequests) { ScopedTestClock clock(10000); FakeTickHistoryProvider provider;