Build 060.17 — Trades Feed Core

This commit is contained in:
2026-07-20 01:55:52 +03:00
parent be6beac560
commit fe308b9e60
11 changed files with 2314 additions and 2 deletions

View File

@@ -1 +1,9 @@
# app/src/market_data/acquisition/adapters/dzengi/__init__.py # app/src/market_data/acquisition/adapters/dzengi/__init__.py
from src.market_data.acquisition.adapters.dzengi.websocket_trade_adapter import (
adapt_websocket_trade_document,
)
__all__ = [
"adapt_websocket_trade_document",
]

View File

@@ -145,6 +145,11 @@ class TradeMappingError(MarketDataAcquisitionError):
pass pass
# Ошибка регистрации или получения Trades Feed.
class TradeFeedRegistryError(MarketDataAcquisitionError):
pass
# Ошибка определения типа входящего WebSocket-сообщения # Ошибка определения типа входящего WebSocket-сообщения
# и выбора специализированного адаптера. # и выбора специализированного адаптера.
class WebSocketMessageRoutingError(MarketDataAcquisitionError): class WebSocketMessageRoutingError(MarketDataAcquisitionError):

View File

@@ -0,0 +1,35 @@
# app/src/market_data/acquisition/feeds/trades_feed.py
from __future__ import annotations
from src.market_data.acquisition.models.trade import Trade
from src.market_data.acquisition.protocol import (
TradeDocumentHandler,
TradeDocumentSource,
)
# Источник готовой внутренней модели Trade.
class TradesFeed:
def __init__(
self,
source: TradeDocumentSource,
handler: TradeDocumentHandler,
) -> None:
self._source = source
self._handler = handler
def load_trade(
self,
symbol: str,
) -> Trade:
"""
Получить последнюю сделку инструмента.
Feed получает сырой документ у источника и делегирует
его преобразование специализированному обработчику.
"""
document = self._source.fetch_trade_document(symbol)
return self._handler.handle_trade_document(document)

View File

@@ -0,0 +1,24 @@
# app/src/market_data/acquisition/handlers/trades_handler.py
# app/src/market_data/acquisition/handlers/trades_handler.py
from __future__ import annotations
from src.market_data.acquisition.adapters.dzengi import (
adapt_websocket_trade_document,
)
from src.market_data.acquisition.models.trade import Trade
from src.market_data.acquisition.validation.schema import (
validate_dzengi_websocket_trade_schema,
)
# Обработчик документа Trades Feed формата Dzengi WebSocket internal.trade.
class DzengiTradeDocumentHandler:
def handle_trade_document(
self,
document: object,
) -> Trade:
validated_document = validate_dzengi_websocket_trade_schema(document)
return adapt_websocket_trade_document(validated_document)

View File

@@ -7,6 +7,7 @@ from typing import Protocol, runtime_checkable
from src.market_data.acquisition.models.candle import Candle from src.market_data.acquisition.models.candle import Candle
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
# Источник сырого документа Instrument Reference Data. # Источник сырого документа Instrument Reference Data.
@@ -87,6 +88,49 @@ class QuoteFeedProtocol(Protocol):
... ...
# Источник сырого документа Trades Feed.
@runtime_checkable
class TradeDocumentSource(Protocol):
def fetch_trade_document(
self,
symbol: str,
) -> object:
"""
Получить декодированный транспортный документ сделок.
Источник не выполняет schema validation, parsing,
value validation или mapping во внутреннюю модель Trade.
"""
...
# Обработчик сырого документа Trades Feed.
@runtime_checkable
class TradeDocumentHandler(Protocol):
def handle_trade_document(
self,
document: object,
) -> Trade:
"""
Преобразовать сырой документ
в проверенную внутреннюю модель Trade.
"""
...
# Источник готовой модели Trade.
@runtime_checkable
class TradesFeedProtocol(Protocol):
def load_trade(
self,
symbol: str,
) -> Trade:
"""
Получить внутреннюю модель последней сделки инструмента.
"""
...
# Источник сырого документа Candles Feed. # Источник сырого документа Candles Feed.
@runtime_checkable @runtime_checkable
class CandlesDocumentSource(Protocol): class CandlesDocumentSource(Protocol):

View File

