Build 060.6: add REST trade mapper

This commit is contained in:
2026-07-19 09:49:01 +03:00
parent a0ae22cb2e
commit 6b0d5badce
4 changed files with 2599 additions and 2 deletions

View File

@@ -13,6 +13,7 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiLotSizeFilter, DzengiLotSizeFilter,
DzengiMinNotionalFilter, DzengiMinNotionalFilter,
DzengiRawNumeric, DzengiRawNumeric,
DzengiRestAggTrade,
DzengiTicker24hrResponse, DzengiTicker24hrResponse,
DzengiWebSocketOhlcEvent, DzengiWebSocketOhlcEvent,
DzengiWebSocketQuoteResponse, DzengiWebSocketQuoteResponse,
@@ -22,11 +23,16 @@ from src.market_data.acquisition.exceptions import (
CandleWebSocketMappingError, CandleWebSocketMappingError,
InstrumentReferenceMappingError, InstrumentReferenceMappingError,
QuoteMappingError, QuoteMappingError,
TradeMappingError,
) )
from src.market_data.acquisition.models.candle import Candle from src.market_data.acquisition.models.candle import Candle
from src.market_data.acquisition.models.candle_close import CandleCloseEvent from src.market_data.acquisition.models.candle_close import CandleCloseEvent
from src.market_data.acquisition.models.instrument import Instrument from src.market_data.acquisition.models.instrument import Instrument
from src.market_data.acquisition.models.quote import Quote from src.market_data.acquisition.models.quote import Quote
from src.market_data.acquisition.models.trade import (
Trade,
TradeAggressorSide,
)
_DZENGI_SOURCE_NAME = "dzengi" _DZENGI_SOURCE_NAME = "dzengi"
@@ -506,4 +512,98 @@ def _required_candle_decimal(
f"Поле {field_name} свечи должно быть конечным числом." f"Поле {field_name} свечи должно быть конечным числом."
) )
return result return result
def map_dzengi_rest_agg_trades_to_trades(
trades: tuple[DzengiRestAggTrade, ...],
*,
symbol: str,
) -> tuple[Trade, ...]:
"""
Преобразовать проверенные transport-модели Dzengi aggTrades
в канонический immutable-набор Trade.
Функция предполагает, что до mapper уже были выполнены:
schema validation, parsing и value validation.
"""
return tuple(
_map_dzengi_rest_agg_trade(
trade,
symbol=symbol,
)
for trade in trades
)
def _map_dzengi_rest_agg_trade(
trade: DzengiRestAggTrade,
*,
symbol: str,
) -> Trade:
return Trade(
symbol=symbol.strip(),
trade_id=trade.aggregate_trade_id,
price=_required_trade_decimal(
trade.price,
field_name="price",
),
quantity=_required_trade_decimal(
trade.quantity,
field_name="quantity",
),
executed_at=_trade_timestamp_ms_to_utc_datetime(
trade.timestamp,
),
aggressor_side=_trade_aggressor_side(
trade.buyer_is_maker,
),
source=_DZENGI_SOURCE_NAME,
)
def _trade_aggressor_side(
buyer_is_maker: bool,
) -> TradeAggressorSide:
if buyer_is_maker:
return TradeAggressorSide.SELL
return TradeAggressorSide.BUY
def _required_trade_decimal(
value: DzengiRawNumeric,
*,
field_name: str,
) -> Decimal:
try:
result = Decimal(str(value))
except (InvalidOperation, ValueError) as exc:
raise TradeMappingError(
f"Поле {field_name} сделки невозможно "
"преобразовать в Decimal."
) from exc
if not result.is_finite():
raise TradeMappingError(
f"Поле {field_name} сделки должно быть "
"конечным числом."
)
return result
def _trade_timestamp_ms_to_utc_datetime(
value: int,
) -> datetime:
try:
return datetime.fromtimestamp(
value / 1000,
tz=timezone.utc,
)
except (OverflowError, OSError, ValueError) as exc:
raise TradeMappingError(
"Поле timestamp сделки невозможно "
"преобразовать в UTC datetime."
) from exc

View File

