Build 060.14 — WebSocket Trade Adapter

This commit is contained in:
2026-07-19 23:08:51 +03:00
parent b7bef85d20
commit bf7645237e
3 changed files with 1162 additions and 0 deletions

View File

@@ -0,0 +1,36 @@
# src/market_data/acquisition/adapters/dzengi/websocket_trade_adapter.py
from __future__ import annotations
from src.market_data.acquisition.adapters.dzengi.mapper import (
map_dzengi_websocket_trade_to_trade,
)
from src.market_data.acquisition.adapters.dzengi.parser import (
parse_dzengi_websocket_trade,
)
from src.market_data.acquisition.models.trade import Trade
from src.market_data.acquisition.validation.schema import (
ValidatedWebSocketTradeDocument,
)
from src.market_data.acquisition.validation.values import (
validate_dzengi_websocket_trade_values,
)
def adapt_websocket_trade_document(
document: ValidatedWebSocketTradeDocument,
) -> Trade:
"""
Преобразовать структурно проверенный документ Dzengi WebSocket Trade
в каноническую immutable-модель Trade.
Функция последовательно выполняет parsing, value validation
и mapping. Schema validation должна быть выполнена вызывающим слоем.
Исключения отдельных этапов не перехватываются и не оборачиваются.
"""
event = parse_dzengi_websocket_trade(document)
validate_dzengi_websocket_trade_values(event)
return map_dzengi_websocket_trade_to_trade(event)

View File

@@ -0,0 +1,162 @@
# tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_trade_adapter.py
from __future__ import annotations
from datetime import datetime, timezone
from decimal import Decimal
from types import MappingProxyType
import pytest
from src.market_data.acquisition.adapters.dzengi import (
websocket_trade_adapter,
)
from src.market_data.acquisition.adapters.dzengi.websocket_trade_adapter import (
adapt_websocket_trade_document,
)
from src.market_data.acquisition.exceptions import (
TradeMappingError,
TradeParseError,
TradeValueError,
)
from src.market_data.acquisition.models.trade import TradeAggressorSide
from src.market_data.acquisition.validation.schema import (
ValidatedWebSocketTradeDocument,
)
def test_adapt_websocket_trade_document_returns_canonical_trade() -> None:
document = ValidatedWebSocketTradeDocument(
payload=MappingProxyType(
{
"buyer": True,
"id": 101,
"orderId": "00a02503-0079-54c4-0000-000081e62b58",
"price": "43210.50",
"size": "0.125",
"symbol": " BTCUSDT ",
"ts": 1_700_000_000_123,
}
),
status="OK",
destination="internal.trade",
correlation_id=None,
)
result = adapt_websocket_trade_document(document)
assert result.symbol == "BTCUSDT"
assert result.trade_id == 101
assert result.price == Decimal("43210.50")
assert result.quantity == Decimal("0.125")
assert result.executed_at == datetime(
2023,
11,
14,
22,
13,
20,
123000,
tzinfo=timezone.utc,
)
assert result.aggressor_side is TradeAggressorSide.BUY
assert result.source == "dzengi_websocket_trade"
def test_adapt_websocket_trade_document_propagates_parse_error(
monkeypatch: pytest.MonkeyPatch,
) -> None:
document = ValidatedWebSocketTradeDocument(
payload=MappingProxyType({}),
status="OK",
destination="internal.trade",
correlation_id=None,
)
expected_error = TradeParseError("test parse error")
def raise_parse_error(
_: ValidatedWebSocketTradeDocument,
) -> object:
raise expected_error
monkeypatch.setattr(
websocket_trade_adapter,
"parse_dzengi_websocket_trade",
raise_parse_error,
)
with pytest.raises(TradeParseError) as exc_info:
adapt_websocket_trade_document(document)
assert exc_info.value is expected_error
def test_adapt_websocket_trade_document_propagates_value_error(
monkeypatch: pytest.MonkeyPatch,
) -> None:
document = ValidatedWebSocketTradeDocument(
payload=MappingProxyType({}),
status="OK",
destination="internal.trade",
correlation_id=None,
)
parsed_event = object()
expected_error = TradeValueError("test value error")
monkeypatch.setattr(
websocket_trade_adapter,
"parse_dzengi_websocket_trade",
lambda _: parsed_event,
)
def raise_value_error(_: object) -> None:
raise expected_error
monkeypatch.setattr(
websocket_trade_adapter,
"validate_dzengi_websocket_trade_values",
raise_value_error,
)
with pytest.raises(TradeValueError) as exc_info:
adapt_websocket_trade_document(document)
assert exc_info.value is expected_error
def test_adapt_websocket_trade_document_propagates_mapping_error(
monkeypatch: pytest.MonkeyPatch,
) -> None:
document = ValidatedWebSocketTradeDocument(
payload=MappingProxyType({}),
status="OK",
destination="internal.trade",
correlation_id=None,
)
parsed_event = object()
expected_error = TradeMappingError("test mapping error")
monkeypatch.setattr(
websocket_trade_adapter,
"parse_dzengi_websocket_trade",
lambda _: parsed_event,
)
monkeypatch.setattr(
websocket_trade_adapter,
"validate_dzengi_websocket_trade_values",
lambda _: None,
)
def raise_mapping_error(_: object) -> object:
raise expected_error
monkeypatch.setattr(
websocket_trade_adapter,
"map_dzengi_websocket_trade_to_trade",
raise_mapping_error,
)
with pytest.raises(TradeMappingError) as exc_info:
adapt_websocket_trade_document(document)
assert exc_info.value is expected_error

