Build 060.4-060.5: add REST trade validation pipeline

This commit is contained in:
2026-07-18 19:32:34 +03:00
parent c0f69ee94a
commit a0ae22cb2e
6 changed files with 1039 additions and 3 deletions

View File

@@ -17,6 +17,7 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiMinNotionalFilter,
DzengiRateLimit,
DzengiRawNumeric,
DzengiRestAggTrade,
DzengiUnknownFilter,
DzengiTicker24hrResponse,
DzengiWebSocketOhlcEvent,
@@ -27,11 +28,13 @@ from src.market_data.acquisition.exceptions import (
CandleWebSocketParseError,
InstrumentReferenceParseError,
QuoteParseError,
TradeParseError,
)
from src.market_data.acquisition.validation.schema import (
ValidatedCandlesDocument,
ValidatedExchangeInfoDocument,
ValidatedQuoteDocument,
ValidatedRestAggTradesDocument,
ValidatedWebSocketOhlcDocument,
ValidatedWebSocketQuoteDocument,
)
@@ -966,3 +969,97 @@ def _candle_required_raw_numeric(
f"получен {type(value).__name__}."
)
return value
# Преобразовать структурно проверенный ответ aggTrades
# в транспортные модели адаптера Dzengi.
def parse_rest_agg_trades(
document: ValidatedRestAggTradesDocument,
) -> tuple[DzengiRestAggTrade, ...]:
"""
Преобразовать элементы проверенного документа aggTrades
в transport-модели Dzengi.
Функция не выполняет schema validation, предметную проверку значений,
преобразование числовых значений в Decimal или mapping
во внутреннюю модель Trade.
"""
return tuple(
_parse_rest_agg_trade_item(
item,
path=f"$[{index}]",
)
for index, item in enumerate(document.items)
)
def _parse_rest_agg_trade_item(
item: Mapping[str, object],
*,
path: str,
) -> DzengiRestAggTrade:
return DzengiRestAggTrade(
aggregate_trade_id=_trade_required_int(
item.get("a"),
path=f"{path}.a",
),
price=_trade_required_raw_numeric(
item.get("p"),
path=f"{path}.p",
),
quantity=_trade_required_raw_numeric(
item.get("q"),
path=f"{path}.q",
),
timestamp=_trade_required_int(
item.get("T"),
path=f"{path}.T",
),
buyer_is_maker=_trade_required_bool(
item.get("m"),
path=f"{path}.m",
),
)
def _trade_required_int(
value: object,
*,
path: str,
) -> int:
if isinstance(value, bool) or not isinstance(value, int):
raise TradeParseError(
f"{path} должен быть целым числом, "
f"получен {type(value).__name__}."
)
return value
def _trade_required_raw_numeric(
value: object,
*,
path: str,
) -> DzengiRawNumeric:
if isinstance(value, bool) or not isinstance(value, (str, int, float)):
raise TradeParseError(
f"{path} должен быть строкой или числом, "
f"получен {type(value).__name__}."
)
return value
def _trade_required_bool(
value: object,
*,
path: str,
) -> bool:
if not isinstance(value, bool):
raise TradeParseError(
f"{path} должен быть логическим значением, "
f"получен {type(value).__name__}."
)
return value

View File

@@ -124,6 +124,16 @@ class TradeSchemaError(MarketDataAcquisitionError):
pass
# Ошибка преобразования проверенного документа в raw-модели сделок.
class TradeParseError(MarketDataAcquisitionError):
pass
# Ошибка проверки допустимости значений raw-моделей сделок.
class TradeValueError(MarketDataAcquisitionError):
pass
# Ошибка определения типа входящего WebSocket-сообщения
# и выбора специализированного адаптера.
class WebSocketMessageRoutingError(MarketDataAcquisitionError):

View File

