From ee8765b7167f884ad38f771dab508268e210a535 Mon Sep 17 00:00:00 2001 From: Sergey Date: Tue, 28 Jul 2026 07:20:09 +0300 Subject: [PATCH] Build 060.23: integrate Trade Stream Acquisition --- .../trade_stream_acquisition_protocol.py | 74 + .../trade_stream_acquisition_service.py | 111 ++ .../trade_stream_message_adapter_protocol.py | 40 + ...gi_trade_stream_acquisition_integration.py | 117 ++ .../test_trade_stream_acquisition_service.py | 383 +++++ docs/migrations/build_060_23.md | 1419 +++++++++++++++++ docs/migrations/build_060_23_architecture.md | 845 ++++++++++ 7 files changed, 2989 insertions(+) create mode 100644 app/src/market_data/acquisition/trade_stream_acquisition_protocol.py create mode 100644 app/src/market_data/acquisition/trade_stream_acquisition_service.py create mode 100644 app/src/market_data/acquisition/trade_stream_message_adapter_protocol.py create mode 100644 app/tests/unit/market_data/acquisition/test_dzengi_trade_stream_acquisition_integration.py create mode 100644 app/tests/unit/market_data/acquisition/test_trade_stream_acquisition_service.py create mode 100644 docs/migrations/build_060_23.md create mode 100644 docs/migrations/build_060_23_architecture.md diff --git a/app/src/market_data/acquisition/trade_stream_acquisition_protocol.py b/app/src/market_data/acquisition/trade_stream_acquisition_protocol.py new file mode 100644 index 0000000..b29e143 --- /dev/null +++ b/app/src/market_data/acquisition/trade_stream_acquisition_protocol.py @@ -0,0 +1,74 @@ +# app/src/market_data/acquisition/trade_stream_acquisition_protocol.py + +from __future__ import annotations + +""" +Публичный контракт интеграции Trade Stream в Acquisition Layer. + +Build 060.23 вводит единую сервисную границу между: + +- Trade Subscription Layer; +- Acquisition Runtime Service; +- Unified WebSocket Adapter; +- Trade Stream Consistency. + +Контракт не определяет production WebSocket transport, receive-loop, +reconnect, recovery, хранение или публикацию канонических Trade. +""" + +from typing import Protocol, runtime_checkable + +from src.market_data.acquisition.models.trade import Trade + + +@runtime_checkable +class TradeStreamAcquisitionServiceProtocol(Protocol): + """ + Контракт сервиса интеграции Trade Stream. + + Сервис отвечает за: + + - передачу команды подписки в Acquisition Runtime; + - преобразование одного входящего WebSocket-документа; + - передачу канонического Trade в Consistency Layer; + - возврат согласованного Trade либо None для дубликата. + + Сервис не владеет WebSocket lifecycle и не запускает receive-loop. + """ + + async def subscribe( + self, + symbols: tuple[str, ...], + *, + correlation_id: str | None = None, + ) -> None: + """ + Передать в Acquisition Runtime команду подписки Trade Stream. + + Args: + symbols: + Символы торговых инструментов для подписки. + + correlation_id: + Необязательный идентификатор транспортного запроса. + При отсутствии значение создаётся Subscription Builder. + """ + ... + + def handle_message( + self, + document: object, + ) -> Trade | None: + """ + Обработать один входящий WebSocket-документ. + + Args: + document: + Сырой документ, полученный из WebSocket Runtime. + + Returns: + Канонический Trade после проверки согласованности; + None, если сообщение не относится к Trade либо является + корректным дубликатом. + """ + ... \ No newline at end of file diff --git a/app/src/market_data/acquisition/trade_stream_acquisition_service.py b/app/src/market_data/acquisition/trade_stream_acquisition_service.py new file mode 100644 index 0000000..394cbd0 --- /dev/null +++ b/app/src/market_data/acquisition/trade_stream_acquisition_service.py @@ -0,0 +1,111 @@ +# app/src/market_data/acquisition/trade_stream_acquisition_service.py + +from __future__ import annotations + +""" +Сервис интеграции Trade Stream с инфраструктурой Acquisition Runtime. + +Build 060.23 вводит сервисную границу между: + +- Trade Subscription Layer; +- Acquisition Runtime Service; +- Unified WebSocket Adapter; +- Trade Stream Consistency. + +На текущем этапе реализована передача команды подписки. +Обработка входящих сообщений будет добавлена следующим шагом Build. +""" + +from src.market_data.acquisition.trade_stream_message_adapter_protocol import ( + TradeStreamMessageAdapterProtocol, +) +from src.market_data.acquisition.consistency.trade_stream_protocol import ( + TradeStreamConsistencyProtocol, +) +from src.market_data.acquisition.models.trade import Trade +from src.market_data.acquisition.runtime.acquisition_runtime_service_protocol import ( + AcquisitionRuntimeServiceProtocol, +) +from src.market_data.acquisition.subscriptions.trades import ( + build_trade_subscribe_command, +) +from src.market_data.acquisition.trade_stream_acquisition_protocol import ( + TradeStreamAcquisitionServiceProtocol, +) + + +class TradeStreamAcquisitionService( + TradeStreamAcquisitionServiceProtocol, +): + """ + Координатор инфраструктуры получения Trade Stream. + + Сервис объединяет: + + - Acquisition Runtime; + - Unified WebSocket Adapter; + - Trade Stream Consistency. + + При этом сам сервис не реализует транспорт, + жизненный цикл WebSocket либо Recovery. + """ + + def __init__( + self, + runtime_service: AcquisitionRuntimeServiceProtocol, + adapter: TradeStreamMessageAdapterProtocol, + consistency_controller: TradeStreamConsistencyProtocol, + ) -> None: + self._runtime_service = runtime_service + self._adapter = adapter + self._consistency_controller = consistency_controller + + async def subscribe( + self, + symbols: tuple[str, ...], + *, + correlation_id: str | None = None, + ) -> None: + """ + Передать в Acquisition Runtime команду подписки Trade Stream. + + Args: + symbols: + Символы торговых инструментов для подписки. + + correlation_id: + Необязательный идентификатор транспортного запроса. + При отсутствии значение создаётся Subscription Builder. + """ + command = build_trade_subscribe_command( + symbols, + correlation_id=correlation_id, + ) + + await self._runtime_service.dispatch(command) + + def handle_message( + self, + document: object, + ) -> Trade | None: + """ + Обработать один входящий WebSocket-документ. + + Документ преобразуется через Unified Adapter. Только канонический + Trade передаётся в Trade Stream Consistency. + + Args: + document: + Сырой WebSocket-документ. + + Returns: + Trade после проверки согласованности; + None, если сообщение не относится к Trade либо является + корректным дубликатом. + """ + result = self._adapter.map_message(document) + + if not isinstance(result, Trade): + return None + + return self._consistency_controller.accept(result) diff --git a/app/src/market_data/acquisition/trade_stream_message_adapter_protocol.py b/app/src/market_data/acquisition/trade_stream_message_adapter_protocol.py new file mode 100644 index 0000000..391a9a0 --- /dev/null +++ b/app/src/market_data/acquisition/trade_stream_message_adapter_protocol.py @@ -0,0 +1,40 @@ +# app/src/market_data/acquisition/trade_stream_message_adapter_protocol.py + +from __future__ import annotations + +""" +Контракт адаптера входящих сообщений Trade Stream. + +Build 060.23 отделяет сервис координации Acquisition +от конкретной реализации WebSocket-протокола биржи. +""" + +from typing import Protocol, runtime_checkable + +from src.market_data.acquisition.models.candle_close import ( + CandleCloseEvent, +) +from src.market_data.acquisition.models.quote import Quote +from src.market_data.acquisition.models.trade import Trade + + +TradeStreamMappedMessage = Quote | CandleCloseEvent | Trade + + +@runtime_checkable +class TradeStreamMessageAdapterProtocol(Protocol): + """ + Контракт преобразования одного входящего WebSocket-документа. + + Реализация может быть exchange-specific, но вызывающий сервис + зависит только от результата преобразования. + """ + + def map_message( + self, + document: object, + ) -> TradeStreamMappedMessage: + """ + Преобразовать один WebSocket-документ в каноническую модель. + """ + ... \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/test_dzengi_trade_stream_acquisition_integration.py b/app/tests/unit/market_data/acquisition/test_dzengi_trade_stream_acquisition_integration.py new file mode 100644 index 0000000..33d1f3c --- /dev/null +++ b/app/tests/unit/market_data/acquisition/test_dzengi_trade_stream_acquisition_integration.py @@ -0,0 +1,117 @@ +# app/tests/unit/market_data/acquisition/test_dzengi_trade_stream_acquisition_integration.py + +from __future__ import annotations + +from datetime import datetime, timezone +from decimal import Decimal + +from src.market_data.acquisition.adapters.dzengi.websocket import ( + DzengiUnifiedWebSocketAdapter, +) +from src.market_data.acquisition.models.trade import ( + Trade, + TradeAggressorSide, +) +from src.market_data.acquisition.runtime.websocket_protocol import ( + AcquisitionRuntimeCommand, +) +from src.market_data.acquisition.trade_stream_acquisition_service import ( + TradeStreamAcquisitionService, +) + + +class FakeRuntimeService: + def __init__(self) -> None: + self.commands: list[AcquisitionRuntimeCommand] = [] + + async def dispatch( + self, + command: AcquisitionRuntimeCommand, + ) -> None: + self.commands.append(command) + + +class RecordingConsistencyController: + def __init__(self) -> None: + self.accepted_trades: list[Trade] = [] + + def accept( + self, + trade: Trade, + ) -> Trade | None: + self.accepted_trades.append(trade) + return trade + + +def _dzengi_trade_document() -> object: + return { + "status": "OK", + "destination": "internal.trade", + "payload": { + "id": 2134857062, + "price": "64497.25", + "size": "0.005", + "ts": 1784218066823, + "symbol": "BTC/USD_LEVERAGE", + "buyer": True, + "orderId": "order-123", + }, + } + + +def test_dzengi_trade_document_passes_complete_acquisition_pipeline() -> None: + runtime = FakeRuntimeService() + consistency = RecordingConsistencyController() + + service = TradeStreamAcquisitionService( + runtime_service=runtime, + adapter=DzengiUnifiedWebSocketAdapter(), + consistency_controller=consistency, + ) + + result = service.handle_message( + _dzengi_trade_document(), + ) + + assert isinstance(result, Trade) + assert consistency.accepted_trades == [result] + + assert result.symbol == "BTC/USD_LEVERAGE" + assert result.trade_id == 2134857062 + assert result.price == Decimal("64497.25") + assert result.quantity == Decimal("0.005") + assert result.executed_at == datetime.fromtimestamp( + 1784218066823 / 1000, + tz=timezone.utc, + ) + assert result.aggressor_side is TradeAggressorSide.BUY + assert result.source == "dzengi_websocket_trade" + + +def test_dzengi_duplicate_trade_is_rejected_by_real_consistency_pipeline() -> None: + from src.market_data.acquisition.consistency.trade_stream_consistency_controller import ( + TradeStreamConsistencyController, + ) + from src.market_data.acquisition.consistency.trade_stream_state_store import ( + TradeStreamStateStore, + ) + + runtime = FakeRuntimeService() + + consistency = TradeStreamConsistencyController( + state_store=TradeStreamStateStore(), + ) + + service = TradeStreamAcquisitionService( + runtime_service=runtime, + adapter=DzengiUnifiedWebSocketAdapter(), + consistency_controller=consistency, + ) + + document = _dzengi_trade_document() + + first_result = service.handle_message(document) + duplicate_result = service.handle_message(document) + + assert isinstance(first_result, Trade) + assert duplicate_result is None \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/test_trade_stream_acquisition_service.py b/app/tests/unit/market_data/acquisition/test_trade_stream_acquisition_service.py new file mode 100644 index 0000000..0676623 --- /dev/null +++ b/app/tests/unit/market_data/acquisition/test_trade_stream_acquisition_service.py @@ -0,0 +1,383 @@ +# app/tests/unit/market_data/acquisition/test_trade_stream_acquisition_service.py + +from __future__ import annotations + +import asyncio +import json +from datetime import datetime, timezone +from decimal import Decimal +from typing import cast + +import pytest + +from src.market_data.acquisition.models.candle_close import ( + CandleCloseEvent, +) +from src.market_data.acquisition.models.quote import Quote +from src.market_data.acquisition.models.trade import ( + Trade, + TradeAggressorSide, +) +from src.market_data.acquisition.runtime.runtime_commands import ( + SubscribeCommand, +) +from src.market_data.acquisition.runtime.transport_messages import ( + TransportTextMessage, +) +from src.market_data.acquisition.runtime.websocket_protocol import ( + AcquisitionRuntimeCommand, +) +from src.market_data.acquisition.trade_stream_acquisition_protocol import ( + TradeStreamAcquisitionServiceProtocol, +) +from src.market_data.acquisition.trade_stream_acquisition_service import ( + TradeStreamAcquisitionService, +) +from src.market_data.acquisition.trade_stream_message_adapter_protocol import ( + TradeStreamMappedMessage, +) + + +class FakeRuntimeService: + def __init__(self) -> None: + self.commands: list[AcquisitionRuntimeCommand] = [] + + async def dispatch( + self, + command: AcquisitionRuntimeCommand, + ) -> None: + self.commands.append(command) + + +class FakeMessageAdapter: + def __init__( + self, + result: TradeStreamMappedMessage, + *, + error: Exception | None = None, + ) -> None: + self._result = result + self._error = error + self.documents: list[object] = [] + + def map_message( + self, + document: object, + ) -> TradeStreamMappedMessage: + self.documents.append(document) + + if self._error is not None: + raise self._error + + return self._result + + +class FakeConsistencyController: + def __init__( + self, + result: Trade | None = None, + *, + use_input_trade: bool = True, + error: Exception | None = None, + ) -> None: + self._result = result + self._use_input_trade = use_input_trade + self._error = error + self.accepted_trades: list[Trade] = [] + + def accept( + self, + trade: Trade, + ) -> Trade | None: + self.accepted_trades.append(trade) + + if self._error is not None: + raise self._error + + if self._use_input_trade: + return trade + + return self._result + + +def _trade( + *, + trade_id: int = 2134857062, +) -> Trade: + return Trade( + symbol="BTC/USD_LEVERAGE", + trade_id=trade_id, + price=Decimal("64497.25"), + quantity=Decimal("0.005"), + executed_at=datetime( + 2026, + 7, + 16, + 11, + 27, + 46, + 823000, + tzinfo=timezone.utc, + ), + aggressor_side=TradeAggressorSide.BUY, + source="dzengi", + ) + + +def _non_trade_quote() -> Quote: + return cast( + Quote, + object.__new__(Quote), + ) + + +def _non_trade_candle() -> CandleCloseEvent: + return cast( + CandleCloseEvent, + object.__new__(CandleCloseEvent), + ) + + +def create_service( + *, + adapter_result: TradeStreamMappedMessage | None = None, + adapter_error: Exception | None = None, + consistency_result: Trade | None = None, + consistency_uses_input_trade: bool = True, + consistency_error: Exception | None = None, +) -> tuple[ + TradeStreamAcquisitionService, + FakeRuntimeService, + FakeMessageAdapter, + FakeConsistencyController, +]: + runtime = FakeRuntimeService() + + adapter = FakeMessageAdapter( + adapter_result if adapter_result is not None else _trade(), + error=adapter_error, + ) + + consistency = FakeConsistencyController( + consistency_result, + use_input_trade=consistency_uses_input_trade, + error=consistency_error, + ) + + service = TradeStreamAcquisitionService( + runtime_service=runtime, + adapter=adapter, + consistency_controller=consistency, + ) + + return ( + service, + runtime, + adapter, + consistency, + ) + + +def test_service_implements_protocol() -> None: + service, _, _, _ = create_service() + + assert isinstance( + service, + TradeStreamAcquisitionServiceProtocol, + ) + + +def test_subscribe_dispatches_subscribe_command() -> None: + service, runtime, _, _ = create_service() + + asyncio.run( + service.subscribe( + ( + "BTCUSDT", + "ETHUSDT", + ) + ) + ) + + assert len(runtime.commands) == 1 + + command = runtime.commands[0] + + assert isinstance( + command, + SubscribeCommand, + ) + + +def test_subscribe_preserves_symbols() -> None: + service, runtime, _, _ = create_service() + + asyncio.run( + service.subscribe( + ( + "BTCUSDT", + "ETHUSDT", + ) + ) + ) + + command = runtime.commands[0] + + assert isinstance( + command, + SubscribeCommand, + ) + + assert "BTCUSDT" in command.subscription_key + assert "ETHUSDT" in command.subscription_key + + +def test_subscribe_accepts_correlation_id() -> None: + service, runtime, _, _ = create_service() + + asyncio.run( + service.subscribe( + ( + "BTCUSDT", + ), + correlation_id="corr-123", + ) + ) + + command = runtime.commands[0] + + assert isinstance(command, SubscribeCommand) + assert isinstance(command.message, TransportTextMessage) + + document = json.loads(command.message.payload) + + assert document["correlationId"] == "corr-123" + + +def test_handle_message_passes_trade_to_consistency() -> None: + trade = _trade() + document = {"destination": "internal.trade"} + + service, _, adapter, consistency = create_service( + adapter_result=trade, + ) + + result = service.handle_message(document) + + assert adapter.documents == [document] + assert consistency.accepted_trades == [trade] + assert result is trade + + +def test_handle_message_preserves_consistency_result_identity() -> None: + adapted_trade = _trade(trade_id=100) + consistency_result = _trade(trade_id=101) + + service, _, _, consistency = create_service( + adapter_result=adapted_trade, + consistency_result=consistency_result, + consistency_uses_input_trade=False, + ) + + result = service.handle_message( + {"destination": "internal.trade"}, + ) + + assert consistency.accepted_trades == [adapted_trade] + assert result is consistency_result + + +def test_handle_message_returns_none_for_duplicate() -> None: + trade = _trade() + + service, _, _, consistency = create_service( + adapter_result=trade, + consistency_result=None, + consistency_uses_input_trade=False, + ) + + result = service.handle_message( + {"destination": "internal.trade"}, + ) + + assert consistency.accepted_trades == [trade] + assert result is None + + +def test_handle_message_ignores_quote() -> None: + quote = _non_trade_quote() + document = {"destination": "quote"} + + service, _, adapter, consistency = create_service( + adapter_result=quote, + ) + + result = service.handle_message(document) + + assert adapter.documents == [document] + assert consistency.accepted_trades == [] + assert result is None + + +def test_handle_message_ignores_candle_close_event() -> None: + candle = _non_trade_candle() + document = {"destination": "ohlc"} + + service, _, adapter, consistency = create_service( + adapter_result=candle, + ) + + result = service.handle_message(document) + + assert adapter.documents == [document] + assert consistency.accepted_trades == [] + assert result is None + + +def test_handle_message_calls_adapter_once() -> None: + document = {"destination": "internal.trade"} + + service, _, adapter, _ = create_service() + + service.handle_message(document) + + assert adapter.documents == [document] + + +def test_handle_message_propagates_adapter_error() -> None: + original = RuntimeError("adapter failed") + + service, _, adapter, consistency = create_service( + adapter_error=original, + ) + + with pytest.raises(RuntimeError, match="adapter failed") as exc_info: + service.handle_message( + {"destination": "internal.trade"}, + ) + + assert exc_info.value is original + assert len(adapter.documents) == 1 + assert consistency.accepted_trades == [] + + +def test_handle_message_propagates_consistency_error() -> None: + original = RuntimeError("consistency failed") + trade = _trade() + + service, _, adapter, consistency = create_service( + adapter_result=trade, + consistency_error=original, + ) + + with pytest.raises( + RuntimeError, + match="consistency failed", + ) as exc_info: + service.handle_message( + {"destination": "internal.trade"}, + ) + + assert exc_info.value is original + assert len(adapter.documents) == 1 + assert consistency.accepted_trades == [trade] \ No newline at end of file diff --git a/docs/migrations/build_060_23.md b/docs/migrations/build_060_23.md new file mode 100644 index 0000000..18a85cf --- /dev/null +++ b/docs/migrations/build_060_23.md @@ -0,0 +1,1419 @@ +# Build 060.23 — Trade Stream Acquisition Integration + +**Engineering Migration Report** + +--- + +# Контроль документа + +| Свойство | Значение | +|----------|----------| +| Build | 060.23 | +| Название | Trade Stream Acquisition Integration | +| Статус | Completed | +| Проект | Dzentra | +| Подсистема | Market Data Acquisition | +| Компонент | Trade Stream Acquisition Integration Layer | +| Версия | 1.0 | + +--- + +# Связанные документы + +- `build_060_23_architecture.md` — архитектурная спецификация Build. +- `build_060_22.md` — Engineering Migration Report предыдущего Build. +- `build_060_22_architecture.md` — спецификация Acquisition Runtime Service. +- `build_060_20_1.md` — Engineering Migration Report корректирующего Build владения состоянием Trade Stream. +- `build_060_19.md` — Engineering Migration Report Trade Recovery. +- `build_060_18.md` — Engineering Migration Report Trade Stream Consistency. + +--- + +# Цель Build + +Build 060.22 завершил создание минимального исполняемого сервисного слоя Acquisition Runtime. + +После предыдущего этапа в системе уже существовали: + +```text +AcquisitionRuntimeServiceProtocol +AcquisitionRuntimeService +AcquisitionRuntimeCommand +AcquisitionRuntimeEvent +``` + +Runtime Service уже умел детерминированно исполнять инфраструктурные команды: + +```text +ConnectCommand +DisconnectCommand +SubscribeCommand +UnsubscribeCommand +SendTextCommand +SendBinaryCommand +``` + +Однако Runtime по-прежнему оставался изолированным от Trade Stream Acquisition Pipeline. + +В системе отсутствовал отдельный компонент, который: + +- формировал команду подписки Trade Stream; +- передавал её в Acquisition Runtime; +- принимал один входящий WebSocket-документ; +- преобразовывал его через Unified Adapter; +- выделял канонический `Trade`; +- передавал `Trade` в Consistency Layer; +- возвращал согласованный результат `Trade | None`. + +Главной задачей Build 060.23 стало создание локального интеграционного слоя: + +```text +TradeStreamAcquisitionService +``` + +и его публичного контракта: + +```text +TradeStreamAcquisitionServiceProtocol +``` + +После завершения Build система получает: + +- единую сервисную точку подписки Trade Stream; +- единый путь обработки входящего Trade-сообщения; +- независимость сервиса от конкретной exchange-specific реализации адаптера; +- интеграцию с `TradeStreamConsistencyProtocol`; +- отдельное unit-покрытие координатора; +- интеграционные тесты полного локального Dzengi pipeline; +- подтверждённую регрессионную совместимость Runtime, Consistency, Recovery, Feeds и Dzengi Adapters. + +При этом Build сознательно не реализует production WebSocket transport, receive-loop, reconnect и Runtime Recovery. + +--- + +# Предпосылки + +К началу настоящего Build были доступны следующие подсистемы. + +## Runtime Layer + +```text +AcquisitionRuntimeServiceProtocol +AcquisitionRuntimeService +AcquisitionRuntimeCommandDispatcherProtocol +AcquisitionRuntimeEventPublisherProtocol +``` + +## Trade Subscription Layer + +```text +build_trade_subscribe_command() +``` + +## WebSocket Adapter Layer + +```text +DzengiUnifiedWebSocketAdapter +DzengiWebSocketTradeAdapter +``` + +## Consistency Layer + +```text +TradeStreamConsistencyProtocol +TradeStreamConsistencyController +TradeStreamStateStore +``` + +## Canonical Model + +```text +Trade +``` + +Каждый компонент уже был реализован и протестирован отдельно. + +Однако их взаимодействие не было оформлено в единый Acquisition Service. + +Архитектура до Build выглядела следующим образом: + +```text +Subscription Builder + │ + ▼ +SubscribeCommand + │ + ▼ +AcquisitionRuntimeService +``` + +и отдельно: + +```text +WebSocket Document + │ + ▼ +DzengiUnifiedWebSocketAdapter + │ + ▼ +Trade +``` + +и отдельно: + +```text +Trade + │ + ▼ +TradeStreamConsistencyController + │ + ▼ +Trade | None +``` + +Build 060.23 должен был соединить эти независимые части без изменения их существующего поведения. + +--- + +# Результаты архитектурного аудита + +Перед реализацией были проанализированы: + +```text +src/market_data/acquisition/protocol.py +src/market_data/acquisition/service.py +src/market_data/acquisition/registry.py +src/market_data/acquisition/feeds/trades_feed.py +src/market_data/acquisition/subscriptions/trades.py +src/market_data/acquisition/adapters/dzengi/websocket.py +src/market_data/acquisition/runtime/acquisition_runtime_service.py +src/market_data/acquisition/runtime/acquisition_runtime_service_protocol.py +src/market_data/acquisition/consistency/trade_stream_protocol.py +src/market_data/acquisition/consistency/trade_stream_consistency_controller.py +src/market_data/acquisition/consistency/trade_stream_state_store.py +``` + +Дополнительно были изучены legacy-компоненты: + +```text +src/integrations/exchange/market_stream.py +src/integrations/exchange/market_data_runner.py +src/integrations/exchange/ws_client.py +src/integrations/exchange/service.py +``` + +и соответствующие unit-тесты. + +Аудит подтвердил следующие выводы. + +--- + +## Существующий TradesFeed остаётся синхронным REST Feed + +Текущий `TradesFeed` работает по модели: + +```text +TradeDocumentSource + │ + ▼ +TradeDocumentHandler + │ + ▼ +Trade +``` + +Он не владеет: + +- WebSocket lifecycle; +- Runtime Commands; +- активными подписками; +- receive-loop; +- reconnect; +- Consistency State. + +Поэтому внедрение Runtime в существующий `TradesFeed` было признано нарушением Single Responsibility. + +Build 060.23 не изменяет `TradesFeed`. + +--- + +## Legacy MarketDataRunner не является точкой интеграции + +`MarketDataRunner` уже содержит собственные: + +- task lifecycle; +- reconnect-циклы; +- REST fallback; +- quote-specific runtime state; +- журналирование; +- status events. + +Подключение нового Trade Stream Acquisition Service внутрь него смешало бы новую Acquisition Architecture с существующим legacy runtime и преждевременно затронуло бы Scope Build 060.24. + +Поэтому `MarketDataRunner` не изменяется. + +--- + +## Legacy market_stream.py не является Composition Root + +`market_stream.py` представляет отдельный quote-specific WebSocket-контур. + +Его запуск в `main.py` не является активной точкой новой Trade Stream Architecture. + +Файл не изменяется и не используется как Composition Root Build 060.23. + +--- + +## ExchangeWebSocketClient не реализует новые Runtime Protocol + +`ExchangeWebSocketClient` является специализированным legacy-клиентом для depth stream. + +Он: + +- самостоятельно открывает соединение; +- самостоятельно отправляет depth request; +- самостоятельно выполняет ping; +- содержит внутренний цикл; +- возвращает raw document; +- жёстко связан со стаканом. + +Следовательно он не является реализацией: + +```text +WebSocketTransportProtocol +WebSocketSessionProtocol +WebSocketSubscriptionManagerProtocol +``` + +и не адаптируется в рамках Build 060.23. + +--- + +## Trade consumer и Trade store отсутствуют + +Repository-wide поиск не выявил утверждённых компонентов: + +```text +TradeStore +TradeConsumer +publish_trade() +on_trade() +consume_trade() +``` + +Поэтому новый сервис не публикует и не сохраняет Trade. + +Его результатом остаётся: + +```text +Trade | None +``` + +--- + +## trades.unsubscribe не поддерживается биржей + +Предыдущий инженерный аудит Dzengi WebSocket API подтвердил отсутствие поддерживаемой операции: + +```text +trades.unsubscribe +``` + +Поэтому Build 060.23 не создаёт фиктивный exchange-specific unsubscribe builder. + +Завершение Trade Stream в будущем должно выполняться через lifecycle соединения либо иной подтверждённый механизм. + +--- + +# Архитектурное решение + +По результатам аудита утверждена следующая схема. + +```text +symbols + │ + ▼ +TradeStreamAcquisitionService.subscribe() + │ + ▼ +build_trade_subscribe_command() + │ + ▼ +SubscribeCommand + │ + ▼ +AcquisitionRuntimeServiceProtocol.dispatch() +``` + +Входящий поток: + +```text +raw WebSocket document + │ + ▼ +TradeStreamMessageAdapterProtocol.map_message() + │ + ▼ +Quote | CandleCloseEvent | Trade + │ + ▼ +isinstance(result, Trade) + │ + ├── нет ──► None + │ + ▼ +TradeStreamConsistencyProtocol.accept() + │ + ▼ +Trade | None +``` + +Новый сервис является stateless-координатором. + +Он не создаёт транспорт, не запускает фоновые задачи и не хранит состояние Trade Stream. + +--- + +# Правило уникальности имён файлов + +Build 060.23 сохраняет обязательное правило проекта Dzentra: + +> Новые файлы должны иметь уникальные имена в масштабе репозитория независимо от каталога. + +Build не создаёт новые общие файлы: + +```text +service.py +protocol.py +test_service.py +``` + +Вместо них добавлены уникальные имена: + +```text +trade_stream_acquisition_protocol.py +trade_stream_acquisition_service.py +trade_stream_message_adapter_protocol.py +test_trade_stream_acquisition_service.py +test_dzengi_trade_stream_acquisition_integration.py +``` + +Имена однозначно отражают назначение файлов и не требуют знания родительского каталога. + +--- + +# TradeStreamAcquisitionServiceProtocol + +В рамках Build создан новый публичный контракт: + +```text +TradeStreamAcquisitionServiceProtocol +``` + +Файл: + +```text +src/market_data/acquisition/ +trade_stream_acquisition_protocol.py +``` + +Protocol определяет две операции. + +## subscribe() + +```python +async def subscribe( + symbols: tuple[str, ...], + *, + correlation_id: str | None = None, +) -> None +``` + +Метод отвечает только за формирование и передачу команды подписки в Runtime. + +## handle_message() + +```python +def handle_message( + document: object, +) -> Trade | None +``` + +Метод отвечает только за обработку одного входящего документа. + +Другие публичные операции не вводятся. + +В частности отсутствуют: + +```text +start() +stop() +run() +receive() +reconnect() +recover() +publish() +``` + +--- + +# TradeStreamMessageAdapterProtocol + +Во время реализации было выявлено, что первоначальная зависимость сервиса от конкретного: + +```text +DzengiUnifiedWebSocketAdapter +``` + +создавала прямую exchange-specific связь внутри координационного слоя. + +Для устранения этой связи введён отдельный контракт: + +```text +TradeStreamMessageAdapterProtocol +``` + +Файл: + +```text +src/market_data/acquisition/ +trade_stream_message_adapter_protocol.py +``` + +Protocol определяет одну операцию: + +```python +def map_message( + document: object, +) -> TradeStreamMappedMessage +``` + +Тип результата: + +```text +TradeStreamMappedMessage +``` + +объединяет: + +```text +Quote +CandleCloseEvent +Trade +``` + +`DzengiUnifiedWebSocketAdapter` структурно удовлетворяет этому Protocol и не требует изменения. + +Благодаря этому `TradeStreamAcquisitionService` зависит от абстракции, а конкретный Dzengi Adapter подключается только при композиции и в интеграционных тестах. + +--- + +# TradeStreamAcquisitionService + +Главным production-компонентом Build является: + +```text +TradeStreamAcquisitionService +``` + +Файл: + +```text +src/market_data/acquisition/ +trade_stream_acquisition_service.py +``` + +Сервис реализует: + +```text +TradeStreamAcquisitionServiceProtocol +``` + +и координирует три зависимости: + +```text +AcquisitionRuntimeServiceProtocol +TradeStreamMessageAdapterProtocol +TradeStreamConsistencyProtocol +``` + +Сервис не создаёт зависимости самостоятельно. + +--- + +# Зависимости сервиса + +Конструктор принимает: + +```text +runtime_service +adapter +consistency_controller +``` + +Полная схема: + +```text +TradeStreamAcquisitionService + │ + ├── AcquisitionRuntimeServiceProtocol + ├── TradeStreamMessageAdapterProtocol + └── TradeStreamConsistencyProtocol +``` + +Это обеспечивает: + +- независимость от конкретного Runtime implementation; +- независимость от конкретной биржи; +- независимость от конкретной Consistency implementation; +- возможность изолированного unit-тестирования. + +--- + +# subscribe() + +Метод `subscribe()` использует существующий builder: + +```text +build_trade_subscribe_command() +``` + +Последовательность: + +```text +symbols + │ + ▼ +build_trade_subscribe_command( + symbols, + correlation_id=... +) + │ + ▼ +SubscribeCommand + │ + ▼ +runtime_service.dispatch(command) +``` + +Сервис не: + +- сериализует JSON; +- создаёт TransportTextMessage вручную; +- строит subscription key самостоятельно; +- нормализует символы; +- хранит активные подписки; +- выполняет connect; +- выполняет retry. + +--- + +# correlation_id + +Параметр: + +```text +correlation_id +``` + +передаётся без изменения существующему Subscription Builder. + +Если он не задан, его создание остаётся ответственностью builder. + +Сервис не генерирует correlation id самостоятельно. + +--- + +# handle_message() + +Метод `handle_message()` обрабатывает один входящий документ. + +Последовательность: + +```text +document + │ + ▼ +adapter.map_message(document) + │ + ▼ +mapped result + │ + ▼ +isinstance(mapped result, Trade) +``` + +Если результат не является `Trade`, сервис возвращает: + +```text +None +``` + +Если результат является `Trade`, сервис вызывает: + +```text +consistency_controller.accept(trade) +``` + +и возвращает результат без изменения. + +--- + +# Поведение для Quote + +Если Adapter возвращает: + +```text +Quote +``` + +сервис: + +- не вызывает Consistency; +- не преобразует Quote; +- не публикует Quote; +- возвращает `None`. + +--- + +# Поведение для CandleCloseEvent + +Если Adapter возвращает: + +```text +CandleCloseEvent +``` + +сервис: + +- не вызывает Consistency; +- не преобразует CandleCloseEvent; +- не публикует событие; +- возвращает `None`. + +--- + +# Поведение для Trade + +Если Adapter возвращает: + +```text +Trade +``` + +сервис передаёт тот же экземпляр в: + +```text +TradeStreamConsistencyProtocol.accept() +``` + +Сервис не копирует и не изменяет модель. + +--- + +# Поведение для дубликата + +Если Consistency Layer возвращает: + +```text +None +``` + +это означает корректный дубликат. + +Сервис сохраняет результат и также возвращает: + +```text +None +``` + +Дополнительная логика не выполняется. + +--- + +# Сохранение идентичности результата Consistency + +Если `TradeStreamConsistencyProtocol.accept()` возвращает отдельный экземпляр `Trade`, сервис возвращает именно этот объект. + +Сервис не заменяет его исходным результатом Adapter и не создаёт копию. + +Данный инвариант подтверждён unit-тестом. + +--- + +# Обработка исключений + +Build не вводит отдельную иерархию ошибок Trade Stream Acquisition Service. + +Исключения распространяются без обёртки. + +## Ошибка Subscription Builder + +Ошибка формирования SubscribeCommand передаётся вызывающей стороне. + +## Ошибка Runtime dispatch + +Ошибка `runtime_service.dispatch()` передаётся вызывающей стороне. + +## Ошибка Adapter + +Schema, parser, value validation или mapper error распространяется без изменения. + +## Ошибка Consistency + +`TradeConsistencyError`, `TradeOrderingError` и иные исключения Consistency Layer распространяются без изменения. + +Сервис не создаёт общий `TradeStreamAcquisitionError`. + +--- + +# Stateless-архитектура + +`TradeStreamAcquisitionService` не хранит операционное состояние. + +В объекте сохраняются только внедрённые зависимости. + +Сервис не хранит: + +- symbols; +- correlation id; +- active subscriptions; +- last trade; +- duplicate window; +- reconnect attempts; +- recovery state; +- runtime status; +- received documents. + +Состояние Consistency принадлежит `TradeStreamStateStore`. + +Состояние Runtime будет принадлежать будущим Session и Subscription Manager. + +--- + +# Изменённые и добавленные файлы + +## Новый Acquisition Service Protocol + +```text +src/market_data/acquisition/ +trade_stream_acquisition_protocol.py +``` + +Добавлен: + +```text +TradeStreamAcquisitionServiceProtocol +``` + +--- + +## Новый Trade Stream Message Adapter Protocol + +```text +src/market_data/acquisition/ +trade_stream_message_adapter_protocol.py +``` + +Добавлены: + +```text +TradeStreamMappedMessage +TradeStreamMessageAdapterProtocol +``` + +--- + +## Новый Trade Stream Acquisition Service + +```text +src/market_data/acquisition/ +trade_stream_acquisition_service.py +``` + +Добавлен: + +```text +TradeStreamAcquisitionService +``` + +--- + +## Unit-тест Acquisition Service + +```text +tests/unit/market_data/acquisition/ +test_trade_stream_acquisition_service.py +``` + +Добавлены Fake-реализации: + +- `FakeRuntimeService`; +- `FakeMessageAdapter`; +- `FakeConsistencyController`. + +Проверены subscribe и handle_message без привязки к Dzengi. + +--- + +## Интеграционный тест Dzengi pipeline + +```text +tests/unit/market_data/acquisition/ +test_dzengi_trade_stream_acquisition_integration.py +``` + +Проверен реальный локальный pipeline: + +```text +Dzengi WebSocket document + │ + ▼ +DzengiUnifiedWebSocketAdapter + │ + ▼ +Canonical Trade + │ + ▼ +TradeStreamConsistencyController + │ + ▼ +Trade | None +``` + +--- + +## Архитектурная документация + +```text +docs/migrations/ +build_060_23_architecture.md +``` + +Документ фиксирует: + +- границы Build; +- ответственность сервиса; +- dependency direction; +- ADR; +- Definition of Done; +- связь с последующими этапами. + +--- + +# Файлы, которые не изменялись + +Build не изменяет: + +```text +src/market_data/acquisition/feeds/trades_feed.py +src/market_data/acquisition/subscriptions/trades.py +src/market_data/acquisition/adapters/dzengi/websocket.py +src/market_data/acquisition/runtime/acquisition_runtime_service.py +src/market_data/acquisition/runtime/acquisition_runtime_service_protocol.py +src/market_data/acquisition/consistency/trade_stream_protocol.py +src/market_data/acquisition/consistency/trade_stream_consistency_controller.py +src/market_data/acquisition/consistency/trade_stream_state_store.py +src/market_data/acquisition/recovery/trade_recovery_controller.py +``` + +Также не изменяются: + +```text +src/integrations/exchange/market_stream.py +src/integrations/exchange/market_data_runner.py +src/integrations/exchange/ws_client.py +src/integrations/exchange/service.py +``` + +Build не затрагивает legacy quote runtime. + +--- + +# Unit-тестирование TradeStreamAcquisitionService + +Новый сервис покрыт отдельным набором unit-тестов. + +Проверены следующие сценарии: + +- сервис соответствует `TradeStreamAcquisitionServiceProtocol`; +- `subscribe()` передаёт одну команду Runtime; +- передаётся именно `SubscribeCommand`; +- symbols сохраняются в сформированной команде; +- `correlation_id` сохраняется в payload; +- Adapter вызывается один раз; +- Trade передаётся в Consistency; +- возвращается тот же объект, который вернула Consistency; +- `None` для корректного дубликата сохраняется; +- Quote игнорируется; +- CandleCloseEvent игнорируется; +- Consistency не вызывается для не-Trade; +- ошибка Adapter распространяется без изменения; +- ошибка Consistency распространяется без изменения. + +Результат совместного локального прогона unit и integration тестов: + +```text +14 passed +``` + +--- + +# Интеграционное тестирование Dzengi pipeline + +Интеграционный тест использует настоящий: + +```text +DzengiUnifiedWebSocketAdapter +``` + +и документ реального WebSocket Trade-контракта: + +```text +status +destination +payload +``` + +Проверяется полный путь: + +```text +schema validation + │ + ▼ +parser + │ + ▼ +value validation + │ + ▼ +mapper + │ + ▼ +Trade + │ + ▼ +Consistency +``` + +Подтверждены канонические поля: + +```text +symbol +trade_id +price +quantity +executed_at +aggressor_side +source +``` + +Для WebSocket Trade подтверждён источник: + +```text +dzengi_websocket_trade +``` + +Также реальный `TradeStreamConsistencyController` подтвердил, что повторная обработка идентичного Trade возвращает: + +```text +None +``` + +--- + +# Расширенное регрессионное тестирование + +После завершения реализации выполнен совместный прогон: + +```text +Runtime +Consistency +Recovery +Feeds +Dzengi Adapters +Trade Stream Acquisition Unit Tests +Dzengi Trade Stream Acquisition Integration Tests +``` + +Охвачены каталоги: + +```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 +tests/unit/market_data/acquisition/adapters/dzengi +``` + +и новые файлы: + +```text +tests/unit/market_data/acquisition/ +test_trade_stream_acquisition_service.py + +tests/unit/market_data/acquisition/ +test_dzengi_trade_stream_acquisition_integration.py +``` + +Результат: + +```text +496 passed +``` + +Тем самым подтверждено: + +- Runtime Service не нарушен; +- Runtime Protocol Layer не нарушен; +- Trade Stream Consistency не нарушена; +- TradeStreamStateStore не нарушен; +- Trade Recovery не нарушена; +- существующие Feeds не нарушены; +- REST Trade Adapter не нарушен; +- WebSocket Quote Adapter не нарушен; +- WebSocket OHLC Adapter не нарушен; +- WebSocket Trade Adapter не нарушен; +- Unified WebSocket Adapter не нарушен; +- новый Acquisition Service корректно интегрируется с реальным Dzengi Adapter; +- duplicate filtering работает через реальный Consistency Controller. + +--- + +# Проверка качества Git diff + +После реализации выполнены: + +```text +git diff --stat +git diff --check +``` + +`git diff --check` не выявил ошибок форматирования в отслеживаемых изменениях. + +Новые production и test-файлы на момент проверки оставались untracked, поэтому не отображались в `git diff --stat`. + +Изменение: + +```text +.gitignore +``` + +не относится к Scope Build 060.23 и не должно включаться в коммит этапа. + +--- + +# Обратная совместимость + +Build сохраняет существующее поведение всех ранее реализованных подсистем. + +Не изменились: + +- `TradesFeedProtocol`; +- `TradesFeed`; +- `TradeDocumentSource`; +- `TradeDocumentHandler`; +- `Trade` model; +- Runtime Commands; +- Runtime Events; +- Acquisition Runtime Service; +- Trade Subscription Builder; +- Unified Dzengi Adapter; +- Trade Stream Consistency API; +- Trade Recovery API; +- Feed Registry; +- legacy Exchange Runtime. + +Новый сервис добавляется как отдельный слой и не заменяет существующие REST Feed. + +--- + +# Производительность + +## subscribe() + +Метод выполняет: + +- один вызов Subscription Builder; +- один асинхронный вызов Runtime dispatch. + +Алгоритмическая сложность зависит от существующего builder и количества symbols. + +Сервис не создаёт дополнительных копий команд и payload. + +## handle_message() + +Метод выполняет: + +- один вызов Adapter; +- одну `isinstance`-проверку; +- не более одного вызова Consistency. + +Координационные накладные расходы постоянны: + +```text +O(1) +``` + +без учёта внутренней стоимости Adapter и Consistency Layer. + +Сервис не вводит: + +- очереди; +- блокировки; +- фоновые задачи; +- retries; +- дополнительную сериализацию; +- storage writes. + +--- + +# Подтверждённые архитектурные инварианты + +## Runtime не знает о Trade + +Новые зависимости направлены из Acquisition Service в Runtime Protocol, а не обратно. + +--- + +## TradeStreamAcquisitionService не зависит от Dzengi + +Production-сервис зависит от: + +```text +TradeStreamMessageAdapterProtocol +``` + +а не от конкретного `DzengiUnifiedWebSocketAdapter`. + +--- + +## Dzengi Adapter проверяется отдельно + +Exchange-specific корректность подтверждена интеграционным тестом. + +--- + +## Consistency остаётся единственной точкой проверки последовательности + +Сервис не реализует duplicate detection или ordering logic самостоятельно. + +--- + +## TradeStreamStateStore остаётся единственным владельцем Consistency State + +Новый сервис не хранит состояния по symbols. + +--- + +## Существующий TradesFeed не изменяется + +REST Feed и WebSocket Acquisition Integration остаются отдельными путями. + +--- + +## Сервис не публикует Trade + +До появления утверждённого Trade Consumer результат возвращается вызывающей стороне. + +--- + +## Сервис не владеет WebSocket lifecycle + +Он не реализует connect, disconnect, receive-loop или reconnect. + +--- + +## trades.unsubscribe не создаётся + +Build не вводит неподдерживаемую exchange-команду. + +--- + +## Уникальность новых имён файлов соблюдена + +Все новые production и test-файлы имеют уникальные имена. + +--- + +# Что не входит в Scope Build + +Build сознательно не реализует: + +- production `WebSocketTransportProtocol`; +- production `WebSocketSessionProtocol`; +- production `WebSocketSubscriptionManagerProtocol`; +- Runtime Event Publisher implementation; +- Runtime Supervisor; +- receive-loop; +- reconnect; +- heartbeat; +- scheduler; +- resubscribe orchestration; +- автоматический вызов Trade Recovery; +- восстановление Trade Stream после разрыва; +- Trade Store; +- Trade Consumer; +- публикацию Trade; +- root Composition; +- изменение `MarketDataRunner`; +- изменение `market_stream.py`; +- изменение `ExchangeWebSocketClient`; +- изменение существующего `TradesFeed`; +- биржевую команду `trades.unsubscribe`. + +Отсутствие перечисленных компонентов является осознанной границей этапа. + +--- + +# Архитектурное значение Build + +Build 060.23 впервые соединяет ранее независимые части Trades Acquisition Architecture: + +```text +Subscription Builder +Runtime Service +Unified Adapter +Canonical Trade +Consistency Controller +``` + +До Build каждый компонент существовал отдельно. + +После завершения Build появляется единая сервисная граница: + +```text +TradeStreamAcquisitionService +``` + +которая поддерживает оба направления интеграции: + +```text +outbound subscription command +``` + +и: + +```text +inbound Trade message processing +``` + +При этом сервис остаётся stateless, exchange-independent и не вмешивается в lifecycle Runtime. + +--- + +# Связь с последующими Build + +## Build 060.24 — Reconnect & Runtime Recovery + +На основе текущего сервиса должны быть построены: + +- Runtime reconnect orchestration; +- автоматическое восстановление подписок; +- восстановление Trade Stream continuity; +- координация Consistency и Recovery после разрыва; +- lifecycle соединения. + +--- + +## Build 060.25 — Integration & Regression + +Будет выполнена полная проверка production graph: + +```text +Runtime +Trade Stream Acquisition +Consistency +Recovery +Trades Feed +``` + +после появления реальных Runtime implementations. + +--- + +## Build 060.26 — Final Documentation + +Будет подготовлена итоговая документация всей ветки Trades Feed Runtime. + +--- + +# Заключение + +Build 060.23 завершает service-level интеграцию Trade Stream с Acquisition Runtime. + +В рамках этапа реализованы: + +```text +TradeStreamAcquisitionServiceProtocol +TradeStreamMessageAdapterProtocol +TradeStreamAcquisitionService +``` + +Новый сервис: + +- формирует Trade SubscribeCommand через существующий builder; +- передаёт команду в `AcquisitionRuntimeServiceProtocol`; +- обрабатывает один WebSocket-документ через Adapter Protocol; +- пропускает только канонический `Trade`; +- делегирует проверку последовательности Consistency Layer; +- возвращает `Trade | None`; +- не хранит состояние; +- не содержит exchange-specific логики; +- не реализует reconnect и Recovery. + +Unit-тесты подтверждают поведение координатора независимо от биржи. + +Интеграционные тесты подтверждают полный локальный Dzengi pipeline и duplicate filtering через реальный Consistency Controller. + +Расширенная регрессия завершена результатом: + +```text +496 passed +``` + +--- + +# Итог Build + +После завершения Build 060.23 система обладает следующими возможностями. + +✓ Создан `TradeStreamAcquisitionServiceProtocol`. + +✓ Создан `TradeStreamMessageAdapterProtocol`. + +✓ Реализован `TradeStreamAcquisitionService`. + +✓ Подписка Trade Stream проходит через существующий Subscription Builder. + +✓ SubscribeCommand передаётся через `AcquisitionRuntimeServiceProtocol`. + +✓ Входящие сообщения проходят через Adapter Protocol. + +✓ `DzengiUnifiedWebSocketAdapter` структурно совместим с новым Protocol. + +✓ Только `Trade` передаётся в Consistency. + +✓ Quote и CandleCloseEvent игнорируются без побочных эффектов. + +✓ Дубликат возвращается как `None`. + +✓ Исключения Adapter и Consistency распространяются без изменения. + +✓ Проверена идентичность результата Consistency. + +✓ Реальный Dzengi WebSocket Trade document проходит полный pipeline. + +✓ Реальный Consistency Controller отклоняет идентичный дубликат. + +✓ Существующий `TradesFeed` не изменён. + +✓ Legacy Exchange Runtime не изменён. + +✓ Все новые файлы имеют уникальные имена. + +✓ Локальный функциональный прогон завершён результатом `14 passed`. + +✓ Расширенный regression завершён результатом `496 passed`. + +Build **060.23 — Trade Stream Acquisition Integration** считается полностью завершённым и готовым к фиксации в Git. diff --git a/docs/migrations/build_060_23_architecture.md b/docs/migrations/build_060_23_architecture.md new file mode 100644 index 0000000..6a4b89e --- /dev/null +++ b/docs/migrations/build_060_23_architecture.md @@ -0,0 +1,845 @@ +# Build 060.23 — Trade Stream Acquisition Integration Architecture + +**Статус:** Architecture Approved + +**Build:** 060.23 + +**Название:** +Trade Stream Acquisition Integration + +--- + +# 1. Назначение Build + +Build 060.23 завершает построение инфраструктурного слоя Acquisition для Trades Feed. + +После Build 060.22 Runtime уже обладает собственной сервисной моделью управления инфраструктурными командами: + +- Runtime Commands; +- Runtime Events; +- Runtime Service; +- Runtime Protocols. + +Однако в настоящий момент Runtime полностью изолирован от существующего Acquisition Layer. + +Trades Feed продолжает работать как синхронный REST Feed: + +``` +REST document + ↓ +TradesHandler + ↓ +Trade +``` + +а Runtime существует отдельно: + +``` +Runtime Commands +Runtime Events +Runtime Service +``` + +В результате отсутствует единая точка, которая соединяет: + +- Runtime; +- Subscription API; +- WebSocket Adapter; +- Consistency Layer. + +Именно эту задачу решает Build 060.23. + +--- + +# 2. Цели Build + +Build вводит единый сервис интеграции Trade Stream. + +Новый сервис становится входной точкой всей Runtime Acquisition Pipeline. + +После завершения Build поток данных приобретает следующий вид: + +``` +Subscribe() + + │ + + ▼ + +Subscription Builder + + │ + + ▼ + +Acquisition Runtime Service + + │ + + ▼ + +WebSocket Runtime + + │ + + ▼ + +Unified Adapter + + │ + + ▼ + +Trade + + │ + + ▼ + +Consistency Controller + + │ + + ▼ + +Trade | None +``` + +При этом Build намеренно не реализует: + +- production WebSocket; +- reconnect; +- heartbeat; +- receive loop; +- runtime scheduler; +- recovery. + +Все перечисленные компоненты остаются предметом следующих Build. + +--- + +# 3. Архитектурная цель + +Главная задача Build — + +полностью отделить инфраструктуру Runtime +от обработки Trade. + +После Build Runtime больше не знает: + +- что такое Trade; +- что такое Feed; +- что такое Consistency; +- что такое Recovery. + +Runtime работает исключительно с инфраструктурными командами. + +Вся логика Trade переносится в отдельный слой Acquisition Integration. + +--- + +# 4. Архитектурные принципы + +Build следует тем же принципам, которые используются во всей новой архитектуре Acquisition. + +## 4.1 Single Responsibility + +Runtime отвечает исключительно за инфраструктуру соединения. + +Trade Integration отвечает исключительно за обработку Trade Stream. + +Consistency отвечает исключительно за целостность последовательности сделок. + +Recovery отвечает исключительно за восстановление пропущенных данных. + +Ни один слой не должен смешивать собственную ответственность с соседними. + +--- + +## 4.2 Dependency Direction + +Все зависимости направлены только вниз. + +``` +Trade Stream Acquisition Service + + │ + + ▼ + +Runtime Service + + │ + + ▼ + +Runtime Protocols +``` + +Обратные зависимости запрещены. + +Runtime не может импортировать Acquisition. + +--- + +## 4.3 Exchange Independence + +Trade Integration не знает ничего о Dzengi. + +Exchange-specific логика уже локализована внутри: + +``` +DzengiUnifiedWebSocketAdapter +``` + +Build не добавляет новых exchange-specific условий. + +--- + +## 4.4 Canonical Pipeline + +После Build существует единственный допустимый путь обработки Trade: + +``` +Raw Message + +↓ + +Unified Adapter + +↓ + +Trade + +↓ + +Consistency + +↓ + +Trade +``` + +Любые альтернативные пути считаются нарушением архитектуры. + +--- + +# 5. Новые компоненты + +Build вводит два новых компонента. + +``` +trade_stream_acquisition_protocol.py + +trade_stream_acquisition_service.py +``` + +Оба файла располагаются внутри: + +``` +src/ +└── market_data/ + └── acquisition/ +``` + +Build намеренно не изменяет существующие Feed. + +Feed остаются независимыми клиентами Acquisition Layer. + +--- + +# 6. Новая ответственность Acquisition Layer + +После Build Acquisition становится полноценным координатором Runtime Pipeline. + +Именно Acquisition теперь отвечает за: + +- отправку Runtime Command; +- получение Trade; +- передачу Trade в Consistency; +- возврат согласованного результата. + +Runtime при этом остаётся полностью инфраструктурным компонентом. + +--- + +--- + +# 7. Архитектура зависимостей + +После завершения Build зависимости между компонентами приобретают следующий вид. + +``` + +------------------------------------+ + | Trade Stream Acquisition Service | + +------------------------------------+ + │ │ + │ │ + ▼ ▼ + +---------------------------+ +---------------------------+ + | Acquisition Runtime | | Dzengi Unified Adapter | + | Service | +---------------------------+ + +---------------------------+ │ + │ │ + │ ▼ + │ Quote | Trade | Candle + │ + ▼ + +---------------------------+ + | Runtime Protocols | + +---------------------------+ + │ + ▼ + Runtime Infrastructure + + +Trade + + │ + + ▼ + +Trade Stream Consistency Controller + + │ + + ▼ + +Trade | None +``` + +Таким образом Runtime перестаёт знать о типах рыночных данных. + +Trade остаётся полностью внутри слоя Acquisition. + +--- + +# 8. Обработка подписок + +Новый сервис становится единственной точкой открытия Trade Stream. + +Внешние компоненты больше не должны самостоятельно создавать Runtime Commands. + +Правильная последовательность выглядит следующим образом. + +``` +symbols + + │ + + ▼ + +build_trade_subscribe_command() + + │ + + ▼ + +SubscribeCommand + + │ + + ▼ + +AcquisitionRuntimeService.dispatch() +``` + +Сам сервис не сериализует сообщения. + +Он использует уже существующий Subscription Builder. + +Build не изменяет формат WebSocket сообщений. + +--- + +# 9. Обработка входящих сообщений + +Второй обязанностью нового сервиса становится обработка входящего транспортного сообщения. + +Полная последовательность обработки выглядит следующим образом. + +``` +Raw WebSocket Document + + │ + + ▼ + +DzengiUnifiedWebSocketAdapter + + │ + + ▼ + +Quote +Trade +CandleCloseEvent + + │ + + ▼ + +isinstance(result, Trade) + + │ + + ├──────────────► нет + │ + │ + ▼ + +Trade Stream Consistency Controller + + │ + + ▼ + +Trade | None +``` + +Если адаптер возвращает: + +``` +Quote +``` + +или + +``` +CandleCloseEvent +``` + +сообщение считается неподходящим для данного сервиса. + +Build не выполняет никаких дополнительных преобразований. + +--- + +# 10. Взаимодействие с Consistency + +Trade Integration не реализует собственную проверку последовательности. + +Вся ответственность делегируется существующему компоненту. + +``` +Trade + + │ + + ▼ + +TradeStreamConsistencyController.accept() + + │ + + ▼ + +Trade | None +``` + +Если Controller возвращает: + +``` +None +``` + +это означает корректный дубликат. + +Никаких дополнительных действий сервис не выполняет. + +Если Controller возбуждает исключение: + +- TradeConsistencyError; +- TradeOrderingError; + +исключение полностью передаётся вызывающей стороне. + +Build не изменяет модель ошибок Consistency Layer. + +--- + +# 11. Взаимодействие с Runtime + +Trade Stream Acquisition Service использует Runtime исключительно как инфраструктурный компонент. + +Допустимыми являются только следующие Runtime Commands. + +``` +ConnectCommand + +DisconnectCommand + +SubscribeCommand + +UnsubscribeCommand + +SendTextCommand + +SendBinaryCommand +``` + +Никакие Runtime Events в Build 060.23 не анализируются. + +Publisher продолжает существовать исключительно как инфраструктурный контракт. + +--- + +# 12. Использование Unified Adapter + +Build намеренно использует существующий: + +``` +DzengiUnifiedWebSocketAdapter +``` + +а не специализированный Trade Adapter. + +Причины данного решения: + +• уже существует единая точка обработки WebSocket сообщений; + +• отсутствует дублирование маршрутизации; + +• Runtime остаётся полностью независимым от типа сообщения; + +• в будущем Runtime сможет одинаково обслуживать: + +- Quotes; +- Trades; +- Candles. + +Таким образом новый сервис использует уже сформированную архитектуру, а не создаёт отдельную ветку обработки Trade. + +--- + +# 13. Публичный API сервиса + +Build вводит минимальный публичный интерфейс. + +```python +class TradeStreamAcquisitionServiceProtocol(Protocol): + + async def subscribe( + self, + symbols: tuple[str, ...], + *, + correlation_id: str | None = None, + ) -> None: + ... + + def handle_message( + self, + document: object, + ) -> Trade | None: + ... +``` + +Оба метода являются частью единой ответственности сервиса. + +Никаких дополнительных методов Build не вводит. + +--- + +--- + +# 14. Architecture Decision Records (ADR) + +## ADR-060.23-001 + +Trade Stream Integration располагается исключительно внутри слоя Acquisition. + +### Причина + +Runtime является инфраструктурным компонентом и не должен знать о типах рыночных данных. + +### Следствие + +Runtime никогда не импортирует: + +- Trade; +- Trades Feed; +- Consistency; +- Recovery; +- Unified Adapter. + +--- + +## ADR-060.23-002 + +Trade обрабатывается только после прохождения Unified Adapter. + +### Причина + +В системе уже существует единая exchange-specific точка преобразования транспортных сообщений. + +``` +DzengiUnifiedWebSocketAdapter +``` + +Build не создаёт альтернативных маршрутов обработки. + +--- + +## ADR-060.23-003 + +Trade Stream Acquisition Service является единственной точкой подписки Trade Runtime. + +### Причина + +Внешние компоненты не должны самостоятельно формировать Runtime Commands. + +Все команды создаются посредством существующих Subscription Builder. + +--- + +## ADR-060.23-004 + +Trade Stream Acquisition Service не хранит состояние Trade Stream. + +### Причина + +Сервис выполняет исключительно координацию. + +Состояние распределено между специализированными компонентами: + +- Runtime Session; +- Subscription Manager; +- Trade Stream Consistency Controller. + +Build не вводит нового состояния. + +--- + +## ADR-060.23-005 + +Trade Stream Acquisition Service не публикует Trade. + +### Причина + +В проекте отсутствует утверждённая инфраструктура публикации канонического Trade Stream. + +До появления такого компонента сервис возвращает: + +``` +Trade | None +``` + +не принимая решений о дальнейшей маршрутизации данных. + +--- + +# 15. Что не входит в Build + +Настоящий Build сознательно не включает: + +- реализацию WebSocket Transport; +- реализацию WebSocket Session; +- production Subscription Manager; +- receive loop; +- reconnect; +- heartbeat; +- scheduler; +- Runtime Recovery; +- публикацию Runtime Events; +- публикацию Trade; +- хранение Trade; +- обработку Quote; +- обработку Candle; +- изменение существующего Trades Feed; +- изменение MarketDataRunner; +- изменение market_stream.py; +- изменение ExchangeWebSocketClient. + +Все перечисленные задачи относятся к следующим Build. + +--- + +# 16. План реализации + +Реализация выполняется в несколько последовательных шагов. + +## Этап 1 + +Создание нового контракта + +``` +trade_stream_acquisition_protocol.py +``` + +Контракт описывает публичный API сервиса. + +--- + +## Этап 2 + +Создание реализации + +``` +trade_stream_acquisition_service.py +``` + +Сервис получает через конструктор зависимости: + +- AcquisitionRuntimeServiceProtocol; +- DzengiUnifiedWebSocketAdapter; +- TradeStreamConsistencyProtocol. + +--- + +## Этап 3 + +Интеграция подписки + +Добавляется метод: + +``` +subscribe(...) +``` + +который: + +- строит SubscribeCommand; +- передаёт его Runtime Service; +- не содержит exchange-specific логики. + +--- + +## Этап 4 + +Интеграция обработки сообщений + +Добавляется метод: + +``` +handle_message(...) +``` + +Последовательность обработки: + +``` +document + +↓ + +Unified Adapter + +↓ + +Trade + +↓ + +Consistency + +↓ + +Trade | None +``` + +--- + +## Этап 5 + +Покрытие тестами + +Добавляются unit-тесты: + +- проверка соответствия Protocol; +- проверка dispatch подписок; +- проверка обработки Trade; +- проверка игнорирования Quote; +- проверка игнорирования Candle; +- проверка передачи Trade в Consistency; +- проверка возврата None для корректного дубликата; +- проверка проброса исключений Consistency; +- проверка отсутствия скрытого состояния. + +--- + +# 17. Definition of Done + +Build считается завершённым после выполнения следующих условий. + +✓ создан Protocol Trade Stream Acquisition Service; + +✓ создана реализация сервиса; + +✓ Runtime интегрирован через AcquisitionRuntimeServiceProtocol; + +✓ Unified Adapter используется как единственная точка преобразования сообщений; + +✓ Consistency используется как единственная точка проверки последовательности сделок; + +✓ сервис не содержит exchange-specific логики; + +✓ сервис не содержит собственного состояния; + +✓ сервис не содержит логики reconnect; + +✓ сервис не содержит логики recovery; + +✓ сервис не взаимодействует с legacy Runtime; + +✓ сервис полностью покрыт unit-тестами; + +✓ существующие тесты Runtime, Feeds, Consistency и Recovery продолжают успешно проходить без изменений. + +--- + +# 18. Архитектурное состояние после Build + +После завершения Build 060.23 инфраструктура Acquisition приобретает завершённую форму. + +``` +Subscription Builder + + │ + + ▼ + +Trade Stream Acquisition Service + + │ + + ▼ + +Acquisition Runtime Service + + │ + + ▼ + +Runtime Protocols + +──────────────────────────────────── + +Raw WebSocket Message + + │ + + ▼ + +DzengiUnifiedWebSocketAdapter + + │ + + ▼ + +Trade + + │ + + ▼ + +Trade Stream Consistency Controller + + │ + + ▼ + +Trade | None +``` + +Таким образом завершается построение сервисного слоя интеграции Runtime и Acquisition. + +Следующий этап — **Build 060.24 — Reconnect & Runtime Recovery**, в рамках которого будет реализовано управление жизненным циклом WebSocket-соединения, автоматическое восстановление подписок и интеграция Recovery Layer с Runtime Infrastructure. \ No newline at end of file