@@ -6,11 +6,13 @@ from src.market_data.acquisition.exceptions import (
CandleFeedRegistryError, CandleFeedRegistryError,
InstrumentFeedRegistryError, InstrumentFeedRegistryError,
QuoteFeedRegistryError, QuoteFeedRegistryError,
TradeFeedRegistryError,
) )
from src.market_data.acquisition.protocol import ( from src.market_data.acquisition.protocol import (
CandlesFeedProtocol, CandlesFeedProtocol,
InstrumentFeedProtocol, InstrumentFeedProtocol,
QuoteFeedProtocol, QuoteFeedProtocol,
TradesFeedProtocol,
) )
@@ -144,6 +146,71 @@ class QuoteFeedRegistry:
return normalized_source_name return normalized_source_name
# Реестр доступных потоков рыночных сделок.
class TradesFeedRegistry:
def __init__(self) -> None:
self._feeds: dict[str, TradesFeedProtocol] = {}
def register(
self,
source_name: str,
feed: TradesFeedProtocol,
) -> None:
"""
Зарегистрировать Trades Feed для указанного источника.
Повторная регистрация того же имени запрещена, чтобы исключить
неявную замену production-зависимости.
"""
normalized_source_name = self._normalize_source_name(source_name)
if not isinstance(feed, TradesFeedProtocol):
raise TradeFeedRegistryError(
f"Объект для источника '{normalized_source_name}' "
"не соответствует TradesFeedProtocol."
)
if normalized_source_name in self._feeds:
raise TradeFeedRegistryError(
f"Trades Feed для источника "
f"'{normalized_source_name}' уже зарегистрирован."
)
self._feeds[normalized_source_name] = feed
def get(
self,
source_name: str,
) -> TradesFeedProtocol:
"""Вернуть зарегистрированный Trades Feed по имени источника."""
normalized_source_name = self._normalize_source_name(source_name)
feed = self._feeds.get(normalized_source_name)
if feed is None:
raise TradeFeedRegistryError(
f"Trades Feed для источника "
f"'{normalized_source_name}' не зарегистрирован."
)
return feed
def _normalize_source_name(
self,
source_name: str,
) -> str:
normalized_source_name = source_name.strip()
if not normalized_source_name:
raise TradeFeedRegistryError(
"Имя источника Trades Feed не должно быть пустым."
)
return normalized_source_name
# Реестр доступных потоков рыночных свечей. # Реестр доступных потоков рыночных свечей.
class CandlesFeedRegistry: class CandlesFeedRegistry:
def __init__(self) -> None: def __init__(self) -> None:

View File

