build 059.8: add unified websocket adapter
This commit is contained in:
@@ -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 (
|
||||
@@ -59,3 +62,43 @@ class DzengiWebSocketOhlcAdapter:
|
||||
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."
|
||||
)
|
||||
@@ -117,3 +117,9 @@ class CandleMappingError(MarketDataAcquisitionError):
|
||||
# Ошибка регистрации или получения Candles Feed.
|
||||
class CandleFeedRegistryError(MarketDataAcquisitionError):
|
||||
pass
|
||||
|
||||
|
||||
# Ошибка определения типа входящего WebSocket-сообщения
|
||||
# и выбора специализированного адаптера.
|
||||
class WebSocketMessageRoutingError(MarketDataAcquisitionError):
|
||||
pass
|
||||
@@ -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)
|
||||
525
docs/migrations/build_059_8.md
Normal file
525
docs/migrations/build_059_8.md
Normal file
@@ -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 с постоянным потоком биржевых сообщений.
|
||||
Reference in New Issue
Block a user