build 059.3: parse websocket OHLC events

This commit is contained in:
2026-07-16 23:21:14 +03:00
parent e0aa04b32a
commit c8c0bb6951
4 changed files with 896 additions and 0 deletions

View File

@@ -19,10 +19,12 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiRawNumeric,
DzengiUnknownFilter,
DzengiTicker24hrResponse,
DzengiWebSocketOhlcEvent,
DzengiWebSocketQuoteResponse,
)
from src.market_data.acquisition.exceptions import (
CandleParseError,
CandleWebSocketParseError,
InstrumentReferenceParseError,
QuoteParseError,
)
@@ -30,6 +32,7 @@ from src.market_data.acquisition.validation.schema import (
ValidatedCandlesDocument,
ValidatedExchangeInfoDocument,
ValidatedQuoteDocument,
ValidatedWebSocketOhlcDocument,
ValidatedWebSocketQuoteDocument,
)
@@ -719,6 +722,101 @@ def _websocket_optional_timestamp(
return _quote_required_int(value, path=path)
# Преобразовать структурно проверенное событие WebSocket OHLC
# в транспортную модель адаптера Dzengi.
def parse_dzengi_websocket_ohlc(
document: ValidatedWebSocketOhlcDocument,
) -> DzengiWebSocketOhlcEvent:
"""
Преобразовать проверенный payload события ``ohlc.event``
в transport-модель Dzengi.
Функция не выполняет предметную проверку значений, не проверяет
допустимость интервала или типа свечи, не преобразует timestamp
в datetime и не выполняет mapping во внутреннюю модель Dzentra.
"""
payload = document.payload
return DzengiWebSocketOhlcEvent(
symbol=_websocket_ohlc_required_string(
payload.get("symbol"),
path="$.payload.symbol",
),
interval=_websocket_ohlc_required_string(
payload.get("interval"),
path="$.payload.interval",
),
candle_type=_websocket_ohlc_required_string(
payload.get("type"),
path="$.payload.type",
),
open_time=_websocket_ohlc_required_int(
payload.get("t"),
path="$.payload.t",
),
open_price=_websocket_ohlc_required_raw_numeric(
payload.get("o"),
path="$.payload.o",
),
high_price=_websocket_ohlc_required_raw_numeric(
payload.get("h"),
path="$.payload.h",
),
low_price=_websocket_ohlc_required_raw_numeric(
payload.get("l"),
path="$.payload.l",
),
close_price=_websocket_ohlc_required_raw_numeric(
payload.get("c"),
path="$.payload.c",
),
)
def _websocket_ohlc_required_string(
value: object,
*,
path: str,
) -> str:
if not isinstance(value, str):
raise CandleWebSocketParseError(
f"{path} должен быть строкой, "
f"получен {type(value).__name__}."
)
return value
def _websocket_ohlc_required_int(
value: object,
*,
path: str,
) -> int:
if isinstance(value, bool) or not isinstance(value, int):
raise CandleWebSocketParseError(
f"{path} должен быть целым числом, "
f"получен {type(value).__name__}."
)
return value
def _websocket_ohlc_required_raw_numeric(
value: object,
*,
path: str,
) -> DzengiRawNumeric:
if isinstance(value, bool) or not isinstance(value, (str, int, float)):
raise CandleWebSocketParseError(
f"{path} должен быть строкой или числом, "
f"получен {type(value).__name__}."
)
return value
# Преобразовать структурно проверенный ответ klines в raw-модели Dzengi.
def parse_candles(
document: ValidatedCandlesDocument,

View File

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

View File

@@ -0,0 +1,200 @@
# app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_parser.py
from __future__ import annotations
from types import MappingProxyType
import pytest
from src.market_data.acquisition.adapters.dzengi.models import (
DzengiWebSocketOhlcEvent,
)
from src.market_data.acquisition.adapters.dzengi.parser import (
parse_dzengi_websocket_ohlc,
)
from src.market_data.acquisition.exceptions import (
CandleWebSocketParseError,
)
from src.market_data.acquisition.validation.schema import (
ValidatedWebSocketOhlcDocument,
)
def _validated_document(
**payload_overrides: object,
) -> ValidatedWebSocketOhlcDocument:
payload: dict[str, object] = {
"symbol": "BTC/USD_LEVERAGE",
"interval": "1m",
"type": "classic",
"t": 1784224740000,
"o": 63992.0,
"h": 64032.55,
"l": 63984.0,
"c": 64032.55,
}
payload.update(payload_overrides)
return ValidatedWebSocketOhlcDocument(
payload=MappingProxyType(payload),
status="OK",
destination="ohlc.event",
correlation_id=None,
)
def test_parse_websocket_ohlc_returns_transport_model() -> None:
result = parse_dzengi_websocket_ohlc(
_validated_document()
)
assert result == DzengiWebSocketOhlcEvent(
symbol="BTC/USD_LEVERAGE",
interval="1m",
candle_type="classic",
open_time=1784224740000,
open_price=63992.0,
high_price=64032.55,
low_price=63984.0,
close_price=64032.55,
)
def test_parse_websocket_ohlc_preserves_heikin_ashi_type() -> None:
result = parse_dzengi_websocket_ohlc(
_validated_document(type="heikin-ashi")
)
assert result.candle_type == "heikin-ashi"
def test_parse_websocket_ohlc_preserves_raw_numeric_types() -> None:
result = parse_dzengi_websocket_ohlc(
_validated_document(
o="63992.00",
h=64033,
l=63984.0,
c="64032.55",
)
)
assert result.open_price == "63992.00"
assert result.high_price == 64033
assert result.low_price == 63984.0
assert result.close_price == "64032.55"
@pytest.mark.parametrize(
("field_name", "value"),
[
("symbol", None),
("symbol", 123),
("interval", None),
("interval", []),
("type", None),
("type", True),
],
)
def test_parse_websocket_ohlc_rejects_non_string_fields(
field_name: str,
value: object,
) -> None:
with pytest.raises(
CandleWebSocketParseError,
match=rf"\$\.payload\.{field_name} должен быть строкой",
):
parse_dzengi_websocket_ohlc(
_validated_document(**{field_name: value})
)
@pytest.mark.parametrize(
"value",
[
None,
True,
1784224740000.0,
"1784224740000",
[],
],
)
def test_parse_websocket_ohlc_rejects_invalid_open_time(
value: object,
) -> None:
with pytest.raises(
CandleWebSocketParseError,
match=r"\$\.payload\.t должен быть целым числом",
):
parse_dzengi_websocket_ohlc(
_validated_document(t=value)
)
@pytest.mark.parametrize(
"field_name",
[
"o",
"h",
"l",
"c",
],
)
@pytest.mark.parametrize(
"value",
[
None,
True,
[],
{},
object(),
],
)
def test_parse_websocket_ohlc_rejects_invalid_raw_numeric(
field_name: str,
value: object,
) -> None:
with pytest.raises(
CandleWebSocketParseError,
match=rf"\$\.payload\.{field_name} должен быть строкой или числом",
):
parse_dzengi_websocket_ohlc(
_validated_document(**{field_name: value})
)
def test_parse_websocket_ohlc_does_not_validate_subject_values() -> None:
result = parse_dzengi_websocket_ohlc(
_validated_document(
symbol="",
interval="unsupported",
type="unknown",
t=-1,
o="NaN",
h=-10,
l=100,
c=0,
)
)
assert result.symbol == ""
assert result.interval == "unsupported"
assert result.candle_type == "unknown"
assert result.open_time == -1
assert result.open_price == "NaN"
assert result.high_price == -10
assert result.low_price == 100
assert result.close_price == 0
def test_parse_websocket_ohlc_ignores_envelope_metadata() -> None:
document = ValidatedWebSocketOhlcDocument(
payload=_validated_document().payload,
status=123,
destination=None,
correlation_id={"unexpected": "value"},
)
result = parse_dzengi_websocket_ohlc(document)
assert result.symbol == "BTC/USD_LEVERAGE"
assert result.interval == "1m"