build 059.1-059.2: add websocket OHLC transport and schema validation

This commit is contained in:
2026-07-16 22:14:25 +03:00
parent ebb578db35
commit e0aa04b32a
6 changed files with 1016 additions and 1 deletions

View File

@@ -148,3 +148,22 @@ class DzengiKline:
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
class DzengiKlinesResponse: class DzengiKlinesResponse:
items: tuple[DzengiKline, ...] 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

View File

@@ -78,6 +78,11 @@ class CandleSchemaError(MarketDataAcquisitionError):
pass pass
# Ошибка структуры события Candles Feed из WebSocket.
class CandleWebSocketSchemaError(MarketDataAcquisitionError):
pass
# Ошибка преобразования проверенного документа в raw-модели свечей. # Ошибка преобразования проверенного документа в raw-модели свечей.
class CandleParseError(MarketDataAcquisitionError): class CandleParseError(MarketDataAcquisitionError):
pass pass

View File

@@ -8,6 +8,7 @@ from typing import Mapping
from src.market_data.acquisition.exceptions import ( from src.market_data.acquisition.exceptions import (
CandleSchemaError, CandleSchemaError,
CandleWebSocketSchemaError,
InstrumentReferenceSchemaError, InstrumentReferenceSchemaError,
QuoteSchemaError, QuoteSchemaError,
) )
@@ -592,3 +593,133 @@ def _require_candle_mapping(
) )
return value 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

View File

@@ -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)

View File

@@ -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,
не затрагивая существующую архитектуру получения свечей.

View File

@@ -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`.