diff --git a/.github/workflows/ubuntu-smoke.yml b/.github/workflows/ubuntu-smoke.yml index da0c279..1e26971 100644 --- a/.github/workflows/ubuntu-smoke.yml +++ b/.github/workflows/ubuntu-smoke.yml @@ -198,6 +198,14 @@ jobs: ./build-linux/market_data_router_test --gtest_brief=1 ./build-linux/market_data_router_example + - name: Build market_data_continuity_test and example + run: cmake --build build-linux --target market_data_continuity_test market_data_continuity_example -j + + - name: Run market_data_continuity_test and example + run: | + ./build-linux/market_data_continuity_test --gtest_brief=1 + ./build-linux/market_data_continuity_example + - 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 850ddba..9f473fb 100644 --- a/.github/workflows/windows-smoke.yml +++ b/.github/workflows/windows-smoke.yml @@ -93,6 +93,11 @@ jobs: - name: Build Telegram live bridge smoke example run: cmake --build build-windows --config Debug --target telegram_live_bridge_smoke + - name: Build market data continuity test and example + run: > + cmake --build build-windows --config Debug --target + market_data_continuity_test market_data_continuity_example + - name: Build TradingView extension bridge smoke example run: cmake --build build-windows --config Debug --target tradingview_extension_bridge_smoke @@ -113,6 +118,8 @@ jobs: .\build-windows\Debug\telegram_signal_bridge_test.exe --gtest_brief=1 .\build-windows\Debug\telegram_worker_source_test.exe --gtest_brief=1 .\build-windows\Debug\trading_view_bridge_test.exe --gtest_brief=1 + .\build-windows\Debug\market_data_continuity_test.exe --gtest_brief=1 + .\build-windows\Debug\market_data_continuity_example.exe .\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 diff --git a/CMakeLists.txt b/CMakeLists.txt index 019ceb2..8be63e4 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -455,6 +455,41 @@ if(OPTIONX_BUILD_EXAMPLES) ) endif() + add_executable(market_data_continuity_example examples/market_data_continuity_example.cpp) + + target_include_directories(market_data_continuity_example PRIVATE + ${EXAMPLE_INCLUDE_DIRS} + ${EXAMPLE_DEPS_INCLUDE_DIRS} + ) + + target_link_directories(market_data_continuity_example PRIVATE ${EXAMPLE_LIBRARY_DIRS}) + target_compile_definitions( + market_data_continuity_example PRIVATE + ${EXAMPLE_DEFINES} + LOGIT_BASE_PATH="${LOGIT_BASE_PATH_FWD}" + ) + target_link_libraries(market_data_continuity_example PRIVATE ${EXAMPLE_LIBS} optionx_cpp) + + if(OPTIONX_BUILD_DEPS) + add_dependencies(market_data_continuity_example mdbx-static AES) + endif() + + foreach(dll ${EXAMPLE_DLL_FILES}) + add_custom_command(TARGET market_data_continuity_example POST_BUILD + COMMAND ${CMAKE_COMMAND} -E copy_if_different + "${dll}" "$" + ) + endforeach() + + if(WIN32) + add_custom_command(TARGET market_data_continuity_example POST_BUILD + COMMAND ${CMAKE_COMMAND} + -DOPTIONX_RUNTIME_DLL_DIR="${EXAMPLE_BUILD_LIBS_DIR}/bin" + -DOPTIONX_RUNTIME_TARGET_DIR="$" + -P "${CMAKE_CURRENT_SOURCE_DIR}/cmake/copy_runtime_dlls.cmake" + ) + endif() + add_executable(lifecycle_stack_example examples/lifecycle_stack_example.cpp) target_include_directories(lifecycle_stack_example PRIVATE diff --git a/examples/README.md b/examples/README.md index 22aab38..bc37e47 100644 --- a/examples/README.md +++ b/examples/README.md @@ -37,6 +37,9 @@ Currently maintained examples: - `market_data_subscriber_base_example.cpp` demonstrates a bot posting subscribe and unsubscribe commands from its own thread by stable provider ID and alias, while provider calls, delivery, and handle cleanup stay in one owner loop. +- `market_data_continuity_example.cpp` demonstrates bar history prefill, + history-before-live ordering, timestamp-gap recovery, and route-scoped + continuity status events with a deterministic local provider. - `trading_condition_hub_example.cpp` demonstrates routing payout, session and expiration-limit changes through `TradingConditionHub`, plus querying the merged current condition snapshot for a concrete symbol. diff --git a/examples/market_data_continuity_example.cpp b/examples/market_data_continuity_example.cpp new file mode 100644 index 0000000..572f861 --- /dev/null +++ b/examples/market_data_continuity_example.cpp @@ -0,0 +1,152 @@ +#include +#include +#include +#include +#include +#include +#include +#include + +#include + +namespace md = optionx::market_data; + +namespace { + +class DemoBarProvider final : public md::BaseMarketDataProvider { +public: + bool subscribe_bars( + md::BarSubscriptionRequest request, + subscription_callback_t callback) override { + m_subscription = md::MarketDataSubscriptionHandle::from_bar_request( + provider_id(), + m_next_subscription_id++, + request); + if (callback) { + callback(md::MarketDataSubscriptionResult::subscribed(m_subscription)); + } + return true; + } + + bool unsubscribe( + md::MarketDataSubscriptionHandle subscription, + subscription_callback_t callback) override { + if (callback) { + callback(md::MarketDataSubscriptionResult::unsubscribed( + std::move(subscription))); + } + m_subscription = {}; + return true; + } + + bool fetch_bar_history( + const optionx::BarHistoryRequest& request, + bar_history_callback_t callback) override { + std::cout << "history request: " << request.symbol + << " [" << request.from_ts << ", " << request.to_ts << "]\n"; + m_history_callbacks.push_back(std::move(callback)); + return true; + } + + void emit_live_bar(std::uint64_t time_ms, double close) { + auto batch = std::make_unique(); + batch->subscription = m_subscription; + batch->type = md::MarketDataType::BARS; + batch->symbol = m_subscription.symbol; + batch->timeframe = m_subscription.timeframe; + batch->items.emplace_back(close - 0.1, close + 0.2, close - 0.3, close, 1.0, time_ms); + batch->items.back().set_flag(optionx::MarketDataFlags::REALTIME); + if (on_bar_data()) on_bar_data()(std::move(batch)); + } + + void complete_history(std::vector bars) { + if (m_history_callbacks.empty()) return; + + auto callback = std::move(m_history_callbacks.front()); + m_history_callbacks.erase(m_history_callbacks.begin()); + + optionx::BarSequence sequence; + sequence.symbol = "EURUSD"; + sequence.provider = "demo-provider"; + sequence.timeframe = 60; + sequence.price_digits = 5; + sequence.volume_digits = 0; + sequence.price_source = optionx::BarPriceSource::MID; + sequence.bars = std::move(bars); + callback(optionx::BarHistoryResult::ok(std::move(sequence))); + } + +private: + md::SubscriptionId m_next_subscription_id = 1; + md::MarketDataSubscriptionHandle m_subscription; + std::vector m_history_callbacks; +}; + +class Chart final : public md::IMarketDataSubscriber { +public: + void on_bar_data(const md::BarDataBatch& batch) override { + for (const auto& bar : batch.items) { + const char* source = "live"; + if (bar.has_flag(optionx::MarketDataFlags::HISTORICAL)) { + source = bar.has_flag(optionx::MarketDataFlags::BACKFILL) + ? "backfill" + : "prefill"; + } + std::cout << "chart bar: route provider subscription #" + << batch.subscription.id + << ", t=" << bar.time_ms + << ", source=" << source << '\n'; + } + } + + void on_market_data_continuity( + const md::MarketDataContinuityUpdate& update) override { + std::cout << "continuity: subscription #" << update.subscription.id + << ", status=" << md::to_str(update.status) + << ", history items=" << update.delivered_items << '\n'; + } +}; + +std::vector make_bars( + std::initializer_list timestamps) { + std::vector bars; + for (const auto timestamp : timestamps) { + bars.emplace_back(100.0, 101.0, 99.0, 100.5, 1.0, timestamp); + } + return bars; +} + +} // namespace + +int main() { + DemoBarProvider provider; + md::MarketDataRouter router; + auto chart = std::make_shared(); + + md::BarSubscriptionRequest request( + "EURUSD", + 60, + optionx::BarPriceSource::MID, + md::MarketDataTransport::WEBSOCKET); + request.continuity.mode = md::MarketDataContinuityMode::PREFILL_AND_RECOVER; + request.continuity.prefill_bars = 2; + request.continuity.max_backfill_bars = 10; + + auto route = router.subscribe_bars(provider, chart, request); + if (!route.active()) { + std::cerr << "could not create a bar route\n"; + return 1; + } + + // This live bar is buffered until the initial historical range is delivered. + provider.emit_live_bar(220000, 100.5); + provider.complete_history(make_bars({100000, 160000})); + + // 280000 is missing, so the Router requests it before releasing 340000. + provider.emit_live_bar(340000, 101.5); + provider.complete_history(make_bars({280000})); + + route.reset(); + router.shutdown(); + return 0; +} diff --git a/guides/api-and-header-contracts.md b/guides/api-and-header-contracts.md index 99116e9..54c22f1 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -171,6 +171,7 @@ Market-data APIs are split into DTO/data types and a provider role: `MarketDataSubscriptionHandle`, `MarketDataSubscriptionResult`, `MarketDataBatch`, `MarketDataHub`, `MarketDataRouter`, `MarketDataSubscriberBase`, `IMarketDataSubscriber`, and + `MarketDataContinuityOptions`, `MarketDataContinuityUpdate`, and `MarketDataContinuityService`. Contract rules: @@ -252,6 +253,17 @@ Contract rules: - `MarketDataContinuityService` is the thin helper for routing recovered history into the same bar batch pipeline. It marks payload bars as `HISTORICAL` and, for gap recovery, `BACKFILL`. +- `BarSubscriptionRequest::continuity` enables Router-owned bar prefill and + optional timestamp-gap recovery. Router buffers live batches until the + corresponding history operation completes and reports route-scoped progress + through `IMarketDataSubscriber::on_market_data_continuity()`. +- Continuity updates carry the concrete provider subscription handle and must + not be confused with stream-level `MarketDataStatusUpdate`. History failure + does not terminate the live route: Router reports `FAILED`, releases buffered + live batches, and returns to `LIVE`. +- Generic history continuity is currently defined for bars only. Tick history, + provider-aware retries, and a universal overlap-deduplication policy remain + separate contracts. `MarketDataRouter` is the subscription-scoped alternative to `MarketDataHub`: @@ -359,7 +371,9 @@ Intrade Bar publishes condition snapshots for every supported symbol/option-type scope after the account context becomes known. `TradingConditionManager` also re-evaluates the time-dependent session, amount, open-trade and sprint-duration limits from the same `AccountInfoData` model used to validate trade requests. -Only scopes whose values changed are emitted. +Only scopes whose values changed are emitted. During platform shutdown it emits +one final `tradable=false` patch for each cached scope before clearing the +manager, which prevents long-lived hubs from retaining stale availability. Intrade Bar intentionally leaves `TradingConditionUpdate::payout` empty. Its payout model depends on the concrete trade amount and duration, but those values diff --git a/guides/implementation-notes.md b/guides/implementation-notes.md index d045eec..c7da43c 100644 --- a/guides/implementation-notes.md +++ b/guides/implementation-notes.md @@ -262,6 +262,8 @@ The manager emits only changed scopes. A scope contains the platform, account type, currency, option type and normalized symbol. When account identity changes, the previous scopes receive a final `tradable=false` patch before the new scopes are published. +The same final unavailable patch is emitted for every cached scope during +platform shutdown before the manager clears its cache. Do not fill `payout` from a made-up reference amount or duration. Current Intrade payout rules are trade-parameter dependent, while `TradingConditionUpdate` does diff --git a/guides/market-data-router.md b/guides/market-data-router.md index 596d787..40a3a6c 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -294,6 +294,100 @@ This is lifecycle replay, not historical market-data recovery. Historical prefill and gap repair belong to `MarketDataContinuityService` and provider history APIs. +## History Prefill And Gap Recovery + +`BarSubscriptionRequest` can ask Router to establish a bar route from +historical data and then continue with live data. The first implementation is +bar-only because `BaseMarketDataProvider` currently has +`fetch_bar_history(...)`, but no generic tick-history operation: + +```cpp +md::BarSubscriptionRequest request( + "EURUSD", + 60, + md::BarPriceSource::MID, + md::MarketDataTransport::WEBSOCKET); +request.continuity.mode = md::MarketDataContinuityMode::PREFILL_AND_RECOVER; +request.continuity.prefill_bars = 100; +request.continuity.max_backfill_bars = 500; + +auto route = router.subscribe_bars("intrade", chart, request); +``` + +The options have these meanings: + +- `LIVE_ONLY` leaves live delivery unchanged. `PREFILL` requests + `prefill_bars` before the first live delivery. `PREFILL_AND_RECOVER` does + both and also repairs timestamp gaps. +- `prefill_bars` is the count-based initial history depth. Router builds an + inclusive timeframe range for that many slots; a provider may still return + fewer bars when part of the range has no data. A `PREFILL` request must + specify a positive count. `PREFILL_AND_RECOVER` may use zero when the + application wants recovery only. +- `max_backfill_bars` bounds one gap request. Zero means that the provider + request is not count-bounded by Router. + +For initial prefill, the delivery order is: + +```text +SUBSCRIBED + -> continuity PREFILLING + -> historical BarDataBatch (HISTORICAL) + -> buffered live BarDataBatch (REALTIME) + -> continuity LIVE +``` + +Live batches arriving while history is in flight are retained by the route. +The history batch is marked `HISTORICAL`; initial prefill is not marked +`BACKFILL`. The buffered live batches are delivered only after the history +operation completes, preserving the history-to-live boundary. Router then +checks the buffered batches in delivery order as well. If one batch contains a +continuous prefix followed by a gap, the prefix is delivered first and only the +suffix is held for recovery. Several queued gaps are recovered serially before +the route emits its final `LIVE` update. + +With `PREFILL_AND_RECOVER`, Router compares increasing bar timestamps with the +requested timeframe. A later bar that skips one or more expected timeframe +buckets produces: + +```text +continuity GAP_DETECTED + -> continuity BACKFILLING + -> historical BarDataBatch (HISTORICAL | BACKFILL) + -> buffered live BarDataBatch + -> continuity LIVE +``` + +`MarketDataContinuityUpdate` is route-scoped. Its `subscription` field is the +concrete provider handle for that route, and `from_time_ms`, `to_time_ms`, +`requested_items`, and `delivered_items` describe the current history request. +Implement `IMarketDataSubscriber::on_market_data_continuity()` when a chart or +bot needs to display or record this state. These updates are separate from +stream-level `MarketDataStatusUpdate` events such as `READY` or `DISCONNECTED`. + +`max_backfill_bars` limits each individual provider request, not the complete +missing interval. When a provider returns a partial but non-empty backfill, +Router keeps draining the queued live batches and can issue another bounded +request for the next remaining gap. An empty successful response for a detected +gap is treated as a failed recovery, so Router releases the queued live data +once instead of retrying the same range forever. + +If a history request fails or is rejected, Router emits `FAILED`, releases any +buffered live batches, and then emits `LIVE`. The route stays usable and live +delivery continues, but the missing historical range is not reconstructed. +Applications that require a complete time series should record the failure and +apply their own retry policy. A provider may return overlapping snapshots for +an in-progress bar; this first continuity layer does not impose a universal +payload deduplication policy, so consumers should correlate bars by stream and +`time_ms` according to their finalized/incomplete-bar policy. + +`MarketDataContinuityService` is the lower-level helper for applications that +want to request history directly. It converts a `BarHistoryResult` into a +`BarDataBatch` and applies `HISTORICAL`/`BACKFILL` flags. Router uses the same +helper for route-owned prefill and recovery. See +`examples/market_data_continuity_example.cpp` for a deterministic chart-like +consumer and provider. + ## Owner Loop And Bot Threads Without an owner dispatcher, Router preserves synchronous behavior. Subscribe, @@ -322,10 +416,13 @@ The dispatcher contract is: Queue acceptance is not an execution guarantee. In particular, `BaseTradingPlatform::post_task()` may cancel accepted tasks that remain queued -when platform shutdown begins. +when platform shutdown begins. A rejected history completion is dropped with +the other provider deliveries; it is never invoked inline on the provider +thread. With a dispatcher configured, Router marshals provider subscription completions, -tick batches, bar batches, and status updates into that loop. Bots use: +history completions, tick batches, bar batches, and status updates into that +loop. Bots use: - `post_subscribe_ticks()` and `post_subscribe_bars()`; - `post_unsubscribe()` and `post_unsubscribe_all()`. diff --git a/guides/market-data-router.ru.md b/guides/market-data-router.ru.md index 73e1ef5..13045a1 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -294,6 +294,102 @@ Router не вызывает пользовательский код, удерж Предварительная история и заполнение разрывов относятся к `MarketDataContinuityService` и history API провайдера. +## Предварительная История И Восстановление Разрывов + +`BarSubscriptionRequest` может попросить Router сначала построить поток баров +из истории, а затем продолжить его live-данными. Первая реализация работает +только с барами: в `BaseMarketDataProvider` сейчас есть +`fetch_bar_history(...)`, но нет общего API истории тиков: + +```cpp +md::BarSubscriptionRequest request( + "EURUSD", + 60, + md::BarPriceSource::MID, + md::MarketDataTransport::WEBSOCKET); +request.continuity.mode = md::MarketDataContinuityMode::PREFILL_AND_RECOVER; +request.continuity.prefill_bars = 100; +request.continuity.max_backfill_bars = 500; + +auto route = router.subscribe_bars("intrade", chart, request); +``` + +Значения options имеют следующий смысл: + +- `LIVE_ONLY` оставляет live-доставку без изменений. `PREFILL` запрашивает + `prefill_bars` перед первой live-доставкой. `PREFILL_AND_RECOVER` делает оба + действия и дополнительно восстанавливает разрывы по timestamp. +- `prefill_bars` задаёт глубину начальной истории в барах. Router строит + inclusive-диапазон на это количество timeframe slots; провайдер всё равно + может вернуть меньше баров, если часть диапазона не содержит данных. Для + `PREFILL` требуется положительное значение. `PREFILL_AND_RECOVER` может + использовать ноль, если приложению нужно только восстановление разрывов. +- `max_backfill_bars` ограничивает один запрос для разрыва. Ноль означает, + что Router не ограничивает provider request количеством баров. + +Для initial prefill порядок доставки такой: + +```text +SUBSCRIBED + -> continuity PREFILLING + -> исторический BarDataBatch (HISTORICAL) + -> накопленный live BarDataBatch (REALTIME) + -> continuity LIVE +``` + +Live batches, пришедшие пока выполняется history request, удерживаются внутри +route. Исторический batch получает `HISTORICAL`; initial prefill не получает +флаг `BACKFILL`. Накопленные live batches доставляются только после завершения +history operation, поэтому граница history-to-live сохраняется. Затем Router +проверяет накопленные batches в порядке доставки. Если один batch содержит +непрерывный prefix, а затем gap, prefix доставляется сразу, а для recovery +удерживается только suffix. Несколько накопленных gap восстанавливаются +последовательно, и только после этого route получает финальный `LIVE`. + +В режиме `PREFILL_AND_RECOVER` Router сравнивает возрастающие timestamps баров +с заданным timeframe. Если более поздний бар перескакивает один или несколько +ожидаемых buckets, возникают события: + +```text +continuity GAP_DETECTED + -> continuity BACKFILLING + -> исторический BarDataBatch (HISTORICAL | BACKFILL) + -> накопленный live BarDataBatch + -> continuity LIVE +``` + +`MarketDataContinuityUpdate` относится к конкретному route. Его поле +`subscription` содержит физический provider handle этого route, а +`from_time_ms`, `to_time_ms`, `requested_items` и `delivered_items` описывают +текущий history request. Переопредели +`IMarketDataSubscriber::on_market_data_continuity()`, если графику или боту +нужно показывать или записывать это состояние. Эти события отделены от +stream-level `MarketDataStatusUpdate`, например `READY` или `DISCONNECTED`. + +`max_backfill_bars` ограничивает каждый отдельный provider request, а не весь +отсутствующий интервал. Если provider вернул неполный, но непустой backfill, +Router продолжит разбирать очередь live batches и может выполнить следующий +ограниченный запрос для оставшегося gap. Успешный пустой ответ для найденного +gap считается failed recovery: Router один раз выпускает накопленные live data и +не зацикливает запрос того же диапазона. + +Если history request завершился ошибкой или provider его отклонил, Router +публикует `FAILED`, выпускает накопленные live batches, затем публикует `LIVE`. +Route остаётся пригодным для работы и live-доставка продолжается, но +отсутствующий исторический диапазон не восстанавливается. Приложение, которому +нужен полный временной ряд, должно записать ошибку и применить собственную +политику повторной попытки. Provider может возвращать пересекающиеся snapshots +для незавершённого бара; этот первый слой continuity не вводит универсальную +политику deduplication, поэтому consumer должен сам сопоставлять бары по +stream и `time_ms` с учётом политики `FINALIZED`/незавершённых баров. + +`MarketDataContinuityService` остаётся низкоуровневым helper для приложений, +которые хотят запрашивать историю напрямую. Он превращает `BarHistoryResult` в +`BarDataBatch` и устанавливает флаги `HISTORICAL`/`BACKFILL`. Router использует +тот же helper для route-owned prefill и recovery. См. подробный +`examples/market_data_continuity_example.cpp` с детерминированным provider и +consumer в стиле графика. + ## Owner Loop И Потоки Ботов Без owner dispatcher Router сохраняет синхронное поведение. Subscribe, @@ -322,10 +418,13 @@ md::MarketDataRouter router( Принятие задачи очередью не гарантирует её выполнение. В частности, `BaseTradingPlatform::post_task()` может отменить принятые задачи, оставшиеся в -очереди к моменту начала shutdown платформы. +очереди к моменту начала shutdown платформы. Отклонённый history completion +отбрасывается вместе с остальными provider deliveries и никогда не вызывается +inline в provider thread. При наличии dispatcher Router переносит provider subscription completions, -tick batches, bar batches и status updates в этот loop. Боты используют: +history completions, tick batches, bar batches и status updates в этот loop. +Боты используют: - `post_subscribe_ticks()` и `post_subscribe_bars()`; - `post_unsubscribe()` и `post_unsubscribe_all()`. diff --git a/guides/platform-api-guide.md b/guides/platform-api-guide.md index 42ea041..de863c9 100644 --- a/guides/platform-api-guide.md +++ b/guides/platform-api-guide.md @@ -65,6 +65,9 @@ snapshots after account context is known and whenever time-dependent limits change. Intrade payout remains absent because it depends on a concrete amount and duration; query that exact trade through `AccountInfoRequest`. +When the platform shuts down, the condition manager publishes a final +`tradable=false` patch for every cached scope before clearing its state, so a +long-lived condition hub does not retain a stale tradable snapshot. ## `market_data::BaseMarketDataProvider` @@ -121,6 +124,14 @@ Subscription rules: timer-based final snapshot to be delivered. - `MarketDataContinuityService` routes recovered historical bars into the same `BarDataBatch` pipeline and marks them as `HISTORICAL`/`BACKFILL`. +- `BarSubscriptionRequest::continuity` lets `MarketDataRouter` request initial + bar history before live delivery and optionally recover timestamp gaps. Live + batches are buffered while history is in flight. Route-scoped progress is + reported through `IMarketDataSubscriber::on_market_data_continuity()`; it is + separate from stream-level `on_market_data_status()`. +- Router continuity is currently bar-only because providers expose + `fetch_bar_history()` but no generic tick-history operation. See the complete + EN/RU Router guides and `market_data_continuity_example.cpp`. - `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 @@ -190,6 +201,13 @@ a cached matching stream status synchronously while a new route is accepted; callbacks should therefore use the subscription carried by the event rather than assume that caller state was already updated after `subscribe_ticks()`. +For a history-first bar route, set `BarSubscriptionRequest::continuity` to +`PREFILL` or `PREFILL_AND_RECOVER`. Router delivers the initial +`HISTORICAL` batch before buffered `REALTIME` batches. In recovery mode it emits +route-scoped `GAP_DETECTED`/`BACKFILLING` updates, delivers recovered bars with +`HISTORICAL | BACKFILL`, then resumes live delivery. A failed history request +emits `FAILED` and releases buffered live data instead of killing the route. + Logical release and physical provider cleanup are separate. A released route stops receiving events immediately. If provider `unsubscribe()` is rejected or completes with failure, the Router keeps the physical handle and callback diff --git a/guides/project-overview.md b/guides/project-overview.md index aaa66da..5155b28 100644 --- a/guides/project-overview.md +++ b/guides/project-overview.md @@ -36,7 +36,7 @@ node lifecycles. See [English](lifecycle-stack.md) or | `optionx_cpp/lifecycle.hpp` | `ILifecycleModule`, `LifecycleStack` | Для необязательной общей обработки и staged shutdown модулей | | `optionx_cpp/data.hpp` | DTO, events, enums, account/symbol/tick/bar/trading data | Для API boundary и сообщений | | `optionx_cpp/components.hpp` | `BaseComponent`, HTTP и trade execution base classes | Для нового manager/component | -| `optionx_cpp/market_data.hpp` | Market-data provider role, subscription DTOs and statuses | Для live tick/bar subscriptions и history API contracts | +| `optionx_cpp/market_data.hpp` | Provider role, subscription DTOs, Router, subscriber base and continuity | Для live tick/bar subscriptions, history prefill and routed delivery | | `optionx_cpp/platforms.hpp` | Base platform и Intrade Bar platform | Для клиентского кода платформ | | `optionx_cpp/storages.hpp` | `ServiceSessionDB` | Для session storage | | `optionx_cpp/bridges.hpp` | `BaseBridge` | Для внешних bridge integrations | @@ -51,7 +51,7 @@ node lifecycles. See [English](lifecycle-stack.md) or | Platform facade | `optionx::platforms`, `include/optionx_cpp/platforms` | Публичный API платформы: connect, auth, trades, account info | | Components/managers | `optionx::components`, platform subnamespaces | Lifecycle-компоненты, которые получают events и выполняют работу | | Trading data | `optionx`, `include/optionx_cpp/data/trading` | `TradeRequest`, `TradeResult`, enums, signals | -| Market data | `optionx::market_data`, `include/optionx_cpp/market_data`, `data/bars`, `data/ticks`, `data/symbol` | Live subscriptions, history requests/results, ticks, bars and symbols | +| Market data | `optionx::market_data`, `include/optionx_cpp/market_data`, `data/bars`, `data/ticks`, `data/symbol` | Provider subscriptions, routed delivery, history prefill/recovery, ticks, bars and symbols | | Events | `optionx::events`, `include/optionx_cpp/data/events` | Pub-sub контракты между components | | Infrastructure | `optionx::utils` | EventBus, tasks, crypto, ids, HTTP helpers | | Storage | `optionx::storage` | AES + mdbx session storage | @@ -79,7 +79,8 @@ node lifecycles. See [English](lifecycle-stack.md) or implementations. Concrete platforms may use HTTP polling, websockets, or both internally; public subscription callbacks report desired-state acceptance, while physical stream readiness is reported through status - callbacks. + callbacks. `MarketDataRouter` adds provider selection, route ownership, + subscription-scoped delivery, and optional bar history continuity. 8. `shutdown()` останавливает tasks, вызывает shutdown у components, затем draining event bus. diff --git a/guides/refactor-backlog.md b/guides/refactor-backlog.md index 5052c44..c78295a 100644 --- a/guides/refactor-backlog.md +++ b/guides/refactor-backlog.md @@ -5,25 +5,22 @@ series. Keep it short and remove items once they are handled. ## Next PR Candidates -- Add `MarketDataRouter` above provider subscriptions. The router should own - provider handles, expose move-only RAII subscription handles, correlate - statuses with concrete subscriptions, and replay the current stream status - to late subscriptions. -- Add `MarketDataSubscriberBase` as optional convenience API for bots and - charts that subscribe from their own methods and retain router handles. -- Publish real broker/platform payout, expiry, amount-limit, and market-open - changes through `TradingConditionUpdate` instead of using the hub only as a - manually populated snapshot cache. +- Extend market-data continuity beyond the first bar-only route implementation: + define provider support for tick history, retries, and a documented + history-to-live boundary for each provider. +- Add robust gap recovery policy with provider-aware retry/backoff, sequence or + timestamp validation, and an explicit deduplication policy for overlapping + historical, backfill, and live bar snapshots. +- Add route-scoped continuity metrics and failure visibility for applications + that need to prove that a chart or strategy has a complete time series. - Add a fuller CMake package/export story for consumers that do not use the project as a direct submodule. The current `optionx_cpp::optionx_cpp` interface target covers build-tree/submodule consumption. ## Explicitly Deferred -- Finalize live bars from platform time/process even when no tick arrives for - the next bar. Tick-driven aggregation alone cannot close an idle stream. -- Remove `SingleTick` from the internal price event path after all remaining - parser/manager consumers use `TickUpdateBatch` directly. +- Add a generic tick-history provider contract. Current continuity support is + intentionally bar-first because providers expose bar history only. - Continue generation-safe lifecycle hardening for legacy bridge transports when their behavior is changed; do not mix that work into market-data API PRs. diff --git a/include/optionx_cpp/market_data.hpp b/include/optionx_cpp/market_data.hpp index 8c935bf..982b23a 100644 --- a/include/optionx_cpp/market_data.hpp +++ b/include/optionx_cpp/market_data.hpp @@ -10,10 +10,13 @@ #include #include #include +#include +#include #include #include #include #include +#include #include #include #include @@ -22,11 +25,14 @@ #include "lifecycle.hpp" #include "utils/fixed_point.hpp" +#include "utils/time_utils.hpp" #include "data/market.hpp" #include "data/bars.hpp" #include "data/ticks.hpp" #include "market_data/enums.hpp" +#include "market_data/MarketDataContinuityOptions.hpp" #include "market_data/MarketDataSubscription.hpp" +#include "market_data/MarketDataContinuity.hpp" #include "market_data/MarketDataBatch.hpp" #include "market_data/BaseMarketDataProvider.hpp" #include "market_data/MarketDataContinuityService.hpp" diff --git a/include/optionx_cpp/market_data/IMarketDataSubscriber.hpp b/include/optionx_cpp/market_data/IMarketDataSubscriber.hpp index e00cc13..309c32f 100644 --- a/include/optionx_cpp/market_data/IMarketDataSubscriber.hpp +++ b/include/optionx_cpp/market_data/IMarketDataSubscriber.hpp @@ -36,6 +36,13 @@ namespace optionx::market_data { virtual void on_market_data_status(const MarketDataStatusUpdate& update) { (void)update; } + + /// \brief Receives history prefill or gap-recovery progress for a route. + /// \param update Route-scoped continuity status and range information. + virtual void on_market_data_continuity( + const MarketDataContinuityUpdate& update) { + (void)update; + } }; } // namespace optionx::market_data diff --git a/include/optionx_cpp/market_data/MarketDataContinuity.hpp b/include/optionx_cpp/market_data/MarketDataContinuity.hpp new file mode 100644 index 0000000..49c2f81 --- /dev/null +++ b/include/optionx_cpp/market_data/MarketDataContinuity.hpp @@ -0,0 +1,61 @@ +#pragma once +#ifndef OPTIONX_HEADER_MARKET_DATA_MARKET_DATA_CONTINUITY_HPP_INCLUDED +#define OPTIONX_HEADER_MARKET_DATA_MARKET_DATA_CONTINUITY_HPP_INCLUDED + +/// \file MarketDataContinuity.hpp +/// \brief Defines route-scoped continuity status updates. + +#include +#include +#include + +namespace optionx::market_data { + + /// \enum MarketDataContinuityStatus + /// \brief Lifecycle state of history-to-live delivery for one route. + enum class MarketDataContinuityStatus { + 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. + LIVE, ///< Live delivery is current, with no pending history work. + FAILED ///< History work failed; live delivery continues without it. + }; + + /// \brief Converts a continuity status to stable text. + inline const char* to_str(MarketDataContinuityStatus status) noexcept { + switch (status) { + case MarketDataContinuityStatus::PREFILLING: + return "PREFILLING"; + case MarketDataContinuityStatus::GAP_DETECTED: + return "GAP_DETECTED"; + case MarketDataContinuityStatus::BACKFILLING: + return "BACKFILLING"; + case MarketDataContinuityStatus::LIVE: + return "LIVE"; + case MarketDataContinuityStatus::FAILED: + return "FAILED"; + case MarketDataContinuityStatus::UNKNOWN: + default: + return "UNKNOWN"; + } + } + + /// \struct MarketDataContinuityUpdate + /// \brief Route-scoped progress or failure information for historical delivery. + struct MarketDataContinuityUpdate { + MarketDataSubscriptionHandle subscription; ///< Concrete provider subscription. + MarketDataType type = MarketDataType::BARS; ///< Continuity payload type. + std::string symbol; ///< Provider symbol. + BarTimeframe timeframe = 0; ///< Bar timeframe in seconds. + 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 delivered_items = 0; ///< Number of history items delivered by the operation. + std::string message; ///< Optional diagnostic text. + }; + +} // namespace optionx::market_data + +#endif // OPTIONX_HEADER_MARKET_DATA_MARKET_DATA_CONTINUITY_HPP_INCLUDED diff --git a/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp new file mode 100644 index 0000000..1605971 --- /dev/null +++ b/include/optionx_cpp/market_data/MarketDataContinuityOptions.hpp @@ -0,0 +1,54 @@ +#pragma once +#ifndef OPTIONX_HEADER_MARKET_DATA_MARKET_DATA_CONTINUITY_OPTIONS_HPP_INCLUDED +#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. + +#include + +namespace optionx::market_data { + + /// \enum MarketDataContinuityMode + /// \brief Selects how a routed bar stream is initialized and recovered. + enum class MarketDataContinuityMode { + LIVE_ONLY = 0, ///< Deliver live provider payloads immediately. + PREFILL, ///< Deliver historical bars before live payloads. + PREFILL_AND_RECOVER ///< Prefill and repair timestamp gaps in live bars. + }; + + /// \struct MarketDataContinuityOptions + /// \brief Configures history prefill and timestamp-gap recovery for a bar route. + struct MarketDataContinuityOptions { + MarketDataContinuityMode mode = MarketDataContinuityMode::LIVE_ONLY; + std::size_t prefill_bars = 0; ///< Number of historical bars requested before live delivery. + std::size_t max_backfill_bars = 1000; ///< Maximum bars per detected gap; zero is unbounded. + + /// \brief Returns true when the option combination is usable. + [[nodiscard]] bool valid() const noexcept { + switch (mode) { + case MarketDataContinuityMode::LIVE_ONLY: + return prefill_bars == 0; + case MarketDataContinuityMode::PREFILL: + return prefill_bars > 0; + case MarketDataContinuityMode::PREFILL_AND_RECOVER: + return true; + default: + return false; + } + } + + /// \brief Returns true when history work is enabled for this route. + [[nodiscard]] bool enabled() const noexcept { + return mode != MarketDataContinuityMode::LIVE_ONLY; + } + + /// \brief Returns true when gap repair is enabled. + [[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/MarketDataContinuityService.hpp b/include/optionx_cpp/market_data/MarketDataContinuityService.hpp index f60caac..1e1f566 100644 --- a/include/optionx_cpp/market_data/MarketDataContinuityService.hpp +++ b/include/optionx_cpp/market_data/MarketDataContinuityService.hpp @@ -3,7 +3,16 @@ #define OPTIONX_HEADER_MARKET_DATA_MARKET_DATA_CONTINUITY_SERVICE_HPP_INCLUDED /// \file MarketDataContinuityService.hpp -/// \brief Defines a small helper for routing historical bars through market-data batches. +/// \brief Defines helpers for routing historical bars through market-data batches. + +#include +#include +#include +#include +#include +#include +#include +#include namespace optionx::market_data { @@ -23,6 +32,66 @@ namespace optionx::market_data { explicit MarketDataContinuityService(BaseMarketDataProvider& provider) : m_provider(provider) {} + /// \brief Converts a non-negative Unix-second value to milliseconds. + /// \param seconds Unix timestamp in seconds. + /// \return Saturated millisecond timestamp, or zero for non-positive input. + static std::uint64_t seconds_to_milliseconds( + std::int64_t seconds) noexcept { + if (seconds <= 0) return 0; + const auto value = static_cast(seconds); + if (value > std::numeric_limits::max() / 1000U) { + return std::numeric_limits::max(); + } + return value * 1000U; + } + + /// \brief Builds a count-based prefill request for a bar subscription. + /// \param request Live bar subscription whose symbol and timeframe are reused. + /// \param now_ms Current Unix timestamp in milliseconds. + /// \param bars Number of inclusive timeframe slots requested before + /// the live boundary. Providers may return fewer bars when + /// data is unavailable for part of the requested range. + /// \return A provider history request using Unix-second boundaries. + static BarHistoryRequest make_prefill_request( + const BarSubscriptionRequest& request, + std::uint64_t now_ms, + std::size_t bars) { + const auto timeframe_ms = safe_multiply( + static_cast(request.timeframe), 1000U); + const auto depth_ms = safe_multiply( + timeframe_ms, + bars > 0 ? bars - 1 : 0); + const auto from_ms = now_ms > depth_ms ? now_ms - depth_ms : 1U; + + return make_history_request(request, from_ms, now_ms); + } + + /// \brief Builds a bounded request for a missing bar range. + /// \param request Live bar subscription whose symbol and price source are reused. + /// \param from_ms Inclusive start of the missing range. + /// \param to_ms Inclusive end of the missing range. + /// \param max_bars Optional upper bound for one backfill operation. + /// \return A provider history request using Unix-second boundaries. + static BarHistoryRequest make_gap_request( + const BarSubscriptionRequest& request, + std::uint64_t from_ms, + std::uint64_t to_ms, + std::size_t max_bars = 0) { + if (max_bars > 0 && to_ms >= from_ms && request.timeframe > 0) { + const auto span_ms = safe_multiply( + safe_multiply( + static_cast(request.timeframe), + 1000U), + max_bars - 1); + const auto bounded_to_ms = from_ms > + std::numeric_limits::max() - span_ms + ? std::numeric_limits::max() + : from_ms + span_ms; + to_ms = std::min(to_ms, bounded_to_ms); + } + return make_history_request(request, from_ms, to_ms); + } + /// \brief Requests historical bars and delivers them as one batch. /// \param request Historical bar range to fetch. /// \param subscription Optional live subscription related to the backfill. @@ -97,6 +166,37 @@ namespace optionx::market_data { } private: + static std::uint64_t safe_multiply( + std::uint64_t left, + std::uint64_t right) noexcept { + if (left == 0 || right == 0) return 0; + if (left > std::numeric_limits::max() / right) { + return std::numeric_limits::max(); + } + return left * right; + } + + static BarHistoryRequest make_history_request( + const BarSubscriptionRequest& request, + std::uint64_t from_ms, + std::uint64_t to_ms) { + BarHistoryRequest history; + history.symbol = request.symbol; + history.timeframe = request.timeframe; + const auto max_seconds = static_cast( + std::numeric_limits::max()); + const auto from_seconds = from_ms / 1000U; + const auto to_seconds = to_ms / 1000U; + history.from_ts = static_cast(std::min( + from_seconds, + max_seconds)); + history.to_ts = static_cast(std::min( + to_seconds, + max_seconds)); + history.price_source = request.price_source; + return history; + } + BaseMarketDataProvider& m_provider; ///< Provider used for history fetches. }; diff --git a/include/optionx_cpp/market_data/MarketDataHub.hpp b/include/optionx_cpp/market_data/MarketDataHub.hpp index 1738fb3..71cd7b5 100644 --- a/include/optionx_cpp/market_data/MarketDataHub.hpp +++ b/include/optionx_cpp/market_data/MarketDataHub.hpp @@ -23,8 +23,8 @@ namespace optionx::market_data { /// /// Status replay is stream-level and happens when a subscriber /// slot is added to the hub. It is not a per-subscription status - /// API; a future router layer should replay status for newly - /// created subscription handles. + /// API; MarketDataRouter adds concrete subscription ownership and + /// replays status for newly created subscription handles. /// /// When MarketDataStatusUpdate::subscription is valid, status /// caching keeps that subscription context distinct from other diff --git a/include/optionx_cpp/market_data/MarketDataRouter.hpp b/include/optionx_cpp/market_data/MarketDataRouter.hpp index 57b5fac..d7048f8 100644 --- a/include/optionx_cpp/market_data/MarketDataRouter.hpp +++ b/include/optionx_cpp/market_data/MarketDataRouter.hpp @@ -54,8 +54,8 @@ namespace optionx::market_data { /// \details The dispatcher must be thread-safe, enqueue tasks in FIFO /// order, and remain available until Router shutdown completes. /// It must not execute posted work inline on a foreign caller. - /// Provider events are dropped if the configured dispatcher - /// rejects them during shutdown. + /// Provider events and history completions are dropped if the + /// configured dispatcher rejects them during shutdown. explicit MarketDataRouter(owner_dispatcher_t owner_dispatcher); /// \brief Copy construction is disabled because provider callbacks are owned. diff --git a/include/optionx_cpp/market_data/MarketDataSubscription.hpp b/include/optionx_cpp/market_data/MarketDataSubscription.hpp index d3e886a..4325606 100644 --- a/include/optionx_cpp/market_data/MarketDataSubscription.hpp +++ b/include/optionx_cpp/market_data/MarketDataSubscription.hpp @@ -5,6 +5,8 @@ /// \file MarketDataSubscription.hpp /// \brief Defines market-data subscription request, handle, and result DTOs. +#include "MarketDataContinuityOptions.hpp" + namespace optionx::market_data { /// \brief Runtime identifier of a market-data provider instance. @@ -55,6 +57,7 @@ namespace optionx::market_data { BarTimeframe timeframe = 0; ///< Bar timeframe in seconds; values <= 0 are invalid. BarPriceSource price_source = BarPriceSource::MID; ///< Price source for bars. MarketDataTransport transport = MarketDataTransport::AUTO; ///< Preferred transport. + MarketDataContinuityOptions continuity; ///< Optional history prefill and gap recovery. /// \brief Default constructor. BarSubscriptionRequest() = default; @@ -64,20 +67,23 @@ namespace optionx::market_data { /// \param timeframe Bar timeframe in seconds. /// \param price_source Price stream used to build OHLC values. /// \param transport Preferred transport for live data. + /// \param continuity Optional history prefill and gap recovery policy. BarSubscriptionRequest( std::string symbol, BarTimeframe timeframe, BarPriceSource price_source = BarPriceSource::MID, - MarketDataTransport transport = MarketDataTransport::AUTO) + MarketDataTransport transport = MarketDataTransport::AUTO, + MarketDataContinuityOptions continuity = {}) : symbol(std::move(symbol)), timeframe(timeframe), price_source(price_source), - transport(transport) {} + transport(transport), + continuity(continuity) {} /// \brief Returns true when the request describes a valid live bar stream. /// \return True when symbol is not empty and timeframe is positive. bool valid() const { - return !symbol.empty() && timeframe > 0; + return !symbol.empty() && timeframe > 0 && continuity.valid(); } }; diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index 592280c..5d2a00f 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -50,6 +50,12 @@ namespace optionx::market_data { std::weak_ptr subscriber; std::shared_ptr control; StreamDescriptor stream; + MarketDataContinuityOptions continuity; + bool continuity_ready = true; + bool continuity_request_in_flight = false; + bool continuity_flushing = false; + std::deque continuity_buffer; + std::uint64_t last_bar_time_ms = 0; MarketDataSubscriptionHandle retained_cleanup_subscription; MarketDataSubscriptionResult unsubscribe_completion; bool subscribe_completion_posted = false; @@ -60,6 +66,20 @@ namespace optionx::market_data { subscription_callback_t release_callback; }; + struct PendingContinuityRequest { + RoutedSubscriptionId router_id; + MarketDataSubscriptionHandle subscription; + BarHistoryRequest request; + MarketDataContinuityStatus status = MarketDataContinuityStatus::UNKNOWN; + std::uint64_t from_time_ms = 0; + std::uint64_t to_time_ms = 0; + std::size_t requested_items = 0; + }; + + struct ContinuityOperation { + bool completed = false; + }; + struct CachedStatus { MarketDataStatusUpdate update; std::uint64_t sequence = 0; @@ -191,6 +211,7 @@ namespace optionx::market_data { std::uint64_t m_next_router_id = 1; bool m_shutdown = false; bool m_shutdown_complete = false; + std::size_t m_continuity_operations_in_flight = 0; MarketDataRouter::owner_dispatcher_t m_owner_dispatcher; static StreamDescriptor stream_from(const TickSubscriptionRequest& request); @@ -220,6 +241,7 @@ namespace optionx::market_data { BaseMarketDataProvider& provider, std::weak_ptr subscriber, StreamDescriptor stream, + MarketDataContinuityOptions continuity, MarketDataProviderId registered_provider_id, std::string& error_message); @@ -234,6 +256,57 @@ namespace optionx::market_data { MarketDataSubscriptionResult result, subscription_callback_t callback); + static MarketDataContinuityOptions continuity_from( + const TickSubscriptionRequest&) noexcept { + return {}; + } + + static MarketDataContinuityOptions continuity_from( + const BarSubscriptionRequest& request) noexcept { + return request.continuity; + } + + void start_continuity( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription); + void request_continuity_history( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription, + BarHistoryRequest request, + MarketDataContinuityStatus status, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::size_t requested_items); + void complete_continuity( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription, + BarHistoryRequest request, + MarketDataContinuityStatus operation_status, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::size_t requested_items, + BarHistoryResult result); + void notify_continuity( + RoutedSubscriptionId router_id, + MarketDataContinuityUpdate update); + static MarketDataContinuityUpdate make_continuity_update( + const MarketDataSubscriptionHandle& subscription, + MarketDataContinuityStatus status, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::size_t requested_items, + std::size_t delivered_items, + std::string message); + bool route_bar_to_entry_no_lock( + const std::shared_ptr& entry, + const BarDataBatch& batch, + std::vector, + BarDataBatch>>& deliveries, + std::vector& continuity_requests, + bool process_buffered = false, + bool allow_gap_recovery = true); + void fail_pending_subscribe( RoutedSubscriptionId router_id, BaseMarketDataProvider& provider, @@ -283,6 +356,8 @@ namespace optionx::market_data { BaseMarketDataProvider& provider, const MarketDataSubscriptionResult& result, bool& shutdown_requested); + bool record_continuity_completion( + const std::shared_ptr& operation); void mark_subscribe_completion_posted( RoutedSubscriptionId router_id, const MarketDataSubscriptionHandle& subscription); @@ -417,6 +492,19 @@ namespace optionx::market_data { }); } + inline bool MarketDataRouterState::record_continuity_completion( + const std::shared_ptr& operation) { + if (!operation) return false; + + std::lock_guard lock(m_mutex); + if (operation->completed) return false; + operation->completed = true; + if (m_continuity_operations_in_flight > 0) { + --m_continuity_operations_in_flight; + } + return true; + } + inline MarketDataRouterState::StreamDescriptor MarketDataRouterState::stream_from(const TickSubscriptionRequest& request) { StreamDescriptor stream; @@ -689,6 +777,7 @@ namespace optionx::market_data { BaseMarketDataProvider& provider, std::weak_ptr subscriber, StreamDescriptor stream, + MarketDataContinuityOptions continuity, MarketDataProviderId registered_provider_id, std::string& error_message) { if (subscriber.expired()) { @@ -743,6 +832,9 @@ namespace optionx::market_data { entry->subscriber = std::move(subscriber); entry->control = control; entry->stream = std::move(stream); + entry->continuity = std::move(continuity); + entry->continuity_ready = + !entry->continuity.enabled() || entry->continuity.prefill_bars == 0; m_entries.emplace(router_id, entry); ++provider_it->second.route_count; } @@ -777,6 +869,7 @@ namespace optionx::market_data { provider, std::move(subscriber), stream_from(request), + continuity_from(request), registered_provider_id, error_message); if (!control) { @@ -1007,6 +1100,7 @@ namespace optionx::market_data { MarketDataStatusUpdate replay; bool has_replay = false; bool release_requested = false; + bool needs_continuity_prefill = false; subscription_callback_t release_callback; BaseMarketDataProvider* unbind = nullptr; @@ -1049,6 +1143,9 @@ namespace optionx::market_data { entry->stream = stream_from(result.subscription); entry->phase = EntryPhase::ACTIVE; set_control_active(entry->control, result.subscription); + needs_continuity_prefill = + entry->stream.type == MarketDataType::BARS && + entry->continuity.prefill_bars > 0; release_requested = entry->release_requested; release_callback = std::move(entry->release_callback); @@ -1074,6 +1171,9 @@ namespace optionx::market_data { if (release_callback && !result.success()) { dispatch_result(std::move(release_callback), result); } + if (needs_continuity_prefill && result.success()) { + start_continuity(router_id, result.subscription); + } if (has_replay && subscriber && result.success()) { bool still_active = false; { @@ -1090,6 +1190,322 @@ namespace optionx::market_data { } } + inline void MarketDataRouterState::start_continuity( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription) { + PendingContinuityRequest pending; + 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 || + entry_it->second->continuity.prefill_bars == 0 || + entry_it->second->continuity_request_in_flight) { + return; + } + + const auto& entry = entry_it->second; + BarSubscriptionRequest bar_request( + entry->stream.symbol, + entry->stream.timeframe, + entry->stream.price_source, + entry->stream.transport); + bar_request.continuity = entry->continuity; + const auto now_ms = static_cast(OPTIONX_TIMESTAMP_MS); + pending.router_id = router_id; + pending.subscription = subscription; + pending.request = MarketDataContinuityService::make_prefill_request( + bar_request, + now_ms, + entry->continuity.prefill_bars); + pending.status = MarketDataContinuityStatus::PREFILLING; + pending.from_time_ms = + MarketDataContinuityService::seconds_to_milliseconds( + pending.request.from_ts); + pending.to_time_ms = + MarketDataContinuityService::seconds_to_milliseconds( + pending.request.to_ts); + pending.requested_items = entry->continuity.prefill_bars; + entry->continuity_ready = false; + entry->continuity_request_in_flight = true; + subscriber = entry->subscriber.lock(); + } + + if (subscriber) { + subscriber->on_market_data_continuity(make_continuity_update( + subscription, + MarketDataContinuityStatus::PREFILLING, + pending.from_time_ms, + pending.to_time_ms, + pending.requested_items, + 0, + "Requesting historical bar prefill.")); + } + request_continuity_history( + pending.router_id, + std::move(pending.subscription), + std::move(pending.request), + pending.status, + pending.from_time_ms, + pending.to_time_ms, + pending.requested_items); + } + + inline void MarketDataRouterState::request_continuity_history( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription, + BarHistoryRequest request, + MarketDataContinuityStatus status, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::size_t requested_items) { + BaseMarketDataProvider* provider = nullptr; + auto operation = std::make_shared(); + { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (entry_it == m_entries.end() || + entry_it->second->phase != EntryPhase::ACTIVE || + !entry_it->second->continuity_request_in_flight) { + return; + } + provider = entry_it->second->provider; + if (provider) ++m_continuity_operations_in_flight; + } + if (!provider) return; + + if (status == MarketDataContinuityStatus::GAP_DETECTED) { + notify_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::GAP_DETECTED, + from_time_ms, + to_time_ms, + requested_items, + 0, + "A gap was detected in the live bar stream.")); + notify_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::BACKFILLING, + from_time_ms, + to_time_ms, + requested_items, + 0, + "Requesting historical bars for the detected gap.")); + } + + const auto state = shared_from_this(); + auto complete_on_owner = [state, + operation, + router_id, + subscription, + request, + status, + from_time_ms, + to_time_ms, + requested_items](BarHistoryResult result) mutable { + if (!state->record_continuity_completion(operation)) return; + + auto task = [state, + router_id, + subscription, + request, + status, + from_time_ms, + to_time_ms, + requested_items, + result = std::move(result)]() mutable { + state->complete_continuity( + router_id, + subscription, + request, + status, + from_time_ms, + to_time_ms, + requested_items, + std::move(result)); + }; + state->dispatch_or_run(std::move(task)); + }; + try { + const bool accepted = provider->fetch_bar_history( + request, + [complete_on_owner](BarHistoryResult result) mutable { + complete_on_owner(std::move(result)); + }); + if (!accepted) { + complete_on_owner(BarHistoryResult::fail( + "Market-data provider did not accept the history request.")); + } + } catch (const std::exception& exception) { + complete_on_owner(BarHistoryResult::fail( + std::string("Market-data provider history request threw: ") + + exception.what())); + } catch (...) { + complete_on_owner(BarHistoryResult::fail( + "Market-data provider history request threw.")); + } + } + + inline void MarketDataRouterState::complete_continuity( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle subscription, + BarHistoryRequest request, + MarketDataContinuityStatus operation_status, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::size_t requested_items, + BarHistoryResult 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->continuity_request_in_flight) { + return; + } + expected_stream = entry_it->second->stream; + entry_it->second->continuity_request_in_flight = false; + entry_it->second->continuity_flushing = true; + } + + const bool history_success = static_cast(result); + std::size_t delivered_history_items = 0; + BarDataBatch history_batch; + bool has_history_batch = false; + bool history_stream_matches = false; + if (history_success) { + history_batch = *MarketDataContinuityService::make_bar_batch( + std::move(result.sequence), + request, + subscription, + operation_status == MarketDataContinuityStatus::GAP_DETECTED); + history_stream_matches = batch_matches_stream( + history_batch, + expected_stream); + if (history_stream_matches) { + delivered_history_items = history_batch.items.size(); + has_history_batch = !history_batch.items.empty(); + } + } + + const bool usable_history = history_success && + history_stream_matches && + (operation_status != MarketDataContinuityStatus::GAP_DETECTED || + delivered_history_items > 0); + if (!usable_history) { + notify_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::FAILED, + from_time_ms, + to_time_ms, + requested_items, + 0, + result.error_desc.empty() + ? (!history_success + ? "Historical market-data continuity request failed." + : !history_stream_matches + ? "Historical bar response does not match the subscribed stream." + : "No bars were returned for the detected gap.") + : result.error_desc)); + } + + if (has_history_batch) { + 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->release_requested; + if (active) { + subscriber = entry_it->second->subscriber.lock(); + for (const auto& bar : history_batch.items) { + if (bar.time_ms > entry_it->second->last_bar_time_ms) { + entry_it->second->last_bar_time_ms = bar.time_ms; + } + } + } + } + if (active && subscriber) { + subscriber->on_bar_data(history_batch); + } + } + + for (;;) { + std::vector, + BarDataBatch>> deliveries; + std::vector continuity_requests; + bool finished = 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) { + return; + } + if (entry_it->second->continuity_buffer.empty()) { + entry_it->second->continuity_ready = true; + entry_it->second->continuity_flushing = false; + finished = true; + } else { + auto batch = std::move( + entry_it->second->continuity_buffer.front()); + entry_it->second->continuity_buffer.pop_front(); + route_bar_to_entry_no_lock( + entry_it->second, + batch, + deliveries, + continuity_requests, + true, + usable_history); + } + } + + for (auto& delivery : deliveries) { + delivery.first->on_bar_data(delivery.second); + } + + if (!continuity_requests.empty()) { + auto pending = std::move(continuity_requests.front()); + request_continuity_history( + pending.router_id, + std::move(pending.subscription), + std::move(pending.request), + pending.status, + pending.from_time_ms, + pending.to_time_ms, + pending.requested_items); + return; + } + + if (finished) { + notify_continuity( + router_id, + make_continuity_update( + subscription, + MarketDataContinuityStatus::LIVE, + from_time_ms, + to_time_ms, + requested_items, + delivered_history_items, + usable_history + ? "Historical market-data continuity is ready." + : "Live delivery continues after history failure.")); + return; + } + } + } + inline void MarketDataRouterState::fail_pending_subscribe( RoutedSubscriptionId router_id, BaseMarketDataProvider& provider, @@ -1435,18 +1851,214 @@ namespace optionx::market_data { }); } + inline MarketDataContinuityUpdate + MarketDataRouterState::make_continuity_update( + const MarketDataSubscriptionHandle& subscription, + MarketDataContinuityStatus status, + std::uint64_t from_time_ms, + std::uint64_t to_time_ms, + std::size_t requested_items, + std::size_t delivered_items, + std::string message) { + MarketDataContinuityUpdate update; + update.subscription = subscription; + update.type = subscription.stream_type; + update.symbol = subscription.symbol; + update.timeframe = subscription.timeframe; + update.status = status; + update.from_time_ms = from_time_ms; + update.to_time_ms = to_time_ms; + update.requested_items = requested_items; + update.delivered_items = delivered_items; + update.message = std::move(message); + return update; + } + + inline void MarketDataRouterState::notify_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; + } + subscriber = entry_it->second->subscriber.lock(); + } + if (subscriber) subscriber->on_market_data_continuity(update); + } + + inline bool MarketDataRouterState::route_bar_to_entry_no_lock( + const std::shared_ptr& entry, + const BarDataBatch& batch, + std::vector, + BarDataBatch>>& deliveries, + std::vector& continuity_requests, + bool process_buffered, + bool allow_gap_recovery) { + 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 (entry->continuity.enabled() && !process_buffered && + (!entry->continuity_ready || entry->continuity_flushing)) { + entry->continuity_buffer.push_back(std::move(routed)); + return false; + } + + if (allow_gap_recovery && entry->continuity.recovers_gaps() && + !entry->continuity_request_in_flight && + entry->stream.timeframe > 0 && + entry->last_bar_time_ms > 0) { + const auto timeframe_ms = + static_cast(entry->stream.timeframe) * 1000U; + std::uint64_t previous_time_ms = entry->last_bar_time_ms; + for (std::size_t index = 0; index < batch.items.size(); ++index) { + const auto& bar = batch.items[index]; + if (bar.time_ms == 0) continue; + const auto expected_time_ms = previous_time_ms > + std::numeric_limits::max() - timeframe_ms + ? std::numeric_limits::max() + : previous_time_ms + timeframe_ms; + if (bar.time_ms > expected_time_ms) { + const auto gap_from_ms = expected_time_ms; + const auto gap_to_ms = bar.time_ms - timeframe_ms; + if (gap_to_ms >= gap_from_ms) { + BarSubscriptionRequest request( + entry->stream.symbol, + entry->stream.timeframe, + entry->stream.price_source, + entry->stream.transport); + request.continuity = entry->continuity; + + auto requested_items = static_cast(0); + const auto gap_items = + ((gap_to_ms - gap_from_ms) / timeframe_ms) + 1U; + requested_items = gap_items > + static_cast( + std::numeric_limits::max()) + ? std::numeric_limits::max() + : static_cast(gap_items); + if (entry->continuity.max_backfill_bars > 0) { + requested_items = std::min( + requested_items, + entry->continuity.max_backfill_bars); + } + + if (index > 0) { + BarDataBatch prefix = routed; + prefix.items.resize(index); + deliveries.emplace_back(subscriber, std::move(prefix)); + for (std::size_t prefix_index = 0; + prefix_index < index; + ++prefix_index) { + entry->last_bar_time_ms = std::max( + entry->last_bar_time_ms, + batch.items[prefix_index].time_ms); + } + } + + routed.items.erase( + routed.items.begin(), + routed.items.begin() + static_cast(index)); + entry->continuity_ready = false; + entry->continuity_request_in_flight = true; + if (process_buffered) { + entry->continuity_buffer.push_front(std::move(routed)); + } else { + entry->continuity_buffer.push_back(std::move(routed)); + } + continuity_requests.push_back(PendingContinuityRequest{ + entry->router_id, + entry->control->provider_subscription, + MarketDataContinuityService::make_gap_request( + request, + gap_from_ms, + gap_to_ms, + entry->continuity.max_backfill_bars), + MarketDataContinuityStatus::GAP_DETECTED, + gap_from_ms, + gap_to_ms, + requested_items}); + return false; + } + } + if (bar.time_ms > previous_time_ms) { + previous_time_ms = bar.time_ms; + } + } + } + + for (const auto& bar : routed.items) { + if (bar.time_ms > entry->last_bar_time_ms) { + entry->last_bar_time_ms = bar.time_ms; + } + } + deliveries.emplace_back(std::move(subscriber), std::move(routed)); + return true; + } + inline void MarketDataRouterState::route_bars( ProviderInstanceId provider_id, std::unique_ptr batch) { - route_batch( - provider_id, - std::move(batch), - [this](const BarDataBatch& data, const StreamDescriptor& stream) { - return batch_matches_stream(data, stream); - }, - [](IMarketDataSubscriber& subscriber, BarDataBatch& data) { - subscriber.on_bar_data(data); - }); + if (!batch) return; + + std::vector, + BarDataBatch>> 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 route_one = [&](const std::shared_ptr& entry) { + route_bar_to_entry_no_lock( + entry, + *batch, + deliveries, + continuity_requests, + false); + }; + + if (batch->subscription.valid()) { + if (batch->subscription.provider_id != provider_id) return; + const auto route_it = provider_it->second.provider_routes.find( + 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()) route_one(entry_it->second); + } else { + for (const auto& [id, entry] : m_entries) { + (void)id; + if (entry->provider_id == provider_id) route_one(entry); + } + } + } + + for (auto& delivery : deliveries) { + delivery.first->on_bar_data(delivery.second); + } + for (auto& request : continuity_requests) { + request_continuity_history( + request.router_id, + std::move(request.subscription), + std::move(request.request), + request.status, + request.from_time_ms, + request.to_time_ms, + request.requested_items); + } } inline void MarketDataRouterState::route_status( @@ -1656,7 +2268,8 @@ namespace optionx::market_data { } } - if (m_entries.empty()) { + if (m_entries.empty() && + m_continuity_operations_in_flight == 0) { m_registered_providers.clear(); m_provider_aliases.clear(); m_registered_provider_ids.clear(); diff --git a/include/optionx_cpp/platforms/IntradeBarPlatform/TradingConditionManager.hpp b/include/optionx_cpp/platforms/IntradeBarPlatform/TradingConditionManager.hpp index 05da6c9..a66f8a0 100644 --- a/include/optionx_cpp/platforms/IntradeBarPlatform/TradingConditionManager.hpp +++ b/include/optionx_cpp/platforms/IntradeBarPlatform/TradingConditionManager.hpp @@ -76,6 +76,19 @@ namespace optionx::platforms::intrade_bar { /// \brief Clears local condition state during platform shutdown. void shutdown() override { + const auto timestamp = current_timestamp_sec(); + for (const auto& previous : m_last_snapshots) { + TradingConditionUpdate retired; + retired.symbol = previous.symbol; + retired.platform_type = previous.platform_type; + retired.account_type = previous.account_type; + retired.currency = previous.currency; + retired.option_type = previous.option_type; + retired.timestamp = timestamp; + retired.tradable = false; + retired.message = "Intrade Bar account condition manager stopped."; + publish(retired); + } m_connected = false; m_last_refresh_sec = 0; m_last_snapshots.clear(); diff --git a/tests/intrade_bar_api/intrade_bar_api_response_test.cpp b/tests/intrade_bar_api/intrade_bar_api_response_test.cpp index 1c963c7..9565608 100644 --- a/tests/intrade_bar_api/intrade_bar_api_response_test.cpp +++ b/tests/intrade_bar_api/intrade_bar_api_response_test.cpp @@ -2258,6 +2258,38 @@ TEST(IntradeBarTradingConditions, OpenTradeLimitUpdatesPlatformBoundHub) { platform.shutdown(); } +TEST(IntradeBarTradingConditions, ShutdownPublishesUnavailableSnapshots) { + IntradeBarPlatform platform; + std::vector updates; + platform.on_trading_condition() = + [&updates](const TradingConditionUpdate& update) { + updates.push_back(update); + }; + + auto account = std::make_shared(); + account->connect = true; + account->account_type = AccountType::DEMO; + account->currency = CurrencyType::USD; + platform.event_bus().notify_async( + std::make_unique( + account, + AccountUpdateStatus::CONNECTED)); + platform.event_bus().drain(); + + const auto published_before_shutdown = updates.size(); + ASSERT_GT(published_before_shutdown, 0U); + + platform.shutdown(); + + ASSERT_EQ(updates.size(), published_before_shutdown * 2U); + for (std::size_t index = published_before_shutdown; + index < updates.size(); + ++index) { + EXPECT_EQ(updates[index].tradable, std::optional(false)); + EXPECT_FALSE(updates[index].symbol.empty()); + } +} + TEST(IntradeBarAccountInfo, AcceptsBtcAliasAndUsesBtcDurationRules) { AccountInfoData account; const int64_t day_timestamp = 1712345600; diff --git a/tests/market_data_continuity_test.cpp b/tests/market_data_continuity_test.cpp new file mode 100644 index 0000000..1c2cfb2 --- /dev/null +++ b/tests/market_data_continuity_test.cpp @@ -0,0 +1,485 @@ +#include + +#include +#include +#include +#include +#include +#include +#include + +#include + +using namespace optionx; +using namespace optionx::market_data; + +namespace { + +class OwnerQueue { +public: + bool post(MarketDataRouter::owner_task_t task) { + if (!m_accepting || !task) return false; + m_tasks.push_back(std::move(task)); + return true; + } + + void drain() { + while (!m_tasks.empty()) { + auto task = std::move(m_tasks.front()); + m_tasks.pop_front(); + task(); + } + } + + void stop_accepting() noexcept { + m_accepting = false; + } + +private: + std::deque m_tasks; + bool m_accepting = true; +}; + +class FakeHistoryProvider final : public BaseMarketDataProvider { +public: + bars_callback_t& on_bar_data() override { + return m_bar_callback; + } + + 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_bars( + BarSubscriptionRequest request, + subscription_callback_t callback) override { + auto subscription = MarketDataSubscriptionHandle::from_bar_request( + provider_id(), + m_next_subscription_id++, + request); + m_active_subscription = subscription; + if (callback) { + callback(MarketDataSubscriptionResult::subscribed( + std::move(subscription))); + } + return true; + } + + bool unsubscribe( + MarketDataSubscriptionHandle subscription, + subscription_callback_t callback) override { + if (callback) { + callback(MarketDataSubscriptionResult::unsubscribed( + std::move(subscription))); + } + return true; + } + + bool fetch_bar_history( + const BarHistoryRequest& request, + bar_history_callback_t callback) override { + history_requests.push_back(request); + if (reject_history) return false; + m_history_callback = std::move(callback); + return true; + } + + void complete_history(BarSequence sequence) { + ASSERT_TRUE(static_cast(m_history_callback)); + auto callback = std::move(m_history_callback); + callback(BarHistoryResult::ok(std::move(sequence))); + } + + void fail_history(std::string message) { + ASSERT_TRUE(static_cast(m_history_callback)); + auto callback = std::move(m_history_callback); + callback(BarHistoryResult::fail(std::move(message))); + } + + void emit_live_bar(std::uint64_t time_ms) { + emit_live_bars({time_ms}); + } + + void emit_live_bars(std::initializer_list times) { + auto batch = std::make_unique(); + batch->subscription = m_active_subscription; + batch->type = MarketDataType::BARS; + batch->symbol = m_active_subscription.symbol; + batch->timeframe = m_active_subscription.timeframe; + for (const auto time_ms : times) { + batch->items.emplace_back(1.0, 2.0, 0.5, 1.5, 1.0, time_ms); + batch->items.back().set_flag(MarketDataFlags::REALTIME); + } + if (m_bar_callback) m_bar_callback(std::move(batch)); + } + + std::vector history_requests; + bool reject_history = false; + +private: + SubscriptionId m_next_subscription_id = 1; + MarketDataSubscriptionHandle m_active_subscription; + bar_history_callback_t m_history_callback; + bars_callback_t m_bar_callback; + ticks_callback_t m_tick_callback; + status_callback_t m_status_callback; +}; + +class RecordingSubscriber final : public IMarketDataSubscriber { +public: + void on_bar_data(const BarDataBatch& batch) override { + bars.push_back(batch); + } + + void on_market_data_continuity( + const MarketDataContinuityUpdate& update) override { + continuity.push_back(update); + } + + std::vector bars; + std::vector continuity; +}; + +BarSequence make_history(std::initializer_list times) { + BarSequence sequence; + sequence.symbol = "EURUSD"; + sequence.provider = "fake"; + sequence.timeframe = 60; + sequence.price_digits = 5; + sequence.volume_digits = 0; + sequence.price_source = BarPriceSource::MID; + for (const auto time_ms : times) { + sequence.bars.emplace_back(1.0, 2.0, 0.5, 1.5, 1.0, time_ms); + } + return sequence; +} + +BarSubscriptionRequest continuity_request(MarketDataContinuityMode mode) { + BarSubscriptionRequest request( + "EURUSD", + 60, + BarPriceSource::MID, + MarketDataTransport::WEBSOCKET); + request.continuity.mode = mode; + request.continuity.prefill_bars = 1; + request.continuity.max_backfill_bars = 10; + return request; +} + +TEST(MarketDataContinuity, PrefillDeliversHistoryBeforeBufferedLiveBars) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL)); + ASSERT_TRUE(route.active()); + ASSERT_EQ(provider.history_requests.size(), 1U); + + provider.emit_live_bar(200000); + EXPECT_TRUE(subscriber->bars.empty()); + + provider.complete_history(make_history({100000})); + + ASSERT_EQ(subscriber->bars.size(), 2U); + EXPECT_TRUE(subscriber->bars[0].items[0].has_flag(MarketDataFlags::HISTORICAL)); + EXPECT_FALSE(subscriber->bars[0].items[0].has_flag(MarketDataFlags::BACKFILL)); + EXPECT_TRUE(subscriber->bars[1].items[0].has_flag(MarketDataFlags::REALTIME)); + ASSERT_EQ(subscriber->continuity.size(), 2U); + EXPECT_EQ(subscriber->continuity[0].status, MarketDataContinuityStatus::PREFILLING); + EXPECT_EQ(subscriber->continuity[1].status, MarketDataContinuityStatus::LIVE); + EXPECT_EQ( + subscriber->continuity[0].subscription.id, + route.provider_subscription().id); + EXPECT_EQ( + subscriber->continuity[1].subscription.id, + route.provider_subscription().id); +} + +TEST(MarketDataContinuity, RecoversTimestampGapBeforeReleasingLiveBatch) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + ASSERT_TRUE(route.active()); + provider.complete_history(make_history({100000})); + ASSERT_EQ(subscriber->bars.size(), 1U); + + provider.emit_live_bar(280000); + EXPECT_EQ(provider.history_requests.size(), 2U); + ASSERT_EQ(subscriber->continuity.size(), 4U); + EXPECT_EQ(subscriber->continuity[2].status, MarketDataContinuityStatus::GAP_DETECTED); + EXPECT_EQ(subscriber->continuity[3].status, MarketDataContinuityStatus::BACKFILLING); + EXPECT_EQ(subscriber->bars.size(), 1U); + + provider.complete_history(make_history({160000, 220000})); + + ASSERT_EQ(subscriber->bars.size(), 3U); + EXPECT_TRUE(subscriber->bars[1].items[0].has_flag(MarketDataFlags::BACKFILL)); + EXPECT_TRUE(subscriber->bars[2].items[0].has_flag(MarketDataFlags::REALTIME)); + ASSERT_EQ(subscriber->continuity.size(), 5U); + EXPECT_EQ(subscriber->continuity.back().status, MarketDataContinuityStatus::LIVE); + EXPECT_EQ( + subscriber->continuity[2].subscription.id, + route.provider_subscription().id); + EXPECT_EQ( + subscriber->continuity[3].subscription.id, + route.provider_subscription().id); +} + +TEST(MarketDataContinuity, ChecksGapsInsideBufferedBatchesInOrder) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + ASSERT_TRUE(route.active()); + + provider.emit_live_bars({160000, 280000}); + provider.emit_live_bar(400000); + provider.complete_history(make_history({100000})); + + ASSERT_EQ(provider.history_requests.size(), 2U); + ASSERT_EQ(subscriber->bars.size(), 2U); + EXPECT_EQ(subscriber->bars[0].items[0].time_ms, 100000U); + EXPECT_EQ(subscriber->bars[1].items[0].time_ms, 160000U); + + provider.complete_history(make_history({220000})); + + ASSERT_EQ(provider.history_requests.size(), 3U); + ASSERT_EQ(subscriber->bars.size(), 4U); + EXPECT_EQ(subscriber->bars[2].items[0].time_ms, 220000U); + EXPECT_EQ(subscriber->bars[3].items[0].time_ms, 280000U); + + provider.complete_history(make_history({340000})); + + ASSERT_EQ(subscriber->bars.size(), 6U); + EXPECT_EQ(subscriber->bars[4].items[0].time_ms, 340000U); + EXPECT_EQ(subscriber->bars[5].items[0].time_ms, 400000U); + EXPECT_TRUE(subscriber->bars[2].items[0].has_flag(MarketDataFlags::BACKFILL)); + EXPECT_TRUE(subscriber->bars[3].items[0].has_flag(MarketDataFlags::REALTIME)); + EXPECT_TRUE(subscriber->bars[4].items[0].has_flag(MarketDataFlags::BACKFILL)); +} + +TEST(MarketDataContinuity, EmptyBackfillFailsOnceAndReleasesLiveBatch) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + ASSERT_TRUE(route.active()); + provider.complete_history(make_history({100000})); + + provider.emit_live_bar(280000); + ASSERT_EQ(provider.history_requests.size(), 2U); + provider.complete_history({}); + + ASSERT_EQ(provider.history_requests.size(), 2U); + ASSERT_EQ(subscriber->bars.size(), 2U); + EXPECT_EQ(subscriber->bars.back().items.front().time_ms, 280000U); + ASSERT_EQ(subscriber->continuity.size(), 6U); + EXPECT_EQ(subscriber->continuity[4].status, MarketDataContinuityStatus::FAILED); + EXPECT_EQ(subscriber->continuity[5].status, MarketDataContinuityStatus::LIVE); +} + +TEST(MarketDataContinuity, HistoryFailureKeepsLiveRouteUsable) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL)); + ASSERT_TRUE(route.active()); + + provider.emit_live_bar(200000); + provider.fail_history("history endpoint unavailable"); + + ASSERT_EQ(subscriber->bars.size(), 1U); + EXPECT_TRUE(subscriber->bars.front().items.front().has_flag( + MarketDataFlags::REALTIME)); + ASSERT_EQ(subscriber->continuity.size(), 3U); + EXPECT_EQ(subscriber->continuity[1].status, MarketDataContinuityStatus::FAILED); + EXPECT_EQ(subscriber->continuity[1].message, "history endpoint unavailable"); + EXPECT_EQ(subscriber->continuity[2].status, MarketDataContinuityStatus::LIVE); + + provider.emit_live_bar(260000); + EXPECT_EQ(subscriber->bars.size(), 2U); +} + +TEST(MarketDataContinuity, ShutdownWaitsForDeferredHistoryOperation) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL)); + ASSERT_TRUE(route.active()); + ASSERT_EQ(provider.history_requests.size(), 1U); + + router.shutdown(); + EXPECT_FALSE(router.is_shutdown_complete()); + + provider.complete_history(make_history({100000})); + router.process(); + + EXPECT_TRUE(router.is_shutdown_complete()); + EXPECT_TRUE(subscriber->bars.empty()); + ASSERT_EQ(subscriber->continuity.size(), 1U); + EXPECT_EQ( + subscriber->continuity.front().status, + MarketDataContinuityStatus::PREFILLING); +} + +TEST(MarketDataContinuity, ShutdownWaitsForFailedHistoryOperation) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL)); + ASSERT_TRUE(route.active()); + + router.shutdown(); + EXPECT_FALSE(router.is_shutdown_complete()); + + provider.fail_history("history endpoint unavailable"); + router.process(); + + EXPECT_TRUE(router.is_shutdown_complete()); + EXPECT_TRUE(subscriber->bars.empty()); + ASSERT_EQ(subscriber->continuity.size(), 1U); + EXPECT_EQ( + subscriber->continuity.front().status, + MarketDataContinuityStatus::PREFILLING); +} + +TEST(MarketDataContinuity, RejectsHistoryResponseFromAnotherStream) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL_AND_RECOVER)); + ASSERT_TRUE(route.active()); + + auto wrong_stream = make_history({100000}); + wrong_stream.symbol = "BTCUSDT"; + wrong_stream.timeframe = 300; + provider.complete_history(std::move(wrong_stream)); + + EXPECT_TRUE(subscriber->bars.empty()); + ASSERT_EQ(subscriber->continuity.size(), 3U); + EXPECT_EQ( + subscriber->continuity[1].status, + MarketDataContinuityStatus::FAILED); + EXPECT_EQ( + subscriber->continuity[1].message, + "Historical bar response does not match the subscribed stream."); + EXPECT_EQ( + subscriber->continuity[2].status, + MarketDataContinuityStatus::LIVE); + + provider.emit_live_bar(280000); + EXPECT_EQ(provider.history_requests.size(), 1U); + ASSERT_EQ(subscriber->bars.size(), 1U); + EXPECT_EQ(subscriber->bars.front().items.front().time_ms, 280000U); +} + +TEST(MarketDataContinuity, PrefillRequestUsesInclusiveBarCountRange) { + BarSubscriptionRequest request( + "EURUSD", + 60, + BarPriceSource::MID, + MarketDataTransport::WEBSOCKET); + + const auto history = MarketDataContinuityService::make_prefill_request( + request, + 600000, + 3); + + EXPECT_EQ(history.from_ts, 480); + EXPECT_EQ(history.to_ts, 600); +} + +TEST(MarketDataContinuity, UnsubscribeDuringHistoryDropsLateDelivery) { + FakeHistoryProvider provider; + MarketDataRouter router; + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL)); + ASSERT_TRUE(route.active()); + route.reset(); + + provider.complete_history(make_history({100000})); + + EXPECT_TRUE(subscriber->bars.empty()); + EXPECT_EQ(router.subscription_count(), 0U); +} + +TEST(MarketDataContinuity, RejectedOwnerDispatchDoesNotDeliverOnProviderThread) { + FakeHistoryProvider provider; + OwnerQueue owner_queue; + MarketDataRouter router( + [&owner_queue](MarketDataRouter::owner_task_t task) { + return owner_queue.post(std::move(task)); + }); + auto subscriber = std::make_shared(); + + auto route = router.subscribe_bars( + provider, + subscriber, + continuity_request(MarketDataContinuityMode::PREFILL)); + EXPECT_TRUE(route.valid()); + EXPECT_FALSE(route.active()); + + owner_queue.drain(); + ASSERT_TRUE(route.active()); + provider.emit_live_bar(200000); + owner_queue.drain(); + + owner_queue.stop_accepting(); + provider.complete_history(make_history({100000})); + + EXPECT_TRUE(subscriber->bars.empty()); + ASSERT_EQ(subscriber->continuity.size(), 1U); + EXPECT_EQ( + subscriber->continuity.front().status, + MarketDataContinuityStatus::PREFILLING); + + router.shutdown(); + EXPECT_TRUE(router.is_shutdown_complete()); +} + +} // namespace + +int main(int argc, char** argv) { + testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +}