diff --git a/app/src/market_data/acquisition/runtime/acquisition_runtime_service.py b/app/src/market_data/acquisition/runtime/acquisition_runtime_service.py new file mode 100644 index 0000000..d6f9b5a --- /dev/null +++ b/app/src/market_data/acquisition/runtime/acquisition_runtime_service.py @@ -0,0 +1,136 @@ +# app/src/market_data/acquisition/runtime/acquisition_runtime_service.py + +from __future__ import annotations + +""" +Сервис исполнения команд Acquisition Runtime. + +Build 060.22 вводит первую производственную реализацию +сервисного слоя Runtime подсистемы Market Data Acquisition. + +Сервис принимает типизированные Runtime-команды и делегирует +их выполнение соответствующим инфраструктурным зависимостям. + +Сервис не содержит знаний о Trades Feed, Consistency или Recovery. +""" + +from src.market_data.acquisition.runtime.acquisition_runtime_service_protocol import ( + AcquisitionRuntimeServiceProtocol, +) +from src.market_data.acquisition.runtime.runtime_commands import ( + ConnectCommand, + DisconnectCommand, + SendBinaryCommand, + SendTextCommand, + SubscribeCommand, + UnsubscribeCommand, +) +from src.market_data.acquisition.runtime.websocket_protocol import ( + AcquisitionRuntimeCommand, + AcquisitionRuntimeCommandDispatcherProtocol, + AcquisitionRuntimeEventPublisherProtocol, + WebSocketSessionProtocol, + WebSocketSubscriptionManagerProtocol, + WebSocketTransportProtocol, +) + + +class AcquisitionRuntimeService( + AcquisitionRuntimeServiceProtocol, + AcquisitionRuntimeCommandDispatcherProtocol, +): + """ + Сервис маршрутизации Acquisition Runtime Commands. + + Каждая команда делегируется ровно одной инфраструктурной + зависимости. + + Сервис не управляет reconnect, heartbeat, scheduler, + Recovery или обработкой рыночных данных. + """ + + def __init__( + self, + session: WebSocketSessionProtocol, + transport: WebSocketTransportProtocol, + subscription_manager: WebSocketSubscriptionManagerProtocol, + event_publisher: AcquisitionRuntimeEventPublisherProtocol, + ) -> None: + """ + Создать Acquisition Runtime Service. + + Args: + session: + Компонент жизненного цикла WebSocket-сессии. + + transport: + Низкоуровневый WebSocket-транспорт. + + subscription_manager: + Компонент управления активными подписками. + + event_publisher: + Компонент публикации инфраструктурных Runtime-событий. + + В Build 060.22 зависимость сохраняется как часть + утверждённой архитектуры, но публикация событий + будет подключена в последующих Build. + """ + self._session = session + self._transport = transport + self._subscription_manager = subscription_manager + self._event_publisher = event_publisher + + async def dispatch( + self, + command: AcquisitionRuntimeCommand, + ) -> None: + """ + Выполнить одну инфраструктурную команду Acquisition Runtime. + + Args: + command: + Типизированная команда Runtime Layer. + + Raises: + TypeError: + Если передан неподдерживаемый тип команды. + + Exception: + Любое исключение инфраструктурной зависимости + распространяется вызывающему компоненту без изменения. + """ + if isinstance(command, ConnectCommand): + await self._session.start() + return + + if isinstance(command, DisconnectCommand): + await self._session.stop() + return + + if isinstance(command, SubscribeCommand): + await self._subscription_manager.subscribe( + command.subscription_key, + command.message, + ) + return + + if isinstance(command, UnsubscribeCommand): + await self._subscription_manager.unsubscribe( + command.subscription_key, + command.message, + ) + return + + if isinstance(command, SendTextCommand): + await self._transport.send(command.message.payload) + return + + if isinstance(command, SendBinaryCommand): + await self._transport.send(command.message.payload) + return + + raise TypeError( + "Неподдерживаемый тип Acquisition Runtime Command: " + f"{type(command).__name__}." + ) \ No newline at end of file diff --git a/app/src/market_data/acquisition/runtime/acquisition_runtime_service_protocol.py b/app/src/market_data/acquisition/runtime/acquisition_runtime_service_protocol.py new file mode 100644 index 0000000..e366835 --- /dev/null +++ b/app/src/market_data/acquisition/runtime/acquisition_runtime_service_protocol.py @@ -0,0 +1,44 @@ +# app/src/market_data/acquisition/runtime/acquisition_runtime_service_protocol.py + +from __future__ import annotations + +""" +Публичный контракт Acquisition Runtime Service. + +Build 060.22 вводит сервисный уровень исполнения типизированных +команд Runtime подсистемы Market Data Acquisition. + +Контракт не определяет устройство WebSocket Session, Transport, +Subscription Manager или механизма публикации событий. +""" + +from typing import Protocol, runtime_checkable + +from src.market_data.acquisition.runtime.websocket_protocol import ( + AcquisitionRuntimeCommand, +) + + +@runtime_checkable +class AcquisitionRuntimeServiceProtocol(Protocol): + """ + Контракт сервиса исполнения Acquisition Runtime Commands. + + Реализация принимает одну типизированную команду и передаёт её + соответствующей инфраструктурной зависимости. + + Сервис не содержит знаний о Trades Feed, Consistency и Recovery. + """ + + async def dispatch( + self, + command: AcquisitionRuntimeCommand, + ) -> None: + """ + Выполнить одну инфраструктурную команду Acquisition Runtime. + + Args: + command: + Типизированная команда Runtime Layer. + """ + ... \ No newline at end of file diff --git a/app/src/market_data/acquisition/runtime/websocket_protocol.py b/app/src/market_data/acquisition/runtime/websocket_protocol.py index 415f058..f3818d2 100644 --- a/app/src/market_data/acquisition/runtime/websocket_protocol.py +++ b/app/src/market_data/acquisition/runtime/websocket_protocol.py @@ -23,6 +23,10 @@ from src.market_data.acquisition.runtime.runtime_events import ( ReconnectFailedEvent, ReconnectStartedEvent, ) +from src.market_data.acquisition.runtime.transport_messages import ( + TransportBinaryMessage, + TransportTextMessage, +) AcquisitionRuntimeCommand = ( @@ -46,6 +50,11 @@ AcquisitionRuntimeEvent = ( | HeartbeatTimeoutEvent ) +AcquisitionSubscriptionMessage = ( + TransportTextMessage + | TransportBinaryMessage +) + @runtime_checkable class WebSocketTransportProtocol(Protocol): @@ -108,10 +117,50 @@ class WebSocketSubscriptionManagerProtocol(Protocol): """ Контракт управления активными WebSocket-подписками. - Конкретные модели подписок и формат сообщений будут добавлены - в последующих Build'ах. + Subscription Manager отвечает за регистрацию и удаление + активных подписок, а также за восстановление зарегистрированных + подписок после повторного подключения. + + Менеджер работает только с универсальными транспортными + сообщениями и не содержит exchange-specific логики. """ + async def subscribe( + self, + subscription_key: str, + message: AcquisitionSubscriptionMessage, + ) -> None: + """ + Зарегистрировать и выполнить WebSocket-подписку. + + Args: + subscription_key: + Уникальный инфраструктурный ключ подписки. + + message: + Полностью сформированное транспортное сообщение + подписки. + """ + ... + + async def unsubscribe( + self, + subscription_key: str, + message: AcquisitionSubscriptionMessage, + ) -> None: + """ + Отменить и удалить существующую WebSocket-подписку. + + Args: + subscription_key: + Уникальный инфраструктурный ключ подписки. + + message: + Полностью сформированное транспортное сообщение + отмены подписки. + """ + ... + async def restore_subscriptions(self) -> None: """Восстановить активные подписки после подключения.""" ... @@ -162,4 +211,4 @@ class AcquisitionRuntimeEventPublisherProtocol(Protocol): """ Опубликовать одно инфраструктурное событие Runtime. """ - ... + ... \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/runtime/test_acquisition_runtime_service.py b/app/tests/unit/market_data/acquisition/runtime/test_acquisition_runtime_service.py new file mode 100644 index 0000000..c187dbc --- /dev/null +++ b/app/tests/unit/market_data/acquisition/runtime/test_acquisition_runtime_service.py @@ -0,0 +1,323 @@ +# app/tests/unit/market_data/acquisition/runtime/test_acquisition_runtime_service.py + +from __future__ import annotations + +import asyncio + +import pytest + +from src.market_data.acquisition.runtime.acquisition_runtime_service import ( + AcquisitionRuntimeService, +) +from src.market_data.acquisition.runtime.acquisition_runtime_service_protocol import ( + AcquisitionRuntimeServiceProtocol, +) +from src.market_data.acquisition.runtime.runtime_commands import ( + ConnectCommand, + DisconnectCommand, + SendBinaryCommand, + SendTextCommand, + SubscribeCommand, + UnsubscribeCommand, +) +from src.market_data.acquisition.runtime.transport_messages import ( + TransportBinaryMessage, + TransportTextMessage, +) +from src.market_data.acquisition.runtime.websocket_protocol import ( + AcquisitionRuntimeCommandDispatcherProtocol, + AcquisitionRuntimeEvent, + AcquisitionSubscriptionMessage, +) + + +class FakeSession: + def __init__(self) -> None: + self.started = 0 + self.stopped = 0 + + @property + def is_connected(self) -> bool: + return False + + async def start(self) -> None: + self.started += 1 + + async def stop(self) -> None: + self.stopped += 1 + + +class FakeTransport: + def __init__(self) -> None: + self.messages: list[str | bytes] = [] + + async def connect(self) -> None: + return None + + async def disconnect(self) -> None: + return None + + async def send( + self, + message: str | bytes, + ) -> None: + self.messages.append(message) + + async def receive(self) -> str | bytes: + raise AssertionError("receive() should not be called.") + + +class FakeSubscriptionManager: + def __init__(self) -> None: + self.subscriptions: list[ + tuple[str, AcquisitionSubscriptionMessage] + ] = [] + self.unsubscriptions: list[ + tuple[str, AcquisitionSubscriptionMessage] + ] = [] + + async def subscribe( + self, + subscription_key: str, + message: AcquisitionSubscriptionMessage, + ) -> None: + self.subscriptions.append( + ( + subscription_key, + message, + ) + ) + + async def unsubscribe( + self, + subscription_key: str, + message: AcquisitionSubscriptionMessage, + ) -> None: + self.unsubscriptions.append( + ( + subscription_key, + message, + ) + ) + + async def restore_subscriptions(self) -> None: + return None + + async def clear_subscriptions(self) -> None: + return None + + +class FakeEventPublisher: + def __init__(self) -> None: + self.events: list[AcquisitionRuntimeEvent] = [] + + async def publish( + self, + event: AcquisitionRuntimeEvent, + ) -> None: + self.events.append(event) + + +def create_service() -> tuple[ + AcquisitionRuntimeService, + FakeSession, + FakeTransport, + FakeSubscriptionManager, + FakeEventPublisher, +]: + session = FakeSession() + transport = FakeTransport() + subscriptions = FakeSubscriptionManager() + publisher = FakeEventPublisher() + + service = AcquisitionRuntimeService( + session=session, + transport=transport, + subscription_manager=subscriptions, + event_publisher=publisher, + ) + + return ( + service, + session, + transport, + subscriptions, + publisher, + ) + + +def test_service_implements_protocol() -> None: + service, *_ = create_service() + + assert isinstance( + service, + AcquisitionRuntimeServiceProtocol, + ) + + +def test_service_is_command_dispatcher() -> None: + service, *_ = create_service() + + assert isinstance( + service, + AcquisitionRuntimeCommandDispatcherProtocol, + ) + + +def test_dispatch_connect_command() -> None: + service, session, *_ = create_service() + + asyncio.run( + service.dispatch( + ConnectCommand(), + ) + ) + + assert session.started == 1 + assert session.stopped == 0 + + +def test_dispatch_disconnect_command() -> None: + service, session, *_ = create_service() + + asyncio.run( + service.dispatch( + DisconnectCommand(), + ) + ) + + assert session.started == 0 + assert session.stopped == 1 + + +def test_dispatch_subscribe_command() -> None: + service, _, _, subscriptions, _ = create_service() + + message = TransportTextMessage( + payload='{"subscribe":true}', + ) + + asyncio.run( + service.dispatch( + SubscribeCommand( + subscription_key="BTCUSDT", + message=message, + ) + ) + ) + + assert subscriptions.subscriptions == [ + ( + "BTCUSDT", + message, + ) + ] + + +def test_dispatch_unsubscribe_command() -> None: + service, _, _, subscriptions, _ = create_service() + + message = TransportTextMessage( + payload='{"unsubscribe":true}', + ) + + asyncio.run( + service.dispatch( + UnsubscribeCommand( + subscription_key="BTCUSDT", + message=message, + ) + ) + ) + + assert subscriptions.unsubscriptions == [ + ( + "BTCUSDT", + message, + ) + ] + + +def test_dispatch_send_text_command() -> None: + service, _, transport, _, _ = create_service() + + asyncio.run( + service.dispatch( + SendTextCommand( + message=TransportTextMessage( + payload="hello", + ) + ) + ) + ) + + assert transport.messages == [ + "hello", + ] + + +def test_dispatch_send_binary_command() -> None: + service, _, transport, _, _ = create_service() + + asyncio.run( + service.dispatch( + SendBinaryCommand( + message=TransportBinaryMessage( + payload=b"\x01\x02", + ) + ) + ) + ) + + assert transport.messages == [ + b"\x01\x02", + ] + + +def test_event_publisher_is_not_used_yet() -> None: + service, _, _, _, publisher = create_service() + + asyncio.run( + service.dispatch( + ConnectCommand(), + ) + ) + + assert publisher.events == [] + + +def test_unsupported_command_raises_type_error() -> None: + class UnsupportedCommand: + pass + + service, *_ = create_service() + + with pytest.raises( + TypeError, + match="Неподдерживаемый тип Acquisition Runtime Command", + ): + asyncio.run( + service.dispatch( + UnsupportedCommand(), # type: ignore[arg-type] + ) + ) + + +def test_session_exception_is_propagated() -> None: + class BrokenSession(FakeSession): + async def start(self) -> None: + raise RuntimeError("boom") + + service = AcquisitionRuntimeService( + session=BrokenSession(), + transport=FakeTransport(), + subscription_manager=FakeSubscriptionManager(), + event_publisher=FakeEventPublisher(), + ) + + with pytest.raises(RuntimeError, match="boom"): + asyncio.run( + service.dispatch( + ConnectCommand(), + ) + ) \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/runtime/test_websocket_protocol.py b/app/tests/unit/market_data/acquisition/runtime/test_websocket_protocol.py index 0a50c38..dfb7e47 100644 --- a/app/tests/unit/market_data/acquisition/runtime/test_websocket_protocol.py +++ b/app/tests/unit/market_data/acquisition/runtime/test_websocket_protocol.py @@ -7,6 +7,7 @@ from src.market_data.acquisition.runtime.websocket_protocol import ( AcquisitionRuntimeCommandDispatcherProtocol, AcquisitionRuntimeEvent, AcquisitionRuntimeEventPublisherProtocol, + AcquisitionSubscriptionMessage, WebSocketSessionProtocol, WebSocketSubscriptionManagerProtocol, WebSocketTransportProtocol, @@ -40,6 +41,20 @@ class FakeWebSocketSession: class FakeWebSocketSubscriptionManager: + async def subscribe( + self, + subscription_key: str, + message: AcquisitionSubscriptionMessage, + ) -> None: + return None + + async def unsubscribe( + self, + subscription_key: str, + message: AcquisitionSubscriptionMessage, + ) -> None: + return None + async def restore_subscriptions(self) -> None: return None @@ -100,6 +115,20 @@ def test_incomplete_transport_does_not_satisfy_protocol() -> None: assert not isinstance(IncompleteTransport(), WebSocketTransportProtocol) +def test_incomplete_subscription_manager_does_not_satisfy_protocol() -> None: + class IncompleteSubscriptionManager: + async def restore_subscriptions(self) -> None: + return None + + async def clear_subscriptions(self) -> None: + return None + + assert not isinstance( + IncompleteSubscriptionManager(), + WebSocketSubscriptionManagerProtocol, + ) + + def test_incomplete_command_dispatcher_does_not_satisfy_protocol() -> None: class IncompleteCommandDispatcher: pass @@ -117,4 +146,4 @@ def test_incomplete_event_publisher_does_not_satisfy_protocol() -> None: assert not isinstance( IncompleteEventPublisher(), AcquisitionRuntimeEventPublisherProtocol, - ) + ) \ No newline at end of file diff --git a/docs/migrations/build_060_22.md b/docs/migrations/build_060_22.md new file mode 100644 index 0000000..6342615 --- /dev/null +++ b/docs/migrations/build_060_22.md @@ -0,0 +1,1224 @@ +# Build 060.22 — Acquisition Runtime Service + +**Engineering Migration Report** + +--- + +# Контроль документа + +| Свойство | Значение | +|----------|----------| +| Build | 060.22 | +| Название | Acquisition Runtime Service | +| Статус | Completed | +| Проект | Dzentra | +| Подсистема | Market Data Acquisition | +| Компонент | Acquisition Runtime Service Layer | +| Версия | 1.0 | + +--- + +# Связанные документы + +- `build_060_22_architecture.md` — архитектурная спецификация Build. +- `build_060_21.md` — Engineering Migration Report предыдущего Build. +- `build_060_21_architecture.md` — спецификация Runtime Integration Contracts. +- `build_060_20_1.md` — Engineering Migration Report корректирующего Build владения состоянием Trade Stream. + +--- + +# Цель Build + +Build 060.21 завершил формирование контрактного слоя Acquisition Runtime. + +После предыдущего этапа в системе уже существовали: + +```text +AcquisitionRuntimeCommand +AcquisitionRuntimeEvent +AcquisitionRuntimeCommandDispatcherProtocol +AcquisitionRuntimeEventPublisherProtocol +``` + +Также были доступны базовые WebSocket-контракты: + +```text +WebSocketTransportProtocol +WebSocketSessionProtocol +WebSocketSubscriptionManagerProtocol +``` + +Однако Runtime Layer оставался декларативным. + +Команды описывали инфраструктурные намерения, но в системе отсутствовал производственный компонент, который: + +- принимал типизированную Runtime-команду; +- определял её конкретный тип; +- выбирал соответствующую инфраструктурную зависимость; +- выполнял одну детерминированную операцию; +- сохранял независимость от Trades Feed, Consistency и Recovery. + +Главной задачей Build 060.22 стало создание первого исполняемого сервисного компонента Acquisition Runtime: + +```text +AcquisitionRuntimeService +``` + +Одновременно был введён отдельный публичный контракт сервисного уровня: + +```text +AcquisitionRuntimeServiceProtocol +``` + +После завершения Build система получает: + +- единый сервис исполнения Acquisition Runtime Commands; +- детерминированную маршрутизацию всех шести существующих команд; +- отдельный публичный Protocol сервиса; +- расширенный контракт управления WebSocket-подписками; +- полное unit-покрытие нового сервисного слоя; +- подтверждённую регрессионную совместимость Runtime, Consistency, Recovery и Feeds. + +При этом Build сознательно не выполняет интеграцию Runtime Service в корневой Acquisition Service или Trades Feed. + +Эта задача относится к Build 060.23. + +--- + +# Предпосылки + +К началу настоящего Build Runtime Layer уже обладал полным набором неизменяемых команд: + +```text +ConnectCommand +DisconnectCommand +SubscribeCommand +UnsubscribeCommand +SendTextCommand +SendBinaryCommand +``` + +Также существовали инфраструктурные события: + +```text +ConnectedEvent +DisconnectedEvent +ConnectFailedEvent +MessageReceivedEvent +MessageSentEvent +ReconnectStartedEvent +ReconnectCompletedEvent +ReconnectFailedEvent +HeartbeatTimeoutEvent +``` + +Build 060.21 определил публичные границы передачи команд и событий, но сознательно не создавал исполняемую реализацию. + +Архитектура выглядела следующим образом: + +```text +AcquisitionRuntimeCommand + │ + ▼ +AcquisitionRuntimeCommandDispatcherProtocol + │ + ▼ + (реализация отсутствует) +``` + +Таким образом следующий этап должен был закрыть именно сервисную границу между типизированными командами и существующими инфраструктурными компонентами. + +--- + +# Результаты архитектурного аудита + +Перед началом реализации был выполнен аудит следующих компонентов: + +```text +src/market_data/acquisition/runtime/websocket_protocol.py +src/market_data/acquisition/runtime/runtime_commands.py +src/market_data/acquisition/runtime/runtime_events.py +src/market_data/acquisition/runtime/transport_messages.py +src/market_data/acquisition/runtime/supervisor.py +src/market_data/acquisition/runtime/reconnect.py +src/market_data/acquisition/runtime/scheduler.py +src/market_data/acquisition/runtime/heartbeat.py +``` + +Дополнительно были проанализированы: + +```text +src/market_data/acquisition/service.py +src/market_data/acquisition/protocol.py +src/market_data/acquisition/registry.py +``` + +а также соответствующие unit-тесты. + +Аудит подтвердил следующее. + +## Runtime orchestration ещё не реализован + +Файлы: + +```text +supervisor.py +reconnect.py +scheduler.py +heartbeat.py +``` + +не содержали производственной реализации. + +Следовательно Build 060.22 не являлся рефакторингом существующего Runtime Service. + +Он создавал первую минимальную сервисную реализацию с нуля. + +--- + +## Supervisor преждевременен + +На раннем этапе рассматривалась возможность реализовать сервис внутри: + +```text +supervisor.py +``` + +Данный вариант был отклонён. + +Supervisor является компонентом более высокого уровня и должен координировать: + +- Runtime Service; +- reconnect; +- heartbeat; +- scheduler; +- lifecycle длительно работающей подсистемы. + +Поскольку перечисленные механизмы ещё не реализованы, использование Supervisor в Build 060.22 создало бы преждевременную orchestration-логику. + +Поэтому был введён отдельный тонкий сервис: + +```text +AcquisitionRuntimeService +``` + +--- + +## Registry для Runtime Service не требуется + +Существующие Feed Registry предназначены для выбора одной реализации из нескольких зарегистрированных источников. + +Runtime Service в текущей архитектуре существует в единственном экземпляре и не представляет семейство взаимозаменяемых Feed-реализаций. + +Создание отдельного Runtime Registry было признано искусственным и исключено из Scope Build. + +--- + +## Корневые Acquisition-файлы не изменяются + +Аудит подтвердил отсутствие необходимости изменять: + +```text +src/market_data/acquisition/service.py +src/market_data/acquisition/protocol.py +src/market_data/acquisition/registry.py +``` + +Интеграция нового Runtime Service с Acquisition Layer относится к следующему этапу: + +```text +Build 060.23 — Acquisition Integration +``` + +--- + +# Архитектурное решение + +По результатам аудита была утверждена следующая модель. + +```text +AcquisitionRuntimeCommand + │ + ▼ +AcquisitionRuntimeServiceProtocol + │ + ▼ +AcquisitionRuntimeService + │ + ┌─────┼──────────────┐ + ▼ ▼ ▼ + Session Transport Subscription Manager +``` + +`AcquisitionRuntimeService` реализует одновременно: + +```text +AcquisitionRuntimeServiceProtocol +AcquisitionRuntimeCommandDispatcherProtocol +``` + +Сервис не реализует WebSocket самостоятельно. + +Он только делегирует каждую команду одной заранее определённой инфраструктурной зависимости. + +--- + +# Правило уникальности имён файлов + +Build 060.22 продолжает обязательное архитектурное правило проекта Dzentra: + +> Во всём репозитории не допускается создание нескольких новых файлов с одинаковым именем независимо от их расположения в каталогах. + +Единственным служебным исключением остаётся: + +```text +__init__.py +``` + +Поэтому Build сознательно не создаёт файлы с общими именами: + +```text +service.py +protocol.py +test_service.py +``` + +Вместо них добавлены уникальные имена: + +```text +acquisition_runtime_service.py +acquisition_runtime_service_protocol.py +test_acquisition_runtime_service.py +``` + +Имена однозначно отражают назначение файлов без необходимости учитывать родительский каталог. + +Repository-wide поиск показал, что в проекте уже существуют старые файлы с общими именами, созданные до введения данного правила. + +Build 060.22 не добавляет новых файлов с такими именами и не выполняет ретроспективное переименование существующего рабочего кода вне согласованного Scope. + +--- + +# AcquisitionRuntimeServiceProtocol + +В рамках Build создан новый публичный контракт: + +```text +AcquisitionRuntimeServiceProtocol +``` + +Файл: + +```text +src/market_data/acquisition/runtime/ +acquisition_runtime_service_protocol.py +``` + +Protocol определяет единственную публичную операцию: + +```python +async def dispatch( + command: AcquisitionRuntimeCommand, +) -> None +``` + +Другие публичные методы не вводятся. + +В частности отсутствуют: + +```text +connect() +disconnect() +subscribe() +unsubscribe() +send() +``` + +Все операции Runtime выражаются исключительно через типизированные команды. + +Это исключает появление второго параллельного API для одной и той же инфраструктурной операции. + +--- + +# Почему Service Protocol существует отдельно + +В Build 060.21 уже был добавлен: + +```text +AcquisitionRuntimeCommandDispatcherProtocol +``` + +Он определяет инфраструктурную способность принимать команды. + +Build 060.22 дополнительно вводит: + +```text +AcquisitionRuntimeServiceProtocol +``` + +Данный контракт фиксирует отдельную сервисную границу Runtime Layer. + +Такое разделение позволяет последующим компонентам зависеть именно от сервисной абстракции, а не от конкретной реализации. + +В будущем под тем же Protocol могут существовать различные реализации: + +```text +Simple Acquisition Runtime Service +Queued Acquisition Runtime Service +Supervised Acquisition Runtime Service +``` + +Build 060.22 реализует только минимальный прямой вариант. + +--- + +# AcquisitionRuntimeService + +Главным production-компонентом Build становится: + +```text +AcquisitionRuntimeService +``` + +Файл: + +```text +src/market_data/acquisition/runtime/ +acquisition_runtime_service.py +``` + +Сервис является тонким stateless-координатором. + +Он: + +- принимает одну типизированную команду; +- определяет её конкретный тип; +- вызывает соответствующую инфраструктурную зависимость; +- завершает выполнение после завершения делегированной операции. + +Сервис не: + +- создаёт фоновые задачи; +- хранит очередь команд; +- выполняет retries; +- управляет reconnect; +- запускает heartbeat; +- хранит подписки; +- обрабатывает Trade; +- вызывает Consistency; +- вызывает Recovery. + +--- + +# Зависимости сервиса + +Все зависимости передаются через конструктор. + +```text +AcquisitionRuntimeService + │ + ├── WebSocketSessionProtocol + ├── WebSocketTransportProtocol + ├── WebSocketSubscriptionManagerProtocol + └── AcquisitionRuntimeEventPublisherProtocol +``` + +Сервис не создаёт зависимости самостоятельно. + +Создание графа объектов относится к Composition Root и не входит в Scope Build 060.22. + +--- + +# Почему Event Publisher внедрён, но пока не используется + +Конструктор Runtime Service принимает: + +```text +AcquisitionRuntimeEventPublisherProtocol +``` + +Однако Build 060.22 не публикует события после выполнения команд. + +Это осознанное решение. + +Publisher включён в окончательную форму конструктора, чтобы последующие Build могли подключить публикацию событий без изменения публичной композиции сервиса. + +При этом преждевременная event-orchestration не добавляется. + +Unit-тест отдельно подтверждает, что Publisher пока остаётся неиспользованным. + +--- + +# Расширение WebSocketSubscriptionManagerProtocol + +Во время реализации было обнаружено расхождение между первоначальной архитектурной матрицей и фактическим контрактом Subscription Manager. + +До Build Protocol содержал только: + +```python +restore_subscriptions() +clear_subscriptions() +``` + +Но Runtime Service должен исполнять: + +```text +SubscribeCommand +UnsubscribeCommand +``` + +Следовательно существующий контракт был недостаточен для детерминированной маршрутизации. + +В Build `WebSocketSubscriptionManagerProtocol` локально расширен методами: + +```python +async def subscribe( + subscription_key: str, + message: AcquisitionSubscriptionMessage, +) -> None +``` + +и: + +```python +async def unsubscribe( + subscription_key: str, + message: AcquisitionSubscriptionMessage, +) -> None +``` + +Существующие методы: + +```text +restore_subscriptions() +clear_subscriptions() +``` + +полностью сохранены. + +Расширение является обратно совместимым с точки зрения существующих моделей данных и не изменяет поведение уже реализованных компонентов. + +--- + +# AcquisitionSubscriptionMessage + +Для операций подписки введён типовой alias: + +```text +AcquisitionSubscriptionMessage +``` + +Он объединяет: + +```text +TransportTextMessage +TransportBinaryMessage +``` + +Благодаря этому Subscription Manager принимает не произвольные `str | bytes`, а уже существующие неизменяемые модели транспортных сообщений. + +Это обеспечивает соответствие моделям: + +```text +SubscribeCommand +UnsubscribeCommand +``` + +и исключает повторную интерпретацию payload внутри Runtime Service. + +--- + +# Официальная матрица маршрутизации + +Build 060.22 реализует следующую детерминированную таблицу. + +| Runtime Command | Действие | +|-----------------|----------| +| `ConnectCommand` | `WebSocketSessionProtocol.start()` | +| `DisconnectCommand` | `WebSocketSessionProtocol.stop()` | +| `SubscribeCommand` | `WebSocketSubscriptionManagerProtocol.subscribe(...)` | +| `UnsubscribeCommand` | `WebSocketSubscriptionManagerProtocol.unsubscribe(...)` | +| `SendTextCommand` | `WebSocketTransportProtocol.send(message.payload)` | +| `SendBinaryCommand` | `WebSocketTransportProtocol.send(message.payload)` | + +Каждый тип команды имеет ровно один маршрут. + +Сервис не выполняет дополнительные действия и не вызывает несколько зависимостей для одной команды. + +--- + +# ConnectCommand + +Маршрут: + +```text +ConnectCommand + │ + ▼ +WebSocketSessionProtocol.start() +``` + +Runtime Service не вызывает `transport.connect()` напрямую. + +Жизненный цикл соединения принадлежит Session. + +--- + +# DisconnectCommand + +Маршрут: + +```text +DisconnectCommand + │ + ▼ +WebSocketSessionProtocol.stop() +``` + +Runtime Service не вызывает `transport.disconnect()` напрямую. + +Завершение сессии остаётся обязанностью Session. + +--- + +# SubscribeCommand + +Маршрут: + +```text +SubscribeCommand + │ + ▼ +WebSocketSubscriptionManagerProtocol.subscribe( + subscription_key, + message, +) +``` + +Runtime Service передаёт исходные значения без изменения. + +Сервис не хранит список подписок и не сериализует сообщение. + +--- + +# UnsubscribeCommand + +Маршрут: + +```text +UnsubscribeCommand + │ + ▼ +WebSocketSubscriptionManagerProtocol.unsubscribe( + subscription_key, + message, +) +``` + +Сервис не изменяет внутреннее состояние Subscription Manager самостоятельно. + +--- + +# SendTextCommand + +Маршрут: + +```text +SendTextCommand + │ + ▼ +WebSocketTransportProtocol.send(str) +``` + +В Transport передаётся: + +```text +command.message.payload +``` + +Runtime Service не вводит отдельный метод `send_text()` и не расширяет Transport Protocol без необходимости. + +--- + +# SendBinaryCommand + +Маршрут: + +```text +SendBinaryCommand + │ + ▼ +WebSocketTransportProtocol.send(bytes) +``` + +В Transport передаётся бинарный payload без преобразования. + +Существующий единый метод: + +```python +send(message: str | bytes) +``` + +полностью покрывает оба типа сообщений. + +--- + +# Обработка неизвестной команды + +Тип `AcquisitionRuntimeCommand` ограничивает штатный публичный вход сервиса. + +Однако Runtime Service также содержит защитную ветку для произвольного неподдерживаемого объекта. + +В этом случае генерируется: + +```text +TypeError +``` + +с диагностическим сообщением, включающим имя фактического типа. + +Это гарантирует, что появление новой Runtime-команды потребует явного обновления маршрутизации и не останется незамеченным. + +--- + +# Обработка исключений + +Build не вводит собственную иерархию ошибок Runtime Service. + +Любое исключение инфраструктурной зависимости распространяется вызывающему компоненту без обёртки. + +Например: + +```text +ConnectCommand + │ + ▼ +AcquisitionRuntimeService + │ + ▼ +WebSocketSession.start() + │ + ▼ +RuntimeError +``` + +Сервис не преобразует ошибку в общий `RuntimeDispatchError`. + +Такое решение сохраняет: + +- исходный тип исключения; +- исходное диагностическое сообщение; +- прозрачность инфраструктурного поведения; +- возможность точной обработки на более высоком уровне. + +--- + +# Stateless-архитектура сервиса + +`AcquisitionRuntimeService` не хранит собственное операционное состояние. + +Внутри объекта сохраняются только внедрённые зависимости. + +Сервис не хранит: + +- состояние соединения; +- выполненные команды; +- активные подписки; +- историю отправленных сообщений; +- reconnect attempts; +- Recovery state; +- Consistency state. + +Все перечисленные данные принадлежат специализированным компонентам. + +--- + +# Изменённые файлы + +## Новый Service Protocol + +```text +src/market_data/acquisition/runtime/ +acquisition_runtime_service_protocol.py +``` + +Добавлен: + +```text +AcquisitionRuntimeServiceProtocol +``` + +Protocol содержит единственный асинхронный метод `dispatch()`. + +--- + +## Новый Runtime Service + +```text +src/market_data/acquisition/runtime/ +acquisition_runtime_service.py +``` + +Добавлен: + +```text +AcquisitionRuntimeService +``` + +Класс реализует: + +```text +AcquisitionRuntimeServiceProtocol +AcquisitionRuntimeCommandDispatcherProtocol +``` + +и маршрутизирует все шесть Runtime Commands. + +--- + +## Расширение WebSocket Protocol Layer + +```text +src/market_data/acquisition/runtime/websocket_protocol.py +``` + +Добавлены: + +```text +AcquisitionSubscriptionMessage +``` + +и методы: + +```text +WebSocketSubscriptionManagerProtocol.subscribe() +WebSocketSubscriptionManagerProtocol.unsubscribe() +``` + +Существующие Protocol и alias сохранены. + +--- + +## Unit-тест Runtime Service + +```text +tests/unit/market_data/acquisition/runtime/ +test_acquisition_runtime_service.py +``` + +Добавлены тестовые реализации: + +- `FakeSession`; +- `FakeTransport`; +- `FakeSubscriptionManager`; +- `FakeEventPublisher`. + +Проверена маршрутизация всех поддерживаемых команд и распространение исключений. + +--- + +## Обновление теста WebSocket Protocol + +```text +tests/unit/market_data/acquisition/runtime/ +test_websocket_protocol.py +``` + +`FakeWebSocketSubscriptionManager` расширен методами: + +```text +subscribe() +unsubscribe() +``` + +Также добавлена негативная проверка неполной реализации обновлённого Protocol. + +--- + +# Файлы, которые не изменялись + +Build не изменяет: + +```text +src/market_data/acquisition/runtime/runtime_commands.py +src/market_data/acquisition/runtime/runtime_events.py +src/market_data/acquisition/runtime/transport_messages.py +src/market_data/acquisition/runtime/supervisor.py +src/market_data/acquisition/runtime/reconnect.py +src/market_data/acquisition/runtime/scheduler.py +src/market_data/acquisition/runtime/heartbeat.py +``` + +Также не изменяются: + +- корневой Acquisition Service; +- Acquisition Registry; +- Acquisition Protocol; +- Trades Feed; +- Quotes Feed; +- Candles Feed; +- Trade Stream Consistency; +- Trade Stream State Store; +- Trade Recovery; +- WebSocket adapters; +- каноническая модель `Trade`. + +--- + +# Unit-тестирование AcquisitionRuntimeService + +Новый сервис покрыт отдельным набором unit-тестов. + +Проверены следующие сценарии: + +- сервис соответствует `AcquisitionRuntimeServiceProtocol`; +- сервис соответствует `AcquisitionRuntimeCommandDispatcherProtocol`; +- `ConnectCommand` вызывает только `session.start()`; +- `DisconnectCommand` вызывает только `session.stop()`; +- `SubscribeCommand` передаёт исходные `subscription_key` и message; +- `UnsubscribeCommand` передаёт исходные `subscription_key` и message; +- `SendTextCommand` передаёт строковый payload в `transport.send()`; +- `SendBinaryCommand` передаёт бинарный payload в `transport.send()`; +- Event Publisher не используется преждевременно; +- неизвестная команда приводит к `TypeError`; +- исключение Session распространяется без изменения. + +Результат: + +```text +11 passed +``` + +--- + +# Unit-тестирование WebSocket Protocol + +После расширения Subscription Manager выполнен локальный прогон обновлённого Protocol Layer. + +Проверены: + +- полная реализация Transport; +- полная реализация Session; +- полная реализация Subscription Manager; +- Command Dispatcher; +- Event Publisher; +- отрицательные сценарии неполных реализаций. + +Результат: + +```text +9 passed +``` + +--- + +# Регрессионное тестирование Runtime Layer + +После добавления Runtime Service выполнен полный прогон: + +```text +tests/unit/market_data/acquisition/runtime +``` + +Проверены: + +- Acquisition Runtime Service; +- Runtime Commands; +- Runtime Events; +- Transport Messages; +- WebSocket Protocol Layer. + +Результат: + +```text +46 passed +``` + +Это подтверждает, что новый сервисный слой не нарушил существующее поведение Runtime. + +--- + +# Расширенное регрессионное тестирование + +Для подтверждения архитектурной изоляции выполнен совместный прогон: + +```text +Runtime +Consistency +Recovery +Feeds +``` + +Команда охватывала каталоги: + +```text +tests/unit/market_data/acquisition/runtime +tests/unit/market_data/acquisition/consistency +tests/unit/market_data/acquisition/recovery +tests/unit/market_data/acquisition/feeds +``` + +Результат: + +```text +180 passed +``` + +Тем самым подтверждено: + +- Runtime Service работает корректно; +- обновление Subscription Manager не нарушило Protocol Layer; +- Trade Stream Consistency не затронута; +- Trade Stream State Store не затронут; +- Trade Recovery не затронута; +- существующие Feed продолжают работать без изменений; +- преждевременная интеграция Build 060.23 не произошла. + +--- + +# Обратная совместимость + +Build сохраняет функциональное поведение существующих подсистем. + +Не изменились: + +- Runtime Command models; +- Runtime Event models; +- Transport Message models; +- WebSocket Transport API; +- WebSocket Session API; +- Canonical Trade Model; +- Trades Feed API; +- Consistency API; +- Recovery API; +- Feed Registry; +- корневые Acquisition Services. + +Единственное расширение существующего контракта относится к `WebSocketSubscriptionManagerProtocol` и добавляет операции, уже необходимые существующим `SubscribeCommand` и `UnsubscribeCommand`. + +--- + +# Производительность + +Runtime Service выполняет конечную цепочку `isinstance`-проверок по шести типам команд и один делегированный асинхронный вызов. + +Сервис не вводит: + +- очереди; +- дополнительные копии payload; +- сериализацию; +- фоновые задачи; +- блокировки; +- retries; +- хранение истории. + +Следовательно накладные расходы ограничены постоянным временем: + +```text +O(1) +``` + +для каждой команды. + +Основная задержка полностью определяется вызываемой инфраструктурной зависимостью. + +--- + +# Подтверждённые архитектурные инварианты + +## Runtime Service не знает предметную область + +Сервис не импортирует: + +```text +Trade +TradesFeed +TradeStreamConsistencyProtocol +TradeRecoveryProtocol +TradeStreamStateStore +``` + +--- + +## Runtime Service stateless + +Сервис хранит только ссылки на внедрённые зависимости. + +--- + +## Одна команда имеет один маршрут + +Ни одна команда не выполняет несколько инфраструктурных операций. + +--- + +## Session владеет lifecycle соединения + +`ConnectCommand` и `DisconnectCommand` делегируются Session, а не Transport напрямую. + +--- + +## Subscription Manager владеет подписками + +Runtime Service не хранит список подписок и не управляет их состоянием самостоятельно. + +--- + +## Transport владеет отправкой payload + +Text и Binary payload передаются через существующий единый метод `send()`. + +--- + +## Event publication не реализована преждевременно + +Publisher внедрён, но не используется до отдельного этапа интеграции событий. + +--- + +## Consistency и Recovery не зависят от Runtime Service + +Новые импорты и зависимости в данных подсистемах отсутствуют. + +--- + +## Уникальность новых имён файлов соблюдена + +Build не добавляет новых `service.py`, `protocol.py` или `test_service.py`. + +--- + +# Что не входит в Scope Build + +Настоящий Build сознательно не реализует: + +- интеграцию Runtime Service в корневой Acquisition Service; +- интеграцию Runtime Service с Trades Feed; +- Composition Root; +- Runtime Supervisor; +- Reconnect orchestration; +- Heartbeat; +- Scheduler; +- очередь команд; +- background worker; +- публикацию Runtime Events; +- обработку входящих WebSocket-сообщений; +- маршрутизацию сообщений к Trade Adapter; +- запуск Trade Stream Consistency; +- автоматический Trade Recovery; +- восстановление подписок после reconnect; +- реальный WebSocket transport implementation; +- сохранение Runtime-состояния между перезапусками. + +Отсутствие перечисленных компонентов является осознанной границей Build и не рассматривается как незавершённость. + +--- + +# Архитектурное значение Build + +Build 060.22 является первым этапом серии Runtime, добавляющим исполняемое поведение. + +До него Runtime Layer содержал: + +- модели сообщений; +- команды; +- события; +- Protocol. + +После завершения Build появляется единая точка выполнения команд: + +```text +AcquisitionRuntimeService +``` + +Это превращает Runtime из исключительно декларативного слоя в минимальную работающую сервисную подсистему. + +При этом сервис остаётся полностью изолированным от предметной области и готов к последующей интеграции без изменения уже реализованных Consistency, Recovery и Feed. + +--- + +# Связь с последующими Build + +## Build 060.23 — Acquisition Integration + +Новый Runtime Service будет подключён к верхнему уровню Acquisition и станет доступен через утверждённую композицию зависимостей. + +--- + +## Build 060.24 — Reconnect & Runtime Recovery + +На основе Runtime Service будут построены: + +- reconnect orchestration; +- resubscribe; +- восстановление Runtime-потока; +- вызов Trade Recovery после нарушения непрерывности. + +--- + +## Build 060.25 — Integration & Regression + +Будет выполнена полная проверка взаимодействия Runtime, Acquisition, Trades Feed, Consistency и Recovery. + +--- + +# Заключение + +Build 060.22 завершает создание минимального исполняемого сервисного слоя Acquisition Runtime. + +В рамках этапа были реализованы: + +```text +AcquisitionRuntimeServiceProtocol +AcquisitionRuntimeService +``` + +а также расширен: + +```text +WebSocketSubscriptionManagerProtocol +``` + +для поддержки существующих команд подписки и отмены подписки. + +Runtime Service детерминированно маршрутизирует все шесть команд к соответствующим инфраструктурным зависимостям и не содержит знаний о рыночных данных. + +Архитектура сохраняет строгие границы: + +- Session управляет lifecycle; +- Transport отправляет payload; +- Subscription Manager управляет подписками; +- Runtime Service только координирует; +- Consistency и Recovery остаются независимыми. + +Все изменения подтверждены локальными и расширенными regression-тестами. + +--- + +# Итог Build + +После завершения Build 060.22 система обладает следующими возможностями. + +✓ Создан отдельный публичный `AcquisitionRuntimeServiceProtocol`. + +✓ Реализован `AcquisitionRuntimeService`. + +✓ Поддерживаются все шесть существующих Acquisition Runtime Commands. + +✓ Каждая команда имеет один детерминированный маршрут. + +✓ `WebSocketSubscriptionManagerProtocol` поддерживает subscribe и unsubscribe. + +✓ Text и Binary payload отправляются через существующий Transport API. + +✓ Исключения инфраструктурных зависимостей распространяются без изменения. + +✓ Event Publisher не используется преждевременно. + +✓ Все новые файлы имеют уникальные имена. + +✓ Runtime regression завершён результатом `46 passed`. + +✓ Расширенный regression завершён результатом `180 passed`. + +Build **060.22 — Acquisition Runtime Service** считается полностью завершённым и готовым к фиксации в Git. diff --git a/docs/migrations/build_060_22_architecture.md b/docs/migrations/build_060_22_architecture.md new file mode 100644 index 0000000..db47a62 --- /dev/null +++ b/docs/migrations/build_060_22_architecture.md @@ -0,0 +1,1424 @@ +# Build 060.22 — Acquisition Runtime Service + +**Статус:** Architecture Specification +**Build:** 060.22 +**Ветка:** Trades Feed (Time & Sales) +**Документ:** `build_060_22_architecture.md` +**Связанные документы:** + +- `build_060_21_architecture.md` — архитектурная спецификация Runtime Integration Contracts; +- `build_060_21.md` — Engineering Migration Report предыдущего Build; +- `build_060_20_1_architecture.md` — спецификация владения состоянием Trade Stream. + +--- + +# Назначение документа + +Настоящий документ является официальной архитектурной спецификацией **Build 060.22 — Acquisition Runtime Service**. + +Документ определяет создание первого исполняемого сервисного компонента Runtime Layer подсистемы **Market Data Acquisition**. + +Build 060.21 сформировал типизированную контрактную границу Runtime: + +```text +AcquisitionRuntimeCommand + +AcquisitionRuntimeEvent + +AcquisitionRuntimeCommandDispatcherProtocol + +AcquisitionRuntimeEventPublisherProtocol +``` + +Однако созданные контракты пока не имеют производственной реализации. + +В системе отсутствует компонент, который: + +- принимает типизированные Runtime-команды; +- маршрутизирует их к соответствующим инфраструктурным зависимостям; +- публикует результирующие Runtime-события; +- формирует единую точку исполнения команд Acquisition Runtime. + +Build 060.22 закрывает именно эту архитектурную задачу. + +Документ служит единственным источником истины при реализации данного Build. + +Изменение зафиксированных решений в процессе написания кода не допускается без отдельного архитектурного пересмотра. + +--- + +# Статус Build + +Build 060.22 является следующим этапом серии Build 060, посвящённой формированию полноценной подсистемы **Trades Feed (Time & Sales)**. + +К моменту начала данного Build в проекте уже существуют: + +- транспортные модели WebSocket; +- Runtime Commands; +- Runtime Events; +- WebSocket transport protocol; +- WebSocket session protocol; +- WebSocket subscription manager protocol; +- контракт диспетчеризации Runtime-команд; +- контракт публикации Runtime-событий; +- Canonical Trade Model; +- Trades Feed; +- Trade Stream Consistency; +- Trade Recovery; +- специализированное хранилище `TradeStreamStateStore`. + +При этом Runtime Layer пока состоит только из контрактов и моделей сообщений. + +Фактический сервис исполнения Runtime-команд отсутствует. + +Следовательно, Runtime ещё не является работающей подсистемой. + +Build 060.22 вводит минимальную исполняемую реализацию сервисного уровня, не подключая её пока к Trades Feed и корневому Acquisition Service. + +--- + +# Контекст + +После завершения Build 060.21 архитектура Runtime Protocol Layer выглядит следующим образом. + +```text +AcquisitionRuntimeCommand + │ + ▼ +AcquisitionRuntimeCommandDispatcherProtocol +``` + +и: + +```text +AcquisitionRuntimeEvent + │ + ▼ +AcquisitionRuntimeEventPublisherProtocol +``` + +Контракты определяют: + +- какие команды допустимы внутри Acquisition Runtime; +- какие события может публиковать Runtime; +- каким образом вызывающий компонент передаёт команду; +- каким образом инфраструктурный результат передаётся потребителям. + +Однако контракт сам по себе не выполняет команду. + +Например: + +```text +ConnectCommand +``` + +уже существует как типизированная команда, но в системе пока отсутствует компонент, который обязан: + +```text +ConnectCommand + │ + ▼ +WebSocketSessionProtocol.start() +``` + +Аналогично: + +```text +DisconnectCommand + │ + ▼ +WebSocketSessionProtocol.stop() +``` + +Для команд управления подписками и отправки сообщений также пока отсутствует единая точка исполнения. + +Таким образом между публичным Runtime Protocol Layer и низкоуровневыми WebSocket-контрактами остаётся незаполненная сервисная граница. + +--- + +# Предпосылки + +Настоящий Build опирается на архитектурные решения, принятые в предыдущих этапах. + +## Build 059.10 + +Сформированы базовые контракты WebSocket Runtime: + +```text +WebSocketTransportProtocol + +WebSocketSessionProtocol + +WebSocketSubscriptionManagerProtocol +``` + +Данные контракты описывают отдельные инфраструктурные возможности, но не объединяют их в единый сервис исполнения команд. + +--- + +## Build 059.11–059.12 + +Сформированы транспортные команды и события Runtime. + +Появились типизированные модели: + +```text +ConnectCommand + +DisconnectCommand + +SubscribeCommand + +UnsubscribeCommand + +SendTextCommand + +SendBinaryCommand +``` + +а также соответствующие инфраструктурные события. + +--- + +## Build 060.18 + +Создана подсистема Trade Stream Consistency. + +Она полностью независима от Runtime Transport и принимает только канонические объекты `Trade`. + +--- + +## Build 060.19 + +Создана подсистема Trade Recovery. + +Recovery остаётся stateless и зависит только от: + +```text +TradeStreamConsistencyProtocol +``` + +--- + +## Build 060.20.1 + +Владение состоянием Trade Stream перенесено в: + +```text +TradeStreamStateStore +``` + +Runtime Transport окончательно перестал рассматриваться как владелец состояния предметной области. + +--- + +## Build 060.21 + +Сформирован интеграционный контракт Runtime Protocol Layer. + +Добавлены: + +```text +AcquisitionRuntimeCommand + +AcquisitionRuntimeEvent + +AcquisitionRuntimeCommandDispatcherProtocol + +AcquisitionRuntimeEventPublisherProtocol +``` + +Имена были намеренно ограничены контекстом `Acquisition`, чтобы исключить конфликт с существующей общесистемной подсистемой: + +```text +src/runtime_events/ +``` + +Build 060.21 определил границы взаимодействия, но сознательно не создавал производственную реализацию. + +--- + +# Проблема + +В текущем состоянии каждая Runtime-команда является только неизменяемым объектом данных. + +Например: + +```python +ConnectCommand() +``` + +не содержит логики подключения. + +Она только выражает намерение вызывающего компонента. + +То же относится к: + +```python +DisconnectCommand() +SubscribeCommand(...) +UnsubscribeCommand(...) +SendTextCommand(...) +SendBinaryCommand(...) +``` + +Для выполнения команд необходим отдельный сервис, который: + +1. принимает объект команды; +2. определяет его конкретный тип; +3. выбирает соответствующую инфраструктурную зависимость; +4. выполняет строго определённую операцию; +5. не содержит знаний о Trade, Consistency или Recovery. + +Без такого компонента Runtime Protocol Layer остаётся декларативным и не может использоваться следующими уровнями системы. + +--- + +# Основная идея Build + +Главная идея Build 060.22 заключается в создании тонкого сервиса маршрутизации Runtime-команд. + +Новый компонент: + +```text +AcquisitionRuntimeService +``` + +становится производственной реализацией: + +```text +AcquisitionRuntimeCommandDispatcherProtocol +``` + +Сервис принимает типизированную команду и делегирует выполнение уже существующим инфраструктурным контрактам. + +Концептуально: + +```text +AcquisitionRuntimeCommand + │ + ▼ +AcquisitionRuntimeService + │ + ┌─────┴─────┐ + ▼ ▼ +WebSocket Subscription +Session Manager +``` + +Сервис не реализует WebSocket самостоятельно. + +Он не содержит сетевого клиента. + +Он не создаёт транспорт. + +Он только координирует существующие зависимости. + +--- + +# Цель Build + +Build обязан создать единственную производственную точку исполнения типизированных Acquisition Runtime Commands. + +После завершения этапа система должна обеспечивать: + +- приём всех команд из `AcquisitionRuntimeCommand`; +- детерминированную маршрутизацию каждой команды; +- делегирование lifecycle-команд WebSocket Session; +- делегирование управления подписками Subscription Manager; +- делегирование отправки транспортных сообщений соответствующей зависимости; +- соответствие `AcquisitionRuntimeCommandDispatcherProtocol`; +- независимое unit-тестирование без реального сетевого соединения. + +Build не должен изменять поведение существующих Runtime Commands и Runtime Events. + +--- + +# Что НЕ входит в Scope Build + +Настоящий Build сознательно не реализует: + +- интеграцию Runtime Service в корневой Acquisition Service; +- интеграцию Runtime Service с Trades Feed; +- Composition Root; +- реальный WebSocket client; +- Runtime Supervisor; +- Reconnect Controller; +- Heartbeat; +- Scheduler; +- автоматическое восстановление подписок; +- автоматический REST Recovery после reconnect; +- обработку входящих Trade-сообщений; +- публикацию Canonical Trade Stream; +- регистрацию Runtime Service в Registry; +- общесистемную шину событий; +- конкурентную очередь команд; +- фоновый worker; +- сохранение Runtime-состояния между перезапусками. + +Эти задачи относятся к Build 060.23–060.25 либо к отдельным последующим этапам. + +Build 060.22 отвечает только за синхронную с точки зрения порядка и асинхронную с точки зрения Python API диспетчеризацию одной команды за один вызов. + +--- + +# Архитектурные принципы + +При реализации Build 060.22 используются следующие обязательные принципы. + +## 1. Unique File Naming + +Во всём репозитории Dzentra запрещено создание нескольких файлов с одинаковым именем независимо от расположения в каталогах. + +Единственное разрешённое исключение: + +```text +__init__.py +``` + +Поэтому в Build не создаются файлы: + +```text +service.py +protocol.py +exceptions.py +test_service.py +``` + +Утверждённые уникальные имена: + +```text +acquisition_runtime_service.py + +acquisition_runtime_service_protocol.py + +test_acquisition_runtime_service.py +``` + +Имя каждого файла должно однозначно определять его назначение без учёта родительского каталога. + +--- + +## 2. Protocol Before Implementation + +Публичный контракт Runtime Service фиксируется отдельно от его реализации. + +Внешние потребители в последующих Build должны зависеть от: + +```text +AcquisitionRuntimeServiceProtocol +``` + +а не от конкретного класса: + +```text +AcquisitionRuntimeService +``` + +--- + +## 3. Thin Application Service + +Runtime Service является тонким координатором. + +Он: + +- принимает команду; +- определяет её тип; +- делегирует операцию; +- завершает вызов. + +Сервис не переносит в себя логику нижележащих компонентов. + +--- + +## 4. No Domain Knowledge + +Runtime Service не должен импортировать: + +```text +Trade +TradeStreamConsistencyProtocol +TradeRecoveryProtocol +TradesFeed +TradeStreamStateStore +``` + +Он работает исключительно с транспортными командами и инфраструктурными контрактами. + +--- + +## 5. Dependency Injection + +Все зависимости передаются в конструктор. + +Сервис не создаёт самостоятельно: + +- WebSocket Session; +- Transport; +- Subscription Manager; +- Publisher; +- Adapter. + +Создание графа зависимостей относится к Composition Root и не входит в Scope Build. + +--- + +## 6. One Command — One Deterministic Route + +Каждый тип команды должен иметь ровно один допустимый маршрут исполнения. + +Команда не может быть обработана несколькими зависимостями одновременно, если это отдельно не определено спецификацией. + +--- + +## 7. No Premature Orchestration + +Runtime Service не является Supervisor. + +Он не запускает фоновые задачи. + +Не выполняет retries. + +Не планирует reconnect. + +Не контролирует heartbeat. + +Перечисленные обязанности появятся в специализированных компонентах позднее. + +--- + +## 8. Existing Contracts Remain Stable + +Build не изменяет публичные модели: + +```text +runtime_commands.py +runtime_events.py +transport_messages.py +websocket_protocol.py +``` + +Если реализация обнаружит недостаточность существующего контракта, изменение должно быть отдельно согласовано до внесения в код. + +--- + +# Acquisition Runtime Service + +## Назначение + +`AcquisitionRuntimeService` является первым исполняемым компонентом Runtime Layer. + +Его единственная задача — выполнение типизированных инфраструктурных Runtime-команд. + +Сервис не содержит собственной бизнес-логики. + +Он выполняет только маршрутизацию команд к уже существующим инфраструктурным зависимостям. + +Концептуально: + +```text +Runtime Command + │ + ▼ +AcquisitionRuntimeService + │ + ├──────────────► WebSocketSession + │ + ├──────────────► WebSocketTransport + │ + └──────────────► SubscriptionManager +``` + +Таким образом сервис становится единственной производственной точкой исполнения Runtime-команд. + +--- + +# Почему появляется отдельный Service + +Во время проектирования были рассмотрены несколько вариантов. + +--- + +## Вариант №1 + +Каждый вызывающий компонент самостоятельно выполняет команды. + +Например: + +```text +TradesFeed + +↓ + +if ConnectCommand + +↓ + +session.start() +``` + +Отклонён. + +Причины: + +- дублирование логики; +- нарушение принципа единственной ответственности; +- отсутствие единой точки маршрутизации. + +--- + +## Вариант №2 + +Перенести выполнение команд внутрь WebSocketSession. + +Например: + +```text +session.dispatch(command) +``` + +Отклонён. + +Причины: + +Session начинает знать о: + +- Subscription; +- Transport; +- Runtime Commands. + +Тем самым нарушается разделение обязанностей. + +--- + +## Вариант №3 + +Создать специализированный Runtime Service. + +Принят. + +Именно Runtime Service становится единственным компонентом, который понимает соответствие между: + +```text +Runtime Command + +↓ + +Infrastructure Action +``` + +--- + +# Граница ответственности + +Runtime Service отвечает исключительно за выполнение Runtime-команд. + +Он НЕ отвечает за: + +- сетевой протокол; +- обработку сообщений; +- Consistency; +- Recovery; +- Trades Feed; +- управление жизненным циклом приложения; +- Supervisor; +- Scheduler; +- Heartbeat; +- Reconnect. + +Все перечисленные обязанности принадлежат другим Build. + +--- + +# AcquisitionRuntimeServiceProtocol + +## Назначение + +Build 060.22 вводит отдельный публичный контракт сервисного уровня. + +```text +AcquisitionRuntimeServiceProtocol +``` + +Он описывает единственную публичную возможность Runtime Service. + +```python +dispatch(...) +``` + +Все последующие Build должны зависеть именно от данного Protocol. + +Конкретная реализация может изменяться без влияния на потребителей. + +--- + +# Почему вводится отдельный Protocol + +На первый взгляд может показаться, что уже существует: + +```text +AcquisitionRuntimeCommandDispatcherProtocol +``` + +Однако данный Protocol описывает инфраструктурную возможность диспетчеризации команд. + +Он не определяет существование самостоятельного сервисного компонента. + +Build 060.22 впервые вводит именно сервисный уровень Runtime. + +Поэтому появляется отдельный сервисный контракт. + +Это позволяет в будущем заменить реализацию Runtime Service без изменения Composition Root и вышестоящих компонентов. + +--- + +# Публичный API + +Сервис предоставляет только один публичный метод. + +```python +async def dispatch( + command: AcquisitionRuntimeCommand, +) -> None +``` + +Других публичных методов Build 060.22 не вводит. + +В частности отсутствуют: + +```python +connect() + +disconnect() + +subscribe() + +unsubscribe() + +send() +``` + +Все операции выражаются исключительно через типизированные команды. + +--- + +# Почему отсутствуют отдельные методы + +Во время проектирования рассматривалась следующая модель. + +```python +runtime.connect() + +runtime.disconnect() + +runtime.subscribe(...) +``` + +Она была отклонена. + +Причина: + +типизированные Runtime Commands уже являются официальным языком взаимодействия Runtime Layer. + +Создание второго API привело бы к существованию двух независимых способов выполнения одной и той же операции. + +Архитектура должна содержать единственный публичный механизм. + +--- + +# Зависимости Runtime Service + +Сервис получает все зависимости через Dependency Injection. + +Минимальный набор зависимостей выглядит следующим образом. + +```text +AcquisitionRuntimeService + + │ + + ├────────► WebSocketSessionProtocol + + │ + + ├────────► WebSocketTransportProtocol + + │ + + ├────────► WebSocketSubscriptionManagerProtocol + + │ + + └────────► AcquisitionRuntimeEventPublisherProtocol +``` + +Никакие другие зависимости Build 060.22 не предусматривает. + +--- + +# Почему внедряется Event Publisher + +Хотя Build 060.22 ещё не реализует полноценную публикацию Runtime-событий, сервис уже принимает зависимость: + +```text +AcquisitionRuntimeEventPublisherProtocol +``` + +Это принципиальное архитектурное решение. + +Причины: + +- исключается изменение конструктора в следующих Build; +- Runtime Service сразу проектируется как источник инфраструктурных событий; +- интеграция публикации становится локальным изменением без перестройки графа зависимостей. + +На данном этапе Publisher допускается не использовать. + +--- + +# Почему внедряется WebSocketTransportProtocol + +Большинство текущих операций выполняются через: + +```text +WebSocketSessionProtocol +``` + +Однако команды: + +```text +SendTextCommand + +SendBinaryCommand +``` + +относятся к транспортному уровню. + +Поэтому Runtime Service сразу получает доступ к Transport. + +Это предотвращает последующее изменение конструктора. + +--- + +# Маршрутизация команд + +Главной обязанностью Runtime Service является детерминированная маршрутизация. + +Каждый тип команды имеет ровно один маршрут исполнения. + +--- + +## ConnectCommand + +```text +ConnectCommand + +↓ + +WebSocketSession.start() +``` + +После успешного завершения управление возвращается вызывающему компоненту. + +Сам Runtime Service не создаёт сетевое соединение. + +--- + +## DisconnectCommand + +```text +DisconnectCommand + +↓ + +WebSocketSession.stop() +``` + +Сервис не закрывает транспорт напрямую. + +Эта ответственность принадлежит Session. + +--- + +## SubscribeCommand + +```text +SubscribeCommand + +↓ + +WebSocketSubscriptionManager.subscribe() +``` + +Runtime Service не хранит список активных подписок. + +Он лишь делегирует выполнение соответствующему компоненту. + +--- + +## UnsubscribeCommand + +```text +UnsubscribeCommand + +↓ + +WebSocketSubscriptionManager.unsubscribe() +``` + +После завершения операции Runtime Service не изменяет собственного состояния. + +--- + +## SendTextCommand + +```text +SendTextCommand + +↓ + +WebSocketTransport.send_text() +``` + +Передаваемое сообщение считается полностью сформированным. + +Runtime Service не сериализует полезную нагрузку повторно. + +--- + +## SendBinaryCommand + +```text +SendBinaryCommand + +↓ + +WebSocketTransport.send_binary() +``` + +Команда делегируется транспортному уровню без каких-либо преобразований. + +--- + +# Официальная матрица маршрутизации + +| Runtime Command | Исполнитель | +|-----------------|-------------| +| ConnectCommand | WebSocketSessionProtocol | +| DisconnectCommand | WebSocketSessionProtocol | +| SubscribeCommand | WebSocketSubscriptionManagerProtocol | +| UnsubscribeCommand | WebSocketSubscriptionManagerProtocol | +| SendTextCommand | WebSocketTransportProtocol | +| SendBinaryCommand | WebSocketTransportProtocol | + +Данная таблица считается официальной спецификацией Build 060.22. + +Любое изменение маршрутов требует отдельного архитектурного решения (ADR). + +--- + +# Детерминированность маршрутизации + +Runtime Service рассматривается как детерминированный маршрутизатор. + +Для каждого входящего объекта существует единственный допустимый маршрут. + +Формально: + +```text +Runtime Command + +↓ + +Dispatch Table + +↓ + +Infrastructure Operation +``` + +Никакие дополнительные проверки, эвристики или выбор стратегии Build 060.22 не предусматривает. + +--- + +# Обработка неизвестной команды + +Build 060.22 не допускает существование неизвестных Runtime-команд. + +Если в сервис поступает объект, не входящий в объединение: + +```text +AcquisitionRuntimeCommand +``` + +это считается внутренней архитектурной ошибкой. + +В таком случае сервис обязан немедленно завершить выполнение исключением. + +Это гарантирует, что добавление новой команды никогда не останется незамеченным и потребует явного обновления таблицы маршрутизации. + +--- + +# Диаграмма взаимодействия + +После реализации Build 060.22 выполнение Runtime-команд приобретает следующий вид. + +```text +Caller + │ + │ dispatch(command) + ▼ +AcquisitionRuntimeService + │ + ├──────────────► WebSocketSessionProtocol + │ + ├──────────────► WebSocketTransportProtocol + │ + ├──────────────► WebSocketSubscriptionManagerProtocol + │ + └──────────────► AcquisitionRuntimeEventPublisherProtocol +``` + +Важно отметить, что Runtime Service остаётся полностью синхронным с точки зрения архитектуры. + +Он: + +- не создаёт фоновых задач; +- не ставит команды в очередь; +- не выполняет повторные попытки; +- не содержит собственного event loop. + +Каждый вызов `dispatch()` завершается только после завершения соответствующей инфраструктурной операции. + +--- + +# Жизненный цикл Runtime Service + +Runtime Service является долгоживущим инфраструктурным сервисом. + +Типичный жизненный цикл выглядит следующим образом. + +```text +Composition Root + +↓ + +создание зависимостей + +↓ + +создание Runtime Service + +↓ + +передача Runtime Service вызывающим компонентам + +↓ + +многократные вызовы dispatch() + +↓ + +завершение приложения +``` + +Сам Runtime Service не требует отдельной инициализации. + +Также отсутствуют методы: + +```python +start() + +stop() + +shutdown() + +dispose() +``` + +Экземпляр полностью готов к работе сразу после создания. + +--- + +# Отношение к состоянию + +Runtime Service является stateless-компонентом. + +Он не хранит: + +- состояние подключения; +- список подписок; +- очередь сообщений; +- информацию о последней выполненной команде; +- счётчики reconnect; +- состояние Recovery; +- состояние Consistency. + +Все перечисленные данные принадлежат специализированным компонентам. + +Следовательно экземпляр Runtime Service можно считать чистым координатором. + +--- + +# Обработка исключений + +Build 060.22 сознательно не вводит собственую иерархию исключений. + +Причина проста. + +Runtime Service не принимает самостоятельных решений. + +Он лишь вызывает инфраструктурные зависимости. + +Следовательно любые исключения должны передаваться вызывающему компоненту без изменения. + +Например: + +```text +Caller + +↓ + +Runtime Service + +↓ + +WebSocket Session + +↓ + +ConnectionError +``` + +Исключение должно пройти обратно без обёртки. + +Это сохраняет прозрачность поведения системы. + +--- + +# Почему исключения не преобразуются + +Во время проектирования рассматривался вариант создания: + +```text +RuntimeDispatchError +``` + +Он был отклонён. + +Причины: + +- потеря информации о первичном исключении; +- необходимость лишнего уровня обработки; +- усложнение диагностики; +- нарушение принципа прозрачной инфраструктуры. + +Runtime Service не должен скрывать происхождение ошибки. + +--- + +# Влияние на существующую архитектуру + +Build 060.22 практически не изменяет существующую систему. + +Новый сервис располагается между контрактами Runtime и инфраструктурными зависимостями. + +До Build: + +```text +Runtime Commands + +↓ + +(нет реализации) +``` + +После Build: + +```text +Runtime Commands + +↓ + +AcquisitionRuntimeService + +↓ + +Infrastructure +``` + +Никакие другие подсистемы не изменяются. + +--- + +# Влияние на Trade Stream Consistency + +Подсистема: + +```text +Trade Stream Consistency +``` + +не получает никаких новых зависимостей. + +Она продолжает работать исключительно с: + +```text +Trade +``` + +и + +```text +TradeStreamConsistencyProtocol +``` + +Build 060.22 не создаёт прямой связи между Runtime и Consistency. + +--- + +# Влияние на Trade Recovery + +Trade Recovery также не изменяется. + +Контроллер восстановления по-прежнему зависит только от: + +```text +TradeStreamConsistencyProtocol +``` + +Runtime Service не импортирует Recovery. + +Recovery не импортирует Runtime Service. + +Архитектурная независимость сохраняется полностью. + +--- + +# Влияние на Trades Feed + +На данном этапе Trades Feed ещё не использует Runtime Service. + +Интеграция будет выполнена отдельным Build: + +```text +060.23 +``` + +Это позволяет протестировать Runtime Service изолированно. + +--- + +# План изменения структуры проекта + +Build добавляет только два новых файла. + +```text +runtime/ + +├── acquisition_runtime_service.py + +└── acquisition_runtime_service_protocol.py +``` + +Также добавляется один файл тестов. + +```text +tests/ + +runtime/ + +└── test_acquisition_runtime_service.py +``` + +Другие файлы Runtime Build 060.22 не изменяет. + +--- + +# Почему используется отдельный Protocol + +Наличие собственного: + +```text +AcquisitionRuntimeServiceProtocol +``` + +позволяет в будущем заменить реализацию. + +Например: + +```text +Simple Runtime Service + +↓ + +Queued Runtime Service + +↓ + +Distributed Runtime Service +``` + +Все перечисленные варианты смогут реализовывать один и тот же публичный контракт. + +Поэтому вышестоящие компоненты не будут зависеть от конкретного класса. + +--- + +# Стратегия тестирования + +Build 060.22 покрывается исключительно Unit Test. + +Никакие интеграционные тесты пока не требуются. + +Каждая Runtime-команда проверяется отдельно. + +Минимальный набор сценариев включает: + +- соответствие `AcquisitionRuntimeServiceProtocol`; +- корректную маршрутизацию `ConnectCommand`; +- корректную маршрутизацию `DisconnectCommand`; +- корректную маршрутизацию `SubscribeCommand`; +- корректную маршрутизацию `UnsubscribeCommand`; +- корректную маршрутизацию `SendTextCommand`; +- корректную маршрутизацию `SendBinaryCommand`; +- отсутствие побочных эффектов; +- прозрачное распространение исключений; +- отсутствие собственного состояния между вызовами. + +Все зависимости заменяются простыми Fake-реализациями. + +Использование реального WebSocket Build 060.22 запрещает. + +--- + +# ADR (Architecture Decision Record) + +## ADR-060.22-01 + +**Решение** + +Создать специализированный сервис: + +```text +AcquisitionRuntimeService +``` + +реализующий выполнение всех типизированных Runtime-команд. + +**Статус** + +Accepted. + +--- + +## ADR-060.22-02 + +**Решение** + +Ввести отдельный публичный контракт: + +```text +AcquisitionRuntimeServiceProtocol +``` + +независимо от уже существующего: + +```text +AcquisitionRuntimeCommandDispatcherProtocol +``` + +**Статус** + +Accepted. + +--- + +## ADR-060.22-03 + +**Решение** + +Runtime Service остаётся stateless. + +**Статус** + +Accepted. + +--- + +## ADR-060.22-04 + +**Решение** + +Все зависимости передаются исключительно через Dependency Injection. + +**Статус** + +Accepted. + +--- + +## ADR-060.22-05 + +**Решение** + +Не создавать новых файлов с именами: + +```text +service.py + +protocol.py + +test_service.py +``` + +Использовать только уникальные имена файлов. + +**Статус** + +Accepted. + +--- + +# Definition of Done + +Build считается завершённым только при выполнении всех условий. + +Обязательно должны существовать: + +- `acquisition_runtime_service_protocol.py`; +- `acquisition_runtime_service.py`; +- `test_acquisition_runtime_service.py`. + +Runtime Service обязан: + +- реализовывать `AcquisitionRuntimeServiceProtocol`; +- принимать зависимости через конструктор; +- поддерживать все типы `AcquisitionRuntimeCommand`; +- выполнять детерминированную маршрутизацию; +- не хранить собственного состояния; +- не создавать зависимости самостоятельно; +- прозрачно распространять исключения. + +Все новые Unit Test должны успешно проходить. + +Build не должен изменять поведение существующих Runtime Commands, Runtime Events, Trade Recovery и Trade Stream Consistency. + +--- + +# Следующий Build + +После завершения Build 060.22 архитектура Runtime впервые получит полноценный исполняемый сервисный слой. + +Следующим этапом станет: + +```text +Build 060.23 + +Acquisition Integration +``` + +На этом этапе `AcquisitionRuntimeService` будет интегрирован в подсистему получения рыночных данных и станет использоваться как единая точка исполнения инфраструктурных Runtime-команд. + +После Build 060.23 Runtime перестанет существовать только как набор контрактов и сервисов и начнёт участвовать в реальном конвейере получения данных от биржи. \ No newline at end of file