@@ -12,6 +12,7 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiMinNotionalFilter,
DzengiRateLimit,
DzengiRawNumeric,
DzengiRestAggTrade,
DzengiUnknownFilter,
DzengiKlinesResponse,
DzengiTicker24hrResponse,
@@ -23,6 +24,7 @@ from src.market_data.acquisition.exceptions import (
CandleWebSocketValueError,
InstrumentReferenceValueError,
QuoteValueError,
TradeValueError,
)
@@ -815,4 +817,91 @@ def _candle_decimal(
f"{path} должно быть конечным числом."
)
return decimal_value
return decimal_value
def validate_rest_agg_trade_values(
trades: tuple[DzengiRestAggTrade, ...],
) -> None:
"""
Проверить допустимость значений transport-моделей Dzengi aggTrades.
Функция не изменяет transport-модели, не преобразует raw numeric
значения в Decimal и не выполняет mapping во внутреннюю модель Trade.
"""
for index, trade in enumerate(trades):
_validate_rest_agg_trade(
trade,
path=f"$[{index}]",
)
def _validate_rest_agg_trade(
trade: DzengiRestAggTrade,
*,
path: str,
) -> None:
_trade_positive_int(
trade.aggregate_trade_id,
path=f"{path}.aggregateTradeId",
)
_trade_positive_decimal(
trade.price,
path=f"{path}.price",
)
_trade_positive_decimal(
trade.quantity,
path=f"{path}.quantity",
)
_trade_positive_int(
trade.timestamp,
path=f"{path}.timestamp",
)
def _trade_positive_int(
value: int,
*,
path: str,
) -> None:
if isinstance(value, bool) or value <= 0:
raise TradeValueError(
f"{path} должно быть целым числом больше нуля."
)
def _trade_positive_decimal(
value: DzengiRawNumeric,
*,
path: str,
) -> None:
decimal_value = _trade_decimal(
value,
path=path,
)
if decimal_value <= 0:
raise TradeValueError(
f"{path} должно быть больше нуля."
)
def _trade_decimal(
value: DzengiRawNumeric,
*,
path: str,
) -> Decimal:
try:
decimal_value = Decimal(str(value))
except (InvalidOperation, ValueError) as exc:
raise TradeValueError(
f"{path} должно быть корректным числом."
) from exc
if not decimal_value.is_finite():
raise TradeValueError(
f"{path} должно быть конечным числом."
)
return decimal_value

View File