View File

@@ -0,0 +1,964 @@
# Build 060.14 — WebSocket Trade Adapter
**Engineering Migration Report**
---
# Контроль документа
| Свойство | Значение |
|----------|----------|
| Build | 060.14 |
| Название | WebSocket Trade Adapter |
| Статус | Completed |
| Проект | Dzentra |
| Подсистема | Market Data Acquisition |
| Компонент | Trades Feed |
| Версия | 1.0 |
---
# Цель Build
После завершения Build 060.13 система получила полностью реализованный уровень **WebSocket Trade Mapper**, выполняющий преобразование транспортной модели
```text
DzengiWebSocketTradeEvent
```
в каноническую модель предметной области
```text
Trade
```
На данном этапе Pipeline уже гарантирует:
- корректность структуры транспортного документа;
- успешное построение транспортной модели;
- корректность всех обязательных значений;
- успешное преобразование транспортной модели в каноническую модель предметной области.
Однако все реализованные уровни Pipeline по-прежнему существовали как независимые компоненты.
Для получения канонической модели сделки вызывающей стороне необходимо было самостоятельно выполнить последовательность операций:
```text
Parser
Value Validation
Mapper
```
Подобный подход противоречит одному из базовых архитектурных принципов Dzentra.
Внутренние компоненты системы не должны знать устройство транспортного Pipeline и порядок вызова его отдельных этапов.
Для решения данной задачи архитектура Dzentra предусматривает следующий обязательный уровень —
**Adapter**.
Именно Adapter объединяет все ранее реализованные уровни обработки в единую точку входа.
Build 060.14 реализует данный уровень для WebSocket Trade.
Основная задача Build — предоставить единственную функцию, принимающую уже прошедший Schema Validation документ
```text
ValidatedWebSocketTradeDocument
```
и возвращающую готовую каноническую модель
```text
Trade
```
выполняя внутри себя весь необходимый конвейер обработки.
При этом Build не затрагивает:
- Schema Validation;
- Parser;
- Value Validation;
- Mapper;
- WebSocket Runtime;
- Routing;
- Trades Feed.
---
# Предпосылки
К началу Build архитектура Market Data Acquisition уже содержала полностью реализованный Adapter для REST Aggregate Trade.
Он объединял несколько архитектурных уровней в единую функцию обработки.
Конвейер REST Trade имел следующий вид.
```text
ValidatedRestAggTradesDocument
Parser
Value Validation
Mapper
tuple[Trade]
```
Для WebSocket Trade после завершения Build 060.13 существовали все необходимые уровни Pipeline.
```text
ValidatedWebSocketTradeDocument
Parser
DzengiWebSocketTradeEvent
Value Validation
DzengiWebSocketTradeEvent
Mapper
Trade
```
Однако единая точка входа отсутствовала.
Каждый уровень необходимо было вызывать отдельно.
Таким образом архитектура WebSocket Trade оставалась незавершённой.
---
# Архитектурное основание
Одним из фундаментальных принципов архитектуры Dzentra является инкапсуляция транспортного Pipeline внутри слоя адаптера.
Внутренние компоненты системы должны работать исключительно с готовыми каноническими моделями предметной области.
Полный Pipeline обработки рыночных данных имеет следующий вид.
```text
Raw Source
Schema Validation
Validated Document
Parser
Transport Model
Value Validation
Transport Model
Mapper
Domain Model
```
Каждый уровень Pipeline отвечает исключительно за собственную область ответственности.
Однако вызывающий код не должен знать внутреннюю структуру данного конвейера.
Эта задача полностью возлагается на Adapter.
Для WebSocket Trade Adapter выполняет следующие действия.
---
```text
Parser
```
создаёт транспортную модель
```text
DzengiWebSocketTradeEvent
```
из прошедшего Schema Validation документа.
---
```text
Value Validation
```
подтверждает корректность всех значений транспортной модели.
---
```text
Mapper
```
преобразует транспортную модель в каноническую модель
```text
Trade
```
---
После завершения Adapter вызывающая сторона получает уже готовую каноническую модель предметной области и не взаимодействует с транспортными моделями либо отдельными уровнями Pipeline.
Adapter принципиально **не выполняет**:
- Schema Validation;
- Runtime-логику;
- маршрутизацию сообщений;
- бизнес-логику;
- принятие торговых решений.
Подобное разделение ответственности делает транспортный Pipeline полностью инкапсулированным и повторно используемым.
---
# Результаты архитектурного аудита
Перед реализацией Build был выполнен аудит существующей архитектуры Adapter.
В ходе анализа подтверждено наличие полностью реализованного Adapter для REST Aggregate Trade.
Его реализация использует единый архитектурный шаблон.
```text
Validated Document
Parser
Value Validation
Mapper
Canonical Model
```
Также аудит подтвердил наличие всех необходимых компонентов WebSocket Trade:
- Schema Validation;
- Parser;
- Value Validation;
- Mapper.
Все перечисленные уровни были полностью реализованы предыдущими Build серии 060.
При этом собственный Adapter для WebSocket Trade отсутствовал.
Таким образом единственным отсутствующим элементом транспортного Pipeline являлась единая точка входа, объединяющая все ранее реализованные уровни.
Build 060.14 полностью устраняет данный пробел и завершает архитектурное построение адаптера обработки WebSocket Trade.
---
# Архитектурное решение
По результатам проведённого аудита было принято решение полностью повторить архитектурный шаблон, уже используемый существующим REST Adapter.
В систему добавлена новая функция
```text
adapt_websocket_trade_document(...)
```
которая принимает документ
```text
ValidatedWebSocketTradeDocument
```
и внутри себя последовательно выполняет:
```text
Parser
Value Validation
Mapper
```
После успешного завершения обработки вызывающая сторона получает единственный результат —
```text
Trade
```
Конвейер WebSocket Trade принимает следующий вид.
```text
ValidatedWebSocketTradeDocument
WebSocket Trade Adapter
├────────► Parser
├────────► Value Validation
├────────► Mapper
Trade
```
Build 060.14 не изменяет архитектуру ранее реализованных компонентов.
Все существующие уровни Pipeline продолжают существовать как самостоятельные независимые компоненты.
Adapter лишь объединяет их в единую точку входа, полностью соответствующую архитектурным принципам Dzentra.
# Реализованный уровень Adapter
В файл
```text
src/market_data/acquisition/adapters/dzengi/websocket_trade_adapter.py
```
добавлена новая функция
```text
adapt_websocket_trade_document(...)
```
Функция принимает объект
```text
ValidatedWebSocketTradeDocument
```
и полностью инкапсулирует транспортный Pipeline обработки WebSocket Trade.
Во время выполнения Adapter последовательно вызывает:
1. Parser;
2. Value Validation;
3. Mapper.
После успешного завершения обработки вызывающая сторона получает готовую каноническую модель
```text
Trade
```
При этом промежуточные транспортные модели не покидают слой адаптера.
---
# Почему Adapter объединяет существующие уровни Pipeline
Во время архитектурного проектирования отдельно рассматривался вопрос о возможности прямого использования Parser, Value Validation и Mapper вызывающими компонентами системы.
По результатам анализа было принято решение отказаться от подобного подхода.
Основные причины:
- вызывающая сторона не должна знать внутреннюю структуру транспортного Pipeline;
- изменение последовательности этапов обработки не должно затрагивать внешний код;
- повторное использование Pipeline должно происходить через единую точку входа;
- транспортные модели должны оставаться внутренней деталью реализации слоя адаптера.
Поэтому Adapter становится единственным публичным способом получения канонической модели
```text
Trade
```
из прошедшего Schema Validation документа.
Подобное решение полностью соответствует базовому архитектурному принципу Dzentra — полной инкапсуляции транспортного Pipeline.
---
# Последовательность обработки
Adapter реализует строго фиксированную последовательность этапов.
```text
ValidatedWebSocketTradeDocument
Parser
DzengiWebSocketTradeEvent
Value Validation
DzengiWebSocketTradeEvent
Mapper
Trade
```
Ни один из этапов не может быть пропущен.
Каждый последующий уровень получает результат предыдущего уровня без изменения общей архитектуры Pipeline.
---
# Использование существующих компонентов
Build 060.14 не вводит новых механизмов обработки данных.
Adapter повторно использует уже существующие компоненты проекта.
Для построения транспортной модели используется:
```text
parse_dzengi_websocket_trade(...)
```
Для проверки корректности значений используется:
```text
validate_dzengi_websocket_trade_values(...)
```
Для построения канонической модели используется:
```text
map_dzengi_websocket_trade_to_trade(...)
```
Таким образом Adapter полностью основан на ранее реализованной инфраструктуре и не дублирует существующую функциональность.
---
# Проброс исключений
Adapter не выполняет собственную обработку ошибок.
Все исключения предыдущих уровней Pipeline пробрасываются вызывающей стороне без изменения.
При ошибке Parser генерируется
```text
TradeParseError
```
При ошибке проверки значений генерируется
```text
TradeValueError
```
При ошибке преобразования транспортной модели генерируется
```text
TradeMappingError
```
Подобный подход сохраняет единообразную архитектуру обработки ошибок во всей подсистеме Market Data Acquisition.
---
# Изменённые файлы
В рамках Build были изменены только два файла.
## Adapter
```text
src/market_data/acquisition/adapters/dzengi/websocket_trade_adapter.py
```
Добавлена функция
```text
adapt_websocket_trade_document(...)
```
Ранее реализованные Parser, Value Validation и Mapper не изменялись.
---
## Unit-тесты
```text
tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_trade_adapter.py
```
Добавлен полный набор unit-тестов нового Adapter.
---
# Добавленные тесты
В рамках Build реализованы четыре unit-теста, полностью покрывающие функциональность Adapter.
---
## Проверка успешной обработки
Подтверждается успешное получение канонической модели
```text
Trade
```
из корректного объекта
```text
ValidatedWebSocketTradeDocument
```
Одновременно подтверждается корректная работа полного Pipeline:
- Parser;
- Value Validation;
- Mapper.
---
## Проверка ошибок Parser
Отдельный тест подтверждает корректный проброс исключения
```text
TradeParseError
```
при ошибке построения транспортной модели.
---
## Проверка ошибок Value Validation
Отдельный тест подтверждает корректный проброс исключения
```text
TradeValueError
```
при обнаружении некорректных значений транспортной модели.
---
## Проверка ошибок Mapper
Отдельный тест подтверждает корректный проброс исключения
```text
TradeMappingError
```
при невозможности построения канонической модели.
---
# Результаты тестирования
После завершения реализации выполнен запуск целевого набора unit-тестов.
```bash
PYTHONPATH=. pytest -q \
tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_trade_adapter.py
```
Результат:
```text
4 passed in 0.02s
```
Все проверки новой функциональности успешно завершены.
Новый набор тестов полностью покрывает:
- успешное выполнение полного Pipeline;
- корректное получение канонической модели `Trade`;
- проброс ошибок Parser;
- проброс ошибок Value Validation;
- проброс ошибок Mapper.
---
# Регрессионное тестирование Adapter
После завершения реализации выполнен полный запуск тестов адаптеров Dzengi.
```bash
PYTHONPATH=. pytest -q \
tests/unit/market_data/acquisition/adapters/dzengi
```
Результат:
```text
301 passed in 0.11s
```
Регрессий существующих Adapter, Parser, Mapper и Validation не обнаружено.
---
# Регрессионное тестирование Market Data
После завершения реализации выполнен запуск полного набора unit-тестов подсистемы Market Data.
```bash
PYTHONPATH=. pytest -q \
tests/unit/market_data
```
Результат:
```text
913 passed in 0.31s
```
Подтверждено отсутствие регрессий во всех ранее реализованных компонентах подсистемы.
# Полное регрессионное тестирование
После завершения реализации выполнен полный запуск unit-тестов проекта.
```bash
PYTHONPATH=. pytest -q
```
Результат:
```text
1234 passed in 2.10s
```
Регрессий существующей функциональности не обнаружено.
Все ранее реализованные Build продолжают работать без каких-либо изменений.
Это подтверждает, что добавленная функциональность полностью изолирована и не влияет на существующие конвейеры обработки Quote, OHLC, REST Trade и остальные подсистемы проекта.
---
# Проверка компиляции
После завершения реализации выполнена полная проверка компиляции проекта.
```bash
python -m compileall src tests
```
Компиляция завершилась успешно.
Ошибок синтаксиса не обнаружено.
Все изменённые файлы успешно компилируются и не нарушают целостность проекта.
---
# Проверка Git diff
После завершения реализации выполнена финальная проверка изменений.
```bash
git diff --check
```
Результат:
```text
без замечаний
```
Проверка подтвердила отсутствие:
- trailing whitespace;
- ошибок окончания строк;
- конфликтов diff;
- нарушений форматирования.
---
# Scope Build 060.14
В рамках данного Build реализован исключительно уровень
```text
WebSocket Trade Adapter
```
Build **не включает**:
- WebSocket Runtime;
- Unified WebSocket Routing;
- Trades Feed;
- Dispatcher;
- Event Bus;
- Runtime Integration.
Подобное ограничение полностью соответствует принятому принципу атомарной реализации Build.
Каждый этап дорожной карты реализует только один архитектурный уровень Pipeline.
---
# Архитектурный результат
После завершения Build система содержит полностью реализованный Adapter WebSocket Trade.
Конвейер обработки принимает следующий вид.
```text
Raw WebSocket Trade
Schema Validation
ValidatedWebSocketTradeDocument
WebSocket Trade Adapter
├────────► Parser
├────────► Value Validation
├────────► Mapper
Trade
```
Таким образом Adapter полностью инкапсулирует транспортный Pipeline.
Внутренние компоненты системы получают уже готовую каноническую модель предметной области и не взаимодействуют с транспортными моделями либо отдельными этапами обработки.
---
# Состояние WebSocket Trade Pipeline
После завершения Build 060.14 конвейер имеет следующий вид.
```text
Raw WebSocket Trade
Schema Validation
ValidatedWebSocketTradeDocument
Adapter
├────────► Parser
├────────► Value Validation
├────────► Mapper
Trade
```
Статус реализации компонентов:
| Компонент | Build | Статус |
|-----------|-------|--------|
| Canonical Trade Model | 060.1 | ✔ Completed |
| WebSocket Trade Transport Model | 060.9 | ✔ Completed |
| WebSocket Trade Schema Validation | 060.10 | ✔ Completed |
| WebSocket Trade Parser | 060.11 | ✔ Completed |
| WebSocket Trade Value Validation | 060.12 | ✔ Completed |
| WebSocket Trade Mapper | 060.13 | ✔ Completed |
| WebSocket Trade Adapter | 060.14 | ✔ Completed |
| Unified WebSocket Routing | 060.15 | Pending |
---
# Соблюдение архитектурных принципов
В рамках Build полностью сохранены архитектурные инварианты Dzentra.
## Локальность изменений
Изменены только:
- `adapters/dzengi/websocket_trade_adapter.py`;
- unit-тесты нового Adapter.
Существующие Parser, Validation и Mapper не изменялись.
---
## Повторное использование архитектуры
Новая реализация полностью повторяет архитектурный шаблон существующего REST Adapter.
Новая архитектура не проектировалась.
Использован уже существующий подход, применяемый для остальных источников рыночных данных.
---
## Повторное использование инфраструктуры
Adapter полностью построен на ранее реализованных компонентах:
- Parser;
- Value Validation;
- Mapper.
Build не вводит новых механизмов обработки и не дублирует существующую функциональность.
Это обеспечивает единообразное поведение всех Adapter проекта.
---
## Разделение ответственности
Adapter отвечает исключительно за объединение ранее реализованных уровней Pipeline.
Build не выполняет:
- Schema Validation;
- Runtime Integration;
- маршрутизацию событий;
- бизнес-логику;
- принятие торговых решений.
Все перечисленные задачи остаются ответственностью других архитектурных уровней.
---
## Обратная совместимость
Существующая обработка:
- Quote;
- OHLC;
- REST Trade;
не изменилась.
Добавленная функциональность полностью изолирована и не оказывает влияния на ранее реализованные компоненты системы.
---
# Архитектурные решения Build (ADR)
## ADR-060.14-001
**Adapter является единственной публичной точкой входа транспортного Pipeline.**
Внутренние компоненты системы не должны самостоятельно вызывать Parser, Value Validation и Mapper.
---
## ADR-060.14-002
**Adapter не содержит собственной бизнес-логики.**
Его задача ограничивается последовательным вызовом существующих уровней Pipeline.
---
## ADR-060.14-003
**Adapter не обрабатывает исключения предыдущих уровней.**
Все ошибки Parser, Value Validation и Mapper пробрасываются вызывающей стороне без изменения.
Это сохраняет единый механизм обработки ошибок во всей подсистеме Market Data Acquisition.
---
## ADR-060.14-004
**Adapter повторно использует существующую инфраструктуру проекта.**
Build не создаёт новых механизмов обработки данных.
Все этапы обработки выполняются уже существующими компонентами системы.
---
# Критерии завершения Build
Build 060.14 считается завершённым, поскольку выполнены все поставленные задачи.
- ✔ реализована функция `adapt_websocket_trade_document()`;
- ✔ Adapter объединяет Parser, Value Validation и Mapper;
- ✔ реализована единая точка входа транспортного Pipeline;
- ✔ транспортный Pipeline полностью инкапсулирован;
- ✔ сохранено разделение ответственности между уровнями;
- ✔ исключения предыдущих уровней корректно пробрасываются;
- ✔ повторно использована существующая инфраструктура проекта;
- ✔ реализовано 4 unit-теста;
- ✔ целевой набор тестов успешно проходит;
- ✔ успешно пройдена регрессия Adapter;
- ✔ успешно пройдена регрессия Market Data;
- ✔ успешно пройдена полная регрессия проекта;
- ✔ проект успешно компилируется;
-`git diff --check` не выявил замечаний;
- ✔ изменения полностью укладываются в согласованный scope Build.
---
# Следующий этап
Следующим этапом дорожной карты является
```text
Build 060.15 — Unified WebSocket Routing
```
Цель следующего Build:
- подключить WebSocket Trade Adapter к существующему WebSocket Runtime;
- реализовать маршрутизацию сообщений по типам событий;
- передавать готовые объекты `Trade` в последующие уровни системы;
- подготовить основу для реализации полноценного Trades Feed.
---
# Итог
Build 060.14 завершил реализацию уровня **WebSocket Trade Adapter** и сделал транспортный Pipeline WebSocket Trade полностью инкапсулированным.
Новая реализация основана на уже существующих архитектурных принципах Dzentra, повторно использует Parser, Value Validation и Mapper, предоставляет единую точку входа для получения канонической модели `Trade` и сохраняет строгое разделение ответственности между архитектурными уровнями.
Build ограничен согласованным scope, успешно прошёл целевое и полное регрессионное тестирование, подтвердил отсутствие регрессий и завершил построение адаптера WebSocket Trade.
Следующим этапом развития серии является **Build 060.15 — Unified WebSocket Routing**, который интегрирует новый Adapter в общий WebSocket Runtime и завершит транспортный конвейер получения сделок в режиме реального времени.