@@ -134,6 +134,12 @@ class TradeValueError(MarketDataAcquisitionError):
pass pass
# Ошибка преобразования raw-модели источника
# во внутреннюю модель Trade.
class TradeMappingError(MarketDataAcquisitionError):
pass
# Ошибка определения типа входящего WebSocket-сообщения # Ошибка определения типа входящего WebSocket-сообщения
# и выбора специализированного адаптера. # и выбора специализированного адаптера.
class WebSocketMessageRoutingError(MarketDataAcquisitionError): class WebSocketMessageRoutingError(MarketDataAcquisitionError):

View File

@@ -3,12 +3,14 @@
from __future__ import annotations from __future__ import annotations
from dataclasses import FrozenInstanceError, replace from dataclasses import FrozenInstanceError, replace
from datetime import datetime, timezone
from decimal import Decimal from decimal import Decimal
import pytest import pytest
from src.market_data.acquisition.adapters.dzengi.mapper import ( from src.market_data.acquisition.adapters.dzengi.mapper import (
map_dzengi_exchange_info_to_instruments, map_dzengi_exchange_info_to_instruments,
map_dzengi_rest_agg_trades_to_trades,
map_dzengi_symbol_to_instrument, map_dzengi_symbol_to_instrument,
) )
from src.market_data.acquisition.adapters.dzengi.models import ( from src.market_data.acquisition.adapters.dzengi.models import (
@@ -17,10 +19,16 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiExchangeInfoSymbol, DzengiExchangeInfoSymbol,
DzengiLotSizeFilter, DzengiLotSizeFilter,
DzengiMinNotionalFilter, DzengiMinNotionalFilter,
DzengiRestAggTrade,
DzengiUnknownFilter, DzengiUnknownFilter,
) )
from src.market_data.acquisition.exceptions import ( from src.market_data.acquisition.exceptions import (
InstrumentReferenceMappingError, InstrumentReferenceMappingError,
TradeMappingError,
)
from src.market_data.acquisition.models.trade import (
Trade,
TradeAggressorSide,
) )
@@ -82,6 +90,23 @@ def _response(
) )
def _rest_agg_trade(
*,
aggregate_trade_id: int = 101,
price: str | int | float = "123.45",
quantity: str | int | float = "0.25",
timestamp: int = 1783537921471,
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_map_complete_dzengi_symbol_to_instrument() -> None: def test_map_complete_dzengi_symbol_to_instrument() -> None:
instrument = map_dzengi_symbol_to_instrument( instrument = map_dzengi_symbol_to_instrument(
_complete_symbol() _complete_symbol()
@@ -404,4 +429,244 @@ def test_mapped_instrument_is_immutable() -> None:
) )
with pytest.raises(FrozenInstanceError): with pytest.raises(FrozenInstanceError):
instrument.status = "BREAK" # type: ignore[misc] instrument.status = "BREAK" # type: ignore[misc]
def test_map_dzengi_rest_agg_trade_to_trade() -> None:
trades = map_dzengi_rest_agg_trades_to_trades(
(_rest_agg_trade(),),
symbol="BTC/USDT",
)
assert isinstance(trades, tuple)
assert len(trades) == 1
trade = trades[0]
assert isinstance(trade, Trade)
assert trade.symbol == "BTC/USDT"
assert trade.trade_id == 101
assert trade.price == Decimal("123.45")
assert trade.quantity == Decimal("0.25")
assert trade.executed_at == datetime.fromtimestamp(
1783537921471 / 1000,
tz=timezone.utc,
)
assert trade.aggressor_side is TradeAggressorSide.BUY
assert trade.source == "dzengi"
def test_trade_mapper_maps_buyer_taker_to_buy_aggressor() -> None:
trades = map_dzengi_rest_agg_trades_to_trades(
(
_rest_agg_trade(
buyer_is_maker=False,
),
),
symbol="BTC/USDT",
)
assert trades[0].aggressor_side is TradeAggressorSide.BUY
def test_trade_mapper_maps_seller_taker_to_sell_aggressor() -> None:
trades = map_dzengi_rest_agg_trades_to_trades(
(
_rest_agg_trade(
buyer_is_maker=True,
),
),
symbol="BTC/USDT",
)
assert trades[0].aggressor_side is TradeAggressorSide.SELL
def test_trade_mapper_preserves_input_order() -> None:
trades = map_dzengi_rest_agg_trades_to_trades(
(
_rest_agg_trade(
aggregate_trade_id=103,
),
_rest_agg_trade(
aggregate_trade_id=101,
),
_rest_agg_trade(
aggregate_trade_id=102,
),
),
symbol="BTC/USDT",
)
assert tuple(
trade.trade_id
for trade in trades
) == (
103,
101,
102,
)
def test_trade_mapper_maps_empty_tuple() -> None:
trades = map_dzengi_rest_agg_trades_to_trades(
(),
symbol="BTC/USDT",
)
assert trades == ()
def test_trade_mapper_strips_symbol() -> None:
trades = map_dzengi_rest_agg_trades_to_trades(
(_rest_agg_trade(),),
symbol=" BTC/USDT ",
)
assert trades[0].symbol == "BTC/USDT"
@pytest.mark.parametrize(
("price", "quantity", "expected_price", "expected_quantity"),
[
(
"123.4500",
"0.2500",
Decimal("123.4500"),
Decimal("0.2500"),
),
(
123,
2,
Decimal("123"),
Decimal("2"),
),
(
123.5,
0.25,
Decimal("123.5"),
Decimal("0.25"),
),
],
)
def test_trade_mapper_converts_numeric_values_to_decimal(
price: str | int | float,
quantity: str | int | float,
expected_price: Decimal,
expected_quantity: Decimal,
) -> None:
trades = map_dzengi_rest_agg_trades_to_trades(
(
_rest_agg_trade(
price=price,
quantity=quantity,
),
),
symbol="BTC/USDT",
)
assert trades[0].price == expected_price
assert trades[0].quantity == expected_quantity
def test_trade_mapper_converts_timestamp_to_utc_datetime() -> None:
timestamp = 1783537921471
trades = map_dzengi_rest_agg_trades_to_trades(
(
_rest_agg_trade(
timestamp=timestamp,
),
),
symbol="BTC/USDT",
)
assert trades[0].executed_at == datetime.fromtimestamp(
timestamp / 1000,
tz=timezone.utc,
)
assert trades[0].executed_at.tzinfo is timezone.utc
def test_mapped_trade_is_immutable() -> None:
trade = map_dzengi_rest_agg_trades_to_trades(
(_rest_agg_trade(),),
symbol="BTC/USDT",
)[0]
with pytest.raises(FrozenInstanceError):
trade.price = Decimal("1") # type: ignore[misc]
@pytest.mark.parametrize(
"field_name",
[
"price",
"quantity",
],
)
def test_trade_mapper_rejects_invalid_decimal_value(
field_name: str,
) -> None:
trade = _rest_agg_trade()
invalid_trade = replace(
trade,
**{field_name: "not-a-number"},
)
with pytest.raises(
TradeMappingError,
match=rf"{field_name}.*невозможно преобразовать в Decimal",
):
map_dzengi_rest_agg_trades_to_trades(
(invalid_trade,),
symbol="BTC/USDT",
)
@pytest.mark.parametrize(
("field_name", "invalid_value"),
[
("price", float("nan")),
("price", float("inf")),
("price", float("-inf")),
("quantity", float("nan")),
("quantity", float("inf")),
("quantity", float("-inf")),
],
)
def test_trade_mapper_rejects_non_finite_decimal_value(
field_name: str,
invalid_value: float,
) -> None:
trade = _rest_agg_trade()
invalid_trade = replace(
trade,
**{field_name: invalid_value},
)
with pytest.raises(
TradeMappingError,
match=rf"{field_name}.*должно быть конечным числом",
):
map_dzengi_rest_agg_trades_to_trades(
(invalid_trade,),
symbol="BTC/USDT",
)
def test_trade_mapper_rejects_unrepresentable_timestamp() -> None:
trade = _rest_agg_trade(
timestamp=10**30,
)
with pytest.raises(
TradeMappingError,
match=r"timestamp.*невозможно преобразовать в UTC datetime",
):
map_dzengi_rest_agg_trades_to_trades(
(trade,),
symbol="BTC/USDT",
)

File diff suppressed because it is too large Load Diff