@@ -9,16 +9,20 @@ import pytest
from src.market_data.acquisition.adapters.dzengi.models import (
DzengiLotSizeFilter,
DzengiMinNotionalFilter,
DzengiRestAggTrade,
DzengiUnknownFilter,
)
from src.market_data.acquisition.adapters.dzengi.parser import (
parse_exchange_info,
parse_rest_agg_trades,
)
from src.market_data.acquisition.exceptions import (
InstrumentReferenceParseError,
TradeParseError,
)
from src.market_data.acquisition.validation.schema import (
ValidatedExchangeInfoDocument,
ValidatedRestAggTradesDocument,
)
@@ -37,6 +41,17 @@ def _validated_document(
)
def _validated_rest_agg_trades_document(
items: list[dict[str, object]],
) -> ValidatedRestAggTradesDocument:
return ValidatedRestAggTradesDocument(
items=tuple(
MappingProxyType(dict(item))
for item in items
),
)
def _complete_symbol() -> dict[str, object]:
return {
"symbol": "ETH/EUR_LEVERAGE",
@@ -402,4 +417,229 @@ def test_parse_result_uses_immutable_sequences() -> None:
assert isinstance(response.payload.symbols, tuple)
assert isinstance(symbol.order_types, tuple)
assert isinstance(symbol.filters, tuple)
assert isinstance(symbol.market_modes, tuple)
assert isinstance(symbol.market_modes, tuple)
def test_parse_rest_agg_trades() -> None:
document = _validated_rest_agg_trades_document(
[
{
"a": 2134857062,
"p": "64497.25",
"q": "0.005",
"T": 1784218066823,
"m": False,
},
{
"a": 2134857063,
"p": 64498.10,
"q": 2,
"T": 1784218067000,
"m": True,
},
]
)
trades = parse_rest_agg_trades(document)
assert trades == (
DzengiRestAggTrade(
aggregate_trade_id=2134857062,
price="64497.25",
quantity="0.005",
timestamp=1784218066823,
buyer_is_maker=False,
),
DzengiRestAggTrade(
aggregate_trade_id=2134857063,
price=64498.10,
quantity=2,
timestamp=1784218067000,
buyer_is_maker=True,
),
)
def test_parse_empty_rest_agg_trades_document() -> None:
document = _validated_rest_agg_trades_document([])
trades = parse_rest_agg_trades(document)
assert trades == ()
def test_rest_agg_trade_parser_preserves_raw_numeric_values() -> None:
document = _validated_rest_agg_trades_document(
[
{
"a": 1,
"p": "1.2300",
"q": 5,
"T": 1000,
"m": False,
},
{
"a": 2,
"p": 1.25,
"q": "0.0100",
"T": 1001,
"m": True,
},
]
)
trades = parse_rest_agg_trades(document)
assert trades[0].price == "1.2300"
assert trades[0].quantity == 5
assert trades[1].price == 1.25
assert trades[1].quantity == "0.0100"
def test_rest_agg_trade_parser_returns_immutable_tuple() -> None:
document = _validated_rest_agg_trades_document(
[
{
"a": 1,
"p": "10",
"q": "2",
"T": 1000,
"m": False,
}
]
)
trades = parse_rest_agg_trades(document)
assert isinstance(trades, tuple)
assert isinstance(trades[0], DzengiRestAggTrade)
@pytest.mark.parametrize(
("field", "value"),
[
("a", None),
("a", "1"),
("a", 1.0),
("a", False),
("T", None),
("T", "1000"),
("T", 1000.0),
("T", True),
],
)
def test_reject_invalid_rest_agg_trade_integer_field(
field: str,
value: object,
) -> None:
item: dict[str, object] = {
"a": 1,
"p": "10",
"q": "2",
"T": 1000,
"m": False,
}
item[field] = value
document = _validated_rest_agg_trades_document([item])
with pytest.raises(
TradeParseError,
match=rf"\$\[0\]\.{field} должен быть целым числом",
):
parse_rest_agg_trades(document)
@pytest.mark.parametrize(
("field", "value"),
[
("p", None),
("p", []),
("p", {}),
("p", False),
("q", None),
("q", []),
("q", {}),
("q", True),
],
)
def test_reject_invalid_rest_agg_trade_raw_numeric_field(
field: str,
value: object,
) -> None:
item: dict[str, object] = {
"a": 1,
"p": "10",
"q": "2",
"T": 1000,
"m": False,
}
item[field] = value
document = _validated_rest_agg_trades_document([item])
with pytest.raises(
TradeParseError,
match=rf"\$\[0\]\.{field} должен быть строкой или числом",
):
parse_rest_agg_trades(document)
@pytest.mark.parametrize(
"value",
[
None,
0,
1,
"false",
[],
{},
],
)
def test_reject_invalid_rest_agg_trade_buyer_is_maker(
value: object,
) -> None:
document = _validated_rest_agg_trades_document(
[
{
"a": 1,
"p": "10",
"q": "2",
"T": 1000,
"m": value,
}
]
)
with pytest.raises(
TradeParseError,
match=r"\$\[0\]\.m должен быть логическим значением",
):
parse_rest_agg_trades(document)
def test_rest_agg_trade_parser_reports_item_index() -> None:
document = _validated_rest_agg_trades_document(
[
{
"a": 1,
"p": "10",
"q": "2",
"T": 1000,
"m": False,
},
{
"a": 2,
"p": [],
"q": "3",
"T": 1001,
"m": True,
},
]
)
with pytest.raises(
TradeParseError,
match=r"\$\[1\]\.p должен быть строкой или числом",
):
parse_rest_agg_trades(document)

View File

@@ -13,13 +13,16 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiLotSizeFilter,
DzengiMinNotionalFilter,
DzengiRateLimit,
DzengiRestAggTrade,
DzengiUnknownFilter,
)
from src.market_data.acquisition.exceptions import (
InstrumentReferenceValueError,
TradeValueError,
)
from src.market_data.acquisition.validation.values import (
validate_exchange_info_values,
validate_rest_agg_trade_values,
)
@@ -86,6 +89,23 @@ def _valid_response(
)
def _valid_rest_agg_trade(
*,
aggregate_trade_id: int = 2134857062,
price: str | int | float = "64497.25",
quantity: str | int | float = "0.005",
timestamp: int = 1784218066823,
buyer_is_maker: bool = False,
) -> DzengiRestAggTrade:
return DzengiRestAggTrade(
aggregate_trade_id=aggregate_trade_id,
price=price,
quantity=quantity,
timestamp=timestamp,
buyer_is_maker=buyer_is_maker,
)
def test_validate_complete_exchange_info_values() -> None:
response = _valid_response(
rate_limits=(
@@ -478,4 +498,190 @@ def test_reject_whitespace_global_exchange_filter_type() -> None:
InstrumentReferenceValueError,
match=r"filterType не должен состоять только из пробелов",
):
validate_exchange_info_values(response)
validate_exchange_info_values(response)
def test_validate_rest_agg_trade_values() -> None:
trades = (
_valid_rest_agg_trade(),
)
assert validate_rest_agg_trade_values(trades) is None
def test_validate_multiple_rest_agg_trade_values() -> None:
trades = (
_valid_rest_agg_trade(
aggregate_trade_id=1,
price="100.2500",
quantity="0.0100",
timestamp=1000,
buyer_is_maker=False,
),
_valid_rest_agg_trade(
aggregate_trade_id=2,
price=101.5,
quantity=2,
timestamp=1001,
buyer_is_maker=True,
),
)
assert validate_rest_agg_trade_values(trades) is None
def test_validate_empty_rest_agg_trade_values() -> None:
assert validate_rest_agg_trade_values(()) is None
@pytest.mark.parametrize(
("aggregate_trade_id", "timestamp"),
[
(0, 1000),
(-1, 1000),
(1, 0),
(1, -1),
],
)
def test_reject_non_positive_rest_agg_trade_integer_value(
aggregate_trade_id: int,
timestamp: int,
) -> None:
trades = (
_valid_rest_agg_trade(
aggregate_trade_id=aggregate_trade_id,
timestamp=timestamp,
),
)
with pytest.raises(
TradeValueError,
match="должно быть целым числом больше нуля",
):
validate_rest_agg_trade_values(trades)
@pytest.mark.parametrize(
("field", "value"),
[
("price", "0"),
("price", 0),
("price", 0.0),
("price", "-0.01"),
("price", -1),
("quantity", "0"),
("quantity", 0),
("quantity", 0.0),
("quantity", "-0.01"),
("quantity", -1),
],
)
def test_reject_non_positive_rest_agg_trade_numeric_value(
field: str,
value: str | int | float,
) -> None:
trade = replace(
_valid_rest_agg_trade(),
**{field: value},
)
with pytest.raises(
TradeValueError,
match=rf"\$\[0\]\.{field} должно быть больше нуля",
):
validate_rest_agg_trade_values((trade,))
@pytest.mark.parametrize(
("field", "value"),
[
("price", "not-a-number"),
("quantity", "not-a-number"),
],
)
def test_reject_invalid_rest_agg_trade_numeric_string(
field: str,
value: str,
) -> None:
trade = replace(
_valid_rest_agg_trade(),
**{field: value},
)
with pytest.raises(
TradeValueError,
match=rf"\$\[0\]\.{field} должно быть корректным числом",
):
validate_rest_agg_trade_values((trade,))
@pytest.mark.parametrize(
("field", "value"),
[
("price", float("nan")),
("price", float("inf")),
("price", float("-inf")),
("price", "NaN"),
("price", "Infinity"),
("price", "-Infinity"),
("quantity", float("nan")),
("quantity", float("inf")),
("quantity", float("-inf")),
("quantity", "NaN"),
("quantity", "Infinity"),
("quantity", "-Infinity"),
],
)
def test_reject_non_finite_rest_agg_trade_numeric_value(
field: str,
value: str | float,
) -> None:
trade = replace(
_valid_rest_agg_trade(),
**{field: value},
)
with pytest.raises(
TradeValueError,
match=rf"\$\[0\]\.{field} должно быть конечным числом",
):
validate_rest_agg_trade_values((trade,))
@pytest.mark.parametrize(
"buyer_is_maker",
[
False,
True,
],
)
def test_accept_rest_agg_trade_buyer_is_maker_values(
buyer_is_maker: bool,
) -> None:
trades = (
_valid_rest_agg_trade(
buyer_is_maker=buyer_is_maker,
),
)
assert validate_rest_agg_trade_values(trades) is None
def test_rest_agg_trade_value_error_reports_item_index() -> None:
trades = (
_valid_rest_agg_trade(
aggregate_trade_id=1,
timestamp=1000,
),
_valid_rest_agg_trade(
aggregate_trade_id=2,
quantity="0",
timestamp=1001,
),
)
with pytest.raises(
TradeValueError,
match=r"\$\[1\]\.quantity должно быть больше нуля",
):
validate_rest_agg_trade_values(trades)