diff --git a/app/src/market_data/acquisition/adapters/dzengi/websocket.py b/app/src/market_data/acquisition/adapters/dzengi/websocket.py index 0f728af..80eb504 100644 --- a/app/src/market_data/acquisition/adapters/dzengi/websocket.py +++ b/app/src/market_data/acquisition/adapters/dzengi/websocket.py @@ -5,16 +5,21 @@ from __future__ import annotations from datetime import datetime, timezone from src.market_data.acquisition.adapters.dzengi.mapper import ( + map_dzengi_websocket_ohlc_to_candle_close_event, map_dzengi_websocket_quote_to_quote, ) from src.market_data.acquisition.adapters.dzengi.parser import ( + parse_dzengi_websocket_ohlc, parse_dzengi_websocket_quote, ) +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 ( + validate_dzengi_websocket_ohlc_schema, validate_dzengi_websocket_quote_schema, ) from src.market_data.acquisition.validation.values import ( + validate_dzengi_websocket_ohlc_values, validate_dzengi_websocket_quote_values, ) @@ -35,3 +40,22 @@ class DzengiWebSocketQuoteAdapter: response, received_at=received_at or datetime.now(timezone.utc), ) + + +# Преобразует одно декодированное сообщение Dzengi WebSocket +# в CandleCloseEvent. +class DzengiWebSocketOhlcAdapter: + def map_message( + self, + document: object, + *, + received_at: datetime | None = None, + ) -> CandleCloseEvent: + validated = validate_dzengi_websocket_ohlc_schema(document) + event = parse_dzengi_websocket_ohlc(validated) + validate_dzengi_websocket_ohlc_values(event) + + return map_dzengi_websocket_ohlc_to_candle_close_event( + event, + received_at=received_at or datetime.now(timezone.utc), + ) \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py b/app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py new file mode 100644 index 0000000..064873d --- /dev/null +++ b/app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py @@ -0,0 +1,112 @@ +# app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py + +from __future__ import annotations + +from datetime import datetime, timezone + +import pytest + +from src.market_data.acquisition.adapters.dzengi.websocket import ( + DzengiWebSocketOhlcAdapter, +) +from src.market_data.acquisition.exceptions import ( + CandleWebSocketValueError, +) + + +def test_adapter_maps_document_to_candle_close_event() -> None: + received_at = datetime(2026, 7, 13, tzinfo=timezone.utc) + + result = DzengiWebSocketOhlcAdapter().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 result.symbol == "BTC/USD" + assert result.interval == "1m" + assert result.candle_type == "classic" + + assert str(result.open_price) == "100" + assert str(result.high_price) == "105" + assert str(result.low_price) == "99" + assert str(result.close_price) == "103" + + assert result.received_at is received_at + assert result.source == "dzengi_websocket_ohlc" + + +def test_adapter_sets_received_at_when_missing() -> None: + result = DzengiWebSocketOhlcAdapter().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 result.received_at.tzinfo is timezone.utc + + +def test_adapter_preserves_value_validation_error() -> None: + with pytest.raises(CandleWebSocketValueError): + DzengiWebSocketOhlcAdapter().map_message( + { + "status": "OK", + "destination": "ohlc.event", + "payload": { + "symbol": "BTC/USD", + "interval": "1m", + "type": "classic", + "t": 1752364800000, + "o": "100", + "h": "90", + "l": "99", + "c": "95", + }, + } + ) + + +def test_adapter_runs_complete_pipeline() -> None: + result = DzengiWebSocketOhlcAdapter().map_message( + { + "status": "OK", + "destination": "ohlc.event", + "payload": { + "symbol": "BTC/USD", + "interval": "5m", + "type": "heikin-ashi", + "t": 1752364800000, + "o": "250", + "h": "260", + "l": "245", + "c": "255", + }, + } + ) + + assert result.interval == "5m" + assert result.candle_type == "heikin-ashi" + assert str(result.close_price) == "255" \ No newline at end of file diff --git a/docs/migrations/build_059_7.md b/docs/migrations/build_059_7.md new file mode 100644 index 0000000..9380ee0 --- /dev/null +++ b/docs/migrations/build_059_7.md @@ -0,0 +1,768 @@ +# Build 059.7 — WebSocket OHLC Adapter + +**Проект:** Dzentra +**Подсистема:** Market Data Acquisition +**Этап:** 059.7 +**Статус:** Completed + +--- + +# Цель + +Добавить верхнеуровневый адаптер обработки WebSocket OHLC-событий Dzengi, объединяющий все ранее реализованные стадии обработки в единый конвейер. + +Результатом работы адаптера должна стать полностью подготовленная внутренняя модель Dzentra: + +```text +CandleCloseEvent +``` + +Адаптер становится единственной точкой входа для преобразования WebSocket-сообщения в событие, пригодное для дальнейшей обработки внутри runtime. + +--- + +# Причина изменения + +После завершения Build 059.6 система уже содержала полный набор независимых компонентов обработки WebSocket OHLC. + +К этому моменту были реализованы: + +```text +Schema Validation +``` + +↓ + +```text +Parser +``` + +↓ + +```text +Value Validation +``` + +↓ + +```text +Mapper +``` + +Каждый компонент решал строго одну задачу и обладал собственным контрактом ответственности. + +Однако использование этих компонентов требовало их последовательного вызова вручную. + +Подобный подход приводил к нескольким архитектурным недостаткам: + +- отсутствовала единая точка входа; +- порядок выполнения этапов мог быть нарушен; +- различные части системы могли использовать разные последовательности обработки; +- возрастал риск появления дублирующего кода. + +Build 059.7 устраняет эти недостатки. + +Теперь вся последовательность обработки инкапсулирована внутри одного адаптера. + +--- + +# Реализовано + +Изменён файл: + +```text +src/market_data/acquisition/adapters/dzengi/websocket.py +``` + +Создан файл: + +```text +tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py +``` + +Добавлена документация: + +```text +docs/migrations/build_059_7.md +``` + +--- + +# Новый адаптер + +Добавлен класс: + +```text +DzengiWebSocketOhlcAdapter +``` + +Основной публичный метод: + +```text +map_message() +``` + +Назначение метода: + +```text +Raw WebSocket document + ↓ +CandleCloseEvent +``` + +Адаптер объединяет все ранее реализованные стадии обработки и предоставляет единый программный интерфейс для остальных компонентов системы. + +--- + +# Конвейер обработки + +После завершения Build 059.7 полный pipeline выглядит следующим образом: + +```text +Raw WebSocket message + ↓ +Schema Validation + ↓ +ValidatedWebSocketOhlcDocument + ↓ +Parser + ↓ +DzengiWebSocketOhlcEvent + ↓ +Value Validation + ↓ +Mapper + ↓ +CandleCloseEvent +``` + +Все стадии выполняются строго последовательно. + +Порядок выполнения не может быть изменён вызывающей стороной. + +--- + +# Последовательность обработки + +Внутри метода `map_message()` выполняются следующие действия. + +## 1. Schema Validation + +Выполняется проверка структуры входящего сообщения. + +Проверяется наличие обязательных полей транспортного контракта Dzengi: + +```text +status +destination +payload +``` + +После успешной проверки строится модель: + +```text +ValidatedWebSocketOhlcDocument +``` + +Ошибки данного этапа приводят к возникновению: + +```text +CandleWebSocketSchemaError +``` + +--- + +## 2. Parser + +Проверенный документ передаётся parser. + +Parser преобразует транспортную модель в специализированную внутреннюю модель интеграционного слоя: + +```text +DzengiWebSocketOhlcEvent +``` + +На этом этапе выполняется преобразование структуры сообщения. + +Проверка предметной корректности значений не производится. + +Ошибки parser приводят к возникновению: + +```text +CandleWebSocketParseError +``` + +--- + +## 3. Value Validation + +Полученная transport-модель проходит предметную проверку. + +Контролируются: + +- корректность временной метки; +- допустимость цен; +- допустимость типа свечи; +- непротиворечивость OHLC. + +При обнаружении ошибок формируется: + +```text +CandleWebSocketValueError +``` + +--- + +## 4. Mapper + +Последней стадией является mapper. + +Он преобразует: + +```text +DzengiWebSocketOhlcEvent +``` + +в: + +```text +CandleCloseEvent +``` + +При этом выполняются: + +- преобразование timestamp; +- преобразование цен в `Decimal`; +- установка `source`; +- перенос `received_at`. + +После завершения mapper транспортный контракт Dzengi полностью исчезает из дальнейшего pipeline. + +Все последующие компоненты работают исключительно с внутренней моделью Dzentra. + +--- + +# Обработка received_at + +Метод `map_message()` принимает дополнительный необязательный параметр: + +```text +received_at +``` + +Тип параметра: + +```text +datetime | None +``` + +Возможны два режима работы. + +--- + +## received_at передан + +Если вызывающая сторона уже зафиксировала момент получения сообщения, адаптер использует именно это значение. + +Никаких дополнительных преобразований времени не выполняется. + +Это позволяет сохранить единый момент получения данных на всём протяжении обработки сообщения. + +--- + +## received_at отсутствует + +Если параметр не передан, адаптер автоматически создаёт значение: + +```python +datetime.now(timezone.utc) +``` + +Таким образом каждое успешно обработанное событие всегда содержит корректное время получения. + +Использование UTC обеспечивает единый формат времени независимо от локальной временной зоны системы. + +--- + +# Контракт адаптера + +Build 059.7 не вводит новых транспортных моделей. + +Адаптер использует исключительно ранее реализованные компоненты. + +Контракт адаптера можно представить следующим образом: + +```text +Input: + +Raw WebSocket document +``` + +↓ + +```text +Output: + +CandleCloseEvent +``` + +При этом адаптер гарантирует строго фиксированную последовательность выполнения всех промежуточных стадий. + +--- + +# Контракт ошибок + +Адаптер намеренно не перехватывает исключения нижележащих компонентов. + +Каждая стадия обработки продолжает использовать собственный тип ошибки. + +Полный контракт выглядит следующим образом: + +```text +Transport + ↓ +CandleTransportError + +Schema Validation + ↓ +CandleWebSocketSchemaError + +Parser + ↓ +CandleWebSocketParseError + +Value Validation + ↓ +CandleWebSocketValueError + +Mapper + ↓ +CandleWebSocketMappingError +``` + +При возникновении любой ошибки выполнение pipeline немедленно прекращается. + +Исключение передаётся вызывающей стороне без изменения типа. + +Такой подход сохраняет независимость отдельных слоёв системы и упрощает диагностику ошибок. + +--- + +# Почему адаптер не перехватывает исключения + +Может показаться, что адаптер должен преобразовывать все ошибки в единое исключение. + +В Build 059.7 это намеренно не реализовано. + +Причины следующие. + +Во-первых, каждая стадия уже обладает собственным контрактом. + +Во-вторых, сохранение исходного типа исключения позволяет точно определить место возникновения ошибки. + +В-третьих, дополнительное оборачивание исключений не несёт новой информации и лишь усложняет диагностику. + +Поэтому адаптер выступает исключительно координатором последовательного выполнения pipeline. + +--- + +# Ответственность адаптера + +После завершения Build 059.7 ответственность адаптера ограничивается следующими задачами. + +Он обязан: + +- принять исходное WebSocket-сообщение; +- выполнить полный pipeline обработки; +- при необходимости создать `received_at`; +- вернуть `CandleCloseEvent`. + +На этом обязанности адаптера заканчиваются. + +--- + +# Что адаптер НЕ делает + +Build 059.7 намеренно не расширяет область ответственности адаптера. + +Адаптер не выполняет: + +- подключение к WebSocket; +- получение сообщений из сети; +- управление соединением; +- подписку на каналы; +- публикацию событий в runtime; +- REST reconciliation; +- получение объёма свечи; +- построение канонической модели `Candle`; +- повторную синхронизацию истории; +- обработку разрыва соединения. + +Все перечисленные задачи относятся к последующим этапам реализации. + +--- + +# Использование существующих компонентов + +Build 059.7 не реализует собственную бизнес-логику обработки свечей. + +Вместо этого адаптер объединяет ранее созданные компоненты: + +```text +Schema Validator +``` + +↓ + +```text +Parser +``` + +↓ + +```text +Value Validator +``` + +↓ + +```text +Mapper +``` + +Таким образом Build 059.7 не дублирует уже реализованную функциональность. + +Он только формирует единый программный интерфейс для использования существующего pipeline. + +--- + +# Архитектурное значение + +После появления адаптера остальные части системы больше не обязаны знать внутреннее устройство обработки WebSocket OHLC. + +Для них существует единственная операция: + +```text +Raw WebSocket message + ↓ +CandleCloseEvent +``` + +Вся остальная логика полностью скрыта внутри адаптера. + +Это соответствует принципу инкапсуляции и уменьшает связанность между подсистемами. + +--- + +# Unit Tests + +Создан файл: + +```text +tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py +``` + +Проверяются следующие сценарии. + +--- + +## Построение CandleCloseEvent + +Проверяется полный pipeline обработки. + +Исходное WebSocket-сообщение должно успешно пройти: + +```text +Schema Validation + ↓ +Parser + ↓ +Value Validation + ↓ +Mapper + ↓ +CandleCloseEvent +``` + +Проверяется корректность всех основных полей: + +- symbol; +- interval; +- candle_type; +- open_price; +- high_price; +- low_price; +- close_price; +- source; +- received_at. + +--- + +## Автоматическая установка received_at + +Если вызывающая сторона не передала параметр: + +```text +received_at +``` + +адаптер обязан автоматически установить текущее UTC-время. + +Тест подтверждает наличие значения и корректность его типа. + +--- + +## Сохранение ошибок Value Validation + +Если входное сообщение содержит недопустимые значения, адаптер не должен скрывать возникшее исключение. + +Проверяется, что наружу передаётся: + +```text +CandleWebSocketValueError +``` + +без изменения типа ошибки. + +--- + +## Проверка полного pipeline + +Отдельный тест подтверждает, что адаптер действительно выполняет все стадии обработки. + +В результате вызывающая сторона получает уже полностью построенный объект: + +```text +CandleCloseEvent +``` + +без необходимости самостоятельно вызывать parser, validator и mapper. + +--- + +Всего выполнено: + +```text +4 tests +``` + +--- + +# Проверка синтаксиса + +Выполнена команда: + +```bash +python -m compileall \ + src/market_data/acquisition/adapters/dzengi/websocket.py \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py +``` + +Результат: + +```text +успешно +``` + +--- + +# Целевые тесты + +Выполнена команда: + +```bash +python -m pytest -q \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py +``` + +Результат: + +```text +4 passed +``` + +--- + +# Регрессия адаптеров + +Выполнена команда: + +```bash +python -m pytest -q \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_quote_adapter.py \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py +``` + +Результат: + +```text +6 passed +``` + +--- + +# Регрессия полного WebSocket pipeline + +Выполнена команда: + +```bash +python -m pytest -q \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_quote_parser.py \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_quote_mapper.py \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_quote_adapter.py \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_parser.py \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_mapper.py \ + tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py +``` + +Результат: + +```text +65 passed +``` + +--- + +# Проверка форматирования + +Выполнена команда: + +```bash +git diff --check +``` + +Ошибок форматирования не обнаружено. + +--- + +# Изменённые файлы + +```text +src/market_data/acquisition/adapters/dzengi/websocket.py +``` + +Создан файл: + +```text +tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_adapter.py +``` + +Документация: + +```text +docs/migrations/build_059_7.md +``` + +--- + +# Совместимость + +Build 059.7 полностью обратно совместим. + +Не изменены: + +- REST Candles Feed; +- REST Candle mapper; +- Quote parser; +- Quote mapper; +- WebSocket Quote Adapter; +- ExchangeService; +- runtime; +- Trading; +- Market Analysis; +- каноническая модель `Candle`. + +Новый адаптер пока не встроен в существующий runtime и не влияет на работу действующего торгового бота. + +--- + +# Архитектурное значение + +Build 059.7 завершает формирование верхнего уровня обработки WebSocket OHLC. + +После предыдущих этапов система уже обладала всеми необходимыми компонентами: + +```text +Schema Validation +Parser +Value Validation +Mapper +``` + +Теперь появился единый компонент, объединяющий их в законченный конвейер. + +После Build 059.7 архитектура обработки выглядит следующим образом: + +```text +Raw WebSocket Message + │ + ▼ +WebSocket OHLC Adapter + │ + ▼ +Schema Validation + │ + ▼ +Parser + │ + ▼ +Value Validation + │ + ▼ +Mapper + │ + ▼ +CandleCloseEvent +``` + +Таким образом внешний код взаимодействует только с адаптером. + +Внутренние детали обработки полностью инкапсулированы. + +Это уменьшает связанность компонентов и соответствует архитектурным принципам Dzentra. + +--- + +# Ограничения этапа + +Build 059.7 намеренно **не реализует**: + +- подписку на WebSocket; +- управление соединением; +- автоматическое получение сообщений; +- публикацию событий в runtime; +- REST reconciliation; +- получение объёма свечи; +- построение канонической модели `Candle`; +- проверку соответствия REST и WebSocket данных; +- восстановление после разрыва соединения; +- повторную синхронизацию истории. + +Все перечисленные задачи относятся к следующим этапам серии Build 059. + +--- + +# Итог + +Build 059.7 добавляет полноценный верхнеуровневый адаптер обработки WebSocket OHLC. + +В результате система получила: + +- единую точку входа для обработки WebSocket OHLC; +- полностью инкапсулированный pipeline обработки; +- автоматическое создание `received_at`; +- единый интерфейс получения `CandleCloseEvent`; +- сохранение независимых контрактов ошибок всех стадий обработки; +- полный набор unit-тестов и успешное прохождение регрессии. + +После завершения Build 059.7 конвейер обработки WebSocket OHLC полностью готов к интеграции с механизмом подписки и дальнейшей публикации событий внутри runtime. + +--- + +# Следующий этап + +```text +Build 059.8 — Unified Adapter +``` + +На этом этапе будет реализован единый адаптер верхнего уровня, способный принимать различные типы сообщений Dzengi WebSocket, определять их тип и направлять в соответствующий специализированный адаптер, формируя унифицированную точку входа для всего WebSocket-потока. \ No newline at end of file