Build 060.15 — Unified WebSocket Routing
This commit is contained in:
@@ -12,14 +12,19 @@ from src.market_data.acquisition.adapters.dzengi.parser import (
|
||||
parse_dzengi_websocket_ohlc,
|
||||
parse_dzengi_websocket_quote,
|
||||
)
|
||||
from src.market_data.acquisition.adapters.dzengi.websocket_trade_adapter import (
|
||||
adapt_websocket_trade_document,
|
||||
)
|
||||
from src.market_data.acquisition.exceptions import (
|
||||
WebSocketMessageRoutingError,
|
||||
)
|
||||
from src.market_data.acquisition.models.candle_close import CandleCloseEvent
|
||||
from src.market_data.acquisition.models.quote import Quote
|
||||
from src.market_data.acquisition.models.trade import Trade
|
||||
from src.market_data.acquisition.validation.schema import (
|
||||
validate_dzengi_websocket_ohlc_schema,
|
||||
validate_dzengi_websocket_quote_schema,
|
||||
validate_dzengi_websocket_trade_schema,
|
||||
)
|
||||
from src.market_data.acquisition.validation.values import (
|
||||
validate_dzengi_websocket_ohlc_values,
|
||||
@@ -64,6 +69,22 @@ class DzengiWebSocketOhlcAdapter:
|
||||
)
|
||||
|
||||
|
||||
# Преобразует одно декодированное сообщение Dzengi WebSocket
|
||||
# internal.trade в каноническую модель Trade.
|
||||
class DzengiWebSocketTradeAdapter:
|
||||
def map_message(
|
||||
self,
|
||||
document: object,
|
||||
) -> Trade:
|
||||
validated = validate_dzengi_websocket_trade_schema(
|
||||
document,
|
||||
)
|
||||
|
||||
return adapt_websocket_trade_document(
|
||||
validated,
|
||||
)
|
||||
|
||||
|
||||
# Определяет тип декодированного сообщения Dzengi WebSocket
|
||||
# и передаёт его соответствующему специализированному адаптеру.
|
||||
class DzengiUnifiedWebSocketAdapter:
|
||||
@@ -72,27 +93,47 @@ class DzengiUnifiedWebSocketAdapter:
|
||||
*,
|
||||
quote_adapter: DzengiWebSocketQuoteAdapter | None = None,
|
||||
ohlc_adapter: DzengiWebSocketOhlcAdapter | None = None,
|
||||
trade_adapter: DzengiWebSocketTradeAdapter | None = None,
|
||||
) -> None:
|
||||
self._quote_adapter = quote_adapter or DzengiWebSocketQuoteAdapter()
|
||||
self._ohlc_adapter = ohlc_adapter or DzengiWebSocketOhlcAdapter()
|
||||
self._quote_adapter = (
|
||||
quote_adapter
|
||||
or DzengiWebSocketQuoteAdapter()
|
||||
)
|
||||
|
||||
self._ohlc_adapter = (
|
||||
ohlc_adapter
|
||||
or DzengiWebSocketOhlcAdapter()
|
||||
)
|
||||
|
||||
self._trade_adapter = (
|
||||
trade_adapter
|
||||
or DzengiWebSocketTradeAdapter()
|
||||
)
|
||||
|
||||
def map_message(
|
||||
self,
|
||||
document: object,
|
||||
*,
|
||||
received_at: datetime | None = None,
|
||||
) -> Quote | CandleCloseEvent:
|
||||
) -> Quote | CandleCloseEvent | Trade:
|
||||
if not isinstance(document, dict):
|
||||
raise WebSocketMessageRoutingError(
|
||||
"Сообщение Dzengi WebSocket должно быть объектом."
|
||||
)
|
||||
|
||||
if document.get("destination") == "ohlc.event":
|
||||
destination = document.get("destination")
|
||||
|
||||
if destination == "ohlc.event":
|
||||
return self._ohlc_adapter.map_message(
|
||||
document,
|
||||
received_at=received_at,
|
||||
)
|
||||
|
||||
if destination == "internal.trade":
|
||||
return self._trade_adapter.map_message(
|
||||
document,
|
||||
)
|
||||
|
||||
if "Payload" in document:
|
||||
return self._quote_adapter.map_message(
|
||||
document,
|
||||
|
||||
Reference in New Issue
Block a user