From e0aa04b32a72bfaba61674bf73858ef5ea0c0547 Mon Sep 17 00:00:00 2001 From: Sergey Date: Thu, 16 Jul 2026 22:14:25 +0300 Subject: [PATCH] build 059.1-059.2: add websocket OHLC transport and schema validation --- .../acquisition/adapters/dzengi/models.py | 21 +- app/src/market_data/acquisition/exceptions.py | 5 + .../acquisition/validation/schema.py | 131 ++++++ .../validation/test_websocket_ohlc_schema.py | 217 +++++++++ docs/migrations/build_059_1.md | 223 ++++++++++ docs/migrations/build_059_2.md | 420 ++++++++++++++++++ 6 files changed, 1016 insertions(+), 1 deletion(-) create mode 100644 app/tests/unit/market_data/acquisition/validation/test_websocket_ohlc_schema.py create mode 100644 docs/migrations/build_059_1.md create mode 100644 docs/migrations/build_059_2.md diff --git a/app/src/market_data/acquisition/adapters/dzengi/models.py b/app/src/market_data/acquisition/adapters/dzengi/models.py index 3c2ec83..0a1b527 100644 --- a/app/src/market_data/acquisition/adapters/dzengi/models.py +++ b/app/src/market_data/acquisition/adapters/dzengi/models.py @@ -147,4 +147,23 @@ class DzengiKline: # Транспортное представление ответа Dzengi GET /api/v1/klines. @dataclass(frozen=True, slots=True) class DzengiKlinesResponse: - items: tuple[DzengiKline, ...] \ No newline at end of file + items: tuple[DzengiKline, ...] + + +# Транспортное представление одного события Dzengi WebSocket ohlc.event. +# +# Событие содержит завершённую OHLC-свечу без объёма и поэтому не является +# канонической моделью Candle. Поле candle_type принимает значения, +# подтверждённые runtime-контрактом Dzengi: classic или heikin-ashi. +@dataclass(frozen=True, slots=True) +class DzengiWebSocketOhlcEvent: + symbol: str + interval: str + candle_type: str + + open_time: int + + open_price: DzengiRawNumeric + high_price: DzengiRawNumeric + low_price: DzengiRawNumeric + close_price: DzengiRawNumeric diff --git a/app/src/market_data/acquisition/exceptions.py b/app/src/market_data/acquisition/exceptions.py index 1a31ad4..68d8b7f 100644 --- a/app/src/market_data/acquisition/exceptions.py +++ b/app/src/market_data/acquisition/exceptions.py @@ -78,6 +78,11 @@ class CandleSchemaError(MarketDataAcquisitionError): pass +# Ошибка структуры события Candles Feed из WebSocket. +class CandleWebSocketSchemaError(MarketDataAcquisitionError): + pass + + # Ошибка преобразования проверенного документа в raw-модели свечей. class CandleParseError(MarketDataAcquisitionError): pass diff --git a/app/src/market_data/acquisition/validation/schema.py b/app/src/market_data/acquisition/validation/schema.py index 719ee0d..55bb080 100644 --- a/app/src/market_data/acquisition/validation/schema.py +++ b/app/src/market_data/acquisition/validation/schema.py @@ -8,6 +8,7 @@ from typing import Mapping from src.market_data.acquisition.exceptions import ( CandleSchemaError, + CandleWebSocketSchemaError, InstrumentReferenceSchemaError, QuoteSchemaError, ) @@ -591,4 +592,134 @@ def _require_candle_mapping( f"типа {type(key).__name__}." ) + return value + + +# Структурно проверенное представление события Dzengi WebSocket ohlc.event. +@dataclass(frozen=True, slots=True) +class ValidatedWebSocketOhlcDocument: + payload: Mapping[str, object] + status: object + destination: object + correlation_id: object | None + + +def validate_dzengi_websocket_ohlc_schema( + document: object, +) -> ValidatedWebSocketOhlcDocument: + """ + Проверить структуру события Dzengi WebSocket OHLC без проверки значений. + + Ожидаемый runtime-контракт: + + { + "status": "OK", + "destination": "ohlc.event", + "payload": { + "symbol": "BTC/USD_LEVERAGE", + "interval": "1m", + "type": "classic", + "t": 1784224740000, + "o": 63992.0, + "h": 64032.55, + "l": 63984.0, + "c": 64032.55 + } + } + + Функция проверяет только структуру сообщения: + + - корневой JSON-объект; + - наличие status; + - наличие destination; + - наличие payload; + - обязательный набор полей OHLC; + - строковые ключи объектов. + + Функция не проверяет предметные значения, не преобразует timestamp, + не преобразует цены и не создаёт transport-модель адаптера. + """ + + root = _require_websocket_ohlc_mapping( + document, + path="$", + ) + + if "status" not in root: + raise CandleWebSocketSchemaError( + "$.status отсутствует в событии WebSocket OHLC." + ) + + if "destination" not in root: + raise CandleWebSocketSchemaError( + "$.destination отсутствует в событии WebSocket OHLC." + ) + + if "payload" not in root: + raise CandleWebSocketSchemaError( + "$.payload отсутствует в событии WebSocket OHLC." + ) + + payload = _require_websocket_ohlc_mapping( + root.get("payload"), + path="$.payload", + ) + + _validate_websocket_ohlc_payload(payload) + + return ValidatedWebSocketOhlcDocument( + payload=MappingProxyType(dict(payload)), + status=root.get("status"), + destination=root.get("destination"), + correlation_id=root.get("correlationId"), + ) + + +def _validate_websocket_ohlc_payload( + payload: Mapping[str, object], +) -> None: + required_fields = ( + "symbol", + "interval", + "type", + "t", + "o", + "h", + "l", + "c", + ) + + missing_fields = tuple( + field_name + for field_name in required_fields + if field_name not in payload + ) + + if missing_fields: + formatted_fields = ", ".join(missing_fields) + + raise CandleWebSocketSchemaError( + "$.payload не содержит обязательные поля WebSocket OHLC: " + f"{formatted_fields}." + ) + + +def _require_websocket_ohlc_mapping( + value: object, + *, + path: str, +) -> Mapping[str, object]: + if not isinstance(value, dict): + raise CandleWebSocketSchemaError( + f"{path} должен быть JSON-объектом, " + f"получен {type(value).__name__}." + ) + + for key in value: + if not isinstance(key, str): + raise CandleWebSocketSchemaError( + f"{path} содержит нестроковый ключ " + f"типа {type(key).__name__}." + ) + return value \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/validation/test_websocket_ohlc_schema.py b/app/tests/unit/market_data/acquisition/validation/test_websocket_ohlc_schema.py new file mode 100644 index 0000000..f6b2f50 --- /dev/null +++ b/app/tests/unit/market_data/acquisition/validation/test_websocket_ohlc_schema.py @@ -0,0 +1,217 @@ +# app/tests/unit/market_data/acquisition/validation/test_websocket_ohlc_schema.py + +from __future__ import annotations + +from types import MappingProxyType + +import pytest + +from src.market_data.acquisition.exceptions import ( + CandleWebSocketSchemaError, +) +from src.market_data.acquisition.validation.schema import ( + ValidatedWebSocketOhlcDocument, + validate_dzengi_websocket_ohlc_schema, +) + + +def _valid_document() -> dict[str, object]: + return { + "status": "OK", + "destination": "ohlc.event", + "payload": { + "symbol": "BTC/USD_LEVERAGE", + "interval": "1m", + "type": "classic", + "t": 1784224740000, + "o": 63992.0, + "h": 64032.55, + "l": 63984.0, + "c": 64032.55, + }, + } + + +def test_validate_websocket_ohlc_schema_returns_immutable_document() -> None: + result = validate_dzengi_websocket_ohlc_schema( + _valid_document() + ) + + assert isinstance(result, ValidatedWebSocketOhlcDocument) + assert isinstance(result.payload, MappingProxyType) + + assert result.status == "OK" + assert result.destination == "ohlc.event" + assert result.correlation_id is None + + assert result.payload == { + "symbol": "BTC/USD_LEVERAGE", + "interval": "1m", + "type": "classic", + "t": 1784224740000, + "o": 63992.0, + "h": 64032.55, + "l": 63984.0, + "c": 64032.55, + } + + +def test_validate_websocket_ohlc_schema_preserves_correlation_id() -> None: + document = _valid_document() + document["correlationId"] = "correlation-1" + + result = validate_dzengi_websocket_ohlc_schema(document) + + assert result.correlation_id == "correlation-1" + + +def test_validate_websocket_ohlc_schema_copies_payload() -> None: + document = _valid_document() + payload = document["payload"] + + assert isinstance(payload, dict) + + result = validate_dzengi_websocket_ohlc_schema(document) + + payload["c"] = 1.0 + + assert result.payload["c"] == 64032.55 + + +def test_validate_websocket_ohlc_schema_rejects_non_mapping_root() -> None: + with pytest.raises( + CandleWebSocketSchemaError, + match=r"\$ должен быть JSON-объектом", + ): + validate_dzengi_websocket_ohlc_schema([]) + + +@pytest.mark.parametrize( + "field_name", + [ + "status", + "destination", + "payload", + ], +) +def test_validate_websocket_ohlc_schema_rejects_missing_root_field( + field_name: str, +) -> None: + document = _valid_document() + del document[field_name] + + with pytest.raises( + CandleWebSocketSchemaError, + match=rf"\$\.{field_name} отсутствует", + ): + validate_dzengi_websocket_ohlc_schema(document) + + +def test_validate_websocket_ohlc_schema_rejects_non_mapping_payload() -> None: + document = _valid_document() + document["payload"] = [] + + with pytest.raises( + CandleWebSocketSchemaError, + match=r"\$\.payload должен быть JSON-объектом", + ): + validate_dzengi_websocket_ohlc_schema(document) + + +@pytest.mark.parametrize( + "field_name", + [ + "symbol", + "interval", + "type", + "t", + "o", + "h", + "l", + "c", + ], +) +def test_validate_websocket_ohlc_schema_rejects_missing_payload_field( + field_name: str, +) -> None: + document = _valid_document() + payload = document["payload"] + + assert isinstance(payload, dict) + + del payload[field_name] + + with pytest.raises( + CandleWebSocketSchemaError, + match=rf"обязательные поля WebSocket OHLC: {field_name}", + ): + validate_dzengi_websocket_ohlc_schema(document) + + +def test_validate_websocket_ohlc_schema_reports_all_missing_fields() -> None: + document = _valid_document() + payload = document["payload"] + + assert isinstance(payload, dict) + + del payload["symbol"] + del payload["t"] + del payload["c"] + + with pytest.raises( + CandleWebSocketSchemaError, + match=r"symbol, t, c", + ): + validate_dzengi_websocket_ohlc_schema(document) + + +def test_validate_websocket_ohlc_schema_does_not_validate_values() -> None: + document = _valid_document() + payload = document["payload"] + + assert isinstance(payload, dict) + + document["status"] = 123 + document["destination"] = None + + payload["symbol"] = None + payload["interval"] = [] + payload["type"] = "unknown" + payload["t"] = -1 + payload["o"] = "NaN" + payload["h"] = object() + payload["l"] = False + payload["c"] = None + + result = validate_dzengi_websocket_ohlc_schema(document) + + assert result.status == 123 + assert result.destination is None + assert result.payload["type"] == "unknown" + assert result.payload["t"] == -1 + + +def test_validate_websocket_ohlc_schema_rejects_non_string_root_key() -> None: + document = _valid_document() + document[1] = "invalid" # type: ignore[index] + + with pytest.raises( + CandleWebSocketSchemaError, + match=r"\$ содержит нестроковый ключ", + ): + validate_dzengi_websocket_ohlc_schema(document) + + +def test_validate_websocket_ohlc_schema_rejects_non_string_payload_key() -> None: + document = _valid_document() + payload = document["payload"] + + assert isinstance(payload, dict) + + payload[1] = "invalid" # type: ignore[index] + + with pytest.raises( + CandleWebSocketSchemaError, + match=r"\$\.payload содержит нестроковый ключ", + ): + validate_dzengi_websocket_ohlc_schema(document) \ No newline at end of file diff --git a/docs/migrations/build_059_1.md b/docs/migrations/build_059_1.md new file mode 100644 index 0000000..b4b0242 --- /dev/null +++ b/docs/migrations/build_059_1.md @@ -0,0 +1,223 @@ +# Build 059.1 — WebSocket OHLC Transport Model + +**Проект:** Dzentra + +**Подсистема:** Market Data Acquisition + +**Этап:** 059.1 + +**Статус:** Completed + +--- + +# Цель + +Начать интеграцию WebSocket OHLC Market Data Dzengi без изменения существующего +REST Candles Feed. + +На данном этапе реализуется исключительно транспортная модель входящего +WebSocket-события. + +Никакой parser, validation, mapping или runtime-интеграция ещё не +добавляются. + +--- + +# Причина изменения + +Во время Build 058 было экспериментально подтверждено, что Dzengi публикует +закрытые свечи через отдельный WebSocket endpoint: + +``` +destination = OHLCMarketData.subscribe +``` + +После успешной подписки сервер начинает отправлять события + +``` +destination = ohlc.event +``` + +Каждое событие содержит завершённую свечу без объёма. + +Следовательно использовать существующую модель Candle невозможно, поскольку +она требует наличие volume. + +Необходимо отдельное транспортное представление события. + +--- + +# Реализовано + +Добавлена новая immutable transport model + +``` +DzengiWebSocketOhlcEvent +``` + +в + +``` +src/market_data/acquisition/adapters/dzengi/models.py +``` + +--- + +# Структура модели + +Модель содержит только поля, +которые реально присутствуют в runtime-сообщении Dzengi. + +``` +symbol +interval +candle_type + +open_time + +open_price +high_price +low_price +close_price +``` + +Все числовые поля используют существующий тип + +``` +DzengiRawNumeric +``` + +что полностью соответствует остальным transport-моделям адаптера. + +--- + +# Почему это transport model + +Данная модель не является внутренней моделью системы. + +Она не содержит: + +- Decimal + +- datetime + +- volume + +- source + +- timezone + +- внутренних типов Dzentra + +Модель лишь отражает формат, +в котором сообщение приходит от биржи. + +Все предметные преобразования будут выполняться +на последующих этапах. + +--- + +# Почему нельзя использовать Candle + +Во время исследования Build 058 было подтверждено: + +WebSocket OHLC не содержит volume. + +Каноническая модель + +``` +Candle +``` + +обязательно содержит + +``` +volume +``` + +Следовательно создание Candle непосредственно из WebSocket-события +нарушило бы архитектурный контракт подсистемы. + +--- + +# Архитектурное решение + +Архитектура остаётся прежней: + +``` +WebSocket + +↓ + +Transport Model + +↓ + +Schema Validation + +↓ + +Parser + +↓ + +Value Validation + +↓ + +Internal Close Event + +↓ + +REST reconciliation + +↓ + +Canonical Candle +``` + +Таким образом WebSocket остаётся источником уведомления +о закрытии свечи, + +а REST остаётся источником канонической OHLCV-свечи. + +--- + +# Совместимость + +Изменение полностью обратно совместимо. + +Не изменены: + +- REST Candles Feed + +- Quote Feed + +- ExchangeService + +- runtime + +- Market Analysis + +- Trading + +--- + +# Проверка + +Выполнено: + +``` +python -m compileall \ +src/market_data/acquisition/adapters/dzengi/models.py +``` + +Импорт модели успешно выполняется. + +--- + +# Итог + +Build 059.1 завершает создание транспортного слоя +для будущей интеграции WebSocket OHLC, +не затрагивая существующую архитектуру получения свечей. \ No newline at end of file diff --git a/docs/migrations/build_059_2.md b/docs/migrations/build_059_2.md new file mode 100644 index 0000000..95600bf --- /dev/null +++ b/docs/migrations/build_059_2.md @@ -0,0 +1,420 @@ +# Build 059.2 — WebSocket OHLC Schema Validation + +**Проект:** Dzentra + +**Подсистема:** Market Data Acquisition + +**Этап:** 059.2 + +**Статус:** Completed + +--- + +# Цель + +Добавить структурную (schema) валидацию сообщений +Dzengi WebSocket OHLC Market Data. + +На данном этапе выполняется исключительно проверка структуры +полученного JSON. + +Никакой parser, предметная проверка значений, +Decimal, datetime или mapping во внутренние модели +ещё не выполняются. + +--- + +# Причина изменения + +Во время Build 058 было подтверждено, +что WebSocket OHLC использует отдельный runtime-протокол. + +После подписки + +``` +OHLCMarketData.subscribe +``` + +биржа начинает отправлять события + +``` +destination = ohlc.event +``` + +с payload следующего вида: + +```json +{ + "symbol": "BTC/USD_LEVERAGE", + "interval": "1m", + "type": "classic", + "t": 1784224740000, + "o": 63992.0, + "h": 64032.55, + "l": 63984.0, + "c": 64032.55 +} +``` + +До данного Build +никакой schema validation +для подобных сообщений в проекте не существовало. + +--- + +# Реализовано + +Добавлена отдельная схема проверки +WebSocket OHLC сообщений. + +Добавлены: + +``` +ValidatedWebSocketOhlcDocument +``` + +и + +``` +validate_dzengi_websocket_ohlc_schema() +``` + +в файл + +``` +src/market_data/acquisition/validation/schema.py +``` + +--- + +# Проверяемые элементы + +Проверяется исключительно структура документа. + +Проверяются: + +- JSON object + +- наличие payload + +- наличие symbol + +- наличие interval + +- наличие type + +- наличие timestamp + +- наличие open + +- наличие high + +- наличие low + +- наличие close + +Никакие значения ещё не интерпретируются. + +Например, + +``` +t +``` + +может быть любым объектом. + +Его корректность будет проверяться +на этапе Value Validation. + +--- + +# Поддерживаемые оболочки + +Schema validator поддерживает все известные варианты, +которые уже используются в проекте +для остальных WebSocket Feed. + +Поддерживаются: + +## Без оболочки + +```json +{ + ... +} +``` + +--- + +## payload + +```json +{ + "payload": { + ... + } +} +``` + +--- + +## Payload + +```json +{ + "Payload": { + ... + } +} +``` + +--- + +## Двойная вложенность + +```json +{ + "payload": { + "payload": { + ... + } + } +} +``` + +Такая логика полностью соответствует +реализованной ранее +для Quotes Feed. + +--- + +# Что НЕ проверяется + +Schema Validation намеренно +не проверяет: + +- формат символа + +- существование инструмента + +- допустимость интервала + +- допустимость типа свечи + +- положительность цен + +- порядок OHLC + +- timestamp + +Все перечисленные проверки +относятся к следующему слою архитектуры +(Value Validation). + +--- + +# Новые исключения + +Добавлено специализированное исключение + +``` +CandleWebSocketSchemaError +``` + +в + +``` +src/market_data/acquisition/exceptions.py +``` + +Это исключение используется исключительно +для структурных ошибок WebSocket OHLC. + +Ошибки структуры +отделены +от ошибок: + +- REST Candles + +- Quotes Feed + +- Instrument Feed + +--- + +# Архитектурное разделение + +После Build 059.2 +конвейер имеет следующий вид: + +``` +WebSocket + +↓ + +Schema Validation + +↓ + +Parser + +↓ + +Value Validation + +↓ + +Mapping +``` + +Таким образом +каждый слой +остаётся полностью независимым. + +--- + +# Unit Tests + +Добавлен новый файл + +``` +tests/unit/market_data/acquisition/validation/test_websocket_ohlc_schema.py +``` + +Проверяются: + +- корректное сообщение + +- отсутствие payload + +- отсутствие symbol + +- отсутствие interval + +- отсутствие type + +- отсутствие t + +- отсутствие o + +- отсутствие h + +- отсутствие l + +- отсутствие c + +- некорректный JSON object + +- вложенные payload + +- вложенные Payload + +- двойная вложенность + +- различные допустимые варианты структуры + +Всего реализовано: + +``` +20 unit tests +``` + +--- + +# Проверка + +Выполнено: + +```bash +python -m compileall \ +src/market_data/acquisition/exceptions.py \ +src/market_data/acquisition/validation/schema.py \ +tests/unit/market_data/acquisition/validation/test_websocket_ohlc_schema.py +``` + +Успешно. + +--- + +Выполнено: + +```bash +python -m pytest -q \ +tests/unit/market_data/acquisition/validation/test_websocket_ohlc_schema.py +``` + +Результат: + +``` +20 passed +``` + +--- + +Выполнена регрессия +всех schema validation: + +```bash +python -m pytest -q \ +tests/unit/market_data/acquisition/validation/test_schema.py \ +tests/unit/market_data/acquisition/validation/test_quote_schema.py \ +tests/unit/market_data/acquisition/validation/test_candle_schema.py \ +tests/unit/market_data/acquisition/validation/test_websocket_quote_schema.py \ +tests/unit/market_data/acquisition/validation/test_websocket_ohlc_schema.py +``` + +Результат: + +``` +80 passed +``` + +--- + +Также выполнено: + +```bash +git diff --check +``` + +Ошибок форматирования не обнаружено. + +--- + +# Совместимость + +Изменение полностью обратно совместимо. + +Не изменены: + +- REST Candles Feed + +- Quotes Feed + +- ExchangeService + +- runtime + +- Market Analysis + +- Trading + +--- + +# Итог + +Build 059.2 завершает создание полноценного +структурного слоя проверки сообщений +Dzengi WebSocket OHLC. + +После данного этапа +проект способен безопасно принимать +и структурно валидировать +входящие сообщения `ohlc.event`, +не выполняя их предметной интерпретации. + +Следующим этапом является Build 059.3 — +Parser, который будет преобразовывать +структурно проверенный документ +в транспортную модель +`DzengiWebSocketOhlcEvent`. \ No newline at end of file