From e5a4efa5d1740e846e85388804ebbb1e731fbdf1 Mon Sep 17 00:00:00 2001 From: Sergey Date: Fri, 17 Jul 2026 13:30:32 +0300 Subject: [PATCH] build 059.8: add unified websocket adapter --- .../acquisition/adapters/dzengi/websocket.py | 43 ++ app/src/market_data/acquisition/exceptions.py | 6 + .../dzengi/test_websocket_unified_adapter.py | 116 ++++ docs/migrations/build_059_8.md | 525 ++++++++++++++++++ 4 files changed, 690 insertions(+) create mode 100644 app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_unified_adapter.py create mode 100644 docs/migrations/build_059_8.md diff --git a/app/src/market_data/acquisition/adapters/dzengi/websocket.py b/app/src/market_data/acquisition/adapters/dzengi/websocket.py index 80eb504..11856d6 100644 --- a/app/src/market_data/acquisition/adapters/dzengi/websocket.py +++ b/app/src/market_data/acquisition/adapters/dzengi/websocket.py @@ -12,6 +12,9 @@ from src.market_data.acquisition.adapters.dzengi.parser import ( parse_dzengi_websocket_ohlc, parse_dzengi_websocket_quote, ) +from src.market_data.acquisition.exceptions import ( + WebSocketMessageRoutingError, +) from src.market_data.acquisition.models.candle_close import CandleCloseEvent from src.market_data.acquisition.models.quote import Quote from src.market_data.acquisition.validation.schema import ( @@ -58,4 +61,44 @@ class DzengiWebSocketOhlcAdapter: return map_dzengi_websocket_ohlc_to_candle_close_event( event, received_at=received_at or datetime.now(timezone.utc), + ) + + +# Определяет тип декодированного сообщения Dzengi WebSocket +# и передаёт его соответствующему специализированному адаптеру. +class DzengiUnifiedWebSocketAdapter: + def __init__( + self, + *, + quote_adapter: DzengiWebSocketQuoteAdapter | None = None, + ohlc_adapter: DzengiWebSocketOhlcAdapter | None = None, + ) -> None: + self._quote_adapter = quote_adapter or DzengiWebSocketQuoteAdapter() + self._ohlc_adapter = ohlc_adapter or DzengiWebSocketOhlcAdapter() + + def map_message( + self, + document: object, + *, + received_at: datetime | None = None, + ) -> Quote | CandleCloseEvent: + if not isinstance(document, dict): + raise WebSocketMessageRoutingError( + "Сообщение Dzengi WebSocket должно быть объектом." + ) + + if document.get("destination") == "ohlc.event": + return self._ohlc_adapter.map_message( + document, + received_at=received_at, + ) + + if "Payload" in document: + return self._quote_adapter.map_message( + document, + received_at=received_at, + ) + + raise WebSocketMessageRoutingError( + "Не удалось определить тип сообщения Dzengi WebSocket." ) \ No newline at end of file diff --git a/app/src/market_data/acquisition/exceptions.py b/app/src/market_data/acquisition/exceptions.py index 5c1effe..f1ae1f9 100644 --- a/app/src/market_data/acquisition/exceptions.py +++ b/app/src/market_data/acquisition/exceptions.py @@ -117,3 +117,9 @@ class CandleMappingError(MarketDataAcquisitionError): # Ошибка регистрации или получения Candles Feed. class CandleFeedRegistryError(MarketDataAcquisitionError): pass + + +# Ошибка определения типа входящего WebSocket-сообщения +# и выбора специализированного адаптера. +class WebSocketMessageRoutingError(MarketDataAcquisitionError): + pass \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_unified_adapter.py b/app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_unified_adapter.py new file mode 100644 index 0000000..091990d --- /dev/null +++ b/app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_unified_adapter.py @@ -0,0 +1,116 @@ +# app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_unified_adapter.py + +from __future__ import annotations + +from datetime import datetime, timezone + +import pytest + +from src.market_data.acquisition.adapters.dzengi.websocket import ( + DzengiUnifiedWebSocketAdapter, +) +from src.market_data.acquisition.exceptions import ( + WebSocketMessageRoutingError, +) +from src.market_data.acquisition.models.candle_close import CandleCloseEvent +from src.market_data.acquisition.models.quote import Quote + + +def test_adapter_routes_quote_message() -> None: + result = DzengiUnifiedWebSocketAdapter().map_message( + { + "Payload": { + "symbolName": "BTC/USD", + "bids": [["10", "1"]], + "asks": [["12", "1"]], + "timestamp": 1000, + } + } + ) + + assert isinstance(result, Quote) + assert result.symbol == "BTC/USD" + assert str(result.last_price) == "11" + + +def test_adapter_routes_ohlc_message() -> None: + result = DzengiUnifiedWebSocketAdapter().map_message( + { + "status": "OK", + "destination": "ohlc.event", + "payload": { + "symbol": "BTC/USD", + "interval": "1m", + "type": "classic", + "t": 1752364800000, + "o": "100", + "h": "105", + "l": "99", + "c": "103", + }, + } + ) + + assert isinstance(result, CandleCloseEvent) + assert result.symbol == "BTC/USD" + assert result.interval == "1m" + assert str(result.close_price) == "103" + + +def test_adapter_passes_received_at_to_quote_adapter() -> None: + received_at = datetime(2026, 7, 17, tzinfo=timezone.utc) + + result = DzengiUnifiedWebSocketAdapter().map_message( + { + "Payload": { + "symbolName": "BTC/USD", + "bids": [["10", "1"]], + "asks": [["12", "1"]], + "timestamp": 1000, + } + }, + received_at=received_at, + ) + + assert isinstance(result, Quote) + assert result.received_at is received_at + + +def test_adapter_passes_received_at_to_ohlc_adapter() -> None: + received_at = datetime(2026, 7, 17, tzinfo=timezone.utc) + + result = DzengiUnifiedWebSocketAdapter().map_message( + { + "status": "OK", + "destination": "ohlc.event", + "payload": { + "symbol": "BTC/USD", + "interval": "1m", + "type": "classic", + "t": 1752364800000, + "o": "100", + "h": "105", + "l": "99", + "c": "103", + }, + }, + received_at=received_at, + ) + + assert isinstance(result, CandleCloseEvent) + assert result.received_at is received_at + + +@pytest.mark.parametrize( + "document", + [ + None, + [], + {}, + {"status": "OK"}, + {"destination": "unknown.event", "payload": {}}, + ], +) +def test_adapter_rejects_unknown_message(document: object) -> None: + with pytest.raises(WebSocketMessageRoutingError): + DzengiUnifiedWebSocketAdapter().map_message(document) \ No newline at end of file diff --git a/docs/migrations/build_059_8.md b/docs/migrations/build_059_8.md new file mode 100644 index 0000000..52d705a --- /dev/null +++ b/docs/migrations/build_059_8.md @@ -0,0 +1,525 @@ +# Build 059.8 — Unified WebSocket Adapter + +**Проект:** Dzentra +**Подсистема:** Market Data Acquisition +**Этап:** 059.8 +**Статус:** Completed + +--- + +# Цель Build + +Завершить формирование слоя WebSocket-адаптеров, добавив единую точку входа для обработки входящих сообщений Dzengi WebSocket без изменения существующих специализированных адаптеров. + +Build 059.8 объединяет ранее реализованные адаптеры Quote и OHLC в единый маршрутизатор сообщений, который определяет тип входящего документа и передаёт его соответствующему обработчику. + +В результате последующие уровни системы больше не должны знать о конкретных форматах сообщений Dzengi WebSocket. + +--- + +# Причина изменения + +После завершения Build 059.7 архитектура содержала два полностью независимых специализированных адаптера: + +- DzengiWebSocketQuoteAdapter +- DzengiWebSocketOhlcAdapter + +Каждый адаптер имел собственный pipeline обработки данных: + +```text +Schema Validation + ↓ +Parser + ↓ +Value Validation + ↓ +Mapper + ↓ +Internal Model +``` + +Оба адаптера обладали одинаковым публичным контрактом: + +```python +map_message(document, received_at=None) +``` + +Однако внешний код должен был самостоятельно определять тип входящего сообщения и выбирать необходимый адаптер. + +Подобная логика выбора типа сообщения не относится к ответственности Runtime или транспортного уровня. + +Она должна существовать непосредственно в слое адаптеров. + +Build 059.8 устраняет данную проблему. + +--- + +# Реализованные изменения + +Изменены файлы: + +```text +src/market_data/acquisition/adapters/dzengi/websocket.py +src/market_data/acquisition/exceptions.py +``` + +Добавлен новый файл: + +```text +tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_unified_adapter.py +``` + +--- + +# Новый компонент + +Добавлен новый класс: + +```text +DzengiUnifiedWebSocketAdapter +``` + +Назначение класса: + +- получение произвольного декодированного сообщения Dzengi WebSocket; +- определение типа сообщения; +- выбор специализированного адаптера; +- возврат внутренней модели системы. + +Unified Adapter не содержит собственной логики обработки рыночных данных. + +Он выполняет исключительно маршрутизацию сообщений. + +--- + +# Архитектура после Build 059.8 + +После завершения Build структура слоя WebSocket имеет следующий вид: + +```text + Raw WebSocket Message + │ + ▼ + DzengiUnifiedWebSocketAdapter + │ + ┌───────────────┴───────────────┐ + │ │ + ▼ ▼ + DzengiWebSocketQuoteAdapter DzengiWebSocketOhlcAdapter + │ │ + ▼ ▼ + Quote CandleCloseEvent +``` + +Каждый специализированный адаптер полностью сохраняет собственный pipeline. + +Unified Adapter лишь выбирает необходимый маршрут обработки. + +--- + +# Маршрутизация сообщений + +Unified Adapter определяет тип сообщения по структуре входящего документа. + +## Quote + +Сообщение котировки определяется по наличию поля: + +```text +Payload +``` + +После определения типа сообщение передаётся в: + +```text +DzengiWebSocketQuoteAdapter +``` + +Дальнейшая обработка полностью выполняется существующим Quote pipeline. + +--- + +## OHLC + +Сообщение закрытия свечи определяется по значению: + +```text +destination == "ohlc.event" +``` + +После определения типа сообщение передаётся в: + +```text +DzengiWebSocketOhlcAdapter +``` + +Дальнейшая обработка полностью выполняется существующим OHLC pipeline. + +--- + +# Отсутствие дублирования логики + +Unified Adapter не выполняет: + +- Schema Validation; +- Parser; +- Value Validation; +- Mapping. + +Вся существующая логика полностью переиспользуется. + +Build 059.8 не изменяет существующие алгоритмы обработки сообщений. + +Добавляется исключительно уровень маршрутизации. + +--- + +# Контракт Unified Adapter + +Публичный метод: + +```python +map_message( + document, + *, + received_at=None, +) +``` + +Возвращаемое значение: + +```text +Quote +или +CandleCloseEvent +``` + +В зависимости от типа входящего сообщения. + +Таким образом внешний код получает уже внутреннюю модель системы и не зависит от структуры сообщений Dzengi. + +--- + +# Передача received_at + +Оба специализированных адаптера уже поддерживали передачу времени получения сообщения через параметр: + +```python +received_at +``` + +Unified Adapter полностью сохраняет данный контракт. + +При вызове: + +```python +map_message( + document, + received_at=received_at, +) +``` + +полученное значение без изменений передаётся выбранному специализированному адаптеру. + +Таким образом единая точка входа не изменяет временные метки и не создаёт новые значения времени самостоятельно. + +Если параметр не передан, соответствующий специализированный адаптер продолжает использовать собственную существующую логику формирования `received_at`. + +Build 059.8 не изменяет поведение предыдущих Build. + +--- + +# Dependency Injection + +Unified Adapter поддерживает передачу специализированных адаптеров через конструктор. + +Используются параметры: + +```python +quote_adapter +ohlc_adapter +``` + +Если адаптеры не переданы, создаются экземпляры по умолчанию: + +```text +DzengiWebSocketQuoteAdapter +DzengiWebSocketOhlcAdapter +``` + +Подобное решение обеспечивает: + +- изоляцию Unit Test; +- возможность последующего расширения; +- отсутствие жёсткой зависимости от конкретных реализаций. + +При этом публичный контракт класса остаётся неизменным. + +--- + +# Новая ошибка маршрутизации + +Build 059.8 вводит новый тип ошибки: + +```text +WebSocketMessageRoutingError +``` + +Назначение ошибки: + +определить ситуацию, когда входящее сообщение невозможно отнести ни к одному поддерживаемому типу. + +Ошибка возникает исключительно на этапе выбора специализированного адаптера. + +Она не относится к: + +- транспортному уровню; +- проверке схемы; +- parser; +- validator; +- mapper. + +Каждый специализированный адаптер продолжает использовать собственные типы ошибок. + +Таким образом уровни ответственности остаются полностью разделёнными. + +--- + +# Контракт обработки ошибок + +При успешном определении типа сообщения Unified Adapter передаёт управление соответствующему адаптеру. + +Все исключения специализированного pipeline проходят наружу без изменений. + +Например: + +```text +QuoteValueError + +QuoteSchemaError + +QuoteMappingError + +CandleWebSocketSchemaError + +CandleWebSocketValueError + +CandleWebSocketMappingError +``` + +не перехватываются и не преобразуются. + +Unified Adapter не изменяет существующую модель обработки ошибок. + +--- + +# Границы ответственности + +Unified Adapter отвечает исключительно за: + +- определение типа сообщения; +- выбор специализированного адаптера; +- передачу параметров; +- возврат внутренней модели. + +Unified Adapter сознательно не выполняет: + +- проверку структуры документа; +- преобразование данных; +- проверку бизнес-ограничений; +- преобразование во внутренние модели; +- взаимодействие с WebSocket-соединением. + +Все перечисленные обязанности остаются в существующих специализированных компонентах. + +Подобное разделение соответствует принципу Single Responsibility. + +--- + +# Поддерживаемые типы сообщений + +После завершения Build 059.8 поддерживаются: + +```text +Quote + +OHLC Close Event +``` + +Архитектура допускает дальнейшее расширение без изменения существующего кода Unified Adapter. + +В последующих Build возможно подключение новых специализированных адаптеров: + +```text +Order Book + +Trades + +Ticker + +Orders + +Positions + +Account Events +``` + +при сохранении единой точки входа. + +--- + +# Unit Test + +Добавлен новый набор тестов: + +```text +test_websocket_unified_adapter.py +``` + +Проверяются следующие сценарии: + +- маршрутизация Quote; +- маршрутизация OHLC; +- передача received_at в Quote Adapter; +- передача received_at в OHLC Adapter; +- обработка неизвестных сообщений. + +Для проверки неизвестного сообщения используется параметризованный тест, охватывающий несколько вариантов входных данных. + +--- + +# Проверка корректности + +Перед завершением Build выполнены проверки. + +## Компиляция + +```text +python -m compileall +``` + +Результат: + +```text +OK +``` + +--- + +## Unit Test Unified Adapter + +```text +9 passed +``` + +--- + +## Совместная регрессия адаптеров + +```text +15 passed +``` + +Проверены: + +- Quote Adapter; +- OHLC Adapter; +- Unified Adapter. + +--- + +## Полная регрессия Market Data Acquisition + +```text +634 passed +``` + +Подтверждено отсутствие регрессий во всей подсистеме получения рыночных данных. + +--- + +## Проверка diff + +Выполнена команда: + +```text +git diff --check +``` + +Замечаний не обнаружено. + +--- + +# Изменённые файлы + +```text +src/market_data/acquisition/adapters/dzengi/websocket.py + +src/market_data/acquisition/exceptions.py + +tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_unified_adapter.py +``` + +--- + +# Совместимость + +Build 059.8 полностью обратно совместим. + +Все существующие специализированные адаптеры продолжают работать без изменений. + +Публичие контракты: + +- Quote Adapter; +- OHLC Adapter; + +остаются неизменными. + +Unified Adapter представляет собой дополнительный уровень маршрутизации и не нарушает существующую архитектуру. + +--- + +# Архитектурное значение Build + +Build 059.8 завершает построение слоя адаптеров Dzengi WebSocket. + +После завершения этапа система получает единую точку входа для обработки входящих сообщений при сохранении независимости специализированных pipeline. + +Это позволяет последующим Build работать уже не с отдельными типами сообщений биржи, а с единым контрактом обработки WebSocket-событий. + +--- + +# Ограничения Build + +Build 059.8 не реализует: + +- управление WebSocket-соединением; +- подписку на каналы; +- повторные подписки; +- восстановление соединения; +- обработку heartbeat; +- синхронизацию состояния. + +Указанные возможности будут реализованы на следующих этапах миграции. + +--- + +# Итог + +Build 059.8 завершил построение слоя маршрутизации сообщений Dzengi WebSocket. + +Архитектура получила единую точку входа без изменения существующих специализированных адаптеров. + +Все проверки успешно пройдены. + +Регрессии отсутствуют. + +--- + +# Следующий этап + +```text +Build 059.9 — Subscription Layer +``` + +Следующий Build посвящён реализации уровня подписок на WebSocket-каналы и подготовке инфраструктуры для интеграции Unified Adapter с постоянным потоком биржевых сообщений. \ No newline at end of file