diff --git a/.github/workflows/ubuntu-smoke.yml b/.github/workflows/ubuntu-smoke.yml index 2636ed7..da0c279 100644 --- a/.github/workflows/ubuntu-smoke.yml +++ b/.github/workflows/ubuntu-smoke.yml @@ -32,6 +32,14 @@ jobs: - name: Configure run: cmake -S . -B build-linux -DOPTIONX_BUILD_DEPS=ON -DOPTIONX_BUILD_TESTS=ON -DOPTIONX_BUILD_EXAMPLES=ON + - name: Build lifecycle stack test and example + run: cmake --build build-linux --target lifecycle_stack_test lifecycle_stack_example -j + + - name: Run lifecycle stack test and example + run: | + ./build-linux/lifecycle_stack_test --gtest_brief=1 + ./build-linux/lifecycle_stack_example + - name: Build trade_record_db_test run: cmake --build build-linux --target trade_record_db_test -j diff --git a/.github/workflows/windows-smoke.yml b/.github/workflows/windows-smoke.yml index 12134a5..850ddba 100644 --- a/.github/workflows/windows-smoke.yml +++ b/.github/workflows/windows-smoke.yml @@ -32,6 +32,18 @@ jobs: -DOPTIONX_BUILD_EXAMPLES=ON -DOPTIONX_LIGHTWEIGHT_BRIDGE_SMOKE_TESTS=ON + - name: Build lifecycle stack test and example + run: > + cmake --build build-windows --config Debug --target + lifecycle_stack_test lifecycle_stack_example + + - name: Run lifecycle stack test and example + shell: pwsh + run: | + $env:PATH = "$PWD\build-windows\bin;$PWD\build-windows\Debug;$env:PATH" + .\build-windows\Debug\lifecycle_stack_test.exe --gtest_brief=1 + .\build-windows\Debug\lifecycle_stack_example.exe + - name: Build MetaTrader path discovery test run: cmake --build build-windows --config Debug --target metatrader_paths_test diff --git a/AGENTS.md b/AGENTS.md index d3af5b7..23aa833 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -18,6 +18,10 @@ контракт Router, provider registry, replay, ownership и bot-thread dispatch. - [Market data router guide RU](guides/market-data-router.ru.md) - русский перевод, обновляемый вместе с каноническим руководством. +- [Lifecycle stack guide](guides/lifecycle-stack.md) - канонический контракт + общего process/shutdown API и staged reverse-order shutdown. +- [Lifecycle stack guide RU](guides/lifecycle-stack.ru.md) - русский перевод, + обновляемый вместе с каноническим руководством. - [Codebase orientation](guides/codebase-orientation.md) - карта проекта, DDD-слои, зависимости, расширение и безопасные точки входа. - [Build and test](guides/build-and-test.md) - CMake options, зависимости, @@ -91,3 +95,7 @@ `guides/market-data-router.md`. Любое смысловое изменение синхронизируй с `guides/market-data-router.ru.md` в том же PR; русский перевод не является источником обратных изменений английского контракта. +- Для общего lifecycle stack каноническая версия - английская + `guides/lifecycle-stack.md`. Любое смысловое изменение синхронизируй с + `guides/lifecycle-stack.ru.md` в том же PR; русский перевод не является + источником обратных изменений английского контракта. diff --git a/CMakeLists.txt b/CMakeLists.txt index c033914..019ceb2 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -455,6 +455,41 @@ if(OPTIONX_BUILD_EXAMPLES) ) endif() + add_executable(lifecycle_stack_example examples/lifecycle_stack_example.cpp) + + target_include_directories(lifecycle_stack_example PRIVATE + ${EXAMPLE_INCLUDE_DIRS} + ${EXAMPLE_DEPS_INCLUDE_DIRS} + ) + + target_link_directories(lifecycle_stack_example PRIVATE ${EXAMPLE_LIBRARY_DIRS}) + target_compile_definitions( + lifecycle_stack_example PRIVATE + ${EXAMPLE_DEFINES} + LOGIT_BASE_PATH="${LOGIT_BASE_PATH_FWD}" + ) + target_link_libraries(lifecycle_stack_example PRIVATE ${EXAMPLE_LIBS} optionx_cpp) + + if(OPTIONX_BUILD_DEPS) + add_dependencies(lifecycle_stack_example mdbx-static AES) + endif() + + foreach(dll ${EXAMPLE_DLL_FILES}) + add_custom_command(TARGET lifecycle_stack_example POST_BUILD + COMMAND ${CMAKE_COMMAND} -E copy_if_different + "${dll}" "$" + ) + endforeach() + + if(WIN32) + add_custom_command(TARGET lifecycle_stack_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(intrade_market_data_example examples/intrade_market_data_example.cpp) target_include_directories(intrade_market_data_example PRIVATE @@ -1032,6 +1067,7 @@ if(OPTIONX_BUILD_TESTS) set(OPTIONX_LIGHTWEIGHT_TESTS bridge_host_test + lifecycle_stack_test metatrader_paths_test telegram_dto_test telegram_signal_parser_test diff --git a/README.md b/README.md index 1974791..5de1339 100644 --- a/README.md +++ b/README.md @@ -5,3 +5,5 @@ OptionX is a C++ library for working with APIs of various brokers and trading pl - Market-data routing, provider selection, replay, and bot threads: [English](guides/market-data-router.md) | [Русский](guides/market-data-router.ru.md) +- Optional common module processing and staged shutdown: + [English](guides/lifecycle-stack.md) | [Русский](guides/lifecycle-stack.ru.md) diff --git a/examples/lifecycle_stack_example.cpp b/examples/lifecycle_stack_example.cpp new file mode 100644 index 0000000..cec1552 --- /dev/null +++ b/examples/lifecycle_stack_example.cpp @@ -0,0 +1,135 @@ +#include +#include + +#include +#include +#include +#include +#include +#include + +namespace { + +namespace lifecycle = optionx::lifecycle; +namespace market_data = optionx::market_data; + +class OwnerLoop final : public lifecycle::ILifecycleModule { +public: + bool post(std::function task) { + if (!task) return false; + std::lock_guard lock(m_mutex); + if (m_shutdown_requested) return false; + m_tasks.push_back(std::move(task)); + return true; + } + + void process() override { + std::vector> tasks; + { + std::lock_guard lock(m_mutex); + tasks.swap(m_tasks); + } + for (auto& task : tasks) task(); + + std::lock_guard lock(m_mutex); + if (m_shutdown_requested && m_tasks.empty()) m_stopped = true; + } + + void shutdown() noexcept override { + std::lock_guard lock(m_mutex); + m_shutdown_requested = true; + if (m_tasks.empty()) m_stopped = true; + } + + [[nodiscard]] bool is_stopped() const noexcept override { + std::lock_guard lock(m_mutex); + return m_stopped; + } + +private: + mutable std::mutex m_mutex; + std::vector> m_tasks; + bool m_shutdown_requested = false; + bool m_stopped = false; +}; + +class DeferredProvider final : public market_data::BaseMarketDataProvider { +public: + explicit DeferredProvider(OwnerLoop& owner_loop) + : m_owner_loop(owner_loop) {} + + bool subscribe_ticks( + market_data::TickSubscriptionRequest request, + subscription_callback_t callback) override { + auto subscription = market_data::MarketDataSubscriptionHandle::from_tick_request( + provider_id(), + m_next_subscription_id++, + request); + if (callback) { + callback(market_data::MarketDataSubscriptionResult::subscribed( + std::move(subscription))); + } + return true; + } + + bool unsubscribe( + market_data::MarketDataSubscriptionHandle subscription, + subscription_callback_t callback) override { + return m_owner_loop.post( + [subscription = std::move(subscription), + callback = std::move(callback)]() mutable { + if (callback) { + callback(market_data::MarketDataSubscriptionResult::unsubscribed( + std::move(subscription))); + } + }); + } + +private: + OwnerLoop& m_owner_loop; + market_data::SubscriptionId m_next_subscription_id = 1; +}; + +class QuoteSink final : public market_data::IMarketDataSubscriber {}; + +} // namespace + +int main() { + OwnerLoop owner_loop; + DeferredProvider provider(owner_loop); + market_data::MarketDataRouter router( + [&owner_loop](market_data::MarketDataRouter::owner_task_t task) { + return owner_loop.post(std::move(task)); + }); + auto subscriber = std::make_shared(); + + lifecycle::LifecycleStack application; + // Registration order declares dependencies: executor first, Router next. + if (!application.add_module(owner_loop) || + !application.add_module(router)) { + std::cerr << "Lifecycle registration failed\n"; + return 1; + } + + auto route = router.subscribe_ticks( + provider, + subscriber, + market_data::TickSubscriptionRequest("EURUSD")); + application.process(); + if (!route.active()) { + std::cerr << "Subscription did not become active\n"; + return 1; + } + + application.shutdown(); + while (!application.is_stopped()) { + application.process(); + } + if (route.valid()) { + std::cerr << "Routed subscription cleanup did not finish\n"; + return 1; + } + + std::cout << "Lifecycle stopped after routed subscription cleanup\n"; + return 0; +} diff --git a/guides/api-and-header-contracts.md b/guides/api-and-header-contracts.md index dd4ceae..553c045 100644 --- a/guides/api-and-header-contracts.md +++ b/guides/api-and-header-contracts.md @@ -108,6 +108,28 @@ etc.) являются implementation detail конкретной платфор метод сначала должен появиться на facade/base contract, а затем делегироваться в manager. +## Common Lifecycle Contract + +`optionx_cpp/lifecycle.hpp` exposes the optional +`lifecycle::ILifecycleModule` and `lifecycle::LifecycleStack` API. + +- A module implements `process()`, idempotent `shutdown() noexcept`, and the + terminal predicate `is_stopped()`. +- The stack stores non-owning references. Registered modules must outlive the + stack and its complete shutdown. +- Registration order is dependency order. Processing runs forward; shutdown is + staged in reverse, one module at a time. +- Dependencies remain processable until the current dependent module stops. +- Stack calls are owner-loop confined. The stack does not create threads. +- Initialization and `run()` remain explicit because existing module startup + contracts differ. +- `MarketDataRouter` and `BaseTradingPlatform` implement the common interface. + Direct lifecycle calls remain supported. + +See [lifecycle-stack.md](lifecycle-stack.md) for the complete contract and +[lifecycle-stack.ru.md](lifecycle-stack.ru.md) for the synchronized Russian +version. + ## Account Info Subscriber Contract `components::AccountInfoHub` is an optional fan-out adapter for the single diff --git a/guides/build-and-test.md b/guides/build-and-test.md index cf89559..3943972 100644 --- a/guides/build-and-test.md +++ b/guides/build-and-test.md @@ -207,6 +207,7 @@ targets. - `examples/event_mediator_test.cpp` - `examples/intrade_bar_api_example.cpp` +- `examples/lifecycle_stack_example.cpp` - `examples/market_data_router_example.cpp` - `examples/market_data_subscriber_base_example.cpp` - `examples/task_manager_example.cpp` diff --git a/guides/codebase-orientation.md b/guides/codebase-orientation.md index 5a17380..869bce3 100644 --- a/guides/codebase-orientation.md +++ b/guides/codebase-orientation.md @@ -24,8 +24,8 @@ ## Public Include Points Главная include-точка - `include/optionx_cpp/optionx.hpp`. Она включает -`utils.hpp`, `data.hpp`, `storages.hpp`, `components.hpp`, `platforms.hpp`, -`bridges.hpp`. +`utils.hpp`, `lifecycle.hpp`, `data.hpp`, `storages.hpp`, `components.hpp`, +`platforms.hpp`, `bridges.hpp`. Aggregate headers в корне `include/optionx_cpp` - часть публичной поверхности. Если добавляешь новый публичный DTO/component/platform, проверь соответствующий diff --git a/guides/implementation-notes.md b/guides/implementation-notes.md index e719c80..3fa89e5 100644 --- a/guides/implementation-notes.md +++ b/guides/implementation-notes.md @@ -50,6 +50,19 @@ Lifecycle: - Registered components хранятся как raw pointers; concrete platform должна владеть ими как fields и гарантировать lifetime. +### Optional Application Lifecycle Stack + +`lifecycle::LifecycleStack` composes top-level modules without replacing their +direct APIs. Register dependencies first. Normal processing runs forward; +shutdown is staged in reverse and does not stop a lower-level executor until +the current dependent module reports `is_stopped()`. + +The stack is non-owning and owner-loop confined. It intentionally does not call +`initialize()` or `run()` because existing platform, component, Router, and bot +startup contracts are not equivalent. Keep startup explicit and use the stack +only when common processing and shutdown are useful. The canonical contract is +in [lifecycle-stack.md](lifecycle-stack.md). + ### IntradeBar Delayed Retry Lifecycle Note Do not report a use-after-free risk for the current IntradeBar settings-switch diff --git a/guides/lifecycle-stack.md b/guides/lifecycle-stack.md new file mode 100644 index 0000000..5868c87 --- /dev/null +++ b/guides/lifecycle-stack.md @@ -0,0 +1,183 @@ +# Lifecycle Stack Guide + +This is the canonical guide for the optional common lifecycle API. Keep +[`lifecycle-stack.ru.md`](lifecycle-stack.ru.md) synchronized when this contract +changes. + +## Purpose + +`LifecycleStack` lets an application drive several modules through one +`process()` call and one staged `shutdown()` request. It is useful when a +platform owner loop, `MarketDataRouter`, bots, node systems, and similar modules +start and stop together. + +The stack is optional. Every module keeps its direct lifecycle API and may be +managed without `LifecycleStack`. + +Public include: + +```cpp +#include +``` + +## Module Contract + +Modules implement `optionx::lifecycle::ILifecycleModule`: + +```cpp +class ILifecycleModule { +public: + virtual void process() = 0; + virtual void shutdown() noexcept = 0; + virtual bool is_stopped() const noexcept = 0; +}; +``` + +The contract is deliberately small: + +- `process()` advances normal work or an in-progress graceful shutdown; +- `shutdown()` is an idempotent request to stop accepting new work; +- `is_stopped()` becomes true only after module-owned work and cleanup finish. + +`shutdown()` does not have to complete asynchronous cleanup before returning. +The owner loop continues calling `process()` until the terminal state is +reached. + +`MarketDataRouter` implements this interface. Its common `is_stopped()` state is +the same as `is_shutdown_complete()`. `BaseTradingPlatform` also implements the +interface and reports its existing terminal lifecycle state. + +## Registration And Ownership + +Register dependencies first and dependents last: + +```text +platform / owner executor + -> provider-facing Router + -> bots or node systems +``` + +```cpp +lifecycle::LifecycleStack application; +application.add_module(platform); +application.add_module(router); +application.add_module(bot); +``` + +The stack stores non-owning pointers. Every registered module must outlive the +stack and its complete shutdown. Duplicate registration, self-registration, and +registration after shutdown starts are rejected. + +Registration order is a dependency declaration, not merely presentation order: + +- normal `process()` runs forward; +- `shutdown()` runs in reverse; +- only one dependent shutdown stage is active at a time; +- lower-level modules keep processing until the current dependent reports + `is_stopped()`. + +Synchronous stages can collapse into one `shutdown()` call. An asynchronous +stage pauses the reverse walk until later `process()` calls complete it. + +## Startup + +`LifecycleStack` does not call `initialize()` or `run()`. Existing objects have +different startup contracts: a platform has `run(bool)`, a component has +`initialize()`, and a bot may require application-specific configuration. Keep +that work explicit: + +```cpp +platform.configure_auth(...); +platform.run(false); +bot.start(); + +lifecycle::LifecycleStack application; +application.add_module(platform); +application.add_module(router); +application.add_module(bot); +``` + +This avoids inventing a lowest-common-denominator startup API. A separate +initialization capability can be added later if multiple real modules share the +same semantics. + +## Owner Loop + +Call `LifecycleStack::process()` and `shutdown()` from the same owner loop. The +stack itself does not create a thread and does not add synchronization around +module methods. + +For a manually driven platform, the host loop becomes: + +```cpp +while (running) { + application.process(); +} + +application.shutdown(); +while (!application.is_stopped()) { + application.process(); +} +``` + +With registration order `platform -> router -> bot`, each tick first lets the +platform execute queued provider callbacks, then lets Router consume retained +lifecycle completions, then processes the bot while it is still active. + +Do not also drive a platform manually when it already owns a worker thread. +For that mode, schedule the common supervisor on the actual owner loop or keep +using the modules' direct lifecycle APIs. + +## Staged Shutdown + +Given this registration order: + +```text +platform -> router -> bot +``` + +the shutdown sequence is: + +```text +bot.shutdown() +while bot is not stopped: + platform.process() + router.process() + bot.process() + +router.shutdown() +while router is not stopped: + platform.process() + router.process() + +platform.shutdown() +``` + +Application code sees only `application.shutdown()` and +`application.process()`. The stack keeps the platform/executor alive while +Router waits for late provider completions and physical unsubscribe results. + +`LifecycleStack` is also an `ILifecycleModule`, so stacks may be nested when a +larger application has independently composed subsystems. + +Nested lifecycle stacks must form an acyclic dependency graph. `add_module()` +rejects self-registration and duplicate references, but it does not detect an +indirect cycle such as `stack_a -> stack_b -> stack_a`; the application must +avoid such registrations. + +## Failures And Limits + +The stack does not invent retry, timeout, or abandon policies. For example, a +failed Router unsubscribe keeps Router and therefore the whole stack in a +non-terminal state. The application may inspect +`failed_unsubscribe_count()` and call `retry_failed_unsubscribes()` with its own +backoff while the provider remains alive. + +`process()` exceptions propagate to the caller. Module shutdown is `noexcept` +by contract. The stack does not own modules, destroy them, or call shutdown from +its destructor. + +The runnable integration example is +[`examples/lifecycle_stack_example.cpp`](../examples/lifecycle_stack_example.cpp). +It combines an owner-loop executor, a deferred market-data provider, and +`MarketDataRouter`, then drains them through the common stack. diff --git a/guides/lifecycle-stack.ru.md b/guides/lifecycle-stack.ru.md new file mode 100644 index 0000000..675968d --- /dev/null +++ b/guides/lifecycle-stack.ru.md @@ -0,0 +1,182 @@ +# Руководство По Lifecycle Stack + +Это русский перевод канонического руководства по необязательному общему +lifecycle API. При изменении контракта синхронизируй его с +[`lifecycle-stack.md`](lifecycle-stack.md). Русская версия не является +источником обратных смысловых правок английского документа. + +## Назначение + +`LifecycleStack` позволяет приложению управлять несколькими модулями через один +вызов `process()` и один staged-запрос `shutdown()`. Это полезно, когда platform +owner loop, `MarketDataRouter`, боты, системы нод и похожие модули запускаются и +останавливаются вместе. + +Stack необязателен. Каждый модуль сохраняет прямой lifecycle API и может +управляться без `LifecycleStack`. + +Публичный include: + +```cpp +#include +``` + +## Контракт Модуля + +Модули реализуют `optionx::lifecycle::ILifecycleModule`: + +```cpp +class ILifecycleModule { +public: + virtual void process() = 0; + virtual void shutdown() noexcept = 0; + virtual bool is_stopped() const noexcept = 0; +}; +``` + +Контракт намеренно мал: + +- `process()` продвигает обычную работу или уже начатый graceful shutdown; +- `shutdown()` является идемпотентным запросом прекратить приём новой работы; +- `is_stopped()` становится true только после завершения принадлежащей модулю + работы и cleanup. + +`shutdown()` не обязан завершать асинхронный cleanup до возврата. Owner loop +продолжает вызывать `process()` до достижения terminal state. + +`MarketDataRouter` реализует этот интерфейс. Его общий `is_stopped()` совпадает +с `is_shutdown_complete()`. `BaseTradingPlatform` также реализует интерфейс и +возвращает своё существующее terminal lifecycle state. + +## Регистрация И Владение + +Сначала регистрируй зависимости, затем зависящие от них модули: + +```text +platform / owner executor + -> provider-facing Router + -> bots или node systems +``` + +```cpp +lifecycle::LifecycleStack application; +application.add_module(platform); +application.add_module(router); +application.add_module(bot); +``` + +Stack хранит non-owning указатели. Каждый зарегистрированный модуль должен жить +дольше stack и полного завершения его shutdown. Повторная регистрация, +self-registration и регистрация после начала shutdown отклоняются. + +Порядок регистрации описывает зависимости, а не только порядок отображения: + +- обычный `process()` идёт вперёд; +- `shutdown()` идёт в обратном порядке; +- одновременно активна только одна стадия shutdown зависимого модуля; +- нижележащие модули продолжают обрабатываться, пока текущий зависимый модуль не + вернёт true из `is_stopped()`. + +Синхронные стадии могут завершиться за один вызов `shutdown()`. Асинхронная +стадия приостанавливает обратный проход до следующих вызовов `process()`. + +## Запуск + +`LifecycleStack` не вызывает `initialize()` или `run()`. У существующих объектов +разные контракты запуска: у платформы есть `run(bool)`, у component есть +`initialize()`, а боту может требоваться прикладная конфигурация. Оставляй эти +действия явными: + +```cpp +platform.configure_auth(...); +platform.run(false); +bot.start(); + +lifecycle::LifecycleStack application; +application.add_module(platform); +application.add_module(router); +application.add_module(bot); +``` + +Так не появляется искусственный startup API по наименьшему общему знаменателю. +Отдельную initialization capability можно добавить позже, если несколько +реальных модулей получат одинаковую семантику. + +## Owner Loop + +Вызывай `LifecycleStack::process()` и `shutdown()` из одного owner loop. Stack не +создаёт поток и не добавляет синхронизацию вокруг методов модулей. + +Для платформы в ручном режиме host loop выглядит так: + +```cpp +while (running) { + application.process(); +} + +application.shutdown(); +while (!application.is_stopped()) { + application.process(); +} +``` + +При порядке регистрации `platform -> router -> bot` каждый tick сначала +позволяет платформе выполнить queued provider callbacks, затем Router забирает +сохранённые lifecycle completions, после чего обрабатывается ещё активный бот. + +Не управляй платформой вручную, если она уже использует собственный worker +thread. В таком режиме запускай общий supervisor в настоящем owner loop или +продолжай использовать прямые lifecycle API модулей. + +## Staged Shutdown + +Для такого порядка регистрации: + +```text +platform -> router -> bot +``` + +остановка выполняется так: + +```text +bot.shutdown() +пока bot не stopped: + platform.process() + router.process() + bot.process() + +router.shutdown() +пока router не stopped: + platform.process() + router.process() + +platform.shutdown() +``` + +Прикладной код видит только `application.shutdown()` и +`application.process()`. Stack сохраняет platform/executor живым, пока Router +ждёт поздние provider completions и результаты физического unsubscribe. + +`LifecycleStack` сам является `ILifecycleModule`, поэтому stacks можно вкладывать, +если большое приложение состоит из независимо собранных подсистем. + +Вложенные lifecycle stacks должны образовывать ацикличный граф зависимостей. +`add_module()` отклоняет саморегистрацию и повторную регистрацию одной ссылки, +но не обнаруживает косвенный цикл вроде `stack_a -> stack_b -> stack_a`; +приложение должно не допускать такие регистрации. + +## Ошибки И Ограничения + +Stack не придумывает retry, timeout или abandon policy. Например, ошибка Router +unsubscribe оставляет Router и поэтому весь stack в non-terminal состоянии. +Приложение может проверить `failed_unsubscribe_count()` и вызвать +`retry_failed_unsubscribes()` со своим backoff, пока provider ещё жив. + +Исключения из `process()` передаются caller. Shutdown модуля имеет контракт +`noexcept`. Stack не владеет модулями, не уничтожает их и не вызывает shutdown +из своего destructor. + +Рабочий интеграционный пример находится в +[`examples/lifecycle_stack_example.cpp`](../examples/lifecycle_stack_example.cpp). +Он объединяет owner-loop executor, deferred market-data provider и +`MarketDataRouter`, а затем дренирует их через общий stack. diff --git a/guides/market-data-router.md b/guides/market-data-router.md index 834fa96..596d787 100644 --- a/guides/market-data-router.md +++ b/guides/market-data-router.md @@ -397,6 +397,10 @@ Applications that compose several process/shutdown modules should put this drain loop in their lifecycle supervisor rather than special-case Router in business code. +The optional [`LifecycleStack`](lifecycle-stack.md) provides that supervisor. +Register the platform/executor before Router; it processes them forward and +keeps the executor alive while stopping Router first in reverse order. + Do not defer subscriber destruction until after the dispatcher is closed. `MarketDataSubscriberBase` normally posts remaining handles as one cleanup task; if posting is no longer possible, its destructor falls back to synchronous diff --git a/guides/market-data-router.ru.md b/guides/market-data-router.ru.md index 2266699..73e1ef5 100644 --- a/guides/market-data-router.ru.md +++ b/guides/market-data-router.ru.md @@ -396,6 +396,10 @@ provider и owner dispatcher Если приложение объединяет несколько process/shutdown modules, этот drain loop должен находиться в lifecycle supervisor, а не в Router-specific business code. +Необязательный [`LifecycleStack`](lifecycle-stack.ru.md) предоставляет такой +supervisor. Зарегистрируй platform/executor перед Router: process пойдёт вперёд, +а при обратном shutdown executor останется жив до полной остановки Router. + Не откладывай уничтожение subscriber до момента, когда dispatcher уже закрыт. Обычно `MarketDataSubscriberBase` отправляет оставшиеся handles одной cleanup задачей; если posting уже невозможен, destructor переходит к синхронному diff --git a/guides/platform-api-guide.md b/guides/platform-api-guide.md index 7c56cba..deba70f 100644 --- a/guides/platform-api-guide.md +++ b/guides/platform-api-guide.md @@ -119,6 +119,19 @@ Subscription rules: websocket feeds may still run for trade lifecycle needs even when there are no public subscribers. +## `lifecycle::LifecycleStack` + +Файлы: `include/optionx_cpp/lifecycle/ILifecycleModule.hpp` и +`include/optionx_cpp/lifecycle/LifecycleStack.hpp`. + +Полное руководство: [English](lifecycle-stack.md) | +[Русский](lifecycle-stack.ru.md). + +Это необязательный application-level supervisor для модулей с `process()`, +`shutdown()` и terminal predicate `is_stopped()`. Зависимости регистрируются +первыми, обычный process идёт вперёд, а shutdown выполняется по одной стадии в +обратном порядке. Startup остаётся явным. + ## `market_data::MarketDataRouter` Файл: `include/optionx_cpp/market_data/MarketDataRouter.hpp`. diff --git a/guides/project-overview.md b/guides/project-overview.md index 965b5a3..aaa66da 100644 --- a/guides/project-overview.md +++ b/guides/project-overview.md @@ -24,10 +24,16 @@ broker API и bridges. Основной сценарий: пользовател ## Public Include Surface +`optionx_cpp/lifecycle.hpp` exposes the optional `ILifecycleModule` and +`LifecycleStack` API for applications that compose platform, Router, bot, and +node lifecycles. See [English](lifecycle-stack.md) or +[Русский](lifecycle-stack.ru.md). + | Include | Что открывает | Когда использовать | |---|---|---| | `optionx_cpp/optionx.hpp` | Все основные subsystems | В приложениях и examples, когда нужен полный API | | `optionx_cpp/utils.hpp` | Pub-sub, tasks, crypto, strings, time, ids | Для инфраструктурного кода и новых components | +| `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 | diff --git a/include/optionx_cpp/lifecycle.hpp b/include/optionx_cpp/lifecycle.hpp new file mode 100644 index 0000000..3c1b6c0 --- /dev/null +++ b/include/optionx_cpp/lifecycle.hpp @@ -0,0 +1,17 @@ +#pragma once +#ifndef OPTIONX_HEADER_LIFECYCLE_HPP_INCLUDED +#define OPTIONX_HEADER_LIFECYCLE_HPP_INCLUDED + +/// \file lifecycle.hpp +/// \brief Includes the optional common lifecycle module API. +/// \note Headers under the lifecycle directory are intended to be included +/// through this aggregate header. + +#include +#include +#include + +#include "lifecycle/ILifecycleModule.hpp" +#include "lifecycle/LifecycleStack.hpp" + +#endif // OPTIONX_HEADER_LIFECYCLE_HPP_INCLUDED diff --git a/include/optionx_cpp/lifecycle/ILifecycleModule.hpp b/include/optionx_cpp/lifecycle/ILifecycleModule.hpp new file mode 100644 index 0000000..ad9cbd4 --- /dev/null +++ b/include/optionx_cpp/lifecycle/ILifecycleModule.hpp @@ -0,0 +1,30 @@ +#pragma once +#ifndef OPTIONX_HEADER_LIFECYCLE_ILIFECYCLE_MODULE_HPP_INCLUDED +#define OPTIONX_HEADER_LIFECYCLE_ILIFECYCLE_MODULE_HPP_INCLUDED + +/// \file ILifecycleModule.hpp +/// \brief Declares the common process and graceful-shutdown module contract. + +namespace optionx::lifecycle { + + /// \class ILifecycleModule + /// \brief Common interface for modules driven by an application owner loop. + /// \details shutdown() requests an idempotent stop. A module may need later + /// process() calls before is_stopped() becomes true. + class ILifecycleModule { + public: + virtual ~ILifecycleModule() noexcept = default; + + /// \brief Advances normal work or an in-progress graceful shutdown. + virtual void process() = 0; + + /// \brief Requests an idempotent graceful shutdown. + virtual void shutdown() noexcept = 0; + + /// \brief Returns true after all module-owned work and cleanup finished. + [[nodiscard]] virtual bool is_stopped() const noexcept = 0; + }; + +} // namespace optionx::lifecycle + +#endif // OPTIONX_HEADER_LIFECYCLE_ILIFECYCLE_MODULE_HPP_INCLUDED diff --git a/include/optionx_cpp/lifecycle/LifecycleStack.hpp b/include/optionx_cpp/lifecycle/LifecycleStack.hpp new file mode 100644 index 0000000..873093f --- /dev/null +++ b/include/optionx_cpp/lifecycle/LifecycleStack.hpp @@ -0,0 +1,114 @@ +#pragma once +#ifndef OPTIONX_HEADER_LIFECYCLE_LIFECYCLE_STACK_HPP_INCLUDED +#define OPTIONX_HEADER_LIFECYCLE_LIFECYCLE_STACK_HPP_INCLUDED + +/// \file LifecycleStack.hpp +/// \brief Defines optional ordered processing and staged shutdown for modules. + +namespace optionx::lifecycle { + + /// \class LifecycleStack + /// \brief Drives non-owning lifecycle modules in dependency order. + /// \details Register lower-level dependencies first and dependents last. + /// process() runs forward. shutdown() stops one module at a time in + /// reverse order, while lower-level dependencies keep processing. + /// All methods must be called from the same owner loop. + class LifecycleStack final : public ILifecycleModule { + public: + LifecycleStack() = default; + + LifecycleStack(const LifecycleStack&) = delete; + LifecycleStack& operator=(const LifecycleStack&) = delete; + LifecycleStack(LifecycleStack&&) = delete; + LifecycleStack& operator=(LifecycleStack&&) = delete; + + /// \brief Registers a non-owning module reference. + /// \param module Module that must outlive this stack and its shutdown. + /// \return False for duplicate/self registration or after shutdown starts. + bool add_module(ILifecycleModule& module) { + if (m_shutdown_requested || &module == this) return false; + if (std::find(m_modules.begin(), m_modules.end(), &module) != + m_modules.end()) { + return false; + } + m_modules.push_back(&module); + return true; + } + + /// \brief Returns the number of registered modules. + [[nodiscard]] std::size_t size() const noexcept { + return m_modules.size(); + } + + /// \brief Returns true when no modules are registered. + [[nodiscard]] bool empty() const noexcept { + return m_modules.empty(); + } + + /// \brief Returns true after staged shutdown has been requested. + [[nodiscard]] bool is_shutdown_requested() const noexcept { + return m_shutdown_requested; + } + + /// \brief Processes modules forward and advances staged shutdown. + void process() override { + if (m_stopped) return; + + const auto process_count = m_shutdown_requested + ? m_shutdown_cursor + : m_modules.size(); + for (std::size_t index = 0; index < process_count; ++index) { + auto* module = m_modules[index]; + if (module && !module->is_stopped()) { + module->process(); + } + } + + if (m_shutdown_requested) advance_shutdown(); + } + + /// \brief Starts staged reverse-order shutdown. + void shutdown() noexcept override { + if (m_stopped || m_shutdown_requested) return; + m_shutdown_requested = true; + m_shutdown_cursor = m_modules.size(); + advance_shutdown(); + } + + /// \brief Returns true after every registered module stopped. + [[nodiscard]] bool is_stopped() const noexcept override { + return m_stopped; + } + + private: + void advance_shutdown() noexcept { + while (m_shutdown_cursor != 0) { + auto* module = m_modules[m_shutdown_cursor - 1]; + if (!module || module->is_stopped()) { + --m_shutdown_cursor; + m_current_shutdown_requested = false; + continue; + } + + if (!m_current_shutdown_requested) { + m_current_shutdown_requested = true; + module->shutdown(); + } + if (!module->is_stopped()) return; + + --m_shutdown_cursor; + m_current_shutdown_requested = false; + } + m_stopped = true; + } + + std::vector m_modules; + std::size_t m_shutdown_cursor = 0; + bool m_shutdown_requested = false; + bool m_current_shutdown_requested = false; + bool m_stopped = false; + }; + +} // namespace optionx::lifecycle + +#endif // OPTIONX_HEADER_LIFECYCLE_LIFECYCLE_STACK_HPP_INCLUDED diff --git a/include/optionx_cpp/market_data.hpp b/include/optionx_cpp/market_data.hpp index 90ab96f..8c935bf 100644 --- a/include/optionx_cpp/market_data.hpp +++ b/include/optionx_cpp/market_data.hpp @@ -20,7 +20,11 @@ #include #include -#include "data.hpp" +#include "lifecycle.hpp" +#include "utils/fixed_point.hpp" +#include "data/market.hpp" +#include "data/bars.hpp" +#include "data/ticks.hpp" #include "market_data/enums.hpp" #include "market_data/MarketDataSubscription.hpp" #include "market_data/MarketDataBatch.hpp" diff --git a/include/optionx_cpp/market_data/MarketDataRouter.hpp b/include/optionx_cpp/market_data/MarketDataRouter.hpp index df2e201..ee61ca3 100644 --- a/include/optionx_cpp/market_data/MarketDataRouter.hpp +++ b/include/optionx_cpp/market_data/MarketDataRouter.hpp @@ -227,7 +227,7 @@ namespace optionx::market_data { /// or fails physical unsubscription, the router retains that provider /// handle, keeps the callback binding, and rejects new routes through the /// affected provider until retry_failed_unsubscribes() succeeds. - class MarketDataRouter { + class MarketDataRouter : public lifecycle::ILifecycleModule { public: using SubscriptionHandle = MarketDataRouterSubscription; ///< Move-only route owner. using subscription_callback_t = BaseMarketDataProvider::subscription_callback_t; ///< Operation callback. @@ -254,8 +254,9 @@ namespace optionx::market_data { /// \brief Move assignment is disabled because handles refer to one router state. MarketDataRouter& operator=(MarketDataRouter&&) = delete; - /// \brief Stops routed subscriptions and releases provider callbacks. - ~MarketDataRouter(); + /// \brief Requests shutdown for any routes still owned by this instance. + /// \details Drain asynchronous provider operations before destruction. + ~MarketDataRouter() override; /// \brief Adds a non-owning provider registration with stable aliases. /// \details Registration does not bind provider callbacks. Aliases are @@ -407,2038 +408,28 @@ namespace optionx::market_data { /// \details Processes provider completions retained during shutdown and /// starts or completes physical subscription cleanup. This method /// does not poll providers or transport data. - void process(); + void process() override; /// \brief Returns true after shutdown drained every provider operation. [[nodiscard]] bool is_shutdown_complete() const noexcept; + /// \brief Returns the common lifecycle terminal state. + [[nodiscard]] bool is_stopped() const noexcept override { + return is_shutdown_complete(); + } + /// \brief Starts an idempotent graceful shutdown. /// \details New routes and user delivery stop immediately. Pending provider /// operations remain owned until later process() calls finish their /// physical cleanup on the owner loop. - void shutdown() noexcept; + void shutdown() noexcept override; private: std::shared_ptr m_state; }; - namespace detail { - - class MarketDataRouterState final - : public std::enable_shared_from_this { - public: - using subscription_callback_t = BaseMarketDataProvider::subscription_callback_t; - - explicit MarketDataRouterState( - MarketDataRouter::owner_dispatcher_t owner_dispatcher = {}) - : m_owner_dispatcher(std::move(owner_dispatcher)) {} - - struct StreamDescriptor { - MarketDataType type = MarketDataType::UNKNOWN; - std::string symbol; - BarTimeframe timeframe = 0; - BarPriceSource price_source = BarPriceSource::MID; - MarketDataTransport transport = MarketDataTransport::AUTO; - }; - - enum class EntryPhase { - PENDING, - ACTIVE, - UNSUBSCRIBING, - CLEANUP_FAILED - }; - - struct Entry { - RoutedSubscriptionId router_id; - ProviderInstanceId provider_id = kInvalidProviderInstanceId; - BaseMarketDataProvider* provider = nullptr; - std::weak_ptr subscriber; - std::shared_ptr control; - StreamDescriptor stream; - MarketDataSubscriptionHandle retained_cleanup_subscription; - MarketDataSubscriptionResult unsubscribe_completion; - bool subscribe_completion_posted = false; - bool subscribe_completion_received = false; - bool unsubscribe_completion_received = false; - EntryPhase phase = EntryPhase::PENDING; - bool release_requested = false; - subscription_callback_t release_callback; - }; - - struct CachedStatus { - MarketDataStatusUpdate update; - std::uint64_t sequence = 0; - }; - - struct ProviderSlot { - BaseMarketDataProvider* provider = nullptr; - std::size_t route_count = 0; - std::unordered_map provider_routes; - std::vector statuses; - std::uint64_t next_status_sequence = 1; - }; - - struct RegisteredProvider { - BaseMarketDataProvider* provider = nullptr; - ProviderInstanceId instance_id = kInvalidProviderInstanceId; - std::vector aliases; - }; - - bool register_provider( - MarketDataProviderId id, - BaseMarketDataProvider& provider, - std::vector aliases); - bool add_provider_alias(MarketDataProviderId id, std::string alias); - bool unregister_provider(MarketDataProviderId id); - [[nodiscard]] std::size_t registered_provider_count() const; - [[nodiscard]] MarketDataProviderId registered_provider_id( - ProviderInstanceId provider_id) const; - [[nodiscard]] std::vector provider_aliases( - MarketDataProviderId id) const; - - MarketDataRouterSubscription subscribe_ticks( - BaseMarketDataProvider& provider, - std::weak_ptr subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback, - MarketDataProviderId registered_provider_id = {}); - - MarketDataRouterSubscription subscribe_ticks( - MarketDataProviderId provider_id, - std::weak_ptr subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback); - - MarketDataRouterSubscription subscribe_ticks( - std::string_view provider_alias, - std::weak_ptr subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback); - - MarketDataRouterSubscription subscribe_bars( - BaseMarketDataProvider& provider, - std::weak_ptr subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback, - MarketDataProviderId registered_provider_id = {}); - - MarketDataRouterSubscription subscribe_bars( - MarketDataProviderId provider_id, - std::weak_ptr subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback); - - MarketDataRouterSubscription subscribe_bars( - std::string_view provider_alias, - std::weak_ptr subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback); - - bool unsubscribe( - const std::shared_ptr& control, - subscription_callback_t callback); - - [[nodiscard]] std::size_t subscription_count() const; - [[nodiscard]] std::size_t failed_unsubscribe_count() const; - std::size_t retry_failed_unsubscribes(); - void process(); - [[nodiscard]] bool is_shutdown_complete() const noexcept; - void shutdown() noexcept; - - bool post_to_owner(MarketDataRouter::owner_task_t task) const; - [[nodiscard]] bool has_owner_dispatcher() const noexcept { - return static_cast(m_owner_dispatcher); - } - - void route_ticks( - ProviderInstanceId provider_id, - std::unique_ptr batch); - void route_bars( - ProviderInstanceId provider_id, - std::unique_ptr batch); - void route_status( - ProviderInstanceId provider_id, - MarketDataStatusUpdate update); - - private: - mutable std::mutex m_mutex; - std::unordered_map< - RoutedSubscriptionId, - std::shared_ptr, - RoutedSubscriptionIdHash> m_entries; - std::unordered_map m_providers; - std::unordered_map< - MarketDataProviderId, - RegisteredProvider, - MarketDataProviderIdHash> m_registered_providers; - std::unordered_map m_provider_aliases; - std::unordered_map - m_registered_provider_ids; - std::uint64_t m_next_router_id = 1; - bool m_shutdown = false; - bool m_shutdown_complete = false; - MarketDataRouter::owner_dispatcher_t m_owner_dispatcher; - - static StreamDescriptor stream_from(const TickSubscriptionRequest& request); - static StreamDescriptor stream_from(const BarSubscriptionRequest& request); - static StreamDescriptor stream_from(const MarketDataSubscriptionHandle& subscription); - - static bool same_status_stream( - const MarketDataStatusUpdate& lhs, - const MarketDataStatusUpdate& rhs) noexcept; - static bool transport_matches( - MarketDataTransport expected, - MarketDataTransport actual) noexcept; - static bool status_matches_stream( - const MarketDataStatusUpdate& update, - const StreamDescriptor& stream) noexcept; - static bool batch_matches_stream( - const TickDataBatch& batch, - const StreamDescriptor& stream) noexcept; - static bool batch_matches_stream( - const BarDataBatch& batch, - const StreamDescriptor& stream) noexcept; - - bool bind_provider(BaseMarketDataProvider& provider, std::string& error_message); - void unbind_provider(BaseMarketDataProvider& provider) const noexcept; - - std::shared_ptr add_pending_entry( - BaseMarketDataProvider& provider, - std::weak_ptr subscriber, - StreamDescriptor stream, - MarketDataProviderId registered_provider_id, - std::string& error_message); - - BaseMarketDataProvider* registered_provider_no_lock( - MarketDataProviderId id) const noexcept; - MarketDataProviderId provider_id_for_alias_no_lock( - std::string_view alias) const; - - void complete_subscribe( - RoutedSubscriptionId router_id, - BaseMarketDataProvider& provider, - MarketDataSubscriptionResult result, - subscription_callback_t callback); - - void fail_pending_subscribe( - RoutedSubscriptionId router_id, - BaseMarketDataProvider& provider, - MarketDataSubscriptionResult result, - subscription_callback_t callback); - - void complete_unsubscribe( - RoutedSubscriptionId router_id, - MarketDataSubscriptionHandle expected_subscription, - MarketDataSubscriptionResult result, - subscription_callback_t callback); - - bool start_unsubscribe( - const std::shared_ptr& entry, - MarketDataSubscriptionHandle subscription, - subscription_callback_t callback); - bool record_unsubscribe_completion( - RoutedSubscriptionId router_id, - const MarketDataSubscriptionHandle& expected_subscription, - const MarketDataSubscriptionResult& result, - bool& shutdown_requested); - void dispatch_unsubscribe_completion( - RoutedSubscriptionId router_id, - MarketDataSubscriptionHandle expected_subscription, - MarketDataSubscriptionResult result, - subscription_callback_t callback); - - BaseMarketDataProvider* remove_entry_no_lock(RoutedSubscriptionId router_id); - void cache_status_no_lock(ProviderSlot& slot, MarketDataStatusUpdate update); - bool replay_status_no_lock( - const ProviderSlot& slot, - const StreamDescriptor& stream, - MarketDataStatusUpdate& update) const; - - static void set_control_active( - const std::shared_ptr& control, - const MarketDataSubscriptionHandle& subscription); - static void set_control_released( - const std::shared_ptr& control); - static void dispatch_result( - subscription_callback_t callback, - MarketDataSubscriptionResult result); - - bool dispatch_or_run(MarketDataRouter::owner_task_t task) const; - bool record_subscribe_completion( - RoutedSubscriptionId router_id, - BaseMarketDataProvider& provider, - const MarketDataSubscriptionResult& result, - bool& shutdown_requested); - void mark_subscribe_completion_posted( - RoutedSubscriptionId router_id, - const MarketDataSubscriptionHandle& subscription); - void dispatch_subscribe_completion( - RoutedSubscriptionId router_id, - BaseMarketDataProvider& provider, - MarketDataSubscriptionResult result, - subscription_callback_t callback); - void request_process(); - - friend class ::optionx::market_data::MarketDataRouterSubscription; - }; - - inline bool MarketDataRouterState::post_to_owner( - MarketDataRouter::owner_task_t task) const { - if (!task || !m_owner_dispatcher) return false; - try { - return m_owner_dispatcher(std::move(task)); - } catch (...) { - return false; - } - } - - inline bool MarketDataRouterState::dispatch_or_run( - MarketDataRouter::owner_task_t task) const { - if (!task) return false; - if (!m_owner_dispatcher) { - task(); - return true; - } - - return post_to_owner(std::move(task)); - } - - inline bool MarketDataRouterState::record_subscribe_completion( - RoutedSubscriptionId router_id, - BaseMarketDataProvider& provider, - const MarketDataSubscriptionResult& result, - bool& shutdown_requested) { - std::lock_guard lock(m_mutex); - const auto entry_it = m_entries.find(router_id); - if (entry_it == m_entries.end() || - entry_it->second->provider != &provider || - entry_it->second->phase != EntryPhase::PENDING || - entry_it->second->subscribe_completion_received) { - return false; - } - - entry_it->second->subscribe_completion_received = true; - if (result.success()) { - entry_it->second->retained_cleanup_subscription = - result.subscription; - } - entry_it->second->subscribe_completion_posted = false; - shutdown_requested = m_shutdown; - return true; - } - - inline void MarketDataRouterState::mark_subscribe_completion_posted( - RoutedSubscriptionId router_id, - const MarketDataSubscriptionHandle& subscription) { - 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::PENDING) { - return; - } - - const auto& reservation = - entry_it->second->retained_cleanup_subscription; - if (reservation.valid() && - reservation.provider_id == subscription.provider_id && - reservation.id == subscription.id) { - entry_it->second->subscribe_completion_posted = true; - } - } - - inline void MarketDataRouterState::dispatch_subscribe_completion( - RoutedSubscriptionId router_id, - BaseMarketDataProvider& provider, - MarketDataSubscriptionResult result, - subscription_callback_t callback) { - if (result.success() && - (!result.subscription.valid() || - result.subscription.provider_id != provider.provider_id())) { - result = MarketDataSubscriptionResult::failed( - std::move(result.subscription), - MarketDataSubscriptionStatus::FAILED, - "Market-data provider returned an invalid subscription handle."); - } - - bool shutdown_requested = false; - if (!record_subscribe_completion( - router_id, - provider, - result, - shutdown_requested)) { - return; - } - - if (shutdown_requested) { - request_process(); - return; - } - - const auto subscription = result.subscription; - auto pending = std::make_shared( - std::move(result)); - const auto state = shared_from_this(); - const bool posted = dispatch_or_run( - [state, - router_id, - provider = &provider, - callback = std::move(callback), - pending]() mutable { - state->complete_subscribe( - router_id, - *provider, - std::move(*pending), - std::move(callback)); - }); - if (posted && subscription.valid()) { - mark_subscribe_completion_posted(router_id, subscription); - } - } - - inline void MarketDataRouterState::request_process() { - if (!m_owner_dispatcher) return; - const auto state = shared_from_this(); - post_to_owner([state]() { - state->process(); - }); - } - - inline MarketDataRouterState::StreamDescriptor - MarketDataRouterState::stream_from(const TickSubscriptionRequest& request) { - StreamDescriptor stream; - stream.type = MarketDataType::TICKS; - stream.symbol = request.symbol; - stream.transport = request.transport; - return stream; - } - - inline MarketDataRouterState::StreamDescriptor - MarketDataRouterState::stream_from(const BarSubscriptionRequest& request) { - StreamDescriptor stream; - stream.type = MarketDataType::BARS; - stream.symbol = request.symbol; - stream.timeframe = request.timeframe; - stream.price_source = request.price_source; - stream.transport = request.transport; - return stream; - } - - inline MarketDataRouterState::StreamDescriptor - MarketDataRouterState::stream_from( - const MarketDataSubscriptionHandle& subscription) { - StreamDescriptor stream; - stream.type = subscription.stream_type; - stream.symbol = subscription.symbol; - stream.timeframe = subscription.timeframe; - stream.price_source = subscription.price_source; - stream.transport = subscription.transport; - return stream; - } - - inline bool MarketDataRouterState::register_provider( - MarketDataProviderId id, - BaseMarketDataProvider& provider, - std::vector aliases) { - if (!id.valid()) return false; - for (std::size_t i = 0; i < aliases.size(); ++i) { - if (aliases[i].empty()) return false; - for (std::size_t j = i + 1; j < aliases.size(); ++j) { - if (aliases[i] == aliases[j]) return false; - } - } - - std::lock_guard lock(m_mutex); - if (m_shutdown || - m_registered_providers.find(id) != m_registered_providers.end() || - m_registered_provider_ids.find(provider.provider_id()) != - m_registered_provider_ids.end()) { - return false; - } - for (const auto& alias : aliases) { - if (m_provider_aliases.find(alias) != m_provider_aliases.end()) { - return false; - } - } - - RegisteredProvider registration; - registration.provider = &provider; - registration.instance_id = provider.provider_id(); - registration.aliases = std::move(aliases); - for (const auto& alias : registration.aliases) { - m_provider_aliases.emplace(alias, id); - } - m_registered_provider_ids.emplace(registration.instance_id, id); - m_registered_providers.emplace(id, std::move(registration)); - return true; - } - - inline bool MarketDataRouterState::add_provider_alias( - MarketDataProviderId id, - std::string alias) { - if (!id.valid() || alias.empty()) return false; - - std::lock_guard lock(m_mutex); - if (m_shutdown) return false; - const auto registration_it = m_registered_providers.find(id); - if (registration_it == m_registered_providers.end()) return false; - - const auto alias_it = m_provider_aliases.find(alias); - if (alias_it != m_provider_aliases.end()) { - return alias_it->second == id; - } - registration_it->second.aliases.push_back(alias); - m_provider_aliases.emplace(std::move(alias), id); - return true; - } - - inline bool MarketDataRouterState::unregister_provider(MarketDataProviderId id) { - if (!id.valid()) return false; - - std::lock_guard lock(m_mutex); - if (m_shutdown) return false; - const auto registration_it = m_registered_providers.find(id); - if (registration_it == m_registered_providers.end()) return false; - - for (const auto& [router_id, entry] : m_entries) { - (void)router_id; - if (entry->provider_id == registration_it->second.instance_id) { - return false; - } - } - for (const auto& alias : registration_it->second.aliases) { - m_provider_aliases.erase(alias); - } - m_registered_provider_ids.erase(registration_it->second.instance_id); - m_registered_providers.erase(registration_it); - return true; - } - - inline std::size_t MarketDataRouterState::registered_provider_count() const { - std::lock_guard lock(m_mutex); - return m_registered_providers.size(); - } - - inline MarketDataProviderId MarketDataRouterState::registered_provider_id( - ProviderInstanceId provider_id) const { - std::lock_guard lock(m_mutex); - const auto it = m_registered_provider_ids.find(provider_id); - return it == m_registered_provider_ids.end() - ? MarketDataProviderId{} - : it->second; - } - - inline std::vector MarketDataRouterState::provider_aliases( - MarketDataProviderId id) const { - std::lock_guard lock(m_mutex); - const auto it = m_registered_providers.find(id); - return it == m_registered_providers.end() - ? std::vector{} - : it->second.aliases; - } - - inline BaseMarketDataProvider* MarketDataRouterState::registered_provider_no_lock( - MarketDataProviderId id) const noexcept { - const auto it = m_registered_providers.find(id); - return it == m_registered_providers.end() ? nullptr : it->second.provider; - } - - inline MarketDataProviderId MarketDataRouterState::provider_id_for_alias_no_lock( - std::string_view alias) const { - const auto it = m_provider_aliases.find(std::string(alias)); - return it == m_provider_aliases.end() ? MarketDataProviderId{} : it->second; - } - - inline bool MarketDataRouterState::same_status_stream( - const MarketDataStatusUpdate& lhs, - const MarketDataStatusUpdate& rhs) noexcept { - return lhs.type == rhs.type && - lhs.symbol == rhs.symbol && - lhs.timeframe == rhs.timeframe && - lhs.transport == rhs.transport; - } - - inline bool MarketDataRouterState::transport_matches( - MarketDataTransport expected, - MarketDataTransport actual) noexcept { - return expected == actual || - expected == MarketDataTransport::AUTO || - expected == MarketDataTransport::HYBRID || - actual == MarketDataTransport::AUTO || - actual == MarketDataTransport::HYBRID; - } - - inline bool MarketDataRouterState::status_matches_stream( - const MarketDataStatusUpdate& update, - const StreamDescriptor& stream) noexcept { - return update.type == stream.type && - update.symbol == stream.symbol && - update.timeframe == stream.timeframe && - transport_matches(stream.transport, update.transport); - } - - inline bool MarketDataRouterState::batch_matches_stream( - const TickDataBatch& batch, - const StreamDescriptor& stream) noexcept { - return stream.type == MarketDataType::TICKS && - batch.type == MarketDataType::TICKS && - batch.symbol == stream.symbol; - } - - inline bool MarketDataRouterState::batch_matches_stream( - const BarDataBatch& batch, - const StreamDescriptor& stream) noexcept { - return stream.type == MarketDataType::BARS && - batch.type == MarketDataType::BARS && - batch.symbol == stream.symbol && - batch.timeframe == stream.timeframe; - } - - inline bool MarketDataRouterState::bind_provider( - BaseMarketDataProvider& provider, - std::string& error_message) { - const auto provider_id = provider.provider_id(); - std::lock_guard lock(m_mutex); - if (m_shutdown) { - error_message = "MarketDataRouter is shut down."; - return false; - } - - const auto existing = m_providers.find(provider_id); - if (existing != m_providers.end()) { - if (existing->second.provider == &provider) return true; - error_message = "Market-data provider runtime ID is already bound."; - return false; - } - - if (provider.on_tick_data() || - provider.on_bar_data() || - provider.on_market_data_status()) { - error_message = - "Market-data provider callbacks are already assigned; " - "unbind the current hub or callback owner first."; - return false; - } - - const auto weak_state = std::weak_ptr(shared_from_this()); - provider.on_tick_data() = - [weak_state, provider_id](std::unique_ptr batch) { - if (const auto state = weak_state.lock()) { - auto pending = std::make_shared>( - std::move(batch)); - state->dispatch_or_run( - [state, provider_id, pending]() mutable { - state->route_ticks(provider_id, std::move(*pending)); - }); - } - }; - provider.on_bar_data() = - [weak_state, provider_id](std::unique_ptr batch) { - if (const auto state = weak_state.lock()) { - auto pending = std::make_shared>( - std::move(batch)); - state->dispatch_or_run( - [state, provider_id, pending]() mutable { - state->route_bars(provider_id, std::move(*pending)); - }); - } - }; - provider.on_market_data_status() = - [weak_state, provider_id](MarketDataStatusUpdate update) { - if (const auto state = weak_state.lock()) { - auto pending = std::make_shared( - std::move(update)); - state->dispatch_or_run( - [state, provider_id, pending]() mutable { - state->route_status(provider_id, std::move(*pending)); - }); - } - }; - - ProviderSlot slot; - slot.provider = &provider; - m_providers.emplace(provider_id, std::move(slot)); - return true; - } - - inline void MarketDataRouterState::unbind_provider( - BaseMarketDataProvider& provider) const noexcept { - try { - provider.on_tick_data() = BaseMarketDataProvider::ticks_callback_t{}; - provider.on_bar_data() = BaseMarketDataProvider::bars_callback_t{}; - provider.on_market_data_status() = BaseMarketDataProvider::status_callback_t{}; - } catch (...) { - } - } - - inline std::shared_ptr - MarketDataRouterState::add_pending_entry( - BaseMarketDataProvider& provider, - std::weak_ptr subscriber, - StreamDescriptor stream, - MarketDataProviderId registered_provider_id, - std::string& error_message) { - if (subscriber.expired()) { - error_message = "Market-data subscriber is null or expired."; - return {}; - } - if (!bind_provider(provider, error_message)) return {}; - - auto control = std::make_shared(); - auto entry = std::make_shared(); - { - std::lock_guard lock(m_mutex); - if (m_shutdown) { - error_message = "MarketDataRouter is shut down."; - return {}; - } - - auto provider_it = m_providers.find(provider.provider_id()); - if (provider_it == m_providers.end() || - provider_it->second.provider != &provider) { - error_message = "Market-data provider binding was released."; - return {}; - } - - for (const auto& [id, existing_entry] : m_entries) { - (void)id; - if (existing_entry->provider_id == provider.provider_id() && - (existing_entry->phase == EntryPhase::UNSUBSCRIBING || - existing_entry->phase == EntryPhase::CLEANUP_FAILED || - (existing_entry->retained_cleanup_subscription.valid() && - !existing_entry->subscribe_completion_posted))) { - error_message = - "Market-data provider has pending or failed physical " - "subscription cleanup; wait for completion or retry failed " - "unsubscriptions before creating new routes."; - return {}; - } - } - - auto router_id_value = m_next_router_id++; - if (router_id_value == 0) { - router_id_value = m_next_router_id++; - } - const RoutedSubscriptionId router_id(router_id_value); - - control->router_id = router_id; - control->registered_provider_id = registered_provider_id; - control->router = shared_from_this(); - entry->router_id = router_id; - entry->provider_id = provider.provider_id(); - entry->provider = &provider; - entry->subscriber = std::move(subscriber); - entry->control = control; - entry->stream = std::move(stream); - m_entries.emplace(router_id, entry); - ++provider_it->second.route_count; - } - return control; - } - - inline MarketDataRouterSubscription MarketDataRouterState::subscribe_ticks( - BaseMarketDataProvider& provider, - std::weak_ptr subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback, - MarketDataProviderId registered_provider_id) { - if (!request.valid()) { - dispatch_result( - std::move(callback), - MarketDataSubscriptionResult::failed( - std::move(request), - MarketDataSubscriptionStatus::INVALID_REQUEST, - "Invalid tick subscription request.")); - return {}; - } - - const auto request_for_failure = request; - std::string error_message; - auto control = add_pending_entry( - provider, - std::move(subscriber), - stream_from(request), - registered_provider_id, - error_message); - if (!control) { - dispatch_result( - std::move(callback), - MarketDataSubscriptionResult::failed( - request_for_failure, - MarketDataSubscriptionStatus::FAILED, - std::move(error_message))); - return {}; - } - - const auto router_id = control->router_id; - const auto state = shared_from_this(); - bool accepted = false; - try { - accepted = provider.subscribe_ticks( - std::move(request), - [state, router_id, &provider, callback]( - MarketDataSubscriptionResult result) mutable { - state->dispatch_subscribe_completion( - router_id, - provider, - std::move(result), - std::move(callback)); - }); - } catch (const std::exception& exception) { - fail_pending_subscribe( - router_id, - provider, - MarketDataSubscriptionResult::failed( - request_for_failure, - MarketDataSubscriptionStatus::FAILED, - std::string("Market-data provider tick subscription threw: ") + - exception.what()), - std::move(callback)); - return MarketDataRouterSubscription(std::move(control)); - } catch (...) { - fail_pending_subscribe( - router_id, - provider, - MarketDataSubscriptionResult::failed( - request_for_failure, - MarketDataSubscriptionStatus::FAILED, - "Market-data provider tick subscription threw."), - std::move(callback)); - return MarketDataRouterSubscription(std::move(control)); - } - - if (!accepted) { - fail_pending_subscribe( - router_id, - provider, - MarketDataSubscriptionResult::failed( - request_for_failure, - MarketDataSubscriptionStatus::FAILED, - "Market-data provider did not accept the tick subscription operation."), - std::move(callback)); - } - return MarketDataRouterSubscription(std::move(control)); - } - - inline MarketDataRouterSubscription MarketDataRouterState::subscribe_ticks( - MarketDataProviderId provider_id, - std::weak_ptr subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback) { - BaseMarketDataProvider* provider = nullptr; - { - std::lock_guard lock(m_mutex); - if (!m_shutdown) provider = registered_provider_no_lock(provider_id); - } - if (!provider) { - dispatch_result( - std::move(callback), - MarketDataSubscriptionResult::failed( - std::move(request), - MarketDataSubscriptionStatus::FAILED, - "Market-data provider ID is not registered.")); - return {}; - } - return subscribe_ticks( - *provider, - std::move(subscriber), - std::move(request), - std::move(callback), - provider_id); - } - - inline MarketDataRouterSubscription MarketDataRouterState::subscribe_ticks( - std::string_view provider_alias, - std::weak_ptr subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback) { - MarketDataProviderId provider_id; - { - std::lock_guard lock(m_mutex); - if (!m_shutdown) { - provider_id = provider_id_for_alias_no_lock(provider_alias); - } - } - if (!provider_id.valid()) { - dispatch_result( - std::move(callback), - MarketDataSubscriptionResult::failed( - std::move(request), - MarketDataSubscriptionStatus::FAILED, - "Market-data provider alias is not registered.")); - return {}; - } - return subscribe_ticks( - provider_id, - std::move(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouterSubscription MarketDataRouterState::subscribe_bars( - BaseMarketDataProvider& provider, - std::weak_ptr subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback, - MarketDataProviderId registered_provider_id) { - if (!request.valid()) { - dispatch_result( - std::move(callback), - MarketDataSubscriptionResult::failed( - std::move(request), - MarketDataSubscriptionStatus::INVALID_REQUEST, - "Invalid bar subscription request.")); - return {}; - } - - const auto request_for_failure = request; - std::string error_message; - auto control = add_pending_entry( - provider, - std::move(subscriber), - stream_from(request), - registered_provider_id, - error_message); - if (!control) { - dispatch_result( - std::move(callback), - MarketDataSubscriptionResult::failed( - request_for_failure, - MarketDataSubscriptionStatus::FAILED, - std::move(error_message))); - return {}; - } - - const auto router_id = control->router_id; - const auto state = shared_from_this(); - bool accepted = false; - try { - accepted = provider.subscribe_bars( - std::move(request), - [state, router_id, &provider, callback]( - MarketDataSubscriptionResult result) mutable { - state->dispatch_subscribe_completion( - router_id, - provider, - std::move(result), - std::move(callback)); - }); - } catch (const std::exception& exception) { - fail_pending_subscribe( - router_id, - provider, - MarketDataSubscriptionResult::failed( - request_for_failure, - MarketDataSubscriptionStatus::FAILED, - std::string("Market-data provider bar subscription threw: ") + - exception.what()), - std::move(callback)); - return MarketDataRouterSubscription(std::move(control)); - } catch (...) { - fail_pending_subscribe( - router_id, - provider, - MarketDataSubscriptionResult::failed( - request_for_failure, - MarketDataSubscriptionStatus::FAILED, - "Market-data provider bar subscription threw."), - std::move(callback)); - return MarketDataRouterSubscription(std::move(control)); - } - - if (!accepted) { - fail_pending_subscribe( - router_id, - provider, - MarketDataSubscriptionResult::failed( - request_for_failure, - MarketDataSubscriptionStatus::FAILED, - "Market-data provider did not accept the bar subscription operation."), - std::move(callback)); - } - return MarketDataRouterSubscription(std::move(control)); - } - - inline MarketDataRouterSubscription MarketDataRouterState::subscribe_bars( - MarketDataProviderId provider_id, - std::weak_ptr subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback) { - BaseMarketDataProvider* provider = nullptr; - { - std::lock_guard lock(m_mutex); - if (!m_shutdown) provider = registered_provider_no_lock(provider_id); - } - if (!provider) { - dispatch_result( - std::move(callback), - MarketDataSubscriptionResult::failed( - std::move(request), - MarketDataSubscriptionStatus::FAILED, - "Market-data provider ID is not registered.")); - return {}; - } - return subscribe_bars( - *provider, - std::move(subscriber), - std::move(request), - std::move(callback), - provider_id); - } - - inline MarketDataRouterSubscription MarketDataRouterState::subscribe_bars( - std::string_view provider_alias, - std::weak_ptr subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback) { - MarketDataProviderId provider_id; - { - std::lock_guard lock(m_mutex); - if (!m_shutdown) { - provider_id = provider_id_for_alias_no_lock(provider_alias); - } - } - if (!provider_id.valid()) { - dispatch_result( - std::move(callback), - MarketDataSubscriptionResult::failed( - std::move(request), - MarketDataSubscriptionStatus::FAILED, - "Market-data provider alias is not registered.")); - return {}; - } - return subscribe_bars( - provider_id, - std::move(subscriber), - std::move(request), - std::move(callback)); - } - - inline void MarketDataRouterState::complete_subscribe( - RoutedSubscriptionId router_id, - BaseMarketDataProvider& provider, - MarketDataSubscriptionResult result, - subscription_callback_t callback) { - std::shared_ptr entry; - std::shared_ptr subscriber; - MarketDataStatusUpdate replay; - bool has_replay = false; - bool release_requested = false; - subscription_callback_t release_callback; - BaseMarketDataProvider* unbind = nullptr; - - if (result.success() && - (!result.subscription.valid() || - result.subscription.provider_id != provider.provider_id())) { - result = MarketDataSubscriptionResult::failed( - std::move(result.subscription), - MarketDataSubscriptionStatus::FAILED, - "Market-data provider returned an invalid subscription handle."); - } - - { - std::lock_guard lock(m_mutex); - const auto entry_it = m_entries.find(router_id); - if (entry_it == m_entries.end() || - m_shutdown || - !entry_it->second->subscribe_completion_received) { - return; - } - - if (!result.success()) { - entry = entry_it->second; - set_control_released(entry->control); - release_callback = std::move(entry->release_callback); - unbind = remove_entry_no_lock(router_id); - } else { - entry = entry_it->second; - const auto& reservation = - entry->retained_cleanup_subscription; - if (!reservation.valid() || - reservation.provider_id != result.subscription.provider_id || - reservation.id != result.subscription.id) { - return; - } - - entry->retained_cleanup_subscription = {}; - entry->subscribe_completion_posted = false; - entry->subscribe_completion_received = false; - entry->stream = stream_from(result.subscription); - entry->phase = EntryPhase::ACTIVE; - set_control_active(entry->control, result.subscription); - release_requested = entry->release_requested; - release_callback = std::move(entry->release_callback); - - auto provider_it = m_providers.find(entry->provider_id); - if (provider_it != m_providers.end() && !release_requested) { - provider_it->second.provider_routes[result.subscription.id] = router_id; - has_replay = replay_status_no_lock( - provider_it->second, - entry->stream, - replay); - subscriber = entry->subscriber.lock(); - } - } - } - - if (unbind) unbind_provider(*unbind); - dispatch_result(callback, result); - - if (release_requested && result.success()) { - unsubscribe(entry->control, std::move(release_callback)); - return; - } - if (release_callback && !result.success()) { - dispatch_result(std::move(release_callback), result); - } - if (has_replay && subscriber && result.success()) { - bool still_active = false; - { - std::lock_guard lock(m_mutex); - const auto current = m_entries.find(router_id); - still_active = current != m_entries.end() && - current->second == entry && - current->second->phase == EntryPhase::ACTIVE && - !current->second->release_requested; - } - if (!still_active) return; - replay.subscription = result.subscription; - subscriber->on_market_data_status(replay); - } - } - - inline void MarketDataRouterState::fail_pending_subscribe( - RoutedSubscriptionId router_id, - BaseMarketDataProvider& provider, - MarketDataSubscriptionResult result, - subscription_callback_t callback) { - dispatch_subscribe_completion( - router_id, - provider, - std::move(result), - std::move(callback)); - } - - inline bool MarketDataRouterState::unsubscribe( - const std::shared_ptr& control, - subscription_callback_t callback) { - if (!control) return false; - - std::shared_ptr entry; - BaseMarketDataProvider* provider = nullptr; - MarketDataSubscriptionHandle subscription; - { - std::lock_guard lock(m_mutex); - const auto it = m_entries.find(control->router_id); - if (it == m_entries.end() || it->second->control != control) { - set_control_released(control); - return false; - } - - entry = it->second; - set_control_released(control); - entry->release_requested = true; - if (entry->phase == EntryPhase::PENDING) { - entry->release_callback = std::move(callback); - return true; - } - if (entry->phase == EntryPhase::UNSUBSCRIBING) return false; - - entry->phase = EntryPhase::UNSUBSCRIBING; - provider = entry->provider; - subscription = control->provider_subscription; - if (!subscription.valid()) { - subscription = entry->retained_cleanup_subscription; - } - const auto provider_it = m_providers.find(entry->provider_id); - if (provider_it != m_providers.end()) { - provider_it->second.provider_routes.erase(subscription.id); - } - } - - if (!provider || !subscription.valid()) { - complete_unsubscribe( - entry->router_id, - subscription, - MarketDataSubscriptionResult::failed( - subscription, - MarketDataSubscriptionStatus::FAILED, - "Routed market-data subscription has no active provider handle."), - std::move(callback)); - return false; - } - - return start_unsubscribe( - entry, - std::move(subscription), - std::move(callback)); - } - - inline bool MarketDataRouterState::start_unsubscribe( - const std::shared_ptr& entry, - MarketDataSubscriptionHandle subscription, - subscription_callback_t callback) { - if (!entry || !entry->provider || !subscription.valid()) return false; - - const auto state = shared_from_this(); - bool accepted = false; - try { - accepted = entry->provider->unsubscribe( - subscription, - [state, - router_id = entry->router_id, - subscription, - callback](MarketDataSubscriptionResult result) mutable { - state->dispatch_unsubscribe_completion( - router_id, - subscription, - std::move(result), - std::move(callback)); - }); - } catch (const std::exception& exception) { - dispatch_unsubscribe_completion( - entry->router_id, - subscription, - MarketDataSubscriptionResult::failed( - subscription, - MarketDataSubscriptionStatus::FAILED, - std::string("Market-data provider unsubscribe threw: ") + - exception.what()), - std::move(callback)); - return false; - } catch (...) { - dispatch_unsubscribe_completion( - entry->router_id, - subscription, - MarketDataSubscriptionResult::failed( - subscription, - MarketDataSubscriptionStatus::FAILED, - "Market-data provider unsubscribe threw."), - std::move(callback)); - return false; - } - - if (!accepted) { - dispatch_unsubscribe_completion( - entry->router_id, - subscription, - MarketDataSubscriptionResult::failed( - subscription, - MarketDataSubscriptionStatus::FAILED, - "Market-data provider did not accept the unsubscribe operation."), - std::move(callback)); - } - return accepted; - } - - inline bool MarketDataRouterState::record_unsubscribe_completion( - RoutedSubscriptionId router_id, - const MarketDataSubscriptionHandle& expected_subscription, - const MarketDataSubscriptionResult& result, - bool& shutdown_requested) { - std::lock_guard lock(m_mutex); - const auto it = m_entries.find(router_id); - if (it == m_entries.end() || - it->second->phase != EntryPhase::UNSUBSCRIBING || - it->second->unsubscribe_completion_received) { - return false; - } - - const auto actual_subscription = result.subscription.valid() - ? result.subscription - : expected_subscription; - if (!actual_subscription.valid() || - actual_subscription.provider_id != expected_subscription.provider_id || - actual_subscription.id != expected_subscription.id) { - return false; - } - - it->second->unsubscribe_completion = result; - if (!it->second->unsubscribe_completion.subscription.valid()) { - it->second->unsubscribe_completion.subscription = expected_subscription; - } - it->second->unsubscribe_completion_received = true; - shutdown_requested = m_shutdown; - return true; - } - - inline void MarketDataRouterState::dispatch_unsubscribe_completion( - RoutedSubscriptionId router_id, - MarketDataSubscriptionHandle expected_subscription, - MarketDataSubscriptionResult result, - subscription_callback_t callback) { - bool shutdown_requested = false; - if (!record_unsubscribe_completion( - router_id, - expected_subscription, - result, - shutdown_requested)) { - return; - } - - if (shutdown_requested) { - request_process(); - return; - } - - auto pending = std::make_shared( - std::move(result)); - const auto state = shared_from_this(); - dispatch_or_run( - [state, - router_id, - expected_subscription, - callback = std::move(callback), - pending]() mutable { - state->complete_unsubscribe( - router_id, - expected_subscription, - std::move(*pending), - std::move(callback)); - }); - } - - inline void MarketDataRouterState::complete_unsubscribe( - RoutedSubscriptionId router_id, - MarketDataSubscriptionHandle expected_subscription, - MarketDataSubscriptionResult result, - subscription_callback_t callback) { - if (!result.subscription.valid()) { - result.subscription = expected_subscription; - } - - BaseMarketDataProvider* unbind = nullptr; - bool handled = false; - bool notify_callback = false; - { - std::lock_guard lock(m_mutex); - const auto it = m_entries.find(router_id); - if (it != m_entries.end() && - it->second->phase == EntryPhase::UNSUBSCRIBING && - (!expected_subscription.valid() || - it->second->unsubscribe_completion_received)) { - it->second->unsubscribe_completion_received = false; - it->second->unsubscribe_completion = {}; - if (result.success()) { - set_control_released(it->second->control); - unbind = remove_entry_no_lock(router_id); - } else { - it->second->phase = EntryPhase::CLEANUP_FAILED; - } - handled = true; - notify_callback = !m_shutdown; - } - } - - if (unbind) unbind_provider(*unbind); - if (handled && notify_callback) { - dispatch_result(std::move(callback), std::move(result)); - } - } - - inline BaseMarketDataProvider* MarketDataRouterState::remove_entry_no_lock( - RoutedSubscriptionId router_id) { - const auto entry_it = m_entries.find(router_id); - if (entry_it == m_entries.end()) return nullptr; - - const auto provider_id = entry_it->second->provider_id; - const auto subscription = entry_it->second->control->provider_subscription; - m_entries.erase(entry_it); - - const auto provider_it = m_providers.find(provider_id); - if (provider_it == m_providers.end()) return nullptr; - provider_it->second.provider_routes.erase(subscription.id); - if (provider_it->second.route_count > 0) { - --provider_it->second.route_count; - } - if (provider_it->second.route_count != 0) return nullptr; - - auto* provider = provider_it->second.provider; - m_providers.erase(provider_it); - return provider; - } - - inline void MarketDataRouterState::cache_status_no_lock( - ProviderSlot& slot, - MarketDataStatusUpdate update) { - update.subscription = {}; - for (auto& cached : slot.statuses) { - if (same_status_stream(cached.update, update)) { - cached.update = std::move(update); - cached.sequence = slot.next_status_sequence++; - return; - } - } - slot.statuses.push_back(CachedStatus{ - std::move(update), - slot.next_status_sequence++}); - } - - inline bool MarketDataRouterState::replay_status_no_lock( - const ProviderSlot& slot, - const StreamDescriptor& stream, - MarketDataStatusUpdate& update) const { - const CachedStatus* latest = nullptr; - for (const auto& cached : slot.statuses) { - if (!status_matches_stream(cached.update, stream)) continue; - if (!latest || cached.sequence > latest->sequence) { - latest = &cached; - } - } - if (!latest) return false; - update = latest->update; - return true; - } - - inline void MarketDataRouterState::route_ticks( - ProviderInstanceId provider_id, - std::unique_ptr batch) { - if (!batch) return; - std::vector, TickDataBatch>> deliveries; - { - std::lock_guard lock(m_mutex); - const auto provider_it = m_providers.find(provider_id); - if (provider_it == m_providers.end()) return; - - 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()) return; - const auto& entry = entry_it->second; - auto subscriber = entry->subscriber.lock(); - if (!subscriber || - entry->phase != EntryPhase::ACTIVE || - !batch_matches_stream(*batch, entry->stream)) { - return; - } - auto routed = *batch; - routed.subscription = entry->control->provider_subscription; - deliveries.emplace_back(std::move(subscriber), std::move(routed)); - } else { - for (const auto& [id, entry] : m_entries) { - (void)id; - if (entry->provider_id != provider_id || - entry->phase != EntryPhase::ACTIVE || - !batch_matches_stream(*batch, entry->stream)) { - continue; - } - auto subscriber = entry->subscriber.lock(); - if (!subscriber) continue; - auto routed = *batch; - routed.subscription = entry->control->provider_subscription; - deliveries.emplace_back(std::move(subscriber), std::move(routed)); - } - } - } - - for (auto& delivery : deliveries) { - delivery.first->on_tick_data(delivery.second); - } - } - - inline void MarketDataRouterState::route_bars( - ProviderInstanceId provider_id, - std::unique_ptr batch) { - if (!batch) return; - std::vector, BarDataBatch>> deliveries; - { - std::lock_guard lock(m_mutex); - const auto provider_it = m_providers.find(provider_id); - if (provider_it == m_providers.end()) return; - - 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()) return; - const auto& entry = entry_it->second; - auto subscriber = entry->subscriber.lock(); - if (!subscriber || - entry->phase != EntryPhase::ACTIVE || - !batch_matches_stream(*batch, entry->stream)) { - return; - } - auto routed = *batch; - routed.subscription = entry->control->provider_subscription; - deliveries.emplace_back(std::move(subscriber), std::move(routed)); - } else { - for (const auto& [id, entry] : m_entries) { - (void)id; - if (entry->provider_id != provider_id || - entry->phase != EntryPhase::ACTIVE || - !batch_matches_stream(*batch, entry->stream)) { - continue; - } - auto subscriber = entry->subscriber.lock(); - if (!subscriber) continue; - auto routed = *batch; - routed.subscription = entry->control->provider_subscription; - deliveries.emplace_back(std::move(subscriber), std::move(routed)); - } - } - } - - for (auto& delivery : deliveries) { - delivery.first->on_bar_data(delivery.second); - } - } - - inline void MarketDataRouterState::route_status( - ProviderInstanceId provider_id, - MarketDataStatusUpdate update) { - std::vector, MarketDataStatusUpdate>> - deliveries; - { - std::lock_guard lock(m_mutex); - const auto provider_it = m_providers.find(provider_id); - if (provider_it == m_providers.end()) return; - if (update.subscription.valid() && - update.subscription.provider_id != provider_id) { - return; - } - - update.provider_id = provider_id; - if (update.subscription.valid()) { - update.type = update.subscription.stream_type; - update.symbol = update.subscription.symbol; - update.timeframe = update.subscription.timeframe; - if (update.transport == MarketDataTransport::AUTO) { - update.transport = update.subscription.transport; - } - } - cache_status_no_lock(provider_it->second, update); - - if (update.subscription.valid()) { - const auto route_it = provider_it->second.provider_routes.find( - update.subscription.id); - if (route_it != provider_it->second.provider_routes.end()) { - const auto entry_it = m_entries.find(route_it->second); - if (entry_it != m_entries.end() && - entry_it->second->phase == EntryPhase::ACTIVE) { - auto subscriber = entry_it->second->subscriber.lock(); - if (subscriber) { - update.subscription = - entry_it->second->control->provider_subscription; - deliveries.emplace_back( - std::move(subscriber), - std::move(update)); - } - } - } - } else { - for (const auto& [id, entry] : m_entries) { - (void)id; - if (entry->provider_id != provider_id || - entry->phase != EntryPhase::ACTIVE || - !status_matches_stream(update, entry->stream)) { - continue; - } - auto subscriber = entry->subscriber.lock(); - if (!subscriber) continue; - auto routed = update; - routed.subscription = entry->control->provider_subscription; - deliveries.emplace_back(std::move(subscriber), std::move(routed)); - } - } - } - - for (auto& delivery : deliveries) { - delivery.first->on_market_data_status(delivery.second); - } - } - - inline void MarketDataRouterState::set_control_active( - const std::shared_ptr& control, - const MarketDataSubscriptionHandle& subscription) { - if (!control) return; - std::lock_guard lock(control->mutex); - control->provider_subscription = subscription; - control->active = true; - } - - inline void MarketDataRouterState::set_control_released( - const std::shared_ptr& control) { - if (!control) return; - std::lock_guard lock(control->mutex); - control->active = false; - control->released = true; - } - - inline void MarketDataRouterState::dispatch_result( - subscription_callback_t callback, - MarketDataSubscriptionResult result) { - if (callback) callback(std::move(result)); - } - - inline std::size_t MarketDataRouterState::subscription_count() const { - std::lock_guard lock(m_mutex); - return m_entries.size(); - } - - inline std::size_t MarketDataRouterState::failed_unsubscribe_count() const { - std::lock_guard lock(m_mutex); - return static_cast(std::count_if( - m_entries.begin(), - m_entries.end(), - [](const auto& item) { - return item.second->phase == EntryPhase::CLEANUP_FAILED; - })); - } - - inline std::size_t MarketDataRouterState::retry_failed_unsubscribes() { - std::vector> controls; - { - std::lock_guard lock(m_mutex); - controls.reserve(m_entries.size()); - for (const auto& [id, entry] : m_entries) { - (void)id; - if (entry->phase == EntryPhase::CLEANUP_FAILED) { - controls.push_back(entry->control); - } - } - } - - std::size_t accepted = 0; - for (const auto& control : controls) { - if (unsubscribe(control, {})) ++accepted; - } - return accepted; - } - - inline void MarketDataRouterState::process() { - struct PendingCompletion { - RoutedSubscriptionId router_id; - MarketDataSubscriptionHandle subscription; - MarketDataSubscriptionResult result; - }; - struct CleanupRequest { - std::shared_ptr entry; - MarketDataSubscriptionHandle subscription; - }; - - for (;;) { - std::vector completions; - { - std::lock_guard lock(m_mutex); - if (!m_shutdown || m_shutdown_complete) return; - - completions.reserve(m_entries.size()); - for (const auto& [id, entry] : m_entries) { - if (entry->phase != EntryPhase::UNSUBSCRIBING || - !entry->unsubscribe_completion_received) { - continue; - } - completions.push_back(PendingCompletion{ - id, - entry->unsubscribe_completion.subscription, - entry->unsubscribe_completion}); - } - } - - for (auto& completion : completions) { - complete_unsubscribe( - completion.router_id, - std::move(completion.subscription), - std::move(completion.result), - {}); - } - - std::vector failed_subscribes; - std::vector cleanup_requests; - std::vector unbind_providers; - bool shutdown_complete = false; - { - std::lock_guard lock(m_mutex); - failed_subscribes.reserve(m_entries.size()); - cleanup_requests.reserve(m_entries.size()); - - for (const auto& [id, entry] : m_entries) { - MarketDataSubscriptionHandle subscription; - if (entry->phase == EntryPhase::PENDING && - entry->subscribe_completion_received) { - if (!entry->retained_cleanup_subscription.valid()) { - failed_subscribes.push_back(id); - continue; - } - subscription = entry->retained_cleanup_subscription; - } else if (entry->phase == EntryPhase::ACTIVE) { - subscription = entry->control->provider_subscription; - } else { - continue; - } - - if (!entry->provider || !subscription.valid()) { - entry->phase = EntryPhase::CLEANUP_FAILED; - continue; - } - entry->phase = EntryPhase::UNSUBSCRIBING; - const auto provider_it = m_providers.find(entry->provider_id); - if (provider_it != m_providers.end()) { - provider_it->second.provider_routes.erase(subscription.id); - } - cleanup_requests.push_back(CleanupRequest{ - entry, - std::move(subscription)}); - } - - for (const auto id : failed_subscribes) { - const auto entry_it = m_entries.find(id); - if (entry_it == m_entries.end()) continue; - set_control_released(entry_it->second->control); - if (auto* provider = remove_entry_no_lock(id)) { - unbind_providers.push_back(provider); - } - } - - if (m_entries.empty()) { - m_registered_providers.clear(); - m_provider_aliases.clear(); - m_registered_provider_ids.clear(); - m_shutdown_complete = true; - shutdown_complete = true; - } - } - - for (auto* provider : unbind_providers) { - if (provider) unbind_provider(*provider); - } - for (auto& request : cleanup_requests) { - start_unsubscribe( - request.entry, - std::move(request.subscription), - {}); - } - - if (shutdown_complete) return; - if (completions.empty() && - failed_subscribes.empty() && - cleanup_requests.empty()) { - return; - } - } - } - - inline bool MarketDataRouterState::is_shutdown_complete() const noexcept { - std::lock_guard lock(m_mutex); - return m_shutdown_complete; - } - - inline void MarketDataRouterState::shutdown() noexcept { - { - std::lock_guard lock(m_mutex); - if (m_shutdown_complete) return; - if (m_shutdown) return; - m_shutdown = true; - for (const auto& [id, entry] : m_entries) { - (void)id; - set_control_released(entry->control); - entry->subscriber.reset(); - entry->release_callback = {}; - if (entry->phase == EntryPhase::CLEANUP_FAILED) { - entry->phase = entry->retained_cleanup_subscription.valid() - ? EntryPhase::PENDING - : EntryPhase::ACTIVE; - } - } - } - - process(); - } - - } // namespace detail - - inline MarketDataRouterSubscription& - MarketDataRouterSubscription::operator=(MarketDataRouterSubscription&& other) noexcept { - if (this == &other) return *this; - reset(); - m_control = std::move(other.m_control); - return *this; - } - - inline RoutedSubscriptionId MarketDataRouterSubscription::router_id() const noexcept { - return m_control ? m_control->router_id : RoutedSubscriptionId{}; - } - - inline MarketDataSubscriptionHandle - MarketDataRouterSubscription::provider_subscription() const { - if (!m_control) return {}; - std::lock_guard lock(m_control->mutex); - return m_control->provider_subscription; - } - - inline MarketDataProviderId - MarketDataRouterSubscription::registered_provider_id() const { - if (!m_control) return {}; - std::lock_guard lock(m_control->mutex); - return m_control->registered_provider_id; - } - - inline bool MarketDataRouterSubscription::valid() const { - if (!m_control) return false; - std::lock_guard lock(m_control->mutex); - return !m_control->released && - m_control->router_id.valid(); - } - - inline bool MarketDataRouterSubscription::active() const { - if (!m_control) return false; - std::lock_guard lock(m_control->mutex); - return !m_control->released && - m_control->active && - m_control->provider_subscription.valid(); - } - - inline bool MarketDataRouterSubscription::unsubscribe( - BaseMarketDataProvider::subscription_callback_t callback) { - if (!m_control) return false; - const auto control = std::move(m_control); - const auto router = control->router.lock(); - if (!router) { - detail::MarketDataRouterState::set_control_released(control); - return false; - } - return router->unsubscribe(control, std::move(callback)); - } - - inline void MarketDataRouterSubscription::reset() noexcept { - if (!m_control) return; - const auto control = std::move(m_control); - try { - if (const auto router = control->router.lock()) { - router->unsubscribe(control, {}); - } else { - detail::MarketDataRouterState::set_control_released(control); - } - } catch (...) { - detail::MarketDataRouterState::set_control_released(control); - } - } - - inline MarketDataRouter::MarketDataRouter() - : m_state(std::make_shared()) {} - - inline MarketDataRouter::MarketDataRouter(owner_dispatcher_t owner_dispatcher) - : m_state(std::make_shared( - std::move(owner_dispatcher))) {} - - inline MarketDataRouter::~MarketDataRouter() { - shutdown(); - } - - inline bool MarketDataRouter::register_provider( - MarketDataProviderId id, - BaseMarketDataProvider& provider, - std::vector aliases) { - return m_state && m_state->register_provider( - id, - provider, - std::move(aliases)); - } - - inline bool MarketDataRouter::add_provider_alias( - MarketDataProviderId id, - std::string alias) { - return m_state && m_state->add_provider_alias(id, std::move(alias)); - } - - inline bool MarketDataRouter::unregister_provider(MarketDataProviderId id) { - return m_state && m_state->unregister_provider(id); - } - - inline std::size_t MarketDataRouter::registered_provider_count() const { - return m_state ? m_state->registered_provider_count() : 0; - } - - inline MarketDataProviderId MarketDataRouter::registered_provider_id( - ProviderInstanceId provider_id) const { - return m_state ? m_state->registered_provider_id(provider_id) : MarketDataProviderId{}; - } - - inline std::vector MarketDataRouter::provider_aliases( - MarketDataProviderId id) const { - return m_state ? m_state->provider_aliases(id) : std::vector{}; - } - - inline bool MarketDataRouter::post_to_owner(owner_task_t task) const { - return m_state && m_state->post_to_owner(std::move(task)); - } - - inline bool MarketDataRouter::has_owner_dispatcher() const noexcept { - return m_state && m_state->has_owner_dispatcher(); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks( - BaseMarketDataProvider& provider, - const std::shared_ptr& subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback) { - return subscribe_ticks_weak( - provider, - std::weak_ptr(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks_weak( - BaseMarketDataProvider& provider, - std::weak_ptr subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback) { - return m_state->subscribe_ticks( - provider, - std::move(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks( - MarketDataProviderId provider_id, - const std::shared_ptr& subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback) { - return subscribe_ticks_weak( - provider_id, - std::weak_ptr(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks_weak( - MarketDataProviderId provider_id, - std::weak_ptr subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback) { - return m_state->subscribe_ticks( - provider_id, - std::move(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks( - std::string_view provider_alias, - const std::shared_ptr& subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback) { - return subscribe_ticks_weak( - provider_alias, - std::weak_ptr(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks_weak( - std::string_view provider_alias, - std::weak_ptr subscriber, - TickSubscriptionRequest request, - subscription_callback_t callback) { - return m_state->subscribe_ticks( - provider_alias, - std::move(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars( - BaseMarketDataProvider& provider, - const std::shared_ptr& subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback) { - return subscribe_bars_weak( - provider, - std::weak_ptr(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars_weak( - BaseMarketDataProvider& provider, - std::weak_ptr subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback) { - return m_state->subscribe_bars( - provider, - std::move(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars( - MarketDataProviderId provider_id, - const std::shared_ptr& subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback) { - return subscribe_bars_weak( - provider_id, - std::weak_ptr(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars_weak( - MarketDataProviderId provider_id, - std::weak_ptr subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback) { - return m_state->subscribe_bars( - provider_id, - std::move(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars( - std::string_view provider_alias, - const std::shared_ptr& subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback) { - return subscribe_bars_weak( - provider_alias, - std::weak_ptr(subscriber), - std::move(request), - std::move(callback)); - } - - inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars_weak( - std::string_view provider_alias, - std::weak_ptr subscriber, - BarSubscriptionRequest request, - subscription_callback_t callback) { - return m_state->subscribe_bars( - provider_alias, - std::move(subscriber), - std::move(request), - std::move(callback)); - } - - inline std::size_t MarketDataRouter::subscription_count() const { - return m_state ? m_state->subscription_count() : 0; - } - - inline std::size_t MarketDataRouter::failed_unsubscribe_count() const { - return m_state ? m_state->failed_unsubscribe_count() : 0; - } - - inline std::size_t MarketDataRouter::retry_failed_unsubscribes() { - return m_state ? m_state->retry_failed_unsubscribes() : 0; - } - - inline void MarketDataRouter::process() { - if (m_state) m_state->process(); - } - - inline bool MarketDataRouter::is_shutdown_complete() const noexcept { - return !m_state || m_state->is_shutdown_complete(); - } - - inline void MarketDataRouter::shutdown() noexcept { - if (m_state) m_state->shutdown(); - } - } // namespace optionx::market_data +#include "detail/MarketDataRouter.ipp" + #endif // OPTIONX_HEADER_MARKET_DATA_MARKET_DATA_ROUTER_HPP_INCLUDED diff --git a/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp new file mode 100644 index 0000000..b9b8ebf --- /dev/null +++ b/include/optionx_cpp/market_data/detail/MarketDataRouter.ipp @@ -0,0 +1,2030 @@ +#pragma once +#ifndef OPTIONX_HEADER_MARKET_DATA_DETAIL_MARKET_DATA_ROUTER_IPP_INCLUDED +#define OPTIONX_HEADER_MARKET_DATA_DETAIL_MARKET_DATA_ROUTER_IPP_INCLUDED + +/// \file MarketDataRouter.ipp +/// \brief Implements subscription-scoped market-data routing utilities. + +namespace optionx::market_data { + + namespace detail { + + class MarketDataRouterState final + : public std::enable_shared_from_this { + public: + using subscription_callback_t = BaseMarketDataProvider::subscription_callback_t; + + explicit MarketDataRouterState( + MarketDataRouter::owner_dispatcher_t owner_dispatcher = {}) + : m_owner_dispatcher(std::move(owner_dispatcher)) {} + + struct StreamDescriptor { + MarketDataType type = MarketDataType::UNKNOWN; + std::string symbol; + BarTimeframe timeframe = 0; + BarPriceSource price_source = BarPriceSource::MID; + MarketDataTransport transport = MarketDataTransport::AUTO; + }; + + enum class EntryPhase { + PENDING, + ACTIVE, + UNSUBSCRIBING, + CLEANUP_FAILED + }; + + struct Entry { + RoutedSubscriptionId router_id; + ProviderInstanceId provider_id = kInvalidProviderInstanceId; + BaseMarketDataProvider* provider = nullptr; + std::weak_ptr subscriber; + std::shared_ptr control; + StreamDescriptor stream; + MarketDataSubscriptionHandle retained_cleanup_subscription; + MarketDataSubscriptionResult unsubscribe_completion; + bool subscribe_completion_posted = false; + bool subscribe_completion_received = false; + bool unsubscribe_completion_received = false; + EntryPhase phase = EntryPhase::PENDING; + bool release_requested = false; + subscription_callback_t release_callback; + }; + + struct CachedStatus { + MarketDataStatusUpdate update; + std::uint64_t sequence = 0; + }; + + struct ProviderSlot { + BaseMarketDataProvider* provider = nullptr; + std::size_t route_count = 0; + std::unordered_map provider_routes; + std::vector statuses; + std::uint64_t next_status_sequence = 1; + }; + + struct RegisteredProvider { + BaseMarketDataProvider* provider = nullptr; + ProviderInstanceId instance_id = kInvalidProviderInstanceId; + std::vector aliases; + }; + + bool register_provider( + MarketDataProviderId id, + BaseMarketDataProvider& provider, + std::vector aliases); + bool add_provider_alias(MarketDataProviderId id, std::string alias); + bool unregister_provider(MarketDataProviderId id); + [[nodiscard]] std::size_t registered_provider_count() const; + [[nodiscard]] MarketDataProviderId registered_provider_id( + ProviderInstanceId provider_id) const; + [[nodiscard]] std::vector provider_aliases( + MarketDataProviderId id) const; + + MarketDataRouterSubscription subscribe_ticks( + BaseMarketDataProvider& provider, + std::weak_ptr subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback, + MarketDataProviderId registered_provider_id = {}); + + MarketDataRouterSubscription subscribe_ticks( + MarketDataProviderId provider_id, + std::weak_ptr subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback); + + MarketDataRouterSubscription subscribe_ticks( + std::string_view provider_alias, + std::weak_ptr subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback); + + MarketDataRouterSubscription subscribe_bars( + BaseMarketDataProvider& provider, + std::weak_ptr subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback, + MarketDataProviderId registered_provider_id = {}); + + MarketDataRouterSubscription subscribe_bars( + MarketDataProviderId provider_id, + std::weak_ptr subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback); + + MarketDataRouterSubscription subscribe_bars( + std::string_view provider_alias, + std::weak_ptr subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback); + + bool unsubscribe( + const std::shared_ptr& control, + subscription_callback_t callback); + + [[nodiscard]] std::size_t subscription_count() const; + [[nodiscard]] std::size_t failed_unsubscribe_count() const; + std::size_t retry_failed_unsubscribes(); + void process(); + [[nodiscard]] bool is_shutdown_complete() const noexcept; + void shutdown() noexcept; + + bool post_to_owner(MarketDataRouter::owner_task_t task) const; + [[nodiscard]] bool has_owner_dispatcher() const noexcept { + return static_cast(m_owner_dispatcher); + } + + void route_ticks( + ProviderInstanceId provider_id, + std::unique_ptr batch); + void route_bars( + ProviderInstanceId provider_id, + std::unique_ptr batch); + void route_status( + ProviderInstanceId provider_id, + MarketDataStatusUpdate update); + + private: + mutable std::mutex m_mutex; + std::unordered_map< + RoutedSubscriptionId, + std::shared_ptr, + RoutedSubscriptionIdHash> m_entries; + std::unordered_map m_providers; + std::unordered_map< + MarketDataProviderId, + RegisteredProvider, + MarketDataProviderIdHash> m_registered_providers; + std::unordered_map m_provider_aliases; + std::unordered_map + m_registered_provider_ids; + std::uint64_t m_next_router_id = 1; + bool m_shutdown = false; + bool m_shutdown_complete = false; + MarketDataRouter::owner_dispatcher_t m_owner_dispatcher; + + static StreamDescriptor stream_from(const TickSubscriptionRequest& request); + static StreamDescriptor stream_from(const BarSubscriptionRequest& request); + static StreamDescriptor stream_from(const MarketDataSubscriptionHandle& subscription); + + static bool same_status_stream( + const MarketDataStatusUpdate& lhs, + const MarketDataStatusUpdate& rhs) noexcept; + static bool transport_matches( + MarketDataTransport expected, + MarketDataTransport actual) noexcept; + static bool status_matches_stream( + const MarketDataStatusUpdate& update, + const StreamDescriptor& stream) noexcept; + static bool batch_matches_stream( + const TickDataBatch& batch, + const StreamDescriptor& stream) noexcept; + static bool batch_matches_stream( + const BarDataBatch& batch, + const StreamDescriptor& stream) noexcept; + + bool bind_provider(BaseMarketDataProvider& provider, std::string& error_message); + void unbind_provider(BaseMarketDataProvider& provider) const noexcept; + + std::shared_ptr add_pending_entry( + BaseMarketDataProvider& provider, + std::weak_ptr subscriber, + StreamDescriptor stream, + MarketDataProviderId registered_provider_id, + std::string& error_message); + + BaseMarketDataProvider* registered_provider_no_lock( + MarketDataProviderId id) const noexcept; + MarketDataProviderId provider_id_for_alias_no_lock( + std::string_view alias) const; + + void complete_subscribe( + RoutedSubscriptionId router_id, + BaseMarketDataProvider& provider, + MarketDataSubscriptionResult result, + subscription_callback_t callback); + + void fail_pending_subscribe( + RoutedSubscriptionId router_id, + BaseMarketDataProvider& provider, + MarketDataSubscriptionResult result, + subscription_callback_t callback); + + void complete_unsubscribe( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle expected_subscription, + MarketDataSubscriptionResult result, + subscription_callback_t callback); + + bool start_unsubscribe( + const std::shared_ptr& entry, + MarketDataSubscriptionHandle subscription, + subscription_callback_t callback); + bool record_unsubscribe_completion( + RoutedSubscriptionId router_id, + const MarketDataSubscriptionHandle& expected_subscription, + const MarketDataSubscriptionResult& result, + bool& shutdown_requested); + void dispatch_unsubscribe_completion( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle expected_subscription, + MarketDataSubscriptionResult result, + subscription_callback_t callback); + + BaseMarketDataProvider* remove_entry_no_lock(RoutedSubscriptionId router_id); + void cache_status_no_lock(ProviderSlot& slot, MarketDataStatusUpdate update); + bool replay_status_no_lock( + const ProviderSlot& slot, + const StreamDescriptor& stream, + MarketDataStatusUpdate& update) const; + + static void set_control_active( + const std::shared_ptr& control, + const MarketDataSubscriptionHandle& subscription); + static void set_control_released( + const std::shared_ptr& control); + static void dispatch_result( + subscription_callback_t callback, + MarketDataSubscriptionResult result); + + bool dispatch_or_run(MarketDataRouter::owner_task_t task) const; + bool record_subscribe_completion( + RoutedSubscriptionId router_id, + BaseMarketDataProvider& provider, + const MarketDataSubscriptionResult& result, + bool& shutdown_requested); + void mark_subscribe_completion_posted( + RoutedSubscriptionId router_id, + const MarketDataSubscriptionHandle& subscription); + void dispatch_subscribe_completion( + RoutedSubscriptionId router_id, + BaseMarketDataProvider& provider, + MarketDataSubscriptionResult result, + subscription_callback_t callback); + void request_process(); + + friend class ::optionx::market_data::MarketDataRouterSubscription; + }; + + inline bool MarketDataRouterState::post_to_owner( + MarketDataRouter::owner_task_t task) const { + if (!task || !m_owner_dispatcher) return false; + try { + return m_owner_dispatcher(std::move(task)); + } catch (...) { + return false; + } + } + + inline bool MarketDataRouterState::dispatch_or_run( + MarketDataRouter::owner_task_t task) const { + if (!task) return false; + if (!m_owner_dispatcher) { + task(); + return true; + } + + return post_to_owner(std::move(task)); + } + + inline bool MarketDataRouterState::record_subscribe_completion( + RoutedSubscriptionId router_id, + BaseMarketDataProvider& provider, + const MarketDataSubscriptionResult& result, + bool& shutdown_requested) { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (entry_it == m_entries.end() || + entry_it->second->provider != &provider || + entry_it->second->phase != EntryPhase::PENDING || + entry_it->second->subscribe_completion_received) { + return false; + } + + entry_it->second->subscribe_completion_received = true; + if (result.success()) { + entry_it->second->retained_cleanup_subscription = + result.subscription; + } + entry_it->second->subscribe_completion_posted = false; + shutdown_requested = m_shutdown; + return true; + } + + inline void MarketDataRouterState::mark_subscribe_completion_posted( + RoutedSubscriptionId router_id, + const MarketDataSubscriptionHandle& subscription) { + 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::PENDING) { + return; + } + + const auto& reservation = + entry_it->second->retained_cleanup_subscription; + if (reservation.valid() && + reservation.provider_id == subscription.provider_id && + reservation.id == subscription.id) { + entry_it->second->subscribe_completion_posted = true; + } + } + + inline void MarketDataRouterState::dispatch_subscribe_completion( + RoutedSubscriptionId router_id, + BaseMarketDataProvider& provider, + MarketDataSubscriptionResult result, + subscription_callback_t callback) { + if (result.success() && + (!result.subscription.valid() || + result.subscription.provider_id != provider.provider_id())) { + result = MarketDataSubscriptionResult::failed( + std::move(result.subscription), + MarketDataSubscriptionStatus::FAILED, + "Market-data provider returned an invalid subscription handle."); + } + + bool shutdown_requested = false; + if (!record_subscribe_completion( + router_id, + provider, + result, + shutdown_requested)) { + return; + } + + if (shutdown_requested) { + request_process(); + return; + } + + const auto subscription = result.subscription; + auto pending = std::make_shared( + std::move(result)); + const auto state = shared_from_this(); + const bool posted = dispatch_or_run( + [state, + router_id, + provider = &provider, + callback = std::move(callback), + pending]() mutable { + state->complete_subscribe( + router_id, + *provider, + std::move(*pending), + std::move(callback)); + }); + if (posted && subscription.valid()) { + mark_subscribe_completion_posted(router_id, subscription); + } + } + + inline void MarketDataRouterState::request_process() { + if (!m_owner_dispatcher) return; + const auto state = shared_from_this(); + post_to_owner([state]() { + state->process(); + }); + } + + inline MarketDataRouterState::StreamDescriptor + MarketDataRouterState::stream_from(const TickSubscriptionRequest& request) { + StreamDescriptor stream; + stream.type = MarketDataType::TICKS; + stream.symbol = request.symbol; + stream.transport = request.transport; + return stream; + } + + inline MarketDataRouterState::StreamDescriptor + MarketDataRouterState::stream_from(const BarSubscriptionRequest& request) { + StreamDescriptor stream; + stream.type = MarketDataType::BARS; + stream.symbol = request.symbol; + stream.timeframe = request.timeframe; + stream.price_source = request.price_source; + stream.transport = request.transport; + return stream; + } + + inline MarketDataRouterState::StreamDescriptor + MarketDataRouterState::stream_from( + const MarketDataSubscriptionHandle& subscription) { + StreamDescriptor stream; + stream.type = subscription.stream_type; + stream.symbol = subscription.symbol; + stream.timeframe = subscription.timeframe; + stream.price_source = subscription.price_source; + stream.transport = subscription.transport; + return stream; + } + + inline bool MarketDataRouterState::register_provider( + MarketDataProviderId id, + BaseMarketDataProvider& provider, + std::vector aliases) { + if (!id.valid()) return false; + for (std::size_t i = 0; i < aliases.size(); ++i) { + if (aliases[i].empty()) return false; + for (std::size_t j = i + 1; j < aliases.size(); ++j) { + if (aliases[i] == aliases[j]) return false; + } + } + + std::lock_guard lock(m_mutex); + if (m_shutdown || + m_registered_providers.find(id) != m_registered_providers.end() || + m_registered_provider_ids.find(provider.provider_id()) != + m_registered_provider_ids.end()) { + return false; + } + for (const auto& alias : aliases) { + if (m_provider_aliases.find(alias) != m_provider_aliases.end()) { + return false; + } + } + + RegisteredProvider registration; + registration.provider = &provider; + registration.instance_id = provider.provider_id(); + registration.aliases = std::move(aliases); + for (const auto& alias : registration.aliases) { + m_provider_aliases.emplace(alias, id); + } + m_registered_provider_ids.emplace(registration.instance_id, id); + m_registered_providers.emplace(id, std::move(registration)); + return true; + } + + inline bool MarketDataRouterState::add_provider_alias( + MarketDataProviderId id, + std::string alias) { + if (!id.valid() || alias.empty()) return false; + + std::lock_guard lock(m_mutex); + if (m_shutdown) return false; + const auto registration_it = m_registered_providers.find(id); + if (registration_it == m_registered_providers.end()) return false; + + const auto alias_it = m_provider_aliases.find(alias); + if (alias_it != m_provider_aliases.end()) { + return alias_it->second == id; + } + registration_it->second.aliases.push_back(alias); + m_provider_aliases.emplace(std::move(alias), id); + return true; + } + + inline bool MarketDataRouterState::unregister_provider(MarketDataProviderId id) { + if (!id.valid()) return false; + + std::lock_guard lock(m_mutex); + if (m_shutdown) return false; + const auto registration_it = m_registered_providers.find(id); + if (registration_it == m_registered_providers.end()) return false; + + for (const auto& [router_id, entry] : m_entries) { + (void)router_id; + if (entry->provider_id == registration_it->second.instance_id) { + return false; + } + } + for (const auto& alias : registration_it->second.aliases) { + m_provider_aliases.erase(alias); + } + m_registered_provider_ids.erase(registration_it->second.instance_id); + m_registered_providers.erase(registration_it); + return true; + } + + inline std::size_t MarketDataRouterState::registered_provider_count() const { + std::lock_guard lock(m_mutex); + return m_registered_providers.size(); + } + + inline MarketDataProviderId MarketDataRouterState::registered_provider_id( + ProviderInstanceId provider_id) const { + std::lock_guard lock(m_mutex); + const auto it = m_registered_provider_ids.find(provider_id); + return it == m_registered_provider_ids.end() + ? MarketDataProviderId{} + : it->second; + } + + inline std::vector MarketDataRouterState::provider_aliases( + MarketDataProviderId id) const { + std::lock_guard lock(m_mutex); + const auto it = m_registered_providers.find(id); + return it == m_registered_providers.end() + ? std::vector{} + : it->second.aliases; + } + + inline BaseMarketDataProvider* MarketDataRouterState::registered_provider_no_lock( + MarketDataProviderId id) const noexcept { + const auto it = m_registered_providers.find(id); + return it == m_registered_providers.end() ? nullptr : it->second.provider; + } + + inline MarketDataProviderId MarketDataRouterState::provider_id_for_alias_no_lock( + std::string_view alias) const { + const auto it = m_provider_aliases.find(std::string(alias)); + return it == m_provider_aliases.end() ? MarketDataProviderId{} : it->second; + } + + inline bool MarketDataRouterState::same_status_stream( + const MarketDataStatusUpdate& lhs, + const MarketDataStatusUpdate& rhs) noexcept { + return lhs.type == rhs.type && + lhs.symbol == rhs.symbol && + lhs.timeframe == rhs.timeframe && + lhs.transport == rhs.transport; + } + + inline bool MarketDataRouterState::transport_matches( + MarketDataTransport expected, + MarketDataTransport actual) noexcept { + return expected == actual || + expected == MarketDataTransport::AUTO || + expected == MarketDataTransport::HYBRID || + actual == MarketDataTransport::AUTO || + actual == MarketDataTransport::HYBRID; + } + + inline bool MarketDataRouterState::status_matches_stream( + const MarketDataStatusUpdate& update, + const StreamDescriptor& stream) noexcept { + return update.type == stream.type && + update.symbol == stream.symbol && + update.timeframe == stream.timeframe && + transport_matches(stream.transport, update.transport); + } + + inline bool MarketDataRouterState::batch_matches_stream( + const TickDataBatch& batch, + const StreamDescriptor& stream) noexcept { + return stream.type == MarketDataType::TICKS && + batch.type == MarketDataType::TICKS && + batch.symbol == stream.symbol; + } + + inline bool MarketDataRouterState::batch_matches_stream( + const BarDataBatch& batch, + const StreamDescriptor& stream) noexcept { + return stream.type == MarketDataType::BARS && + batch.type == MarketDataType::BARS && + batch.symbol == stream.symbol && + batch.timeframe == stream.timeframe; + } + + inline bool MarketDataRouterState::bind_provider( + BaseMarketDataProvider& provider, + std::string& error_message) { + const auto provider_id = provider.provider_id(); + std::lock_guard lock(m_mutex); + if (m_shutdown) { + error_message = "MarketDataRouter is shut down."; + return false; + } + + const auto existing = m_providers.find(provider_id); + if (existing != m_providers.end()) { + if (existing->second.provider == &provider) return true; + error_message = "Market-data provider runtime ID is already bound."; + return false; + } + + if (provider.on_tick_data() || + provider.on_bar_data() || + provider.on_market_data_status()) { + error_message = + "Market-data provider callbacks are already assigned; " + "unbind the current hub or callback owner first."; + return false; + } + + const auto weak_state = std::weak_ptr(shared_from_this()); + provider.on_tick_data() = + [weak_state, provider_id](std::unique_ptr batch) { + if (const auto state = weak_state.lock()) { + auto pending = std::make_shared>( + std::move(batch)); + state->dispatch_or_run( + [state, provider_id, pending]() mutable { + state->route_ticks(provider_id, std::move(*pending)); + }); + } + }; + provider.on_bar_data() = + [weak_state, provider_id](std::unique_ptr batch) { + if (const auto state = weak_state.lock()) { + auto pending = std::make_shared>( + std::move(batch)); + state->dispatch_or_run( + [state, provider_id, pending]() mutable { + state->route_bars(provider_id, std::move(*pending)); + }); + } + }; + provider.on_market_data_status() = + [weak_state, provider_id](MarketDataStatusUpdate update) { + if (const auto state = weak_state.lock()) { + auto pending = std::make_shared( + std::move(update)); + state->dispatch_or_run( + [state, provider_id, pending]() mutable { + state->route_status(provider_id, std::move(*pending)); + }); + } + }; + + ProviderSlot slot; + slot.provider = &provider; + m_providers.emplace(provider_id, std::move(slot)); + return true; + } + + inline void MarketDataRouterState::unbind_provider( + BaseMarketDataProvider& provider) const noexcept { + try { + provider.on_tick_data() = BaseMarketDataProvider::ticks_callback_t{}; + provider.on_bar_data() = BaseMarketDataProvider::bars_callback_t{}; + provider.on_market_data_status() = BaseMarketDataProvider::status_callback_t{}; + } catch (...) { + } + } + + inline std::shared_ptr + MarketDataRouterState::add_pending_entry( + BaseMarketDataProvider& provider, + std::weak_ptr subscriber, + StreamDescriptor stream, + MarketDataProviderId registered_provider_id, + std::string& error_message) { + if (subscriber.expired()) { + error_message = "Market-data subscriber is null or expired."; + return {}; + } + if (!bind_provider(provider, error_message)) return {}; + + auto control = std::make_shared(); + auto entry = std::make_shared(); + { + std::lock_guard lock(m_mutex); + if (m_shutdown) { + error_message = "MarketDataRouter is shut down."; + return {}; + } + + auto provider_it = m_providers.find(provider.provider_id()); + if (provider_it == m_providers.end() || + provider_it->second.provider != &provider) { + error_message = "Market-data provider binding was released."; + return {}; + } + + for (const auto& [id, existing_entry] : m_entries) { + (void)id; + if (existing_entry->provider_id == provider.provider_id() && + (existing_entry->phase == EntryPhase::UNSUBSCRIBING || + existing_entry->phase == EntryPhase::CLEANUP_FAILED || + (existing_entry->retained_cleanup_subscription.valid() && + !existing_entry->subscribe_completion_posted))) { + error_message = + "Market-data provider has pending or failed physical " + "subscription cleanup; wait for completion or retry failed " + "unsubscriptions before creating new routes."; + return {}; + } + } + + auto router_id_value = m_next_router_id++; + if (router_id_value == 0) { + router_id_value = m_next_router_id++; + } + const RoutedSubscriptionId router_id(router_id_value); + + control->router_id = router_id; + control->registered_provider_id = registered_provider_id; + control->router = shared_from_this(); + entry->router_id = router_id; + entry->provider_id = provider.provider_id(); + entry->provider = &provider; + entry->subscriber = std::move(subscriber); + entry->control = control; + entry->stream = std::move(stream); + m_entries.emplace(router_id, entry); + ++provider_it->second.route_count; + } + return control; + } + + inline MarketDataRouterSubscription MarketDataRouterState::subscribe_ticks( + BaseMarketDataProvider& provider, + std::weak_ptr subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback, + MarketDataProviderId registered_provider_id) { + if (!request.valid()) { + dispatch_result( + std::move(callback), + MarketDataSubscriptionResult::failed( + std::move(request), + MarketDataSubscriptionStatus::INVALID_REQUEST, + "Invalid tick subscription request.")); + return {}; + } + + const auto request_for_failure = request; + std::string error_message; + auto control = add_pending_entry( + provider, + std::move(subscriber), + stream_from(request), + registered_provider_id, + error_message); + if (!control) { + dispatch_result( + std::move(callback), + MarketDataSubscriptionResult::failed( + request_for_failure, + MarketDataSubscriptionStatus::FAILED, + std::move(error_message))); + return {}; + } + + const auto router_id = control->router_id; + const auto state = shared_from_this(); + bool accepted = false; + try { + accepted = provider.subscribe_ticks( + std::move(request), + [state, router_id, &provider, callback]( + MarketDataSubscriptionResult result) mutable { + state->dispatch_subscribe_completion( + router_id, + provider, + std::move(result), + std::move(callback)); + }); + } catch (const std::exception& exception) { + fail_pending_subscribe( + router_id, + provider, + MarketDataSubscriptionResult::failed( + request_for_failure, + MarketDataSubscriptionStatus::FAILED, + std::string("Market-data provider tick subscription threw: ") + + exception.what()), + std::move(callback)); + return MarketDataRouterSubscription(std::move(control)); + } catch (...) { + fail_pending_subscribe( + router_id, + provider, + MarketDataSubscriptionResult::failed( + request_for_failure, + MarketDataSubscriptionStatus::FAILED, + "Market-data provider tick subscription threw."), + std::move(callback)); + return MarketDataRouterSubscription(std::move(control)); + } + + if (!accepted) { + fail_pending_subscribe( + router_id, + provider, + MarketDataSubscriptionResult::failed( + request_for_failure, + MarketDataSubscriptionStatus::FAILED, + "Market-data provider did not accept the tick subscription operation."), + std::move(callback)); + } + return MarketDataRouterSubscription(std::move(control)); + } + + inline MarketDataRouterSubscription MarketDataRouterState::subscribe_ticks( + MarketDataProviderId provider_id, + std::weak_ptr subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback) { + BaseMarketDataProvider* provider = nullptr; + { + std::lock_guard lock(m_mutex); + if (!m_shutdown) provider = registered_provider_no_lock(provider_id); + } + if (!provider) { + dispatch_result( + std::move(callback), + MarketDataSubscriptionResult::failed( + std::move(request), + MarketDataSubscriptionStatus::FAILED, + "Market-data provider ID is not registered.")); + return {}; + } + return subscribe_ticks( + *provider, + std::move(subscriber), + std::move(request), + std::move(callback), + provider_id); + } + + inline MarketDataRouterSubscription MarketDataRouterState::subscribe_ticks( + std::string_view provider_alias, + std::weak_ptr subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback) { + MarketDataProviderId provider_id; + { + std::lock_guard lock(m_mutex); + if (!m_shutdown) { + provider_id = provider_id_for_alias_no_lock(provider_alias); + } + } + if (!provider_id.valid()) { + dispatch_result( + std::move(callback), + MarketDataSubscriptionResult::failed( + std::move(request), + MarketDataSubscriptionStatus::FAILED, + "Market-data provider alias is not registered.")); + return {}; + } + return subscribe_ticks( + provider_id, + std::move(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouterSubscription MarketDataRouterState::subscribe_bars( + BaseMarketDataProvider& provider, + std::weak_ptr subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback, + MarketDataProviderId registered_provider_id) { + if (!request.valid()) { + dispatch_result( + std::move(callback), + MarketDataSubscriptionResult::failed( + std::move(request), + MarketDataSubscriptionStatus::INVALID_REQUEST, + "Invalid bar subscription request.")); + return {}; + } + + const auto request_for_failure = request; + std::string error_message; + auto control = add_pending_entry( + provider, + std::move(subscriber), + stream_from(request), + registered_provider_id, + error_message); + if (!control) { + dispatch_result( + std::move(callback), + MarketDataSubscriptionResult::failed( + request_for_failure, + MarketDataSubscriptionStatus::FAILED, + std::move(error_message))); + return {}; + } + + const auto router_id = control->router_id; + const auto state = shared_from_this(); + bool accepted = false; + try { + accepted = provider.subscribe_bars( + std::move(request), + [state, router_id, &provider, callback]( + MarketDataSubscriptionResult result) mutable { + state->dispatch_subscribe_completion( + router_id, + provider, + std::move(result), + std::move(callback)); + }); + } catch (const std::exception& exception) { + fail_pending_subscribe( + router_id, + provider, + MarketDataSubscriptionResult::failed( + request_for_failure, + MarketDataSubscriptionStatus::FAILED, + std::string("Market-data provider bar subscription threw: ") + + exception.what()), + std::move(callback)); + return MarketDataRouterSubscription(std::move(control)); + } catch (...) { + fail_pending_subscribe( + router_id, + provider, + MarketDataSubscriptionResult::failed( + request_for_failure, + MarketDataSubscriptionStatus::FAILED, + "Market-data provider bar subscription threw."), + std::move(callback)); + return MarketDataRouterSubscription(std::move(control)); + } + + if (!accepted) { + fail_pending_subscribe( + router_id, + provider, + MarketDataSubscriptionResult::failed( + request_for_failure, + MarketDataSubscriptionStatus::FAILED, + "Market-data provider did not accept the bar subscription operation."), + std::move(callback)); + } + return MarketDataRouterSubscription(std::move(control)); + } + + inline MarketDataRouterSubscription MarketDataRouterState::subscribe_bars( + MarketDataProviderId provider_id, + std::weak_ptr subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback) { + BaseMarketDataProvider* provider = nullptr; + { + std::lock_guard lock(m_mutex); + if (!m_shutdown) provider = registered_provider_no_lock(provider_id); + } + if (!provider) { + dispatch_result( + std::move(callback), + MarketDataSubscriptionResult::failed( + std::move(request), + MarketDataSubscriptionStatus::FAILED, + "Market-data provider ID is not registered.")); + return {}; + } + return subscribe_bars( + *provider, + std::move(subscriber), + std::move(request), + std::move(callback), + provider_id); + } + + inline MarketDataRouterSubscription MarketDataRouterState::subscribe_bars( + std::string_view provider_alias, + std::weak_ptr subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback) { + MarketDataProviderId provider_id; + { + std::lock_guard lock(m_mutex); + if (!m_shutdown) { + provider_id = provider_id_for_alias_no_lock(provider_alias); + } + } + if (!provider_id.valid()) { + dispatch_result( + std::move(callback), + MarketDataSubscriptionResult::failed( + std::move(request), + MarketDataSubscriptionStatus::FAILED, + "Market-data provider alias is not registered.")); + return {}; + } + return subscribe_bars( + provider_id, + std::move(subscriber), + std::move(request), + std::move(callback)); + } + + inline void MarketDataRouterState::complete_subscribe( + RoutedSubscriptionId router_id, + BaseMarketDataProvider& provider, + MarketDataSubscriptionResult result, + subscription_callback_t callback) { + std::shared_ptr entry; + std::shared_ptr subscriber; + MarketDataStatusUpdate replay; + bool has_replay = false; + bool release_requested = false; + subscription_callback_t release_callback; + BaseMarketDataProvider* unbind = nullptr; + + if (result.success() && + (!result.subscription.valid() || + result.subscription.provider_id != provider.provider_id())) { + result = MarketDataSubscriptionResult::failed( + std::move(result.subscription), + MarketDataSubscriptionStatus::FAILED, + "Market-data provider returned an invalid subscription handle."); + } + + { + std::lock_guard lock(m_mutex); + const auto entry_it = m_entries.find(router_id); + if (entry_it == m_entries.end() || + m_shutdown || + !entry_it->second->subscribe_completion_received) { + return; + } + + if (!result.success()) { + entry = entry_it->second; + set_control_released(entry->control); + release_callback = std::move(entry->release_callback); + unbind = remove_entry_no_lock(router_id); + } else { + entry = entry_it->second; + const auto& reservation = + entry->retained_cleanup_subscription; + if (!reservation.valid() || + reservation.provider_id != result.subscription.provider_id || + reservation.id != result.subscription.id) { + return; + } + + entry->retained_cleanup_subscription = {}; + entry->subscribe_completion_posted = false; + entry->subscribe_completion_received = false; + entry->stream = stream_from(result.subscription); + entry->phase = EntryPhase::ACTIVE; + set_control_active(entry->control, result.subscription); + release_requested = entry->release_requested; + release_callback = std::move(entry->release_callback); + + auto provider_it = m_providers.find(entry->provider_id); + if (provider_it != m_providers.end() && !release_requested) { + provider_it->second.provider_routes[result.subscription.id] = router_id; + has_replay = replay_status_no_lock( + provider_it->second, + entry->stream, + replay); + subscriber = entry->subscriber.lock(); + } + } + } + + if (unbind) unbind_provider(*unbind); + dispatch_result(callback, result); + + if (release_requested && result.success()) { + unsubscribe(entry->control, std::move(release_callback)); + return; + } + if (release_callback && !result.success()) { + dispatch_result(std::move(release_callback), result); + } + if (has_replay && subscriber && result.success()) { + bool still_active = false; + { + std::lock_guard lock(m_mutex); + const auto current = m_entries.find(router_id); + still_active = current != m_entries.end() && + current->second == entry && + current->second->phase == EntryPhase::ACTIVE && + !current->second->release_requested; + } + if (!still_active) return; + replay.subscription = result.subscription; + subscriber->on_market_data_status(replay); + } + } + + inline void MarketDataRouterState::fail_pending_subscribe( + RoutedSubscriptionId router_id, + BaseMarketDataProvider& provider, + MarketDataSubscriptionResult result, + subscription_callback_t callback) { + dispatch_subscribe_completion( + router_id, + provider, + std::move(result), + std::move(callback)); + } + + inline bool MarketDataRouterState::unsubscribe( + const std::shared_ptr& control, + subscription_callback_t callback) { + if (!control) return false; + + std::shared_ptr entry; + BaseMarketDataProvider* provider = nullptr; + MarketDataSubscriptionHandle subscription; + { + std::lock_guard lock(m_mutex); + const auto it = m_entries.find(control->router_id); + if (it == m_entries.end() || it->second->control != control) { + set_control_released(control); + return false; + } + + entry = it->second; + set_control_released(control); + entry->release_requested = true; + if (entry->phase == EntryPhase::PENDING) { + entry->release_callback = std::move(callback); + return true; + } + if (entry->phase == EntryPhase::UNSUBSCRIBING) return false; + + entry->phase = EntryPhase::UNSUBSCRIBING; + provider = entry->provider; + subscription = control->provider_subscription; + if (!subscription.valid()) { + subscription = entry->retained_cleanup_subscription; + } + const auto provider_it = m_providers.find(entry->provider_id); + if (provider_it != m_providers.end()) { + provider_it->second.provider_routes.erase(subscription.id); + } + } + + if (!provider || !subscription.valid()) { + complete_unsubscribe( + entry->router_id, + subscription, + MarketDataSubscriptionResult::failed( + subscription, + MarketDataSubscriptionStatus::FAILED, + "Routed market-data subscription has no active provider handle."), + std::move(callback)); + return false; + } + + return start_unsubscribe( + entry, + std::move(subscription), + std::move(callback)); + } + + inline bool MarketDataRouterState::start_unsubscribe( + const std::shared_ptr& entry, + MarketDataSubscriptionHandle subscription, + subscription_callback_t callback) { + if (!entry || !entry->provider || !subscription.valid()) return false; + + const auto state = shared_from_this(); + bool accepted = false; + try { + accepted = entry->provider->unsubscribe( + subscription, + [state, + router_id = entry->router_id, + subscription, + callback](MarketDataSubscriptionResult result) mutable { + state->dispatch_unsubscribe_completion( + router_id, + subscription, + std::move(result), + std::move(callback)); + }); + } catch (const std::exception& exception) { + dispatch_unsubscribe_completion( + entry->router_id, + subscription, + MarketDataSubscriptionResult::failed( + subscription, + MarketDataSubscriptionStatus::FAILED, + std::string("Market-data provider unsubscribe threw: ") + + exception.what()), + std::move(callback)); + return false; + } catch (...) { + dispatch_unsubscribe_completion( + entry->router_id, + subscription, + MarketDataSubscriptionResult::failed( + subscription, + MarketDataSubscriptionStatus::FAILED, + "Market-data provider unsubscribe threw."), + std::move(callback)); + return false; + } + + if (!accepted) { + dispatch_unsubscribe_completion( + entry->router_id, + subscription, + MarketDataSubscriptionResult::failed( + subscription, + MarketDataSubscriptionStatus::FAILED, + "Market-data provider did not accept the unsubscribe operation."), + std::move(callback)); + } + return accepted; + } + + inline bool MarketDataRouterState::record_unsubscribe_completion( + RoutedSubscriptionId router_id, + const MarketDataSubscriptionHandle& expected_subscription, + const MarketDataSubscriptionResult& result, + bool& shutdown_requested) { + std::lock_guard lock(m_mutex); + const auto it = m_entries.find(router_id); + if (it == m_entries.end() || + it->second->phase != EntryPhase::UNSUBSCRIBING || + it->second->unsubscribe_completion_received) { + return false; + } + + const auto actual_subscription = result.subscription.valid() + ? result.subscription + : expected_subscription; + if (!actual_subscription.valid() || + actual_subscription.provider_id != expected_subscription.provider_id || + actual_subscription.id != expected_subscription.id) { + return false; + } + + it->second->unsubscribe_completion = result; + if (!it->second->unsubscribe_completion.subscription.valid()) { + it->second->unsubscribe_completion.subscription = expected_subscription; + } + it->second->unsubscribe_completion_received = true; + shutdown_requested = m_shutdown; + return true; + } + + inline void MarketDataRouterState::dispatch_unsubscribe_completion( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle expected_subscription, + MarketDataSubscriptionResult result, + subscription_callback_t callback) { + bool shutdown_requested = false; + if (!record_unsubscribe_completion( + router_id, + expected_subscription, + result, + shutdown_requested)) { + return; + } + + if (shutdown_requested) { + request_process(); + return; + } + + auto pending = std::make_shared( + std::move(result)); + const auto state = shared_from_this(); + dispatch_or_run( + [state, + router_id, + expected_subscription, + callback = std::move(callback), + pending]() mutable { + state->complete_unsubscribe( + router_id, + expected_subscription, + std::move(*pending), + std::move(callback)); + }); + } + + inline void MarketDataRouterState::complete_unsubscribe( + RoutedSubscriptionId router_id, + MarketDataSubscriptionHandle expected_subscription, + MarketDataSubscriptionResult result, + subscription_callback_t callback) { + if (!result.subscription.valid()) { + result.subscription = expected_subscription; + } + + BaseMarketDataProvider* unbind = nullptr; + bool handled = false; + bool notify_callback = false; + { + std::lock_guard lock(m_mutex); + const auto it = m_entries.find(router_id); + if (it != m_entries.end() && + it->second->phase == EntryPhase::UNSUBSCRIBING && + (!expected_subscription.valid() || + it->second->unsubscribe_completion_received)) { + it->second->unsubscribe_completion_received = false; + it->second->unsubscribe_completion = {}; + if (result.success()) { + set_control_released(it->second->control); + unbind = remove_entry_no_lock(router_id); + } else { + it->second->phase = EntryPhase::CLEANUP_FAILED; + } + handled = true; + notify_callback = !m_shutdown; + } + } + + if (unbind) unbind_provider(*unbind); + if (handled && notify_callback) { + dispatch_result(std::move(callback), std::move(result)); + } + } + + inline BaseMarketDataProvider* MarketDataRouterState::remove_entry_no_lock( + RoutedSubscriptionId router_id) { + const auto entry_it = m_entries.find(router_id); + if (entry_it == m_entries.end()) return nullptr; + + const auto provider_id = entry_it->second->provider_id; + const auto subscription = entry_it->second->control->provider_subscription; + m_entries.erase(entry_it); + + const auto provider_it = m_providers.find(provider_id); + if (provider_it == m_providers.end()) return nullptr; + provider_it->second.provider_routes.erase(subscription.id); + if (provider_it->second.route_count > 0) { + --provider_it->second.route_count; + } + if (provider_it->second.route_count != 0) return nullptr; + + auto* provider = provider_it->second.provider; + m_providers.erase(provider_it); + return provider; + } + + inline void MarketDataRouterState::cache_status_no_lock( + ProviderSlot& slot, + MarketDataStatusUpdate update) { + update.subscription = {}; + for (auto& cached : slot.statuses) { + if (same_status_stream(cached.update, update)) { + cached.update = std::move(update); + cached.sequence = slot.next_status_sequence++; + return; + } + } + slot.statuses.push_back(CachedStatus{ + std::move(update), + slot.next_status_sequence++}); + } + + inline bool MarketDataRouterState::replay_status_no_lock( + const ProviderSlot& slot, + const StreamDescriptor& stream, + MarketDataStatusUpdate& update) const { + const CachedStatus* latest = nullptr; + for (const auto& cached : slot.statuses) { + if (!status_matches_stream(cached.update, stream)) continue; + if (!latest || cached.sequence > latest->sequence) { + latest = &cached; + } + } + if (!latest) return false; + update = latest->update; + return true; + } + + inline void MarketDataRouterState::route_ticks( + ProviderInstanceId provider_id, + std::unique_ptr batch) { + if (!batch) return; + std::vector, TickDataBatch>> deliveries; + { + std::lock_guard lock(m_mutex); + const auto provider_it = m_providers.find(provider_id); + if (provider_it == m_providers.end()) return; + + 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()) return; + const auto& entry = entry_it->second; + auto subscriber = entry->subscriber.lock(); + if (!subscriber || + entry->phase != EntryPhase::ACTIVE || + !batch_matches_stream(*batch, entry->stream)) { + return; + } + auto routed = *batch; + routed.subscription = entry->control->provider_subscription; + deliveries.emplace_back(std::move(subscriber), std::move(routed)); + } else { + for (const auto& [id, entry] : m_entries) { + (void)id; + if (entry->provider_id != provider_id || + entry->phase != EntryPhase::ACTIVE || + !batch_matches_stream(*batch, entry->stream)) { + continue; + } + auto subscriber = entry->subscriber.lock(); + if (!subscriber) continue; + auto routed = *batch; + routed.subscription = entry->control->provider_subscription; + deliveries.emplace_back(std::move(subscriber), std::move(routed)); + } + } + } + + for (auto& delivery : deliveries) { + delivery.first->on_tick_data(delivery.second); + } + } + + inline void MarketDataRouterState::route_bars( + ProviderInstanceId provider_id, + std::unique_ptr batch) { + if (!batch) return; + std::vector, BarDataBatch>> deliveries; + { + std::lock_guard lock(m_mutex); + const auto provider_it = m_providers.find(provider_id); + if (provider_it == m_providers.end()) return; + + 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()) return; + const auto& entry = entry_it->second; + auto subscriber = entry->subscriber.lock(); + if (!subscriber || + entry->phase != EntryPhase::ACTIVE || + !batch_matches_stream(*batch, entry->stream)) { + return; + } + auto routed = *batch; + routed.subscription = entry->control->provider_subscription; + deliveries.emplace_back(std::move(subscriber), std::move(routed)); + } else { + for (const auto& [id, entry] : m_entries) { + (void)id; + if (entry->provider_id != provider_id || + entry->phase != EntryPhase::ACTIVE || + !batch_matches_stream(*batch, entry->stream)) { + continue; + } + auto subscriber = entry->subscriber.lock(); + if (!subscriber) continue; + auto routed = *batch; + routed.subscription = entry->control->provider_subscription; + deliveries.emplace_back(std::move(subscriber), std::move(routed)); + } + } + } + + for (auto& delivery : deliveries) { + delivery.first->on_bar_data(delivery.second); + } + } + + inline void MarketDataRouterState::route_status( + ProviderInstanceId provider_id, + MarketDataStatusUpdate update) { + std::vector, MarketDataStatusUpdate>> + deliveries; + { + std::lock_guard lock(m_mutex); + const auto provider_it = m_providers.find(provider_id); + if (provider_it == m_providers.end()) return; + if (update.subscription.valid() && + update.subscription.provider_id != provider_id) { + return; + } + + update.provider_id = provider_id; + if (update.subscription.valid()) { + update.type = update.subscription.stream_type; + update.symbol = update.subscription.symbol; + update.timeframe = update.subscription.timeframe; + if (update.transport == MarketDataTransport::AUTO) { + update.transport = update.subscription.transport; + } + } + cache_status_no_lock(provider_it->second, update); + + if (update.subscription.valid()) { + const auto route_it = provider_it->second.provider_routes.find( + update.subscription.id); + if (route_it != provider_it->second.provider_routes.end()) { + const auto entry_it = m_entries.find(route_it->second); + if (entry_it != m_entries.end() && + entry_it->second->phase == EntryPhase::ACTIVE) { + auto subscriber = entry_it->second->subscriber.lock(); + if (subscriber) { + update.subscription = + entry_it->second->control->provider_subscription; + deliveries.emplace_back( + std::move(subscriber), + std::move(update)); + } + } + } + } else { + for (const auto& [id, entry] : m_entries) { + (void)id; + if (entry->provider_id != provider_id || + entry->phase != EntryPhase::ACTIVE || + !status_matches_stream(update, entry->stream)) { + continue; + } + auto subscriber = entry->subscriber.lock(); + if (!subscriber) continue; + auto routed = update; + routed.subscription = entry->control->provider_subscription; + deliveries.emplace_back(std::move(subscriber), std::move(routed)); + } + } + } + + for (auto& delivery : deliveries) { + delivery.first->on_market_data_status(delivery.second); + } + } + + inline void MarketDataRouterState::set_control_active( + const std::shared_ptr& control, + const MarketDataSubscriptionHandle& subscription) { + if (!control) return; + std::lock_guard lock(control->mutex); + control->provider_subscription = subscription; + control->active = true; + } + + inline void MarketDataRouterState::set_control_released( + const std::shared_ptr& control) { + if (!control) return; + std::lock_guard lock(control->mutex); + control->active = false; + control->released = true; + } + + inline void MarketDataRouterState::dispatch_result( + subscription_callback_t callback, + MarketDataSubscriptionResult result) { + if (callback) callback(std::move(result)); + } + + inline std::size_t MarketDataRouterState::subscription_count() const { + std::lock_guard lock(m_mutex); + return m_entries.size(); + } + + inline std::size_t MarketDataRouterState::failed_unsubscribe_count() const { + std::lock_guard lock(m_mutex); + return static_cast(std::count_if( + m_entries.begin(), + m_entries.end(), + [](const auto& item) { + return item.second->phase == EntryPhase::CLEANUP_FAILED; + })); + } + + inline std::size_t MarketDataRouterState::retry_failed_unsubscribes() { + std::vector> controls; + { + std::lock_guard lock(m_mutex); + controls.reserve(m_entries.size()); + for (const auto& [id, entry] : m_entries) { + (void)id; + if (entry->phase == EntryPhase::CLEANUP_FAILED) { + controls.push_back(entry->control); + } + } + } + + std::size_t accepted = 0; + for (const auto& control : controls) { + if (unsubscribe(control, {})) ++accepted; + } + return accepted; + } + + inline void MarketDataRouterState::process() { + struct PendingCompletion { + RoutedSubscriptionId router_id; + MarketDataSubscriptionHandle subscription; + MarketDataSubscriptionResult result; + }; + struct CleanupRequest { + std::shared_ptr entry; + MarketDataSubscriptionHandle subscription; + }; + + for (;;) { + std::vector completions; + { + std::lock_guard lock(m_mutex); + if (!m_shutdown || m_shutdown_complete) return; + + completions.reserve(m_entries.size()); + for (const auto& [id, entry] : m_entries) { + if (entry->phase != EntryPhase::UNSUBSCRIBING || + !entry->unsubscribe_completion_received) { + continue; + } + completions.push_back(PendingCompletion{ + id, + entry->unsubscribe_completion.subscription, + entry->unsubscribe_completion}); + } + } + + for (auto& completion : completions) { + complete_unsubscribe( + completion.router_id, + std::move(completion.subscription), + std::move(completion.result), + {}); + } + + std::vector failed_subscribes; + std::vector cleanup_requests; + std::vector unbind_providers; + bool shutdown_complete = false; + { + std::lock_guard lock(m_mutex); + failed_subscribes.reserve(m_entries.size()); + cleanup_requests.reserve(m_entries.size()); + + for (const auto& [id, entry] : m_entries) { + MarketDataSubscriptionHandle subscription; + if (entry->phase == EntryPhase::PENDING && + entry->subscribe_completion_received) { + if (!entry->retained_cleanup_subscription.valid()) { + failed_subscribes.push_back(id); + continue; + } + subscription = entry->retained_cleanup_subscription; + } else if (entry->phase == EntryPhase::ACTIVE) { + subscription = entry->control->provider_subscription; + } else { + continue; + } + + if (!entry->provider || !subscription.valid()) { + entry->phase = EntryPhase::CLEANUP_FAILED; + continue; + } + entry->phase = EntryPhase::UNSUBSCRIBING; + const auto provider_it = m_providers.find(entry->provider_id); + if (provider_it != m_providers.end()) { + provider_it->second.provider_routes.erase(subscription.id); + } + cleanup_requests.push_back(CleanupRequest{ + entry, + std::move(subscription)}); + } + + for (const auto id : failed_subscribes) { + const auto entry_it = m_entries.find(id); + if (entry_it == m_entries.end()) continue; + set_control_released(entry_it->second->control); + if (auto* provider = remove_entry_no_lock(id)) { + unbind_providers.push_back(provider); + } + } + + if (m_entries.empty()) { + m_registered_providers.clear(); + m_provider_aliases.clear(); + m_registered_provider_ids.clear(); + m_shutdown_complete = true; + shutdown_complete = true; + } + } + + for (auto* provider : unbind_providers) { + if (provider) unbind_provider(*provider); + } + for (auto& request : cleanup_requests) { + start_unsubscribe( + request.entry, + std::move(request.subscription), + {}); + } + + if (shutdown_complete) return; + if (completions.empty() && + failed_subscribes.empty() && + cleanup_requests.empty()) { + return; + } + } + } + + inline bool MarketDataRouterState::is_shutdown_complete() const noexcept { + std::lock_guard lock(m_mutex); + return m_shutdown_complete; + } + + inline void MarketDataRouterState::shutdown() noexcept { + { + std::lock_guard lock(m_mutex); + if (m_shutdown_complete) return; + if (m_shutdown) return; + m_shutdown = true; + for (const auto& [id, entry] : m_entries) { + (void)id; + set_control_released(entry->control); + entry->subscriber.reset(); + entry->release_callback = {}; + if (entry->phase == EntryPhase::CLEANUP_FAILED) { + entry->phase = entry->retained_cleanup_subscription.valid() + ? EntryPhase::PENDING + : EntryPhase::ACTIVE; + } + } + } + + process(); + } + + } // namespace detail + + inline MarketDataRouterSubscription& + MarketDataRouterSubscription::operator=(MarketDataRouterSubscription&& other) noexcept { + if (this == &other) return *this; + reset(); + m_control = std::move(other.m_control); + return *this; + } + + inline RoutedSubscriptionId MarketDataRouterSubscription::router_id() const noexcept { + return m_control ? m_control->router_id : RoutedSubscriptionId{}; + } + + inline MarketDataSubscriptionHandle + MarketDataRouterSubscription::provider_subscription() const { + if (!m_control) return {}; + std::lock_guard lock(m_control->mutex); + return m_control->provider_subscription; + } + + inline MarketDataProviderId + MarketDataRouterSubscription::registered_provider_id() const { + if (!m_control) return {}; + std::lock_guard lock(m_control->mutex); + return m_control->registered_provider_id; + } + + inline bool MarketDataRouterSubscription::valid() const { + if (!m_control) return false; + std::lock_guard lock(m_control->mutex); + return !m_control->released && + m_control->router_id.valid(); + } + + inline bool MarketDataRouterSubscription::active() const { + if (!m_control) return false; + std::lock_guard lock(m_control->mutex); + return !m_control->released && + m_control->active && + m_control->provider_subscription.valid(); + } + + inline bool MarketDataRouterSubscription::unsubscribe( + BaseMarketDataProvider::subscription_callback_t callback) { + if (!m_control) return false; + const auto control = std::move(m_control); + const auto router = control->router.lock(); + if (!router) { + detail::MarketDataRouterState::set_control_released(control); + return false; + } + return router->unsubscribe(control, std::move(callback)); + } + + inline void MarketDataRouterSubscription::reset() noexcept { + if (!m_control) return; + const auto control = std::move(m_control); + try { + if (const auto router = control->router.lock()) { + router->unsubscribe(control, {}); + } else { + detail::MarketDataRouterState::set_control_released(control); + } + } catch (...) { + detail::MarketDataRouterState::set_control_released(control); + } + } + + inline MarketDataRouter::MarketDataRouter() + : m_state(std::make_shared()) {} + + inline MarketDataRouter::MarketDataRouter(owner_dispatcher_t owner_dispatcher) + : m_state(std::make_shared( + std::move(owner_dispatcher))) {} + + inline MarketDataRouter::~MarketDataRouter() { + shutdown(); + } + + inline bool MarketDataRouter::register_provider( + MarketDataProviderId id, + BaseMarketDataProvider& provider, + std::vector aliases) { + return m_state && m_state->register_provider( + id, + provider, + std::move(aliases)); + } + + inline bool MarketDataRouter::add_provider_alias( + MarketDataProviderId id, + std::string alias) { + return m_state && m_state->add_provider_alias(id, std::move(alias)); + } + + inline bool MarketDataRouter::unregister_provider(MarketDataProviderId id) { + return m_state && m_state->unregister_provider(id); + } + + inline std::size_t MarketDataRouter::registered_provider_count() const { + return m_state ? m_state->registered_provider_count() : 0; + } + + inline MarketDataProviderId MarketDataRouter::registered_provider_id( + ProviderInstanceId provider_id) const { + return m_state ? m_state->registered_provider_id(provider_id) : MarketDataProviderId{}; + } + + inline std::vector MarketDataRouter::provider_aliases( + MarketDataProviderId id) const { + return m_state ? m_state->provider_aliases(id) : std::vector{}; + } + + inline bool MarketDataRouter::post_to_owner(owner_task_t task) const { + return m_state && m_state->post_to_owner(std::move(task)); + } + + inline bool MarketDataRouter::has_owner_dispatcher() const noexcept { + return m_state && m_state->has_owner_dispatcher(); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks( + BaseMarketDataProvider& provider, + const std::shared_ptr& subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback) { + return subscribe_ticks_weak( + provider, + std::weak_ptr(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks_weak( + BaseMarketDataProvider& provider, + std::weak_ptr subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback) { + return m_state->subscribe_ticks( + provider, + std::move(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks( + MarketDataProviderId provider_id, + const std::shared_ptr& subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback) { + return subscribe_ticks_weak( + provider_id, + std::weak_ptr(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks_weak( + MarketDataProviderId provider_id, + std::weak_ptr subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback) { + return m_state->subscribe_ticks( + provider_id, + std::move(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks( + std::string_view provider_alias, + const std::shared_ptr& subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback) { + return subscribe_ticks_weak( + provider_alias, + std::weak_ptr(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_ticks_weak( + std::string_view provider_alias, + std::weak_ptr subscriber, + TickSubscriptionRequest request, + subscription_callback_t callback) { + return m_state->subscribe_ticks( + provider_alias, + std::move(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars( + BaseMarketDataProvider& provider, + const std::shared_ptr& subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback) { + return subscribe_bars_weak( + provider, + std::weak_ptr(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars_weak( + BaseMarketDataProvider& provider, + std::weak_ptr subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback) { + return m_state->subscribe_bars( + provider, + std::move(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars( + MarketDataProviderId provider_id, + const std::shared_ptr& subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback) { + return subscribe_bars_weak( + provider_id, + std::weak_ptr(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars_weak( + MarketDataProviderId provider_id, + std::weak_ptr subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback) { + return m_state->subscribe_bars( + provider_id, + std::move(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars( + std::string_view provider_alias, + const std::shared_ptr& subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback) { + return subscribe_bars_weak( + provider_alias, + std::weak_ptr(subscriber), + std::move(request), + std::move(callback)); + } + + inline MarketDataRouter::SubscriptionHandle MarketDataRouter::subscribe_bars_weak( + std::string_view provider_alias, + std::weak_ptr subscriber, + BarSubscriptionRequest request, + subscription_callback_t callback) { + return m_state->subscribe_bars( + provider_alias, + std::move(subscriber), + std::move(request), + std::move(callback)); + } + + inline std::size_t MarketDataRouter::subscription_count() const { + return m_state ? m_state->subscription_count() : 0; + } + + inline std::size_t MarketDataRouter::failed_unsubscribe_count() const { + return m_state ? m_state->failed_unsubscribe_count() : 0; + } + + inline std::size_t MarketDataRouter::retry_failed_unsubscribes() { + return m_state ? m_state->retry_failed_unsubscribes() : 0; + } + + inline void MarketDataRouter::process() { + if (m_state) m_state->process(); + } + + inline bool MarketDataRouter::is_shutdown_complete() const noexcept { + return !m_state || m_state->is_shutdown_complete(); + } + + inline void MarketDataRouter::shutdown() noexcept { + if (m_state) m_state->shutdown(); + } + + +} // namespace optionx::market_data + +#endif // OPTIONX_HEADER_MARKET_DATA_DETAIL_MARKET_DATA_ROUTER_IPP_INCLUDED diff --git a/include/optionx_cpp/optionx.hpp b/include/optionx_cpp/optionx.hpp index d6db941..6eef746 100644 --- a/include/optionx_cpp/optionx.hpp +++ b/include/optionx_cpp/optionx.hpp @@ -6,6 +6,7 @@ /// \brief Includes core headers for the OptionX library. #include "utils.hpp" +#include "lifecycle.hpp" #include "data.hpp" #include "market_data.hpp" #include "storages.hpp" diff --git a/include/optionx_cpp/platforms.hpp b/include/optionx_cpp/platforms.hpp index acad157..64c98b7 100644 --- a/include/optionx_cpp/platforms.hpp +++ b/include/optionx_cpp/platforms.hpp @@ -18,6 +18,7 @@ #include "config.hpp" #include "utils.hpp" +#include "lifecycle.hpp" #include "data.hpp" #include "market_data.hpp" #include "storages.hpp" diff --git a/include/optionx_cpp/platforms/common/BaseTradingPlatform.hpp b/include/optionx_cpp/platforms/common/BaseTradingPlatform.hpp index a6833ab..98b6b94 100644 --- a/include/optionx_cpp/platforms/common/BaseTradingPlatform.hpp +++ b/include/optionx_cpp/platforms/common/BaseTradingPlatform.hpp @@ -9,7 +9,10 @@ namespace optionx::platforms { /// \class BaseTradingPlatform /// \brief Base endpoint facade for trading platforms, account data, lifecycle, and connection state. - class BaseTradingPlatform : public BaseEndpoint, public BaseTradingApi { + class BaseTradingPlatform + : public BaseEndpoint, + public BaseTradingApi, + public lifecycle::ILifecycleModule { public: BaseTradingPlatform(std::shared_ptr account_info) : m_account_info(std::move(account_info)), @@ -241,6 +244,11 @@ namespace optionx::platforms { m_stopped.store(true, std::memory_order_release); }; + /// \brief Returns true after the platform lifecycle has stopped. + [[nodiscard]] bool is_stopped() const noexcept override { + return m_stopped.load(std::memory_order_acquire); + } + /// \brief Returns a reference to the event bus. utils::EventBus& event_bus() { return m_event_bus; } diff --git a/tests/lifecycle_event_safety_test.cpp b/tests/lifecycle_event_safety_test.cpp index d0943eb..472e49f 100644 --- a/tests/lifecycle_event_safety_test.cpp +++ b/tests/lifecycle_event_safety_test.cpp @@ -62,6 +62,16 @@ class TestPlatform final : public optionx::platforms::BaseTradingPlatform { } // namespace +TEST(BaseTradingPlatformLifecycle, ImplementsCommonLifecycleModule) { + TestPlatform platform; + optionx::lifecycle::ILifecycleModule& module = platform; + + EXPECT_FALSE(module.is_stopped()); + module.process(); + module.shutdown(); + EXPECT_TRUE(module.is_stopped()); +} + TEST(BaseTradingPlatformLifecycle, RepeatedRunDoesNotDuplicateLifecycleTasks) { TestPlatform platform; diff --git a/tests/lifecycle_stack_test.cpp b/tests/lifecycle_stack_test.cpp new file mode 100644 index 0000000..a424f20 --- /dev/null +++ b/tests/lifecycle_stack_test.cpp @@ -0,0 +1,132 @@ +#include + +#include + +#include +#include +#include + +namespace { + +using optionx::lifecycle::ILifecycleModule; +using optionx::lifecycle::LifecycleStack; + +class RecordingModule final : public ILifecycleModule { +public: + RecordingModule( + std::string name, + std::size_t shutdown_process_count, + std::vector& events) + : m_name(std::move(name)), + m_shutdown_process_count(shutdown_process_count), + m_events(events) {} + + void process() override { + m_events.push_back("process:" + m_name); + if (!m_shutdown_requested || m_shutdown_process_count == 0) return; + --m_shutdown_process_count; + if (m_shutdown_process_count == 0) m_stopped = true; + } + + void shutdown() noexcept override { + m_events.push_back("shutdown:" + m_name); + m_shutdown_requested = true; + if (m_shutdown_process_count == 0) m_stopped = true; + } + + [[nodiscard]] bool is_stopped() const noexcept override { + return m_stopped; + } + +private: + std::string m_name; + std::size_t m_shutdown_process_count = 0; + std::vector& m_events; + bool m_shutdown_requested = false; + bool m_stopped = false; +}; + +TEST(LifecycleStack, ProcessesForwardAndShutsDownDependentsOneAtATime) { + std::vector events; + RecordingModule platform("platform", 0, events); + RecordingModule router("router", 1, events); + RecordingModule bot("bot", 1, events); + LifecycleStack stack; + + ASSERT_TRUE(stack.add_module(platform)); + ASSERT_TRUE(stack.add_module(router)); + ASSERT_TRUE(stack.add_module(bot)); + EXPECT_FALSE(stack.add_module(router)); + EXPECT_FALSE(stack.add_module(stack)); + + stack.process(); + EXPECT_EQ(events, (std::vector{ + "process:platform", + "process:router", + "process:bot"})); + events.clear(); + + stack.shutdown(); + EXPECT_TRUE(stack.is_shutdown_requested()); + EXPECT_FALSE(stack.is_stopped()); + EXPECT_EQ(events, (std::vector{"shutdown:bot"})); + events.clear(); + + stack.process(); + EXPECT_FALSE(stack.is_stopped()); + EXPECT_EQ(events, (std::vector{ + "process:platform", + "process:router", + "process:bot", + "shutdown:router"})); + events.clear(); + + stack.process(); + EXPECT_TRUE(stack.is_stopped()); + EXPECT_EQ(events, (std::vector{ + "process:platform", + "process:router", + "shutdown:platform"})); + + events.clear(); + stack.process(); + stack.shutdown(); + EXPECT_TRUE(events.empty()); + EXPECT_FALSE(stack.add_module(bot)); +} + +TEST(LifecycleStack, StopsSynchronousModulesInReverseOrder) { + std::vector events; + RecordingModule first("first", 0, events); + RecordingModule second("second", 0, events); + LifecycleStack stack; + + ASSERT_TRUE(stack.add_module(first)); + ASSERT_TRUE(stack.add_module(second)); + stack.shutdown(); + + EXPECT_TRUE(stack.is_stopped()); + EXPECT_EQ(events, (std::vector{ + "shutdown:second", + "shutdown:first"})); +} + +TEST(LifecycleStack, EmptyStackStopsOnRequest) { + LifecycleStack stack; + + EXPECT_TRUE(stack.empty()); + EXPECT_EQ(stack.size(), 0u); + EXPECT_FALSE(stack.is_stopped()); + + stack.shutdown(); + + EXPECT_TRUE(stack.is_shutdown_requested()); + EXPECT_TRUE(stack.is_stopped()); +} + +} // namespace + +int main(int argc, char** argv) { + testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +} diff --git a/tests/market_data_router_test.cpp b/tests/market_data_router_test.cpp index 36498c6..86bab1a 100644 --- a/tests/market_data_router_test.cpp +++ b/tests/market_data_router_test.cpp @@ -274,6 +274,16 @@ MarketDataStatusUpdate ready_status( } // namespace +TEST(MarketDataRouter, ImplementsCommonLifecycleModule) { + MarketDataRouter router; + optionx::lifecycle::ILifecycleModule& module = router; + + EXPECT_FALSE(module.is_stopped()); + module.process(); + module.shutdown(); + EXPECT_TRUE(module.is_stopped()); +} + TEST(MarketDataRouter, UsesStrongRoutedSubscriptionIds) { const RoutedSubscriptionId empty; EXPECT_FALSE(empty.valid());