@@ -0,0 +1,194 @@
# app/tests/unit/market_data/acquisition/feeds/test_trades_feed.py
from __future__ import annotations
from datetime import datetime, timezone
from decimal import Decimal
import pytest
from src.market_data.acquisition.feeds.trades_feed import TradesFeed
from src.market_data.acquisition.models.trade import (
Trade,
TradeAggressorSide,
)
from src.market_data.acquisition.protocol import (
TradeDocumentHandler,
TradeDocumentSource,
TradesFeedProtocol,
)
def _trade() -> Trade:
return Trade(
symbol="BTCUSDT",
trade_id=101,
price=Decimal("43210.50"),
quantity=Decimal("0.125"),
executed_at=datetime(
2023,
11,
14,
22,
13,
20,
123000,
tzinfo=timezone.utc,
),
aggressor_side=TradeAggressorSide.BUY,
source="dzengi_websocket_trade",
)
class StubTradeDocumentSource:
def __init__(self, document: object) -> None:
self.document = document
self.requested_symbols: list[str] = []
def fetch_trade_document(
self,
symbol: str,
) -> object:
self.requested_symbols.append(symbol)
return self.document
class StubTradeDocumentHandler:
def __init__(self, trade: Trade) -> None:
self.trade = trade
self.received_documents: list[object] = []
def handle_trade_document(
self,
document: object,
) -> Trade:
self.received_documents.append(document)
return self.trade
def test_source_conforms_to_trade_document_source_protocol() -> None:
source = StubTradeDocumentSource(document={})
assert isinstance(source, TradeDocumentSource)
def test_handler_conforms_to_trade_document_handler_protocol() -> None:
handler = StubTradeDocumentHandler(trade=_trade())
assert isinstance(handler, TradeDocumentHandler)
def test_feed_conforms_to_trades_feed_protocol() -> None:
feed = TradesFeed(
source=StubTradeDocumentSource(document={}),
handler=StubTradeDocumentHandler(trade=_trade()),
)
assert isinstance(feed, TradesFeedProtocol)
def test_load_trade_requests_document_for_symbol() -> None:
document = {
"status": "OK",
"destination": "internal.trade",
"payload": {},
}
source = StubTradeDocumentSource(document=document)
handler = StubTradeDocumentHandler(trade=_trade())
feed = TradesFeed(
source=source,
handler=handler,
)
feed.load_trade("BTCUSDT")
assert source.requested_symbols == ["BTCUSDT"]
def test_load_trade_passes_document_to_handler() -> None:
document = {
"status": "OK",
"destination": "internal.trade",
"payload": {},
}
source = StubTradeDocumentSource(document=document)
handler = StubTradeDocumentHandler(trade=_trade())
feed = TradesFeed(
source=source,
handler=handler,
)
feed.load_trade("BTCUSDT")
assert handler.received_documents == [document]
def test_load_trade_returns_handler_result_without_modification() -> None:
expected_trade = _trade()
source = StubTradeDocumentSource(document={})
handler = StubTradeDocumentHandler(trade=expected_trade)
feed = TradesFeed(
source=source,
handler=handler,
)
result = feed.load_trade("BTCUSDT")
assert result is expected_trade
def test_load_trade_does_not_normalize_symbol() -> None:
source = StubTradeDocumentSource(document={})
handler = StubTradeDocumentHandler(trade=_trade())
feed = TradesFeed(
source=source,
handler=handler,
)
feed.load_trade(" BTCUSDT ")
assert source.requested_symbols == [" BTCUSDT "]
def test_load_trade_propagates_source_error() -> None:
expected_error = RuntimeError("source error")
class FailingSource:
def fetch_trade_document(
self,
symbol: str,
) -> object:
raise expected_error
feed = TradesFeed(
source=FailingSource(),
handler=StubTradeDocumentHandler(trade=_trade()),
)
with pytest.raises(RuntimeError) as exc_info:
feed.load_trade("BTCUSDT")
assert exc_info.value is expected_error
def test_load_trade_propagates_handler_error() -> None:
expected_error = RuntimeError("handler error")
class FailingHandler:
def handle_trade_document(
self,
document: object,
) -> Trade:
raise expected_error
feed = TradesFeed(
source=StubTradeDocumentSource(document={}),
handler=FailingHandler(),
)
with pytest.raises(RuntimeError) as exc_info:
feed.load_trade("BTCUSDT")
assert exc_info.value is expected_error

View File

