From be6beac560a441f4563400ae7b3abf76d40f52b8 Mon Sep 17 00:00:00 2001 From: Sergey Date: Mon, 20 Jul 2026 00:49:45 +0300 Subject: [PATCH] Build 060.16: add Trade Subscription Layer and subscription builders --- app/scripts/check_trades_unsubscribe.py | 214 +++ .../acquisition/adapters/dzengi/__init__.py | 1 + .../acquisition/subscriptions/__init__.py | 6 + .../acquisition/subscriptions/trades.py | 111 ++ .../acquisition/subscriptions/test_trades.py | 178 +++ docs/migrations/build_060_16.md | 1283 +++++++++++++++++ 6 files changed, 1793 insertions(+) create mode 100644 app/scripts/check_trades_unsubscribe.py create mode 100644 app/src/market_data/acquisition/subscriptions/__init__.py create mode 100644 app/src/market_data/acquisition/subscriptions/trades.py create mode 100644 app/tests/unit/market_data/acquisition/subscriptions/test_trades.py create mode 100644 docs/migrations/build_060_16.md diff --git a/app/scripts/check_trades_unsubscribe.py b/app/scripts/check_trades_unsubscribe.py new file mode 100644 index 0000000..664c914 --- /dev/null +++ b/app/scripts/check_trades_unsubscribe.py @@ -0,0 +1,214 @@ +from __future__ import annotations + +import asyncio +import json +import os +from typing import Any +from uuid import uuid4 + +import websockets +from websockets.typing import Subprotocol + + +SYMBOL = os.getenv( + "TRADE_TEST_SYMBOL", + "BTC/USD_LEVERAGE", +) + +OBSERVE_BEFORE_UNSUBSCRIBE_SECONDS = 30 +OBSERVE_AFTER_UNSUBSCRIBE_SECONDS = 60 +RECEIVE_TIMEOUT_SECONDS = 10 + + +def build_ws_url() -> str: + raw_url = os.getenv( + "EXCHANGE_WS_URL", + "wss://api-adapter.dzengi.com", + ).rstrip("/") + + if raw_url.startswith("https://"): + raw_url = raw_url.replace("https://", "wss://", 1) + elif raw_url.startswith("http://"): + raw_url = raw_url.replace("http://", "ws://", 1) + + if raw_url.endswith("/connect"): + return raw_url + + return f"{raw_url}/connect" + + +def build_headers() -> dict[str, str]: + base_url = os.getenv( + "EXCHANGE_BASE_URL", + "https://api-adapter.dzengi.com", + ).rstrip("/") + + headers = { + "Origin": base_url, + "Content-Type": "application/json", + } + + api_key = os.getenv("EXCHANGE_API_KEY", "").strip() + + if api_key: + headers["X-MBX-APIKEY"] = api_key + + return headers + + +def build_request(destination: str) -> dict[str, Any]: + return { + "correlationId": str(uuid4()), + "destination": destination, + "payload": { + "symbols": [SYMBOL], + }, + } + + +def decode_message(raw_message: str | bytes) -> Any: + if isinstance(raw_message, bytes): + raw_message = raw_message.decode("utf-8") + + try: + return json.loads(raw_message) + except json.JSONDecodeError: + return raw_message + + +def print_message(label: str, message: Any) -> None: + print() + print("=" * 80) + print(label) + print("=" * 80) + + if isinstance(message, (dict, list)): + print(json.dumps(message, indent=2, ensure_ascii=False)) + else: + print(message) + + +def is_trade_event(message: Any) -> bool: + return ( + isinstance(message, dict) + and message.get("destination") == "internal.trade" + ) + + +async def receive_until( + websocket: Any, + *, + duration_seconds: int, + phase: str, +) -> int: + loop = asyncio.get_running_loop() + deadline = loop.time() + duration_seconds + trade_count = 0 + + while loop.time() < deadline: + remaining = deadline - loop.time() + timeout = min(RECEIVE_TIMEOUT_SECONDS, remaining) + + try: + raw_message = await asyncio.wait_for( + websocket.recv(), + timeout=timeout, + ) + except asyncio.TimeoutError: + print( + f"[{phase}] За последние {timeout:.0f} секунд " + "сообщений не получено; соединение остаётся открытым." + ) + continue + + message = decode_message(raw_message) + print_message(f"[{phase}] Получено сообщение", message) + + if is_trade_event(message): + trade_count += 1 + + return trade_count + + +async def main() -> None: + ws_url = build_ws_url() + headers = build_headers() + + subscribe_request = build_request("trades.subscribe") + unsubscribe_request = build_request("trades.unsubscribe") + + print(f"WebSocket URL: {ws_url}") + print(f"Symbol: {SYMBOL}") + + async with websockets.connect( + ws_url, + extra_headers=headers, + subprotocols=[Subprotocol("json")], + ping_interval=20, + ping_timeout=20, + open_timeout=20, + ) as websocket: + print_message( + "Отправляется trades.subscribe", + subscribe_request, + ) + + await websocket.send( + json.dumps(subscribe_request), + ) + + before_count = await receive_until( + websocket, + duration_seconds=OBSERVE_BEFORE_UNSUBSCRIBE_SECONDS, + phase="BEFORE UNSUBSCRIBE", + ) + + print() + print( + "Количество Trade Event до unsubscribe: " + f"{before_count}" + ) + + print_message( + "Отправляется trades.unsubscribe", + unsubscribe_request, + ) + + await websocket.send( + json.dumps(unsubscribe_request), + ) + + after_count = await receive_until( + websocket, + duration_seconds=OBSERVE_AFTER_UNSUBSCRIBE_SECONDS, + phase="AFTER UNSUBSCRIBE", + ) + + print() + print("=" * 80) + print("РЕЗУЛЬТАТ НАБЛЮДЕНИЯ") + print("=" * 80) + print(f"Trade Event до unsubscribe: {before_count}") + print(f"Trade Event после unsubscribe: {after_count}") + + if before_count == 0: + print( + "До unsubscribe не было получено ни одной сделки. " + "Эксперимент нельзя считать доказательным." + ) + elif after_count == 0: + print( + "После trades.unsubscribe новые сделки не поступили. " + "Команда, вероятно, поддерживается, но необходимо " + "проверить ACK или повторить тест при высокой активности." + ) + else: + print( + "После trades.unsubscribe сделки продолжили поступать. " + "Команда либо не поддерживается, либо была отклонена, " + "либо имеет другой формат." + ) + + +if __name__ == "__main__": + asyncio.run(main()) \ No newline at end of file diff --git a/app/src/market_data/acquisition/adapters/dzengi/__init__.py b/app/src/market_data/acquisition/adapters/dzengi/__init__.py index e69de29..e4aab50 100644 --- a/app/src/market_data/acquisition/adapters/dzengi/__init__.py +++ b/app/src/market_data/acquisition/adapters/dzengi/__init__.py @@ -0,0 +1 @@ +# app/src/market_data/acquisition/adapters/dzengi/__init__.py \ No newline at end of file diff --git a/app/src/market_data/acquisition/subscriptions/__init__.py b/app/src/market_data/acquisition/subscriptions/__init__.py new file mode 100644 index 0000000..24a1adc --- /dev/null +++ b/app/src/market_data/acquisition/subscriptions/__init__.py @@ -0,0 +1,6 @@ +""" +Формирование транспортных подписок Market Data Acquisition. + +Модули пакета преобразуют параметры предметной области в универсальные +Runtime-команды и не выполняют сетевые операции самостоятельно. +""" \ No newline at end of file diff --git a/app/src/market_data/acquisition/subscriptions/trades.py b/app/src/market_data/acquisition/subscriptions/trades.py new file mode 100644 index 0000000..b2b4fae --- /dev/null +++ b/app/src/market_data/acquisition/subscriptions/trades.py @@ -0,0 +1,111 @@ +# app/src/market_data/acquisition/subscriptions/trades.py + +from __future__ import annotations + +import json +from collections.abc import Sequence +from uuid import uuid4 + +from src.market_data.acquisition.runtime.runtime_commands import ( + SubscribeCommand, +) +from src.market_data.acquisition.runtime.transport_messages import ( + TransportTextMessage, +) + + +TRADE_SUBSCRIPTION_DESTINATION = "trades.subscribe" +TRADE_SUBSCRIPTION_KEY_PREFIX = "trades" + + +def _normalize_trade_subscription_symbols( + symbols: Sequence[str], +) -> tuple[str, ...]: + """ + Нормализовать символы Trade-подписки. + + Удаляет внешние пробелы, пустые значения и дубликаты. + Возвращает символы в стабильном лексикографическом порядке. + """ + + normalized_symbols = { + symbol.strip() + for symbol in symbols + if symbol.strip() + } + + if not normalized_symbols: + raise ValueError( + "Trade subscription requires at least one non-empty symbol." + ) + + return tuple(sorted(normalized_symbols)) + + +def build_trade_subscription_key( + symbols: Sequence[str], +) -> str: + """ + Построить стабильный идентификатор Trade-подписки. + + Идентификатор не зависит от порядка входных символов и используется + Runtime для регистрации и последующего восстановления подписки. + """ + + normalized_symbols = _normalize_trade_subscription_symbols(symbols) + + return ( + f"{TRADE_SUBSCRIPTION_KEY_PREFIX}:" + f"{','.join(normalized_symbols)}" + ) + + +def build_trade_subscribe_message( + symbols: Sequence[str], + *, + correlation_id: str | None = None, +) -> TransportTextMessage: + """ + Построить текстовое сообщение подписки Dzengi Trade WebSocket. + + correlation_id может быть передан вызывающим кодом для детерминированных + тестов и сопоставления ACK. Если значение не передано, создаётся новый UUID. + """ + + normalized_symbols = _normalize_trade_subscription_symbols(symbols) + + document = { + "correlationId": correlation_id or str(uuid4()), + "destination": TRADE_SUBSCRIPTION_DESTINATION, + "payload": { + "symbols": list(normalized_symbols), + }, + } + + return TransportTextMessage( + payload=json.dumps( + document, + separators=(",", ":"), + ) + ) + + +def build_trade_subscribe_command( + symbols: Sequence[str], + *, + correlation_id: str | None = None, +) -> SubscribeCommand: + """ + Построить универсальную Runtime-команду подписки на Trade Feed. + + Runtime получает стабильный ключ и готовое транспортное сообщение, + не интерпретируя Dzengi-specific JSON. + """ + + return SubscribeCommand( + subscription_key=build_trade_subscription_key(symbols), + message=build_trade_subscribe_message( + symbols, + correlation_id=correlation_id, + ), + ) \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/subscriptions/test_trades.py b/app/tests/unit/market_data/acquisition/subscriptions/test_trades.py new file mode 100644 index 0000000..5cafaba --- /dev/null +++ b/app/tests/unit/market_data/acquisition/subscriptions/test_trades.py @@ -0,0 +1,178 @@ +# app/tests/unit/market_data/acquisition/subscriptions/test_trades.py + +from __future__ import annotations + +import json +from uuid import UUID + +import pytest + +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.subscriptions.trades import ( + TRADE_SUBSCRIPTION_DESTINATION, + build_trade_subscribe_command, + build_trade_subscribe_message, + build_trade_subscription_key, +) + + +def test_build_trade_subscription_key_for_single_symbol() -> None: + result = build_trade_subscription_key( + ["BTC/USD_LEVERAGE"], + ) + + assert result == "trades:BTC/USD_LEVERAGE" + + +def test_build_trade_subscription_key_normalizes_symbols() -> None: + result = build_trade_subscription_key( + [ + " ETH/USD_LEVERAGE ", + "BTC/USD_LEVERAGE", + "ETH/USD_LEVERAGE", + "", + " ", + ], + ) + + assert result == ( + "trades:BTC/USD_LEVERAGE,ETH/USD_LEVERAGE" + ) + + +def test_build_trade_subscription_key_rejects_empty_symbols() -> None: + with pytest.raises( + ValueError, + match="at least one non-empty symbol", + ): + build_trade_subscription_key( + ["", " "], + ) + + +def test_build_trade_subscribe_message_returns_transport_message() -> None: + result = build_trade_subscribe_message( + ["BTC/USD_LEVERAGE"], + correlation_id="trade-subscription-1", + ) + + assert isinstance(result, TransportTextMessage) + + +def test_build_trade_subscribe_message_builds_dzengi_document() -> None: + result = build_trade_subscribe_message( + [ + "ETH/USD_LEVERAGE", + "BTC/USD_LEVERAGE", + ], + correlation_id="trade-subscription-1", + ) + + document = json.loads(result.payload) + + assert document == { + "correlationId": "trade-subscription-1", + "destination": TRADE_SUBSCRIPTION_DESTINATION, + "payload": { + "symbols": [ + "BTC/USD_LEVERAGE", + "ETH/USD_LEVERAGE", + ], + }, + } + + +def test_build_trade_subscribe_message_generates_uuid() -> None: + result = build_trade_subscribe_message( + ["BTC/USD_LEVERAGE"], + ) + + document = json.loads(result.payload) + + generated_id = UUID(document["correlationId"]) + + assert str(generated_id) == document["correlationId"] + + +def test_build_trade_subscribe_message_normalizes_symbols() -> None: + result = build_trade_subscribe_message( + [ + " BTC/USD_LEVERAGE ", + "BTC/USD_LEVERAGE", + "", + ], + correlation_id="trade-subscription-1", + ) + + document = json.loads(result.payload) + + assert document["payload"]["symbols"] == [ + "BTC/USD_LEVERAGE", + ] + + +def test_build_trade_subscribe_message_rejects_empty_symbols() -> None: + with pytest.raises( + ValueError, + match="at least one non-empty symbol", + ): + build_trade_subscribe_message( + [], + correlation_id="trade-subscription-1", + ) + + +def test_build_trade_subscribe_command_returns_runtime_command() -> None: + result = build_trade_subscribe_command( + ["BTC/USD_LEVERAGE"], + correlation_id="trade-subscription-1", + ) + + assert isinstance(result, SubscribeCommand) + assert result.subscription_key == "trades:BTC/USD_LEVERAGE" + assert isinstance(result.message, TransportTextMessage) + + +def test_build_trade_subscribe_command_contains_dzengi_message() -> None: + result = build_trade_subscribe_command( + ["BTC/USD_LEVERAGE"], + correlation_id="trade-subscription-1", + ) + + document = json.loads(result.message.payload) + + assert document == { + "correlationId": "trade-subscription-1", + "destination": "trades.subscribe", + "payload": { + "symbols": [ + "BTC/USD_LEVERAGE", + ], + }, + } + + +def test_build_trade_subscribe_command_uses_same_normalized_symbols() -> None: + result = build_trade_subscribe_command( + [ + "ETH/USD_LEVERAGE", + " BTC/USD_LEVERAGE ", + "ETH/USD_LEVERAGE", + ], + correlation_id="trade-subscription-1", + ) + + document = json.loads(result.message.payload) + + assert result.subscription_key == ( + "trades:BTC/USD_LEVERAGE,ETH/USD_LEVERAGE" + ) + assert document["payload"]["symbols"] == [ + "BTC/USD_LEVERAGE", + "ETH/USD_LEVERAGE", + ] \ No newline at end of file diff --git a/docs/migrations/build_060_16.md b/docs/migrations/build_060_16.md new file mode 100644 index 0000000..1a2baf7 --- /dev/null +++ b/docs/migrations/build_060_16.md @@ -0,0 +1,1283 @@ +# Build 060.16 — Trade Subscription Layer + +**Engineering Migration Report** + +--- + +# Контроль документа + +| Свойство | Значение | +|----------|----------| +| Build | 060.16 | +| Название | Trade Subscription Layer | +| Статус | Completed | +| Проект | Dzentra | +| Подсистема | Market Data Acquisition | +| Компонент | Trades Feed | +| Версия | 1.0 | + +--- + +# Цель Build + +После завершения Build 060.15 система получила полностью интегрированный механизм маршрутизации входящих WebSocket-сообщений. + +К этому моменту архитектура уже обеспечивала: + +- единый WebSocket Runtime; +- единый механизм транспортной маршрутизации; +- полноценный Trade Adapter; +- преобразование транспортных сообщений в каноническую модель `Trade`. + +Однако один важный уровень инфраструктуры всё ещё отсутствовал. + +Система уже умела принимать сообщения о сделках, но ещё не обладала единым механизмом формирования исходящих WebSocket-подписок. + +Подписка на поток Trade оставалась частью низкоуровневого протокола биржи и не была представлена отдельным архитектурным уровнем. + +Это приводило к нескольким ограничениям. + +Во-первых, детали протокола Dzengi могли начать проникать в Runtime. + +Во-вторых, отсутствовал единый механизм построения подписок для различных типов рыночных данных. + +В-третьих, будущие Feed были бы вынуждены самостоятельно формировать JSON-документы подписки. + +Подобная архитектура противоречила одному из базовых принципов Dzentra — строгому разделению ответственности между слоями системы. + +Основная задача Build заключается в создании отдельного уровня **Trade Subscription Layer**, который полностью инкапсулирует знания о протоколе подписки Dzengi и предоставляет Runtime исключительно универсальные транспортные команды. + +После завершения Build система должна уметь: + +- построить стабильный идентификатор подписки; +- сформировать корректный WebSocket-документ `trades.subscribe`; +- преобразовать транспортное сообщение в универсальную `SubscribeCommand`; +- полностью скрыть детали протокола Dzengi от Runtime. + +При этом Build принципиально не затрагивает: + +- WebSocket Runtime; +- WebSocket Protocol; +- Unified Router; +- Parser; +- Value Validation; +- Mapper; +- Trade Adapter; +- Trades Feed; +- Registry подписок; +- механизм восстановления подписок после reconnect. + +--- + +# Предпосылки + +К началу Build архитектура Market Data Acquisition уже содержала полностью реализованный транспортный уровень Runtime. + +Он обеспечивал выполнение всех операций, связанных с жизненным циклом WebSocket-соединения. + +Конвейер управления соединением выглядел следующим образом. + +```text +SubscribeCommand + │ + ▼ +Runtime + │ + ▼ +WebSocket Transport +``` + +Runtime уже обладал универсальной инфраструктурой транспортных команд. + +В частности, были реализованы: + +```text +ConnectCommand + +DisconnectCommand + +SubscribeCommand + +UnsubscribeCommand + +SendTextCommand + +SendBinaryCommand +``` + +При этом Runtime принципиально не содержал информации о конкретной бирже. + +Он не должен был знать: + +- структуру JSON-документов; +- названия WebSocket destination; +- формат сообщений подписки; +- особенности протокола Dzengi. + +Одновременно с этим Build 060.15 завершил интеграцию Trade Pipeline в Unified WebSocket Routing. + +Таким образом система уже умела получать входящие сообщения: + +```text +internal.trade +``` + +и преобразовывать их в каноническую модель: + +```text +Trade +``` + +Однако обратная часть жизненного цикла — формирование исходящей подписки — всё ещё отсутствовала как самостоятельный архитектурный уровень. + +В результате знания о протоколе подписки неизбежно начали бы распространяться по нескольким компонентам системы. + +--- + +# Архитектурное основание + +Одним из фундаментальных принципов архитектуры Dzentra является изоляция транспортного протокола биржи от внутренних компонентов системы. + +Каждый уровень должен знать исключительно ту информацию, которая относится к его собственной зоне ответственности. + +Runtime отвечает исключительно за передачу транспортных сообщений. + +Feed отвечает исключительно за организацию потока данных. + +Parser отвечает исключительно за разбор транспортных моделей. + +Mapper отвечает исключительно за построение канонических моделей. + +Следовательно, знания о формате WebSocket-подписки также должны быть сосредоточены в отдельном специализированном слое. + +После завершения Build архитектурная схема формирования подписки принимает следующий вид. + +```text +Trade Symbols + │ + ▼ +Trade Subscription Builder + │ + ▼ +TransportTextMessage + │ + ▼ +SubscribeCommand + │ + ▼ +Runtime + │ + ▼ +WebSocket +``` + +Каждый уровень отвечает только за собственную задачу. + +Builder знает исключительно протокол Dzengi. + +Runtime знает исключительно транспортные команды. + +WebSocket знает исключительно способ передачи сообщения. + +Подобное разделение полностью соответствует общей архитектуре подсистемы Market Data Acquisition. + +--- + +# Результаты архитектурного аудита + +Перед началом реализации Build был выполнен полный аудит существующей инфраструктуры Runtime. + +Анализ подтвердил наличие полностью сформированного транспортного уровня. + +В частности, компонент + +```text +runtime/runtime_commands.py +``` + +уже содержал универсальную команду + +```text +SubscribeCommand +``` + +которая полностью отделяет вызывающий код от конкретного способа передачи сообщения. + +Дополнительно был выполнен аудит транспортных сообщений. + +Компонент + +```text +runtime/transport_messages.py +``` + +уже содержал универсальную модель + +```text +TransportTextMessage +``` + +предназначенную для передачи текстовых WebSocket-документов. + +Также был выполнен аудит протокола Runtime. + +Интерфейс + +```text +WebSocketSubscriptionManagerProtocol +``` + +уже предусматривал поддержку восстановления подписок после переподключения посредством методов: + +```text +restore_subscriptions() + +clear_subscriptions() +``` + +Тем самым было подтверждено, что архитектура Runtime уже подготовлена к будущей реализации Registry подписок. + +Вмешательство в Runtime в рамках Build 060.16 не требуется. + +Единственным отсутствующим уровнем являлся специализированный механизм построения транспортных подписок. + +Именно этот уровень и становится предметом настоящего Build. + +--- + +# Исследование протокола Dzengi + +Перед реализацией Subscription Layer был выполнен отдельный анализ публичного WebSocket-протокола Dzengi. + +Изучение Swagger и существующей производственной интеграции подтвердило использование следующего формата подписки. + +```json +{ + "correlationId": "...", + "destination": "trades.subscribe", + "payload": { + "symbols": [ + "BTC/USD_LEVERAGE" + ] + } +} +``` + +После успешной обработки сервер возвращает подтверждение подписки. + +```json +{ + "status": "OK", + "destination": "trades.subscribe", + "payload": { + "subscriptions": { + "BTC/USD_LEVERAGE": "PROCESSED" + } + } +} +``` + +После этого сервер начинает передавать сообщения + +```text +destination = "internal.trade" +``` + +которые поступают в ранее реализованный Trade Pipeline. + +Проведённый анализ также подтвердил отсутствие документированной операции выборочной отмены подписки. + +В Swagger отсутствует описание команды + +```text +trades.unsubscribe +``` + +что потребовало проведения дополнительного инженерного исследования. + +--- + +# Экспериментальная проверка `trades.unsubscribe` + +В рамках Build был проведён отдельный эксперимент, целью которого являлась проверка существования недокументированной операции выборочной отмены подписки. + +Для этого был подготовлен диагностический сценарий, выполняющий следующую последовательность действий. + +```text +Открытие WebSocket + │ + ▼ +trades.subscribe + │ + ▼ +Получение ACK + │ + ▼ +trades.unsubscribe + │ + ▼ +Анализ ответа сервера +``` + +В ходе эксперимента сервер успешно принял команду + +```text +trades.subscribe +``` + +и подтвердил регистрацию подписки. + +После этого была отправлена команда + +```text +destination = "trades.unsubscribe" +``` + +Ответ сервера оказался следующим. + +```json +{ + "status": "ERROR", + "payload": { + "errorCode": "BAD_REQUEST" + } +} +``` + +Таким образом было подтверждено: + +- сервер распознаёт транспортный запрос; +- предложенный формат команды не поддерживается; +- документированного механизма выборочной отмены подписки не существует. + +Полученный результат стал одним из ключевых архитектурных оснований настоящего Build. + +Dzentra не реализует неподтверждённые возможности внешнего API. + +Subscription Layer поддерживает исключительно подтверждённую операцию + +```text +trades.subscribe +``` + +а прекращение получения потока данных остаётся частью жизненного цикла самого WebSocket-соединения. + +--- + +# Архитектурное решение + +По результатам проведённого аудита было принято решение не изменять существующую инфраструктуру Runtime. + +Вместо добавления новой логики в транспортный уровень был реализован отдельный слой построения подписок. + +Данный слой полностью изолирует знания о протоколе Dzengi от остальных компонентов системы. + +Новая архитектура принимает следующий вид. + +```text +Trade Symbols + │ + ▼ +Trade Subscription Layer + │ + ├────────► build_trade_subscription_key() + │ + ├────────► build_trade_subscribe_message() + │ + └────────► build_trade_subscribe_command() + │ + ▼ + SubscribeCommand + │ + ▼ + Runtime + │ + ▼ + WebSocket +``` + +Таким образом Runtime продолжает работать исключительно с универсальными транспортными командами. + +Ни один компонент Runtime не знает: + +- название WebSocket destination; +- структуру JSON-документа; +- список символов; +- особенности протокола Dzengi. + +Все эти знания полностью сосредоточены внутри нового Subscription Layer. + +Подобное решение сохраняет существующее разделение ответственности и обеспечивает возможность дальнейшего расширения системы без изменения транспортной инфраструктуры. + +--- + +# Новый пакет Subscription Layer + +Для реализации нового архитектурного уровня создан отдельный пакет. + +```text +src/ +└── market_data/ + └── acquisition/ + └── subscriptions/ + ├── __init__.py + └── trades.py +``` + +Появление отдельного пакета является осознанным архитектурным решением. + +Во время проектирования рассматривались несколько альтернатив. + +Размещение Builder внутри Runtime было отклонено, поскольку Runtime не должен содержать знания о конкретной бирже. + +Размещение Builder внутри Adapter также было признано неудачным. + +Adapter отвечает исключительно за преобразование входящих транспортных документов. + +Построение исходящих подписок относится к другой области ответственности. + +В результате был выбран самостоятельный пакет + +```text +acquisition/subscriptions +``` + +который становится специализированным уровнем формирования транспортных подписок. + +Данная архитектура обеспечивает единое место хранения логики построения подписок для всех будущих потоков рыночных данных. + +В дальнейшем аналогичным образом могут быть реализованы: + +```text +quotes.py + +candles.py + +orderbook.py + +status.py +``` + +Каждый модуль будет отвечать исключительно за построение транспортных подписок собственного типа данных. + +--- + +# Реализованные Builder-функции + +В рамках Build реализованы три чистые функции. + +Они полностью покрывают процесс формирования универсальной Runtime-команды. + +--- + +## Построение ключа подписки + +Функция + +```text +build_trade_subscription_key() +``` + +формирует стабильный идентификатор подписки. + +При построении ключа выполняется: + +- удаление пустых элементов; +- удаление дубликатов; +- нормализация пробелов; +- каноническая сортировка символов. + +Например + +```text +BTC/USD_LEVERAGE +``` + +преобразуется в + +```text +trades:BTC/USD_LEVERAGE +``` + +а набор + +```text +ETH/USD_LEVERAGE +BTC/USD_LEVERAGE +ETH/USD_LEVERAGE +``` + +преобразуется в + +```text +trades:BTC/USD_LEVERAGE,ETH/USD_LEVERAGE +``` + +Использование стабильного идентификатора позволяет в дальнейшем однозначно регистрировать подписки внутри Runtime Registry. + +--- + +## Построение транспортного сообщения + +Функция + +```text +build_trade_subscribe_message() +``` + +формирует полностью готовый WebSocket-документ протокола Dzengi. + +Результатом работы функции является универсальная транспортная модель + +```text +TransportTextMessage +``` + +содержащая сериализованный JSON-документ. + +При отсутствии внешнего + +```text +correlationId +``` + +Builder автоматически создаёт новый UUID. + +При передаче значения извне оно сохраняется без изменений. + +Подобный механизм обеспечивает возможность детерминированного тестирования и сопоставления серверного ACK. + +--- + +## Построение Runtime-команды + +Функция + +```text +build_trade_subscribe_command() +``` + +является верхним уровнем Subscription Layer. + +Она объединяет результаты двух предыдущих функций. + +В результате вызывающая сторона получает полностью сформированную + +```text +SubscribeCommand +``` + +которая уже содержит: + +- стабильный идентификатор подписки; +- готовое транспортное сообщение. + +После построения команды никакая дополнительная обработка более не требуется. + +Runtime получает полностью готовую транспортную команду. + +--- + +# Последовательность формирования SubscribeCommand + +После завершения Build полный процесс построения подписки выглядит следующим образом. + +```text +Symbols + │ + ▼ +Normalization + │ + ▼ +Subscription Key + │ + ▼ +Trade JSON Builder + │ + ▼ +TransportTextMessage + │ + ▼ +SubscribeCommand + │ + ▼ +Runtime + │ + ▼ +WebSocket +``` + +Каждый этап отвечает исключительно за собственную задачу. + +Нормализация символов не знает о Runtime. + +Runtime не знает о JSON. + +JSON Builder не знает о транспортном соединении. + +Подобное разделение обеспечивает высокую повторную используемость компонентов. + +--- + +# Использование существующей Runtime-инфраструктуры + +Одним из ключевых требований Build являлось максимальное повторное использование уже реализованных компонентов. + +Subscription Layer не вводит новых транспортных моделей. + +Не создаёт новых Runtime-команд. + +Не изменяет существующий протокол передачи сообщений. + +Вместо этого используются уже реализованные компоненты проекта. + +Для передачи транспортного документа применяется + +```text +TransportTextMessage +``` + +Для взаимодействия с Runtime используется + +```text +SubscribeCommand +``` + +Передача данных продолжает выполняться существующим WebSocket Runtime. + +Таким образом новый Build полностью повторно использует уже реализованную инфраструктуру проекта. + +--- + +# Почему Subscription Layer реализован функциями + +Во время проектирования рассматривался вариант реализации отдельного класса + +```text +TradeSubscriptionManager +``` + +Однако после проведения архитектурного аудита данный вариант был отклонён. + +Причина заключается в отсутствии собственного состояния. + +Subscription Layer: + +- не хранит активные подписки; +- не выполняет регистрацию; +- не взаимодействует с Runtime Registry; +- не выполняет восстановление после reconnect. + +Единственной задачей слоя является построение транспортных объектов. + +Подобная логика является полностью детерминированной. + +Она зависит исключительно от входных параметров. + +Поэтому реализация в виде чистых функций полностью соответствует принятому в проекте архитектурному правилу. + +Если компонент не хранит состояние, он должен быть реализован функциями. + +Создание класса в данном случае привело бы лишь к искусственному усложнению архитектуры. + +--- + +# Почему Runtime не знает о `trades.subscribe` + +Во время проектирования отдельно рассматривался вариант переноса формирования JSON-документа непосредственно в Runtime. + +После анализа архитектуры данный вариант был отклонён. + +Основные причины: + +- Runtime не должен зависеть от конкретной биржи; +- Runtime не должен знать транспортный протокол; +- Runtime не должен содержать названия WebSocket destination; +- Runtime должен работать исключительно с универсальными командами. + +После завершения Build Runtime продолжает получать только + +```text +SubscribeCommand +``` + +которая уже содержит полностью сформированное транспортное сообщение. + +Таким образом достигается полная независимость Runtime от конкретной реализации протокола Dzengi. + +--- + +# Почему `trades.unsubscribe` не реализован + +Во время подготовки Build отдельно рассматривалась возможность реализации симметричной операции отмены подписки. + +На первый взгляд подобное решение выглядело естественным продолжением механизма + +```text +trades.subscribe +``` + +Однако архитектура Dzentra строится исключительно на подтверждённых возможностях внешнего API. + +Поэтому перед реализацией был проведён отдельный инженерный эксперимент. + +В результате сервер вернул ответ + +```text +status = ERROR +errorCode = BAD_REQUEST +``` + +Тем самым было подтверждено, что используемый формат + +```text +trades.unsubscribe +``` + +не поддерживается. + +Отсутствие документированного контракта означает невозможность гарантировать корректную работу подобной функции. + +Поэтому Build сознательно ограничен исключительно подтверждённой возможностью + +```text +trades.subscribe +``` + +Данное решение полностью соответствует одному из инженерных принципов проекта. + +Dzentra никогда не реализует функциональность, существование которой не подтверждено документацией либо экспериментально. + +Удаление подписки будет реализовано исключительно посредством управления жизненным циклом WebSocket-соединения и будущего Runtime Registry. + +--- + +# Изменённые файлы + +В рамках Build добавлены два новых файла. + +## Subscription Layer + +```text +src/market_data/acquisition/subscriptions/__init__.py +``` + +Создан новый пакет формирования транспортных подписок. + +Пакет становится единым местом хранения Builder-функций для различных типов рыночных данных. + +--- + +## Trade Subscription Builder + +```text +src/market_data/acquisition/subscriptions/trades.py +``` + +Реализованы: + +```text +build_trade_subscription_key() + +build_trade_subscribe_message() + +build_trade_subscribe_command() +``` + +Никакие существующие компоненты Runtime изменены не были. + +--- + +## Unit-тесты + +Добавлен новый набор тестов. + +```text +tests/unit/market_data/acquisition/subscriptions/test_trades.py +``` + +Все существующие тесты Runtime сохранены без изменений. + +--- + +# Добавленные тесты + +В рамках Build реализирован полный набор unit-тестов нового Subscription Layer. + +Проверяются следующие сценарии. + +--- + +## Построение ключа подписки + +Подтверждается корректная генерация: + +```text +trades:BTC/USD_LEVERAGE +``` + +для одиночного символа. + +--- + +## Нормализация символов + +Подтверждается: + +- удаление дубликатов; +- удаление пустых строк; +- удаление пробелов; +- стабильная сортировка. + +Получаемый ключ не зависит от порядка входных данных. + +--- + +## Проверка пустого списка + +Подтверждается генерация + +```text +ValueError +``` + +при отсутствии корректных символов. + +--- + +## Построение транспортного сообщения + +Подтверждается возврат объекта + +```text +TransportTextMessage +``` + +с корректным JSON-документом Dzengi. + +--- + +## Генерация UUID + +Подтверждается автоматическое создание + +```text +correlationId +``` + +при отсутствии значения. + +--- + +## Использование внешнего correlationId + +Подтверждается сохранение значения, + +переданного вызывающей стороной. + +Это обеспечивает возможность детерминированного тестирования. + +--- + +## Построение SubscribeCommand + +Подтверждается возврат корректной + +```text +SubscribeCommand +``` + +с правильным: + +- subscription key; +- TransportTextMessage. + +--- + +## Согласованность Subscription Builder + +Подтверждается использование одинакового набора нормализованных символов как при построении ключа подписки, так и при построении JSON-документа. + +Это гарантирует отсутствие рассогласования между Runtime Registry и транспортным протоколом. + +--- + +# Результаты тестирования + +После завершения реализации выполнен запуск нового набора unit-тестов. + +```bash +PYTHONPATH="$PWD" python -m pytest \ +tests/unit/market_data/acquisition/subscriptions/test_trades.py -v +``` + +Результат: + +```text +11 passed +``` + +Подтверждена корректная работа всех Builder-функций. + +--- + +# Регрессионное тестирование Runtime + +После завершения реализации выполнен запуск полного набора Runtime-тестов. + +```bash +PYTHONPATH="$PWD" python -m pytest \ +tests/unit/market_data/acquisition/runtime +``` + +Результат: + +```text +30 passed +``` + +Подтверждено отсутствие регрессий существующей транспортной инфраструктуры. + +Ни один существующий Runtime-компонент не изменился. + +--- + +# Проверка компиляции + +После завершения реализации выполнена проверка компиляции новых компонентов. + +```bash +python -m compileall \ +src/market_data/acquisition/subscriptions \ +tests/unit/market_data/acquisition/subscriptions +``` + +Компиляция завершилась успешно. + +Ошибок синтаксиса не обнаружено. + +Все новые файлы успешно компилируются. + +--- + +# Проверка Git diff + +После завершения реализации выполнена финальная проверка изменений. + +```bash +git diff --check +``` + +Результат: + +```text +без замечаний +``` + +Проверка подтвердила отсутствие: + +- trailing whitespace; +- ошибок окончания строк; +- конфликтов diff; +- нарушений форматирования. + +--- + +# Scope Build 060.16 + +В рамках настоящего Build реализован исключительно уровень + +```text +Trade Subscription Layer +``` + +Build **не включает**: + +- Trades Feed; +- Runtime Registry; +- Runtime Integration; +- Reconnect; +- восстановление подписок; +- обработку входящих Trade; +- сортировку сделок; +- дедупликацию; +- REST Backfill. + +Подобное ограничение полностью соответствует принятому принципу атомарной реализации Build. + +Каждый этап дорожной карты реализует только один самостоятельный архитектурный уровень. + +--- + +# Архитектурный результат + +После завершения Build система получила полностью самостоятельный слой формирования Trade-подписок. + +Архитектура приобретает следующий вид. + +```text +Trade Symbols + │ + ▼ +Trade Subscription Builder + │ + ├────────► Subscription Key + │ + ├────────► TransportTextMessage + │ + └────────► SubscribeCommand + │ + ▼ + Runtime + │ + ▼ + WebSocket +``` + +Subscription Layer полностью изолирует знания о протоколе Dzengi. + +Runtime продолжает работать исключительно с универсальными транспортными командами. + +Это позволяет в дальнейшем добавлять новые типы подписок без каких-либо изменений транспортной инфраструктуры. + +--- + +# Соблюдение архитектурных принципов + +В рамках Build полностью сохранены архитектурные инварианты Dzentra. + +## Локальность изменений + +Добавлены только: + +- `src/market_data/acquisition/subscriptions/__init__.py`; +- `src/market_data/acquisition/subscriptions/trades.py`; +- `tests/unit/market_data/acquisition/subscriptions/test_trades.py`. + +Существующие Runtime-компоненты не изменялись. + +--- + +## Повторное использование инфраструктуры + +Subscription Layer полностью использует уже существующие компоненты: + +- `TransportTextMessage`; +- `SubscribeCommand`; +- WebSocket Runtime. + +Новая транспортная инфраструктура не создавалась. + +--- + +## Разделение ответственности + +Subscription Builder отвечает исключительно за построение транспортных подписок. + +Runtime отвечает исключительно за передачу сообщений. + +Feed будет отвечать исключительно за организацию потока данных. + +Подобное разделение полностью соответствует архитектуре Market Data Acquisition. + +--- + +## Повторное использование архитектуры + +Build не вводит нового транспортного уровня. + +Не изменяет существующий Runtime. + +Не изменяет WebSocket Protocol. + +Не изменяет механизм передачи сообщений. + +Вместо этого новый уровень полностью использует уже существующую архитектуру Market Data Acquisition. + +Фактически Build является логическим продолжением ранее реализованной Runtime-инфраструктуры. + +Это позволяет последующим этапам дорожной карты использовать единый механизм формирования подписок независимо от конкретного типа рыночных данных. + +--- + +## Минимальность изменений + +Одним из ключевых требований Build являлось минимальное вмешательство в существующий код. + +По результатам архитектурного аудита было принято решение отказаться от изменения Runtime. + +Все изменения локализованы внутри нового пакета + +```text +acquisition/subscriptions +``` + +Подобная локализация существенно снижает риск возникновения регрессий и позволяет независимо развивать Subscription Layer без влияния на остальные компоненты системы. + +--- + +## Функциональный подход + +Все компоненты Subscription Layer реализованы в виде чистых функций. + +Каждая функция: + +- имеет одну область ответственности; +- не зависит от внешнего состояния; +- не хранит собственных данных; +- детерминированно преобразует входные параметры в результат. + +Подобный подход полностью соответствует принятому в проекте правилу: + +> Stateless-компоненты реализуются функциями. Stateful-компоненты реализуются классами. + +Subscription Layer относится к первой категории. + +--- + +## Подготовка к масштабированию + +Несмотря на то что Build реализует только Trade Subscription Layer, выбранная архитектура изначально ориентирована на дальнейшее расширение. + +В дальнейшем аналогичный механизм может использоваться для: + +```text +Quote Subscription Layer + +Candles Subscription Layer + +Order Book Subscription Layer + +Market Status Subscription Layer +``` + +Каждый новый Builder будет полностью независимым. + +При этом Runtime останется неизменным. + +Таким образом уже на этапе Build 060.16 закладывается единый архитектурный шаблон для всех будущих потоков рыночных данных. + +--- + +# Архитектурные решения Build (ADR) + +## ADR-060.16-001 + +**Runtime остаётся полностью транспортно-независимым.** + +Runtime не знает: + +- структуру JSON; +- `destination`; +- особенности протокола Dzengi; +- типы рыночных данных. + +Он работает исключительно с универсальными транспортными командами. + +--- + +## ADR-060.16-002 + +**Построение подписок выполняется специализированным Subscription Layer.** + +Все знания о WebSocket-протоколе Dzengi сосредоточены внутри Builder-функций. + +Feed и Runtime не взаимодействуют с транспортным протоколом напрямую. + +--- + +## ADR-060.16-003 + +**Неподтверждённые возможности внешнего API не реализуются.** + +Во время Build выполнена экспериментальная проверка команды + +```text +trades.unsubscribe +``` + +Полученный ответ + +```text +BAD_REQUEST +``` + +подтвердил отсутствие поддерживаемого контракта. + +В результате Subscription Layer реализует исключительно документированную операцию + +```text +trades.subscribe +``` + +--- + +## ADR-060.16-004 + +**Subscription Builder не хранит собственного состояния.** + +Все функции являются полностью детерминированными. + +Registry подписок, управление жизненным циклом подписок и механизм восстановления после reconnect относятся к следующим Build и не входят в область ответственности настоящего слоя. + +--- + +# Критерии завершения Build + +Build 060.16 считается завершённым, поскольку выполнены все поставленные задачи. + +- ✔ создан отдельный пакет `acquisition/subscriptions`; +- ✔ реализована функция `build_trade_subscription_key()`; +- ✔ реализована функция `build_trade_subscribe_message()`; +- ✔ реализована функция `build_trade_subscribe_command()`; +- ✔ выполнена нормализация списка символов; +- ✔ реализована генерация стабильного ключа подписки; +- ✔ реализовано построение `TransportTextMessage`; +- ✔ реализовано построение `SubscribeCommand`; +- ✔ Runtime не изменён; +- ✔ подтверждено отсутствие поддержки `trades.unsubscribe`; +- ✔ реализованы unit-тесты нового Subscription Layer; +- ✔ успешно пройдены новые тесты (`11 passed`); +- ✔ успешно пройдена регрессия Runtime (`30 passed`); +- ✔ проект успешно компилируется; +- ✔ `git diff --check` не выявил замечаний; +- ✔ изменения полностью укладываются в согласованный scope Build. + +--- + +# Следующий этап + +Следующим этапом дорожной карты является + +```text +Build 060.17 — Trades Feed Core +``` + +Цель следующего Build: + +- реализовать специализированный Feed обработки сделок; +- подключить поток `Trade` к инфраструктуре Market Data Acquisition; +- организовать приём канонических моделей `Trade`; +- подготовить основу для последующих этапов: + - упорядочивания сделок; + - дедупликации; + - восстановления истории после переподключения; + - интеграции с Runtime. + +Subscription Layer, реализованный в Build 060.16, станет источником формирования подписок для нового Feed. + +--- + +# Итог + +Build 060.16 завершил формирование самостоятельного уровня **Trade Subscription Layer** в подсистеме Market Data Acquisition. + +Новая реализация полностью повторно использует существующую инфраструктуру Runtime, не изменяет транспортный протокол и не нарушает архитектурные инварианты проекта. + +Все знания о WebSocket-протоколе Dzengi теперь сосредоточены в одном специализированном пакете, тогда как Runtime продолжает работать исключительно с универсальными транспортными командами. + +В ходе Build была не только реализована новая функциональность, но и проведено инженерное исследование публичного WebSocket API Dzengi, подтвердившее отсутствие поддерживаемой операции `trades.unsubscribe`. Полученные результаты легли в основу принятых архитектурных решений и закреплены в ADR настоящего Build. + +Реализация ограничена согласованным scope, успешно прошла целевое и регрессионное тестирование, не потребовала изменений существующей Runtime-инфраструктуры и сформировала единый архитектурный шаблон для будущих Subscription Layer всех типов рыночных данных. + +После завершения Build система получила завершённый механизм формирования подписок `Trade`, который станет фундаментом для реализации **Trades Feed Core** в следующем этапе дорожной карты. \ No newline at end of file