diff --git a/guides/market-data-router.md b/guides/market-data-router.md index e6b8595..c32e116 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -721,6 +721,15 @@ unrelated successful range cannot hide the earlier unresolved watermark. History requests are bounded by `max_backfill_ms` and are scheduled by `process()`, so tick continuity does not create a timer thread. +The same bound applies to the initial prefill. When `max_backfill_ms` is +non-zero, Router splits the lookback into inclusive chunks, starts the next +chunk at the previous end point when overlap can advance the range, and keeps +the live tail buffered until the final chunk reaches the prefill boundary. +With a provider history grid narrower than the bound, the resulting requests +are `5000..6000`, `6000..7000`, and so on; if the bound cannot contain one +grid step, Router advances to the next provider boundary to avoid repeating a +request. A zero bound keeps the provider range unbounded. + On reconnect, tick continuity reports `STALE`, waits for `READY`, and requests the unresolved range through the latest observed time. History overlap is removed according to the selected tick identity policy. The default provider @@ -730,8 +739,9 @@ make an otherwise identical observation distinct. If the continuity buffer excee its batch or item limit, Router releases the held live data, reports `FAILED`/`DEGRADED`, disables continuity for that route, and resumes ordinary live delivery. If transport is interrupted during the initial prefill, Router -restarts from the original lookback start after `READY` and extends the request -through the current time, so the interrupted interval is not silently skipped. +restarts from the original lookback start after `READY` and extends the +chunked request sequence through the current time, so the interrupted interval +is not silently skipped. ## Owner Loop And Bot Threads diff --git a/guides/market-data-router.ru.md b/guides/market-data-router.ru.md index 7816f89..0af130f 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -877,6 +877,14 @@ Router сначала отправляет historical ticks с флагом `HIS ограничивает каждый history request, а `process()` обслуживает retries без создания отдельного timer thread. +То же ограничение действует для initial prefill: если значение не равно +нулю, Router делит lookback на inclusive chunks и удерживает live tail в +buffer до завершения последнего chunk. При provider grid и лимите, который +допускает overlap, запросы идут как `5000..6000`, `6000..7000` и так далее. +Если лимит меньше одного шага grid, Router переходит к следующей provider +boundary, чтобы не повторять тот же запрос. Нулевой лимит оставляет диапазон +provider без ограничения. + После reconnect tick continuity публикует `STALE`, ждёт `READY` и запрашивает unresolved range до последнего observed time. History overlap удаляется по выбранной tick identity policy. Provider default для Intrade использует только @@ -886,5 +894,6 @@ snapshot новым. При переполнении buffer Router освобо `FAILED`/`DEGRADED`, отключает continuity для этого route и возобновляет обычную live delivery. Если transport прервался во время initial prefill, после `READY` Router -начинает повторный запрос с исходного начала lookback и расширяет его до -текущего времени, поэтому прерванный интервал не пропускается молча. +начинает повторную chunked-последовательность с исходного начала lookback и +расширяет её до текущего времени, поэтому прерванный интервал не пропускается +молча. diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp index 05457e1..c7c420d 100644 --- a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -615,6 +615,13 @@ namespace optionx::market_data { std::uint64_t target_time_ms, std::uint64_t max_backfill_ms, std::uint64_t interval_ms) noexcept; + static bool next_tick_history_range( + std::uint64_t previous_to_time_ms, + std::uint64_t target_time_ms, + std::uint64_t max_backfill_ms, + std::uint64_t interval_ms, + std::uint64_t& next_from_time_ms, + std::uint64_t& next_to_time_ms) noexcept; static std::uint64_t duration_to_milliseconds( std::chrono::steady_clock::duration duration) noexcept; static MarketDataContinuitySnapshot make_continuity_snapshot_no_lock( @@ -1736,6 +1743,7 @@ namespace optionx::market_data { const auto current_boundary_time_ms = std::max( original_boundary_time_ms, current_time_ms); + continuity.reconnect_target_time_ms = current_boundary_time_ms; const auto requested_from_time_ms = original_boundary_time_ms > lookback ? original_boundary_time_ms - lookback : 1U; @@ -1746,7 +1754,16 @@ namespace optionx::market_data { current_boundary_time_ms >= history_interval_ms) { pending.from_time_ms = history_interval_ms; } - pending.to_time_ms = current_boundary_time_ms; + pending.to_time_ms = bounded_tick_history_to_time_ms( + pending.from_time_ms, + current_boundary_time_ms, + entry->tick_continuity.max_backfill_ms, + history_interval_ms); + if (pending.to_time_ms < pending.from_time_ms && + history_interval_ms > 0) { + pending.to_time_ms = pending.from_time_ms; + } + if (pending.to_time_ms < pending.from_time_ms) return; pending.request = TickHistoryRequest( entry->stream.symbol, pending.from_time_ms, @@ -2663,6 +2680,41 @@ namespace optionx::market_data { return aligned_bound >= from_time_ms ? aligned_bound : from_time_ms; } + inline bool MarketDataRouterState::next_tick_history_range( + std::uint64_t previous_to_time_ms, + std::uint64_t target_time_ms, + std::uint64_t max_backfill_ms, + std::uint64_t interval_ms, + std::uint64_t& next_from_time_ms, + std::uint64_t& next_to_time_ms) noexcept { + if (target_time_ms <= previous_to_time_ms) return false; + + const bool can_overlap = max_backfill_ms == 0 || + (interval_ms == 0 + ? max_backfill_ms > 1U + : max_backfill_ms > interval_ms); + next_from_time_ms = previous_to_time_ms; + if (!can_overlap) { + const auto step = interval_ms > 0 ? interval_ms : 1U; + if (previous_to_time_ms > + std::numeric_limits::max() - step) { + return false; + } + next_from_time_ms = previous_to_time_ms + step; + } + if (next_from_time_ms > target_time_ms) return false; + + next_to_time_ms = bounded_tick_history_to_time_ms( + next_from_time_ms, + target_time_ms, + max_backfill_ms, + interval_ms); + if (next_to_time_ms < next_from_time_ms && interval_ms > 0) { + next_to_time_ms = next_from_time_ms; + } + return next_to_time_ms >= next_from_time_ms; + } + inline bool MarketDataRouterState::tick_history_covers_range( const TickDataBatch& batch, std::uint64_t from_time_ms, @@ -3298,7 +3350,9 @@ namespace optionx::market_data { const auto& entry = entry_it->second; auto& continuity = entry->tick_continuity_state; if (kind == ContinuityRequestKind::PREFILL) { - continuity.initial_prefill_pending = false; + if (!usable_history) { + continuity.initial_prefill_pending = false; + } } if (!usable_history) { mark_tick_unverified_no_lock(entry, from_time_ms); @@ -3405,6 +3459,7 @@ namespace optionx::market_data { } bool schedule_next = false; + bool notify_next_prefill = false; PendingTickContinuityRequest next_request; { std::lock_guard lock(m_mutex); @@ -3422,39 +3477,55 @@ namespace optionx::market_data { align_tick_history_end_time_ms( continuity.reconnect_target_time_ms, history_interval_ms); - if (usable_history && - kind != ContinuityRequestKind::PREFILL && - history_target_time_ms > to_time_ms) { - const auto max_backfill_ms = entry->tick_continuity.max_backfill_ms; - // Inclusive ranges overlap at the boundary whenever the - // configured limit can contain that overlap. If a - // provider grid is wider than the limit, advance to the - // next representable boundary to avoid repeating a range. - const bool can_overlap = max_backfill_ms == 0 || - (history_interval_ms == 0 - ? max_backfill_ms > 1U - : max_backfill_ms > history_interval_ms); - std::uint64_t next_from = to_time_ms; - if (!can_overlap) { - const auto step = history_interval_ms > 0 - ? history_interval_ms - : 1U; - if (to_time_ms > - std::numeric_limits::max() - step) { - return; + if (usable_history && kind == ContinuityRequestKind::PREFILL) { + const auto prefill_target_time_ms = history_target_time_ms > 0 + ? history_target_time_ms + : continuity.initial_prefill_boundary_time_ms; + if (prefill_target_time_ms > to_time_ms) { + std::uint64_t next_from = 0; + std::uint64_t next_to = 0; + if (next_tick_history_range( + to_time_ms, + prefill_target_time_ms, + entry->tick_continuity.max_backfill_ms, + history_interval_ms, + next_from, + next_to)) { + next_request.router_id = router_id; + next_request.subscription = + entry->control->provider_subscription; + next_request.request = TickHistoryRequest( + entry->stream.symbol, + next_from, + next_to); + next_request.kind = ContinuityRequestKind::PREFILL; + next_request.from_time_ms = next_from; + next_request.to_time_ms = next_to; + next_request.requested_items = 0; + next_request.attempt = 1; + continuity.phase = + MarketDataContinuityPhase::PREFILLING; + continuity.request_in_flight = true; + schedule_next = true; + notify_next_prefill = true; } - next_from = to_time_ms + step; } - if (next_from > history_target_time_ms) { - return; + if (!schedule_next) { + continuity.initial_prefill_pending = false; } - auto next_to = bounded_tick_history_to_time_ms( - next_from, - history_target_time_ms, - max_backfill_ms, - history_interval_ms); - if (next_to < next_from && history_interval_ms > 0) { - next_to = next_from; + } else if (usable_history && + kind != ContinuityRequestKind::PREFILL && + history_target_time_ms > to_time_ms) { + std::uint64_t next_from = 0; + std::uint64_t next_to = 0; + if (!next_tick_history_range( + to_time_ms, + history_target_time_ms, + entry->tick_continuity.max_backfill_ms, + history_interval_ms, + next_from, + next_to)) { + return; } if (next_to < next_from) return; next_request.router_id = router_id; @@ -3474,6 +3545,18 @@ namespace optionx::market_data { } } if (schedule_next) { + if (notify_next_prefill) { + notify_tick_continuity( + router_id, + make_continuity_update( + next_request.subscription, + MarketDataContinuityStatus::PREFILLING, + next_request.from_time_ms, + next_request.to_time_ms, + next_request.requested_items, + 0, + "Requesting the next historical tick prefill chunk.")); + } request_tick_continuity_history( next_request.router_id, std::move(next_request.subscription), diff --git a/tests/market_data_tick_continuity_test.cpp b/tests/market_data_tick_continuity_test.cpp index c963277..8c871cb 100644 --- a/tests/market_data_tick_continuity_test.cpp +++ b/tests/market_data_tick_continuity_test.cpp @@ -273,6 +273,72 @@ TEST(MarketDataTickContinuity, CompletesPrefillAndDeliversDirectRealtime) { EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); } +TEST(MarketDataTickContinuity, BoundsInitialPrefillByMaxBackfill) { + ScopedTestClock clock(10000); + 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, + 5000, + 1500)); + ASSERT_TRUE(route.valid()); + ASSERT_EQ(provider.history_requests.size(), 1U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 5000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 6000U); + + provider.emit_ticks({make_tick(10001)}); + EXPECT_TRUE(subscriber->ticks.empty()); + + provider.complete_history(make_history({5000, 6000})); + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 6000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 7000U); + provider.complete_history(make_history({6000, 7000})); + ASSERT_EQ(provider.history_requests.size(), 3U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 7000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 8000U); + provider.complete_history(make_history({7000, 8000})); + ASSERT_EQ(provider.history_requests.size(), 4U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 8000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 9000U); + provider.complete_history(make_history({8000, 9000})); + ASSERT_EQ(provider.history_requests.size(), 5U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 9000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 10000U); + + EXPECT_EQ(std::count_if( + subscriber->ticks.begin(), + subscriber->ticks.end(), + [](const TickDataBatch& batch) { + return std::any_of( + batch.items.begin(), + batch.items.end(), + [](const Tick& tick) { return tick.time_ms == 10001U; }); + }), 0); + + provider.complete_history(make_history({9000, 10000})); + ASSERT_EQ(subscriber->ticks.size(), 6U); + ASSERT_EQ(subscriber->ticks.back().items.size(), 1U); + EXPECT_EQ(subscriber->ticks.back().items.front().time_ms, 10001U); + EXPECT_TRUE(subscriber->ticks.back().items.front().has_flag( + MarketDataFlags::LIVE_SOURCE)); + EXPECT_TRUE(subscriber->ticks.back().items.front().has_flag( + MarketDataFlags::CATCHUP)); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 1U); + + const auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::LIVE); + EXPECT_EQ(snapshot->verified_through_time_ms, 10001U); + EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); +} + TEST(MarketDataTickContinuity, UsesProviderAlignedClockForPrefill) { ScopedTestClock clock(3123); FakeTickHistoryProvider provider; @@ -358,6 +424,86 @@ TEST(MarketDataTickContinuity, RestartsInterruptedPrefillFromOriginalRange) { EXPECT_EQ(snapshot->unverified_from_time_ms, 0U); } +TEST(MarketDataTickContinuity, RestartsChunkedPrefillThroughCurrentBoundary) { + 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()); + ASSERT_EQ(provider.history_requests.size(), 1U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 2000U); + + provider.emit_status(MarketDataStreamStatus::DISCONNECTED); + provider.emit_ticks({make_tick(6001)}); + market_data_tick_continuity_test_clock::now_ms = 6000; + provider.emit_status(MarketDataStreamStatus::READY); + + ASSERT_EQ(provider.history_requests.size(), 2U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 1000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 2000U); + + // The completion from before transport loss is obsolete. The restarted + // request must remain in flight and keep the buffered tick held. + provider.complete_history(make_history({1000, 2000})); + EXPECT_TRUE(subscriber->ticks.empty()); + auto snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_TRUE(snapshot->request_in_flight); + + provider.complete_history(make_history({1000, 2000})); + ASSERT_EQ(provider.history_requests.size(), 3U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 2000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 3000U); + provider.complete_history(make_history({2000, 3000})); + ASSERT_EQ(provider.history_requests.size(), 4U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 3000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 4000U); + provider.complete_history(make_history({3000, 4000})); + ASSERT_EQ(provider.history_requests.size(), 5U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 4000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 5000U); + provider.complete_history(make_history({4000, 5000})); + ASSERT_EQ(provider.history_requests.size(), 6U); + EXPECT_EQ(provider.history_requests.back().from_time_ms, 5000U); + EXPECT_EQ(provider.history_requests.back().to_time_ms, 6000U); + + EXPECT_TRUE(std::none_of( + subscriber->ticks.begin(), + subscriber->ticks.end(), + [](const TickDataBatch& batch) { + return std::any_of( + batch.items.begin(), + batch.items.end(), + [](const Tick& tick) { return tick.time_ms == 6001U; }); + })); + + provider.complete_history(make_history({5000, 6000})); + ASSERT_FALSE(subscriber->ticks.empty()); + const auto& catchup_batch = subscriber->ticks.back(); + ASSERT_EQ(catchup_batch.items.size(), 1U); + EXPECT_EQ(catchup_batch.items.front().time_ms, 6001U); + EXPECT_TRUE(catchup_batch.items.front().has_flag( + MarketDataFlags::LIVE_SOURCE)); + EXPECT_TRUE(catchup_batch.items.front().has_flag( + MarketDataFlags::CATCHUP)); + EXPECT_EQ(count_status(*subscriber, MarketDataContinuityStatus::LIVE), 1U); + + snapshot = router.continuity_snapshot(route.router_id()); + ASSERT_TRUE(snapshot.has_value()); + EXPECT_EQ(snapshot->phase, MarketDataContinuityPhase::LIVE); + EXPECT_EQ(snapshot->verified_through_time_ms, 6001U); +} + TEST(MarketDataTickContinuity, ChecksGapsInsidePrefillBacklog) { ScopedTestClock clock(3000); FakeTickHistoryProvider provider; @@ -556,31 +702,24 @@ TEST(MarketDataTickContinuity, DeduplicatesInclusiveHistoryAcrossDeliveredBatche 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)})); + provider.emit_ticks({make_tick(4501, 5.0)}); - ASSERT_EQ(provider.history_requests.size(), 4U); + ASSERT_EQ(provider.history_requests.size(), 3U); 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(3000, 3.1), make_tick(4000, 4.0)})); - ASSERT_EQ(provider.history_requests.size(), 4U); + ASSERT_EQ(provider.history_requests.size(), 3U); std::vector delivered; for (const auto& batch : subscriber->ticks) { delivered.insert(delivered.end(), batch.items.begin(), batch.items.end());