@@ -0,0 +1,124 @@
# app/tests/unit/market_data/acquisition/handlers/test_trades_handler.py
from __future__ import annotations
from datetime import datetime, timezone
from decimal import Decimal
import pytest
from src.market_data.acquisition.exceptions import TradeSchemaError
from src.market_data.acquisition.handlers import trades_handler as handler_module
from src.market_data.acquisition.handlers.trades_handler import (
DzengiTradeDocumentHandler,
)
from src.market_data.acquisition.models.trade import (
Trade,
TradeAggressorSide,
)
from src.market_data.acquisition.protocol import TradeDocumentHandler
from src.market_data.acquisition.validation.schema import (
ValidatedWebSocketTradeDocument,
)
def _trade_document() -> dict[str, object]:
return {
"status": "OK",
"destination": "internal.trade",
"payload": {
"buyer": True,
"id": 101,
"orderId": "00a02503-0079-54c4-0000-000081e62b58",
"price": "43210.50",
"size": "0.125",
"symbol": " BTCUSDT ",
"ts": 1_700_000_000_123,
},
}
def _trade() -> Trade:
return Trade(
symbol="BTCUSDT",
trade_id=101,
price=Decimal("43210.50"),
quantity=Decimal("0.125"),
executed_at=datetime(
2023,
11,
14,
22,
13,
20,
123000,
tzinfo=timezone.utc,
),
aggressor_side=TradeAggressorSide.BUY,
source="dzengi_websocket_trade",
)
def test_handler_conforms_to_trade_document_handler_protocol() -> None:
handler = DzengiTradeDocumentHandler()
assert isinstance(handler, TradeDocumentHandler)
def test_handler_returns_canonical_trade() -> None:
handler = DzengiTradeDocumentHandler()
result = handler.handle_trade_document(_trade_document())
assert result == _trade()
def test_handler_validates_schema_before_calling_adapter(
monkeypatch: pytest.MonkeyPatch,
) -> None:
handler = DzengiTradeDocumentHandler()
expected_trade = _trade()
captured_document: ValidatedWebSocketTradeDocument | None = None
def fake_adapter(
document: ValidatedWebSocketTradeDocument,
) -> Trade:
nonlocal captured_document
captured_document = document
return expected_trade
monkeypatch.setattr(
handler_module,
"adapt_websocket_trade_document",
fake_adapter,
)
result = handler.handle_trade_document(_trade_document())
assert result is expected_trade
assert isinstance(
captured_document,
ValidatedWebSocketTradeDocument,
)
assert captured_document.status == "OK"
assert captured_document.destination == "internal.trade"
assert captured_document.correlation_id is None
assert captured_document.payload["id"] == 101
assert captured_document.payload["symbol"] == " BTCUSDT "
def test_handler_propagates_schema_error() -> None:
handler = DzengiTradeDocumentHandler()
document = {
"status": "OK",
"destination": "internal.trade",
"payload": {},
}
with pytest.raises(
TradeSchemaError,
match=r"\$\.payload не содержит обязательные поля WebSocket Trade",
):
handler.handle_trade_document(document)

View File

