Build 060.13 — WebSocket Trade Mapper
This commit is contained in:
@@ -17,6 +17,7 @@ from src.market_data.acquisition.adapters.dzengi.models import (
|
|||||||
DzengiTicker24hrResponse,
|
DzengiTicker24hrResponse,
|
||||||
DzengiWebSocketOhlcEvent,
|
DzengiWebSocketOhlcEvent,
|
||||||
DzengiWebSocketQuoteResponse,
|
DzengiWebSocketQuoteResponse,
|
||||||
|
DzengiWebSocketTradeEvent,
|
||||||
)
|
)
|
||||||
from src.market_data.acquisition.exceptions import (
|
from src.market_data.acquisition.exceptions import (
|
||||||
CandleMappingError,
|
CandleMappingError,
|
||||||
@@ -37,6 +38,7 @@ from src.market_data.acquisition.models.trade import (
|
|||||||
|
|
||||||
_DZENGI_SOURCE_NAME = "dzengi"
|
_DZENGI_SOURCE_NAME = "dzengi"
|
||||||
_DZENGI_WEBSOCKET_OHLC_SOURCE_NAME = "dzengi_websocket_ohlc"
|
_DZENGI_WEBSOCKET_OHLC_SOURCE_NAME = "dzengi_websocket_ohlc"
|
||||||
|
_DZENGI_WEBSOCKET_TRADE_SOURCE_NAME = "dzengi_websocket_trade"
|
||||||
|
|
||||||
|
|
||||||
def map_dzengi_symbol_to_instrument(
|
def map_dzengi_symbol_to_instrument(
|
||||||
@@ -537,6 +539,41 @@ def map_dzengi_rest_agg_trades_to_trades(
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def map_dzengi_websocket_trade_to_trade(
|
||||||
|
event: DzengiWebSocketTradeEvent,
|
||||||
|
) -> Trade:
|
||||||
|
"""
|
||||||
|
Преобразовать проверенную transport-модель
|
||||||
|
WebSocket Trade в каноническую модель Trade.
|
||||||
|
|
||||||
|
Предполагается, что ранее успешно выполнены:
|
||||||
|
|
||||||
|
- Schema Validation
|
||||||
|
- Parser
|
||||||
|
- Value Validation
|
||||||
|
"""
|
||||||
|
|
||||||
|
return Trade(
|
||||||
|
symbol=event.symbol.strip(),
|
||||||
|
trade_id=event.trade_id,
|
||||||
|
price=_required_trade_decimal(
|
||||||
|
event.price,
|
||||||
|
field_name="price",
|
||||||
|
),
|
||||||
|
quantity=_required_trade_decimal(
|
||||||
|
event.size,
|
||||||
|
field_name="size",
|
||||||
|
),
|
||||||
|
executed_at=_trade_timestamp_ms_to_utc_datetime(
|
||||||
|
event.timestamp,
|
||||||
|
),
|
||||||
|
aggressor_side=_websocket_trade_aggressor_side(
|
||||||
|
event.buyer,
|
||||||
|
),
|
||||||
|
source=_DZENGI_WEBSOCKET_TRADE_SOURCE_NAME,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def _map_dzengi_rest_agg_trade(
|
def _map_dzengi_rest_agg_trade(
|
||||||
trade: DzengiRestAggTrade,
|
trade: DzengiRestAggTrade,
|
||||||
*,
|
*,
|
||||||
@@ -572,6 +609,24 @@ def _trade_aggressor_side(
|
|||||||
return TradeAggressorSide.BUY
|
return TradeAggressorSide.BUY
|
||||||
|
|
||||||
|
|
||||||
|
def _websocket_trade_aggressor_side(
|
||||||
|
buyer: bool,
|
||||||
|
) -> TradeAggressorSide:
|
||||||
|
"""
|
||||||
|
Преобразовать направление WebSocket Trade
|
||||||
|
в каноническую сторону агрессора.
|
||||||
|
|
||||||
|
WebSocket:
|
||||||
|
buyer=True -> BUY
|
||||||
|
buyer=False -> SELL
|
||||||
|
"""
|
||||||
|
|
||||||
|
if buyer:
|
||||||
|
return TradeAggressorSide.BUY
|
||||||
|
|
||||||
|
return TradeAggressorSide.SELL
|
||||||
|
|
||||||
|
|
||||||
def _required_trade_decimal(
|
def _required_trade_decimal(
|
||||||
value: DzengiRawNumeric,
|
value: DzengiRawNumeric,
|
||||||
*,
|
*,
|
||||||
|
|||||||
@@ -0,0 +1,347 @@
|
|||||||
|
# app/tests/unit/market_data/acquisition/adapters/dzengi/
|
||||||
|
# test_websocket_trade_mapper.py
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from dataclasses import FrozenInstanceError
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
from decimal import Decimal
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from src.market_data.acquisition.adapters.dzengi.mapper import (
|
||||||
|
map_dzengi_websocket_trade_to_trade,
|
||||||
|
)
|
||||||
|
from src.market_data.acquisition.adapters.dzengi.models import (
|
||||||
|
DzengiWebSocketTradeEvent,
|
||||||
|
)
|
||||||
|
from src.market_data.acquisition.exceptions import (
|
||||||
|
TradeMappingError,
|
||||||
|
)
|
||||||
|
from src.market_data.acquisition.models.trade import (
|
||||||
|
Trade,
|
||||||
|
TradeAggressorSide,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _websocket_trade_event(
|
||||||
|
*,
|
||||||
|
trade_id: int = 101,
|
||||||
|
price: str | int | float = "123.45",
|
||||||
|
size: str | int | float = "0.25",
|
||||||
|
timestamp: int = 1783537921471,
|
||||||
|
symbol: str = "BTC/USDT",
|
||||||
|
buyer: bool = True,
|
||||||
|
order_id: str = "order-101",
|
||||||
|
) -> DzengiWebSocketTradeEvent:
|
||||||
|
return DzengiWebSocketTradeEvent(
|
||||||
|
trade_id=trade_id,
|
||||||
|
price=price,
|
||||||
|
size=size,
|
||||||
|
timestamp=timestamp,
|
||||||
|
symbol=symbol,
|
||||||
|
buyer=buyer,
|
||||||
|
order_id=order_id,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_map_dzengi_websocket_trade_to_trade() -> None:
|
||||||
|
event = _websocket_trade_event()
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
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_websocket_trade"
|
||||||
|
|
||||||
|
|
||||||
|
def test_websocket_trade_mapper_maps_size_to_quantity() -> None:
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
size="1.75",
|
||||||
|
)
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
assert trade.quantity == Decimal("1.75")
|
||||||
|
|
||||||
|
|
||||||
|
def test_websocket_trade_mapper_maps_buyer_to_buy_aggressor() -> None:
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
buyer=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
assert trade.aggressor_side is TradeAggressorSide.BUY
|
||||||
|
|
||||||
|
|
||||||
|
def test_websocket_trade_mapper_maps_seller_to_sell_aggressor() -> None:
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
buyer=False,
|
||||||
|
)
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
assert trade.aggressor_side is TradeAggressorSide.SELL
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
(
|
||||||
|
"price",
|
||||||
|
"size",
|
||||||
|
"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_websocket_trade_mapper_converts_numeric_values_to_decimal(
|
||||||
|
price: str | int | float,
|
||||||
|
size: str | int | float,
|
||||||
|
expected_price: Decimal,
|
||||||
|
expected_quantity: Decimal,
|
||||||
|
) -> None:
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
price=price,
|
||||||
|
size=size,
|
||||||
|
)
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
assert trade.price == expected_price
|
||||||
|
assert trade.quantity == expected_quantity
|
||||||
|
|
||||||
|
|
||||||
|
def test_websocket_trade_mapper_converts_timestamp_to_utc_datetime() -> None:
|
||||||
|
timestamp = 1783537921471
|
||||||
|
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
timestamp=timestamp,
|
||||||
|
)
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
assert trade.executed_at == datetime.fromtimestamp(
|
||||||
|
timestamp / 1000,
|
||||||
|
tz=timezone.utc,
|
||||||
|
)
|
||||||
|
assert trade.executed_at.tzinfo is timezone.utc
|
||||||
|
|
||||||
|
|
||||||
|
def test_websocket_trade_mapper_strips_symbol() -> None:
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
symbol=" BTC/USDT ",
|
||||||
|
)
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
assert trade.symbol == "BTC/USDT"
|
||||||
|
|
||||||
|
|
||||||
|
def test_websocket_trade_mapper_sets_websocket_source() -> None:
|
||||||
|
event = _websocket_trade_event()
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
assert trade.source == "dzengi_websocket_trade"
|
||||||
|
|
||||||
|
|
||||||
|
def test_websocket_trade_mapper_does_not_expose_order_id() -> None:
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
order_id="exchange-order-999",
|
||||||
|
)
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
assert not hasattr(trade, "order_id")
|
||||||
|
|
||||||
|
|
||||||
|
def test_websocket_trade_mapper_does_not_modify_transport_event() -> None:
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
trade_id=202,
|
||||||
|
price="456.78",
|
||||||
|
size="3.5",
|
||||||
|
timestamp=1783537921999,
|
||||||
|
symbol=" ETH/USDT ",
|
||||||
|
buyer=False,
|
||||||
|
order_id="order-202",
|
||||||
|
)
|
||||||
|
|
||||||
|
original_values = (
|
||||||
|
event.trade_id,
|
||||||
|
event.price,
|
||||||
|
event.size,
|
||||||
|
event.timestamp,
|
||||||
|
event.symbol,
|
||||||
|
event.buyer,
|
||||||
|
event.order_id,
|
||||||
|
)
|
||||||
|
|
||||||
|
map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
assert (
|
||||||
|
event.trade_id,
|
||||||
|
event.price,
|
||||||
|
event.size,
|
||||||
|
event.timestamp,
|
||||||
|
event.symbol,
|
||||||
|
event.buyer,
|
||||||
|
event.order_id,
|
||||||
|
) == original_values
|
||||||
|
|
||||||
|
|
||||||
|
def test_mapped_websocket_trade_is_immutable() -> None:
|
||||||
|
event = _websocket_trade_event()
|
||||||
|
|
||||||
|
trade = map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
with pytest.raises(FrozenInstanceError):
|
||||||
|
trade.price = Decimal("1") # type: ignore[misc]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
(
|
||||||
|
"field_name",
|
||||||
|
"invalid_value",
|
||||||
|
"expected_mapper_field",
|
||||||
|
),
|
||||||
|
[
|
||||||
|
(
|
||||||
|
"price",
|
||||||
|
"not-a-number",
|
||||||
|
"price",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"size",
|
||||||
|
"not-a-number",
|
||||||
|
"size",
|
||||||
|
),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_websocket_trade_mapper_rejects_invalid_decimal_value(
|
||||||
|
field_name: str,
|
||||||
|
invalid_value: str,
|
||||||
|
expected_mapper_field: str,
|
||||||
|
) -> None:
|
||||||
|
values = {
|
||||||
|
"price": "123.45",
|
||||||
|
"size": "0.25",
|
||||||
|
}
|
||||||
|
values[field_name] = invalid_value
|
||||||
|
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
price=values["price"],
|
||||||
|
size=values["size"],
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(
|
||||||
|
TradeMappingError,
|
||||||
|
match=(
|
||||||
|
rf"{expected_mapper_field}.*"
|
||||||
|
r"невозможно преобразовать в Decimal"
|
||||||
|
),
|
||||||
|
):
|
||||||
|
map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
(
|
||||||
|
"field_name",
|
||||||
|
"invalid_value",
|
||||||
|
"expected_mapper_field",
|
||||||
|
),
|
||||||
|
[
|
||||||
|
(
|
||||||
|
"price",
|
||||||
|
float("nan"),
|
||||||
|
"price",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"price",
|
||||||
|
float("inf"),
|
||||||
|
"price",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"price",
|
||||||
|
float("-inf"),
|
||||||
|
"price",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"size",
|
||||||
|
float("nan"),
|
||||||
|
"size",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"size",
|
||||||
|
float("inf"),
|
||||||
|
"size",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"size",
|
||||||
|
float("-inf"),
|
||||||
|
"size",
|
||||||
|
),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_websocket_trade_mapper_rejects_non_finite_decimal_value(
|
||||||
|
field_name: str,
|
||||||
|
invalid_value: float,
|
||||||
|
expected_mapper_field: str,
|
||||||
|
) -> None:
|
||||||
|
values: dict[str, str | int | float] = {
|
||||||
|
"price": "123.45",
|
||||||
|
"size": "0.25",
|
||||||
|
}
|
||||||
|
values[field_name] = invalid_value
|
||||||
|
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
price=values["price"],
|
||||||
|
size=values["size"],
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(
|
||||||
|
TradeMappingError,
|
||||||
|
match=(
|
||||||
|
rf"{expected_mapper_field}.*"
|
||||||
|
r"должно быть конечным числом"
|
||||||
|
),
|
||||||
|
):
|
||||||
|
map_dzengi_websocket_trade_to_trade(event)
|
||||||
|
|
||||||
|
|
||||||
|
def test_websocket_trade_mapper_rejects_unrepresentable_timestamp() -> None:
|
||||||
|
event = _websocket_trade_event(
|
||||||
|
timestamp=10**30,
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(
|
||||||
|
TradeMappingError,
|
||||||
|
match=r"timestamp.*невозможно преобразовать в UTC datetime",
|
||||||
|
):
|
||||||
|
map_dzengi_websocket_trade_to_trade(event)
|
||||||
1169
docs/migrations/build_060_13.md
Normal file
1169
docs/migrations/build_060_13.md
Normal file
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user