@@ -3,6 +3,17 @@
from __future__ import annotations from __future__ import annotations
from decimal import Decimal from decimal import Decimal
from datetime import datetime, timezone
from src.market_data.acquisition.models.trade import (
Trade,
TradeAggressorSide,
)
from src.market_data.acquisition.protocol import (
TradeDocumentHandler,
TradeDocumentSource,
TradesFeedProtocol,
)
from src.market_data.acquisition.exceptions import ( from src.market_data.acquisition.exceptions import (
InstrumentReferenceMappingError, InstrumentReferenceMappingError,
@@ -67,6 +78,57 @@ class StubInstrumentFeed:
return (_instrument(),) return (_instrument(),)
def _trade() -> Trade:
return Trade(
symbol="BTCUSDT",
trade_id=101,
price=Decimal("43210.50"),
quantity=Decimal("0.125"),
executed_at=datetime(
2023,
11,
14,
22,
13,
20,
123000,
tzinfo=timezone.utc,
),
aggressor_side=TradeAggressorSide.BUY,
source="dzengi_websocket_trade",
)
class StubTradeDocumentSource:
def fetch_trade_document(
self,
symbol: str,
) -> object:
return {
"symbol": symbol,
}
class StubTradeDocumentHandler:
def handle_trade_document(
self,
document: object,
) -> Trade:
del document
return _trade()
class StubTradesFeed:
def load_trade(
self,
symbol: str,
) -> Trade:
del symbol
return _trade()
class InvalidSource: class InvalidSource:
pass pass
@@ -156,4 +218,45 @@ def test_all_instrument_reference_errors_share_base_type() -> None:
assert all( assert all(
isinstance(error, MarketDataAcquisitionError) isinstance(error, MarketDataAcquisitionError)
for error in errors for error in errors
) )
def test_trade_document_source_satisfies_protocol() -> None:
source = StubTradeDocumentSource()
assert isinstance(source, TradeDocumentSource)
assert source.fetch_trade_document("BTCUSDT") == {
"symbol": "BTCUSDT",
}
def test_trade_document_handler_satisfies_protocol() -> None:
handler = StubTradeDocumentHandler()
assert isinstance(handler, TradeDocumentHandler)
trade = handler.handle_trade_document({})
assert trade.symbol == "BTCUSDT"
def test_trades_feed_satisfies_protocol() -> None:
feed = StubTradesFeed()
assert isinstance(feed, TradesFeedProtocol)
trade = feed.load_trade("BTCUSDT")
assert trade.symbol == "BTCUSDT"
def test_trade_protocols_support_structural_typing() -> None:
source: TradeDocumentSource = StubTradeDocumentSource()
handler: TradeDocumentHandler = StubTradeDocumentHandler()
feed: TradesFeedProtocol = StubTradesFeed()
document = source.fetch_trade_document("BTCUSDT")
assert handler.handle_trade_document(document).symbol == "BTCUSDT"
assert feed.load_trade("BTCUSDT").symbol == "BTCUSDT"

View File

@@ -3,9 +3,28 @@
from __future__ import annotations from __future__ import annotations
from decimal import Decimal from decimal import Decimal
from datetime import datetime, timezone
from typing import cast
import pytest import pytest
from src.market_data.acquisition.exceptions import (
TradeFeedRegistryError,
)
from src.market_data.acquisition.models.trade import (
Trade,
TradeAggressorSide,
)
from src.market_data.acquisition.protocol import (
TradesFeedProtocol,
)
from src.market_data.acquisition.registry import (
TradesFeedRegistry,
)
from src.market_data.acquisition.exceptions import ( from src.market_data.acquisition.exceptions import (
InstrumentFeedRegistryError, InstrumentFeedRegistryError,
MarketDataAcquisitionError, MarketDataAcquisitionError,
@@ -58,6 +77,45 @@ class StubInstrumentFeed:
return self.instruments return self.instruments
def _trade() -> Trade:
return Trade(
symbol="BTCUSDT",
trade_id=101,
price=Decimal("43210.50"),
quantity=Decimal("0.125"),
executed_at=datetime(
2023,
11,
14,
22,
13,
20,
123000,
tzinfo=timezone.utc,
),
aggressor_side=TradeAggressorSide.BUY,
source="dzengi_websocket_trade",
)
class StubTradesFeed:
def __init__(
self,
*,
trade: Trade | None = None,
) -> None:
self.trade = trade or _trade()
self.load_call_count = 0
def load_trade(
self,
symbol: str,
) -> Trade:
self.load_call_count += 1
return self.trade
class InvalidFeed: class InvalidFeed:
pass pass
@@ -688,3 +746,83 @@ def test_candles_registry_error_inherits_acquisition_error() -> None:
assert isinstance(error, MarketDataAcquisitionError) assert isinstance(error, MarketDataAcquisitionError)
assert str(error) == "Registry error." assert str(error) == "Registry error."
def test_register_and_get_trades_feed() -> None:
registry = TradesFeedRegistry()
feed = StubTradesFeed()
registry.register("dzengi", feed)
assert registry.get("dzengi") is feed
def test_trades_registry_accepts_protocol() -> None:
registry = TradesFeedRegistry()
feed = StubTradesFeed()
assert isinstance(feed, TradesFeedProtocol)
registry.register("dzengi", feed)
assert registry.get("dzengi") is feed
def test_trades_registry_preserves_identity() -> None:
registry = TradesFeedRegistry()
feed = StubTradesFeed()
registry.register("dzengi", feed)
assert registry.get("dzengi") is feed
def test_trades_registry_strips_outer_whitespace() -> None:
registry = TradesFeedRegistry()
feed = StubTradesFeed()
registry.register(" dzengi ", feed)
assert registry.get("dzengi") is feed
assert registry.get(" dzengi ") is feed
def test_trades_registry_rejects_invalid_feed() -> None:
registry = TradesFeedRegistry()
with pytest.raises(
TradeFeedRegistryError,
match=r"не соответствует TradesFeedProtocol",
):
registry.register(
"dzengi",
cast(TradesFeedProtocol, object()),
)
def test_trades_registry_rejects_duplicate_source() -> None:
registry = TradesFeedRegistry()
registry.register("dzengi", StubTradesFeed())
with pytest.raises(
TradeFeedRegistryError,
match=r"уже зарегистрирован",
):
registry.register("dzengi", StubTradesFeed())
def test_trades_registry_rejects_unknown_source() -> None:
registry = TradesFeedRegistry()
with pytest.raises(
TradeFeedRegistryError,
match=r"не зарегистрирован",
):
registry.get("dzengi")
def test_trades_registry_error_inherits_market_data_error() -> None:
error = TradeFeedRegistryError("registry error")
assert isinstance(error, MarketDataAcquisitionError)

File diff suppressed because it is too large Load Diff