From fe308b9e60fc65b25fb6a896ea3d6a6307dbf216 Mon Sep 17 00:00:00 2001 From: Sergey Date: Mon, 20 Jul 2026 01:55:52 +0300 Subject: [PATCH] =?UTF-8?q?Build=20060.17=20=E2=80=94=20Trades=20Feed=20Co?= =?UTF-8?q?re?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../acquisition/adapters/dzengi/__init__.py | 10 +- app/src/market_data/acquisition/exceptions.py | 5 + .../acquisition/feeds/trades_feed.py | 35 + .../acquisition/handlers/trades_handler.py | 24 + app/src/market_data/acquisition/protocol.py | 44 + app/src/market_data/acquisition/registry.py | 67 + .../acquisition/feeds/test_trades_feed.py | 194 ++ .../handlers/test_trades_handler.py | 124 ++ .../market_data/acquisition/test_protocol.py | 105 +- .../market_data/acquisition/test_registry.py | 138 ++ docs/migrations/build_060_17.md | 1570 +++++++++++++++++ 11 files changed, 2314 insertions(+), 2 deletions(-) create mode 100644 app/tests/unit/market_data/acquisition/feeds/test_trades_feed.py create mode 100644 app/tests/unit/market_data/acquisition/handlers/test_trades_handler.py create mode 100644 docs/migrations/build_060_17.md diff --git a/app/src/market_data/acquisition/adapters/dzengi/__init__.py b/app/src/market_data/acquisition/adapters/dzengi/__init__.py index e4aab50..c17d89b 100644 --- a/app/src/market_data/acquisition/adapters/dzengi/__init__.py +++ b/app/src/market_data/acquisition/adapters/dzengi/__init__.py @@ -1 +1,9 @@ -# app/src/market_data/acquisition/adapters/dzengi/__init__.py \ No newline at end of file +# 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", +] \ No newline at end of file diff --git a/app/src/market_data/acquisition/exceptions.py b/app/src/market_data/acquisition/exceptions.py index 9f79046..e5fb153 100644 --- a/app/src/market_data/acquisition/exceptions.py +++ b/app/src/market_data/acquisition/exceptions.py @@ -145,6 +145,11 @@ class TradeMappingError(MarketDataAcquisitionError): pass +# Ошибка регистрации или получения Trades Feed. +class TradeFeedRegistryError(MarketDataAcquisitionError): + pass + + # Ошибка определения типа входящего WebSocket-сообщения # и выбора специализированного адаптера. class WebSocketMessageRoutingError(MarketDataAcquisitionError): diff --git a/app/src/market_data/acquisition/feeds/trades_feed.py b/app/src/market_data/acquisition/feeds/trades_feed.py index e69de29..9f779e9 100644 --- a/app/src/market_data/acquisition/feeds/trades_feed.py +++ b/app/src/market_data/acquisition/feeds/trades_feed.py @@ -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) \ No newline at end of file diff --git a/app/src/market_data/acquisition/handlers/trades_handler.py b/app/src/market_data/acquisition/handlers/trades_handler.py index e69de29..3697b01 100644 --- a/app/src/market_data/acquisition/handlers/trades_handler.py +++ b/app/src/market_data/acquisition/handlers/trades_handler.py @@ -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) \ No newline at end of file diff --git a/app/src/market_data/acquisition/protocol.py b/app/src/market_data/acquisition/protocol.py index 3003530..f4bb5c9 100644 --- a/app/src/market_data/acquisition/protocol.py +++ b/app/src/market_data/acquisition/protocol.py @@ -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.instrument import Instrument from src.market_data.acquisition.models.quote import Quote +from src.market_data.acquisition.models.trade import Trade # Источник сырого документа 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. @runtime_checkable class CandlesDocumentSource(Protocol): diff --git a/app/src/market_data/acquisition/registry.py b/app/src/market_data/acquisition/registry.py index ec6f650..396fed8 100644 --- a/app/src/market_data/acquisition/registry.py +++ b/app/src/market_data/acquisition/registry.py @@ -6,11 +6,13 @@ from src.market_data.acquisition.exceptions import ( CandleFeedRegistryError, InstrumentFeedRegistryError, QuoteFeedRegistryError, + TradeFeedRegistryError, ) from src.market_data.acquisition.protocol import ( CandlesFeedProtocol, InstrumentFeedProtocol, QuoteFeedProtocol, + TradesFeedProtocol, ) @@ -144,6 +146,71 @@ class QuoteFeedRegistry: 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: def __init__(self) -> None: diff --git a/app/tests/unit/market_data/acquisition/feeds/test_trades_feed.py b/app/tests/unit/market_data/acquisition/feeds/test_trades_feed.py new file mode 100644 index 0000000..02212c9 --- /dev/null +++ b/app/tests/unit/market_data/acquisition/feeds/test_trades_feed.py @@ -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 \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/handlers/test_trades_handler.py b/app/tests/unit/market_data/acquisition/handlers/test_trades_handler.py new file mode 100644 index 0000000..7797398 --- /dev/null +++ b/app/tests/unit/market_data/acquisition/handlers/test_trades_handler.py @@ -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) \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/test_protocol.py b/app/tests/unit/market_data/acquisition/test_protocol.py index de46ee1..47d0c5d 100644 --- a/app/tests/unit/market_data/acquisition/test_protocol.py +++ b/app/tests/unit/market_data/acquisition/test_protocol.py @@ -3,6 +3,17 @@ from __future__ import annotations 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 ( InstrumentReferenceMappingError, @@ -67,6 +78,57 @@ class StubInstrumentFeed: 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: pass @@ -156,4 +218,45 @@ def test_all_instrument_reference_errors_share_base_type() -> None: assert all( isinstance(error, MarketDataAcquisitionError) for error in errors - ) \ No newline at end of file + ) + + +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" \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/test_registry.py b/app/tests/unit/market_data/acquisition/test_registry.py index ae36b98..1aae36b 100644 --- a/app/tests/unit/market_data/acquisition/test_registry.py +++ b/app/tests/unit/market_data/acquisition/test_registry.py @@ -3,9 +3,28 @@ from __future__ import annotations from decimal import Decimal +from datetime import datetime, timezone +from typing import cast 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 ( InstrumentFeedRegistryError, MarketDataAcquisitionError, @@ -58,6 +77,45 @@ class StubInstrumentFeed: 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: pass @@ -688,3 +746,83 @@ def test_candles_registry_error_inherits_acquisition_error() -> None: assert isinstance(error, MarketDataAcquisitionError) 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) diff --git a/docs/migrations/build_060_17.md b/docs/migrations/build_060_17.md new file mode 100644 index 0000000..3dce793 --- /dev/null +++ b/docs/migrations/build_060_17.md @@ -0,0 +1,1570 @@ +# Build 060.17 — Trades Feed Core + +**Engineering Migration Report** + +--- + +# Контроль документа + +| Свойство | Значение | +|----------|----------| +| Build | 060.17 | +| Название | Trades Feed Core | +| Статус | Completed | +| Проект | Dzentra | +| Подсистема | Market Data Acquisition | +| Компонент | Trades Feed | +| Версия | 1.0 | + +--- + +# Цель Build + +После завершения Build 060.16 подсистема Market Data Acquisition получила полностью самостоятельный слой формирования подписок на поток сделок. + +К этому моменту архитектура уже обеспечивала: + +- универсальный WebSocket Runtime; +- транспортную маршрутизацию входящих сообщений; +- Subscription Layer; +- полноценный Trade Adapter; +- преобразование транспортных сообщений в каноническую модель `Trade`. + +Несмотря на это, один важный архитектурный уровень ещё отсутствовал. + +Система уже умела: + +- подписываться на поток сделок; +- получать WebSocket-документы; +- выполнять их валидацию; +- преобразовывать транспортную модель в каноническую модель `Trade`. + +Однако отсутствовал компонент, объединяющий эти части в единый поток получения данных. + +Между источником транспортных сообщений и канонической моделью не существовало специализированного Feed. + +В результате будущие компоненты системы были бы вынуждены самостоятельно координировать получение документов, их обработку и преобразование в бизнес-модель. + +Подобное решение противоречило одному из фундаментальных принципов архитектуры Dzentra. + +Каждый уровень системы должен обладать собственной областью ответственности и инкапсулировать соответствующую ей логику. + +Главной задачей Build 060.17 становится создание специализированного слоя **Trades Feed**, который объединяет источник транспортных документов и существующий механизм их преобразования в единую точку получения канонических моделей `Trade`. + +После завершения Build система должна получить следующие архитектурные компоненты: + +- специализированный `TradeDocumentSource`; +- специализированный `TradeDocumentHandler`; +- специализированный `TradesFeed`; +- специализированный `TradesFeedRegistry`; +- отдельный тип ошибки регистрации Feed; +- набор Protocol-интерфейсов для всех перечисленных компонентов. + +При этом Build принципиально не затрагивает: + +- WebSocket Runtime; +- Subscription Layer; +- Unified Router; +- Parser; +- Value Validation; +- Mapper; +- Trade Adapter; +- Runtime Registry; +- восстановление подписок; +- обработку истории сделок; +- дедупликацию; +- сортировку сделок; +- интеграцию Feed с Runtime. + +Все перечисленные задачи относятся к следующим этапам дорожной карты. + +--- + +# Предпосылки + +К началу Build архитектура обработки сделок уже была практически полностью сформирована. + +Существовала завершённая вертикаль преобразования транспортного сообщения в каноническую модель. + +Она включала следующие этапы. + +```text +WebSocket Document + │ + ▼ +Schema Validation + │ + ▼ +Parser + │ + ▼ +Value Validation + │ + ▼ +Mapper + │ + ▼ +Trade +``` + +Каждый из перечисленных компонентов обладал строго определённой зоной ответственности. + +Schema Validation проверяла соответствие транспортного документа контракту WebSocket API. + +Parser извлекал необходимые поля. + +Value Validation проверяла корректность полученных значений. + +Mapper строил внутреннюю модель `Trade`. + +Подобная архитектура уже полностью соответствовала общим принципам Market Data Acquisition. + +Однако вся цепочка существовала исключительно как последовательность независимых компонентов. + +Отсутствовал слой, который связывал бы получение транспортного документа и его преобразование в единую операцию. + +Именно эту роль должен был взять на себя новый Feed. + +--- + +# Архитектурное основание + +Одним из фундаментальных принципов архитектуры Dzentra является разделение компонентов по уровням ответственности. + +Каждый уровень отвечает исключительно за собственную задачу. + +Источник данных отвечает только за получение транспортного документа. + +Handler отвечает только за преобразование документа. + +Feed отвечает только за организацию потока данных. + +Ни один из этих компонентов не должен брать на себя ответственность другого. + +До начала Build подобное разделение отсутствовало. + +Trade Adapter уже существовал, однако он представлял собой исключительно преобразователь транспортной модели. + +Он не должен был: + +- получать документы; +- выбирать источник данных; +- координировать последовательность обработки. + +Эти обязанности относятся к уровню Feed. + +После завершения Build архитектура приобретает следующий вид. + +```text +TradeDocumentSource + │ + ▼ +TradeDocumentHandler + │ + ▼ +TradesFeed + │ + ▼ +Trade +``` + +При этом каждый уровень продолжает выполнять только одну архитектурную функцию. + +Такое разделение полностью повторяет ранее реализованную архитектуру других потоков рыночных данных и формирует единый шаблон построения Feed внутри подсистемы Market Data Acquisition. + +--- + +# Результаты архитектурного аудита + +Перед началом реализации был выполнен полный аудит существующей структуры подсистемы Market Data Acquisition. + +Проверка показала, что большая часть необходимой инфраструктуры уже присутствует в проекте. + +В частности, были обнаружены: + +```text +feeds/ + quotes_feed.py + instrument_feed.py +``` + +а также соответствующие обработчики + +```text +handlers/ + quotes_handler.py + instrument_handler.py +``` + +Дополнительно в структуре проекта уже существовали пустые заготовки файлов + +```text +feeds/trades_feed.py + +handlers/trades_handler.py +``` + +что подтверждало первоначальный архитектурный замысел по созданию отдельного Trade Feed. + +Отдельный аудит был проведён для существующего Trade Adapter. + +Проверка показала, что функция + +```text +adapt_websocket_trade_document() +``` + +уже реализует полный внутренний конвейер обработки. + +Она выполняет: + +- Parser; +- Value Validation; +- Mapper. + +При этом предварительная Schema Validation в состав Adapter не входит. + +Данное архитектурное решение оказалось принципиально важным. + +Это означало, что новый Handler должен выполнять только одну дополнительную обязанность — проверку транспортного документа перед передачей его Adapter. + +Таким образом существующий Adapter не потребовал никаких изменений. + +Все архитектурные обязанности удалось распределить между новыми компонентами без нарушения существующей реализации. + +--- + +# Архитектурное решение + +По результатам проведённого аудита было принято решение не изменять существующую цепочку преобразования транспортных сообщений. + +Parser, Value Validation и Mapper уже были реализованы, протестированы и использовались Trade Adapter. + +Поэтому вместо изменения существующих компонентов было принято решение построить недостающий уровень координации. + +Новая архитектура получила следующий вид. + +```text +TradeDocumentSource + │ + ▼ +TradeDocumentHandler + │ + ▼ +TradesFeed + │ + ▼ +Trade +``` + +При этом каждый уровень отвечает исключительно за собственную задачу. + +Источник знает только способ получения транспортного документа. + +Handler знает только правила преобразования документа. + +Feed знает только последовательность взаимодействия между двумя предыдущими компонентами. + +Подобное разделение полностью соответствует принятой архитектуре Market Data Acquisition. + +--- + +# Новые Protocol-интерфейсы + +Одной из целей Build являлось формирование полноценного контрактного уровня для нового Feed. + +До начала Build соответствующие интерфейсы отсутствовали. + +В рамках реализации были добавлены три новых Protocol. + +```text +TradeDocumentSource + +TradeDocumentHandler + +TradesFeedProtocol +``` + +Появление Protocol имеет сразу несколько архитектурных преимуществ. + +Во-первых, новый Feed перестаёт зависеть от конкретной реализации источника данных. + +Во-вторых, Feed перестаёт зависеть от конкретной реализации Handler. + +В-третьих, появляется возможность независимого тестирования каждого компонента посредством Stub-реализаций. + +Вся дальнейшая работа Trade Feed строится исключительно через данные контракты. + +Это полностью соответствует принципу Dependency Inversion, принятому в архитектуре Dzentra. + +--- + +## TradeDocumentSource + +Protocol + +```text +TradeDocumentSource +``` + +определяет единственную обязанность. + +Получение транспортного документа сделки. + +Контракт имеет следующий вид. + +```text +symbol + │ + ▼ +fetch_trade_document() + │ + ▼ +raw document +``` + +Источник не знает: + +- структуру модели `Trade`; +- Parser; +- Mapper; +- Feed; +- Runtime. + +Он отвечает исключительно за получение транспортного документа. + +--- + +## TradeDocumentHandler + +Вторым новым контрактом стал + +```text +TradeDocumentHandler +``` + +Его обязанностью является преобразование транспортного документа в каноническую модель. + +Контракт выглядит следующим образом. + +```text +raw document + │ + ▼ +handle_trade_document() + │ + ▼ +Trade +``` + +Принципиальным архитектурным решением стало использование типа + +```text +object +``` + +в качестве входного параметра Protocol. + +Это означает, что контракт полностью независим от конкретной реализации транспортного документа. + +Handler сам определяет, каким образом необходимо проверить и преобразовать входящие данные. + +Подобное решение предотвращает проникновение знаний о транспортной модели Dzengi в общий слой Protocol. + +--- + +## TradesFeedProtocol + +Последним новым интерфейсом стал + +```text +TradesFeedProtocol +``` + +Он определяет публичный контракт самого Feed. + +Feed предоставляет единственную операцию. + +```text +load_trade(symbol) + │ + ▼ +Trade +``` + +Никаких дополнительных обязанностей Protocol не содержит. + +Он не определяет: + +- очереди; +- кеш; +- подписки; +- историю; +- Runtime. + +Все перечисленные возможности будут добавляться в следующих Build без изменения публичного интерфейса Feed. + +Именно поэтому уже сейчас выбран максимально простой и устойчивый контракт. + +--- + +# Новый TradeFeedRegistry + +Следующим отсутствующим архитектурным уровнем являлся Registry. + +До начала Build инфраструктура уже содержала аналогичные Registry для других потоков рыночных данных. + +Поэтому было принято решение полностью повторить существующий архитектурный шаблон. + +Новый компонент получил название + +```text +TradesFeedRegistry +``` + +Его обязанностью является регистрация и получение экземпляров Feed по имени источника. + +Архитектура Registry имеет следующий вид. + +```text +Source name + │ + ▼ +TradesFeedRegistry + │ + ▼ +TradesFeed +``` + +Registry не создаёт Feed. + +Не управляет жизненным циклом. + +Не выполняет ленивую инициализацию. + +Не знает ничего о Runtime. + +Он лишь хранит зарегистрированные реализации. + +--- + +## Нормализация имени источника + +При регистрации Feed выполняется нормализация имени источника. + +Удаляются внешние пробелы. + +Это обеспечивает идентичное поведение независимо от способа передачи строки. + +Например + +```text +"dzengi" +``` + +и + +```text +" dzengi " +``` + +рассматриваются как один и тот же источник. + +Подобное решение полностью повторяет поведение остальных Registry подсистемы. + +--- + +## Проверка соответствия Protocol + +Во время регистрации выполняется обязательная проверка + +```text +isinstance(..., TradesFeedProtocol) +``` + +Таким образом Registry гарантирует хранение исключительно корректных реализаций Feed. + +Попытка зарегистрировать произвольный объект приводит к генерации специализированного исключения. + +--- + +## Новый тип ошибки + +Для Trade Feed введён собственный тип исключения. + +```text +TradeFeedRegistryError +``` + +Выделение отдельного класса ошибки обеспечивает независимую обработку ошибок регистрации различных типов Feed. + +При этом новый класс наследуется от общей иерархии исключений Market Data Acquisition и полностью соответствует существующей архитектуре проекта. + +--- + +# Новый TradeDocumentHandler + +Одним из ключевых компонентов Build становится + +```text +DzengiTradeDocumentHandler +``` + +Во время проектирования рассматривалось несколько вариантов распределения обязанностей между Handler и Adapter. + +Первоначально предполагалось передавать в Handler уже предварительно проверенный документ. + +Однако проведённый архитектурный аудит показал, что подобный подход приводит к утечке транспортной модели Dzengi в общий слой Protocol. + +Поэтому было принято другое решение. + +Handler принимает обычный объект. + +```text +object +``` + +После чего самостоятельно выполняет Schema Validation. + +Лишь затем документ передаётся существующему Adapter. + +Итоговая последовательность выглядит следующим образом. + +```text +raw document + │ + ▼ +Schema Validation + │ + ▼ +Trade Adapter + │ + ▼ +Trade +``` + +Таким образом удалось сохранить полную независимость общего Protocol от конкретной реализации транспортной модели. + +--- + +## Почему Schema Validation выполняет Handler + +Данное решение было принято осознанно. + +Schema Validation относится к уровню обработки транспортного документа. + +Feed не должен знать устройство WebSocket-сообщения. + +Adapter не должен заниматься проверкой структуры документа. + +Он получает уже корректную транспортную модель. + +Следовательно именно Handler становится естественным местом выполнения Schema Validation. + +Такое распределение обязанностей обеспечивает строгое разделение ответственности между всеми уровнями архитектуры. + +--- + +# Новый TradesFeed + +Центральным компонентом настоящего Build становится + +```text +TradesFeed +``` + +Именно он завершает формирование базовой инфраструктуры получения сделок внутри подсистемы Market Data Acquisition. + +До начала Build существовали все необходимые строительные блоки. + +Система уже умела: + +- получать транспортные сообщения; +- выполнять их валидацию; +- преобразовывать транспортную модель в каноническую модель `Trade`. + +Однако отсутствовал компонент, который объединял бы эти операции в единый сценарий получения данных. + +После завершения Build данную роль выполняет `TradesFeed`. + +--- + +## Архитектура Feed + +Конструкция Feed намеренно сделана максимально простой. + +Она состоит только из двух зависимостей. + +```text +TradeDocumentSource + +TradeDocumentHandler +``` + +При создании Feed оба компонента передаются через конструктор. + +```text +TradesFeed + │ + ├── TradeDocumentSource + │ + └── TradeDocumentHandler +``` + +Подобная схема полностью соответствует принципу Dependency Injection. + +Feed ничего не знает о конкретной реализации источника данных. + +Он также не знает, каким образом реализован Handler. + +Ему известны исключительно Protocol-контракты. + +Это обеспечивает независимость компонентов и значительно упрощает тестирование. + +--- + +## Последовательность обработки + +Работа Feed состоит всего из двух операций. + +Сначала выполняется получение транспортного документа. + +Затем этот документ передаётся Handler. + +Полученная каноническая модель возвращается вызывающему компоненту. + +Полный конвейер выглядит следующим образом. + +```text +symbol + │ + ▼ +TradeDocumentSource + │ + ▼ +raw document + │ + ▼ +TradeDocumentHandler + │ + ▼ +Trade +``` + +Feed не выполняет никаких дополнительных действий. + +Он не изменяет данные. + +Не валидирует значения. + +Не преобразует транспортную модель. + +Не выполняет кеширование. + +Не хранит историю. + +Не осуществляет повторные запросы. + +Он лишь организует взаимодействие двух независимых компонентов. + +--- + +# Почему TradesFeed является Orchestration Layer + +Во время проектирования отдельно рассматривался вопрос о распределении логики между Feed и Handler. + +Существовало два возможных варианта. + +Первый вариант предполагал перенос части логики обработки непосредственно в Feed. + +В этом случае Feed должен был бы: + +- выполнять Schema Validation; +- вызывать Parser; +- вызывать Mapper; +- контролировать отдельные этапы обработки. + +После анализа архитектуры данный подход был отклонён. + +Feed переставал быть координатором и начинал выполнять функции обработчика. + +Это нарушало бы принцип единственной ответственности. + +Поэтому был выбран второй вариант. + +Feed не содержит бизнес-логики. + +Он лишь организует последовательность вызовов. + +Именно поэтому архитектурно он относится к категории **Orchestration Layer**. + +--- + +## Преимущества выбранного решения + +Подобная архитектура обладает рядом важных преимуществ. + +Во-первых, Feed остаётся исключительно лёгким координатором. + +Во-вторых, обработка документа полностью сосредоточена внутри Handler. + +В-третьих, изменение внутреннего устройства Adapter не требует изменения Feed. + +В-четвёртых, в дальнейшем Handler может быть заменён другой реализацией без изменения Feed. + +Таким образом достигается максимально слабая связанность компонентов. + +--- + +# Почему Adapter не изменялся + +Во время реализации Build отдельно анализировалась возможность расширения существующего Adapter. + +На первый взгляд могло показаться логичным добавить в него Schema Validation. + +После детального анализа было принято решение этого не делать. + +Trade Adapter уже обладает строго определённой зоной ответственности. + +Он отвечает исключительно за преобразование транспортной модели. + +В его состав входят: + +```text +Parser + +Value Validation + +Mapper +``` + +Schema Validation относится к предыдущему уровню обработки. + +Она проверяет корректность транспортного документа ещё до появления транспортной модели. + +Следовательно включение Schema Validation внутрь Adapter привело бы к смешению двух различных архитектурных уровней. + +Поэтому существующий Adapter был оставлен без изменений. + +Это позволило полностью сохранить обратную совместимость всей ранее реализованной инфраструктуры. + +--- + +# Полный конвейер обработки сделки + +После завершения Build архитектура обработки Trade принимает окончательный вид. + +```text +TradeDocumentSource + │ + ▼ +DzengiTradeDocumentHandler + │ + ▼ +Schema Validation + │ + ▼ +Trade Adapter + │ + ▼ +Parser + │ + ▼ +Value Validation + │ + ▼ +Mapper + │ + ▼ +Trade +``` + +Таким образом каждый уровень отвечает исключительно за собственную область ответственности. + +Источник отвечает за получение данных. + +Handler отвечает за подготовку документа. + +Adapter отвечает за преобразование транспортной модели. + +Feed отвечает за организацию взаимодействия между компонентами. + +Подобное разделение полностью соответствует архитектурным принципам Dzentra. + +--- + +# Использование существующей инфраструктуры + +Одним из важнейших требований настоящего Build являлось максимальное повторное использование уже существующих компонентов проекта. + +В ходе реализации не создавались новые Parser. + +Не создавались новые Mapper. + +Не создавались новые модели данных. + +Не создавались новые транспортные структуры. + +Вместо этого новый Feed полностью использует уже существующую инфраструктуру. + +Повторно используются: + +- Schema Validation; +- Trade Adapter; +- Parser; +- Value Validation; +- Mapper; +- каноническая модель `Trade`. + +Это существенно снижает риск появления регрессий и обеспечивает единый механизм обработки Trade во всех компонентах системы. + +--- + +# Почему Feed не взаимодействует с Runtime + +Во время проектирования отдельно рассматривалась возможность непосредственного подключения Feed к WebSocket Runtime. + +На первый взгляд подобная схема позволяла уменьшить количество промежуточных компонентов. + +Однако после анализа архитектуры данный вариант был отклонён. + +Runtime отвечает исключительно за транспорт. + +Feed отвечает исключительно за получение канонических моделей. + +Эти два уровня не должны знать друг о друге. + +Связь между ними будет реализована отдельным интеграционным слоем в одном из следующих Build. + +Благодаря такому решению текущий Feed остаётся полностью независимым от способа получения транспортных документов. + +Он может работать: + +- с WebSocket; +- с REST; +- с тестовыми источниками; +- с Mock-реализациями. + +Без каких-либо изменений собственного кода. + +--- + +# Подготовка к дальнейшему развитию + +Несмотря на минимальный объём собственной логики, настоящий Build закладывает фундамент для последующих этапов развития Trade Pipeline. + +Именно поверх нового Feed будут реализованы: + +- непрерывный поток сделок; +- интеграция с WebSocket Runtime; +- восстановление подписок после reconnect; +- дедупликация сообщений; +- упорядочивание сделок; +- обработка пропусков; +- REST Backfill; +- синхронизация истории. + +Таким образом Build 060.17 завершает создание базовой архитектуры Trades Feed и формирует стабильную основу для последующего развития подсистемы. + +--- + +# Изменённые файлы + +В рамках Build были добавлены и расширены несколько компонентов подсистемы Market Data Acquisition. + +Все изменения являются локальными и полностью соответствуют согласованному scope настоящего Build. + +Существующая архитектура Parser, Mapper, Runtime и Subscription Layer изменена не была. + +--- + +## Protocol + +```text +src/market_data/acquisition/protocol.py +``` + +Добавлены три новых контракта. + +```text +TradeDocumentSource + +TradeDocumentHandler + +TradesFeedProtocol +``` + +Новые Protocol завершают контрактный уровень Trade Feed и позволяют всем последующим компонентам взаимодействовать исключительно через абстракции. + +--- + +## Exceptions + +```text +src/market_data/acquisition/exceptions.py +``` + +Добавлен новый специализированный тип ошибки. + +```text +TradeFeedRegistryError +``` + +Он используется исключительно инфраструктурой регистрации Feed. + +Выделение собственного класса позволяет независимо обрабатывать ошибки различных Registry. + +--- + +## Registry + +```text +src/market_data/acquisition/registry.py +``` + +Добавлен новый компонент + +```text +TradesFeedRegistry +``` + +Реализованы: + +- регистрация Feed; +- получение Feed; +- проверка соответствия Protocol; +- защита от повторной регистрации; +- нормализация имени источника; +- генерация специализированных ошибок. + +Архитектура полностью повторяет существующий шаблон остальных Registry проекта. + +--- + +## Handler + +```text +src/market_data/acquisition/handlers/trades_handler.py +``` + +Реализован + +```text +DzengiTradeDocumentHandler +``` + +В его обязанности входят: + +- Schema Validation; +- вызов существующего Trade Adapter; +- возврат канонической модели `Trade`. + +Handler не содержит собственной бизнес-логики и не изменяет существующий Adapter. + +--- + +## Feed + +```text +src/market_data/acquisition/feeds/trades_feed.py +``` + +Реализован новый + +```text +TradesFeed +``` + +Feed выполняет исключительно координацию взаимодействия между: + +- TradeDocumentSource; +- TradeDocumentHandler. + +Feed не содержит логики обработки данных. + +--- + +## Adapter Export + +```text +src/market_data/acquisition/adapters/dzengi/__init__.py +``` + +Добавлен экспорт + +```text +adapt_websocket_trade_document +``` + +Изменение носит исключительно инфраструктурный характер и обеспечивает единообразное использование Adapter другими компонентами системы. + +--- + +# Добавленные unit-тесты + +Настоящий Build сопровождается расширением существующего набора unit-тестов. + +Все новые компоненты получили независимое покрытие. + +Кроме того, были расширены тесты существующей инфраструктуры Protocol и Registry. + +--- + +## Handler + +Добавлен новый файл. + +```text +tests/unit/market_data/acquisition/handlers/test_trades_handler.py +``` + +Проверяются следующие сценарии. + +--- + +### Корректная обработка документа + +Подтверждается успешное преобразование транспортного документа в модель `Trade`. + +--- + +### Вызов Schema Validation + +Подтверждается, что Handler выполняет предварительную проверку документа перед передачей его Adapter. + +--- + +### Использование Adapter + +Подтверждается, что после успешной проверки документ передаётся существующему Trade Adapter. + +--- + +### Возврат канонической модели + +Подтверждается, что результатом работы Handler всегда является объект `Trade`. + +--- + +### Обработка ошибок + +Проверяется корректное распространение исключений при нарушении структуры входного документа. + +--- + +## Feed + +Добавлен новый файл. + +```text +tests/unit/market_data/acquisition/feeds/test_trades_feed.py +``` + +Проверяются следующие сценарии. + +--- + +### Получение документа + +Подтверждается вызов метода + +```text +fetch_trade_document() +``` + +у источника данных. + +--- + +### Передача документа Handler + +Подтверждается, что Feed передаёт полученный документ обработчику без изменений. + +--- + +### Возврат модели Trade + +Подтверждается возврат канонической модели вызывающему компоненту. + +--- + +### Корректная последовательность вызовов + +Проверяется, что сначала вызывается Source, а затем Handler. + +--- + +## Protocol + +Расширен существующий файл. + +```text +tests/unit/market_data/acquisition/test_protocol.py +``` + +Добавлены проверки новых контрактов. + +Подтверждается: + +- соответствие `TradeDocumentSource`; +- соответствие `TradeDocumentHandler`; +- соответствие `TradesFeedProtocol`; +- корректность Structural Typing. + +--- + +## Registry + +Расширен существующий файл. + +```text +tests/unit/market_data/acquisition/test_registry.py +``` + +Добавлены проверки нового + +```text +TradesFeedRegistry +``` + +Проверяются: + +- успешная регистрация Feed; +- получение зарегистрированного Feed; +- проверка Protocol; +- защита от повторной регистрации; +- нормализация имени источника; +- обработка неизвестного источника; +- генерация `TradeFeedRegistryError`; +- наследование нового исключения от общей иерархии ошибок Market Data Acquisition. + +--- + +# Результаты тестирования + +После завершения реализации был выполнен запуск полного набора unit-тестов подсистемы Market Data Acquisition. + +Использовалась команда. + +```bash +python -m pytest tests/unit/market_data/acquisition -q +``` + +Результат выполнения. + +```text +950 passed +``` + +Все существующие тесты проекта успешно пройдены. + +Новые компоненты не вызвали регрессий ранее реализованной функциональности. + +--- + +# Регрессионное тестирование + +Дополнительно были выполнены целевые проверки новых компонентов. + +На ранних этапах разработки тестирование проводилось отдельно для: + +- Trade Handler; +- Trades Feed. + +После завершения реализации был выполнен полный прогон всех тестов подсистемы. + +Результат подтвердил полную совместимость нового Feed с существующей архитектурой. + +Ни один ранее реализованный компонент не потребовал изменений. + +--- + +# Проверка компиляции + +После завершения реализации выполнена проверка корректности компиляции новых компонентов. + +Ошибок синтаксиса обнаружено не было. + +Все новые файлы успешно проходят статическую проверку интерпретатора Python. + +--- + +# Проверка Git diff + +После завершения Build рекомендуется выполнить стандартную инженерную проверку. + +```bash +git diff --check +``` + +Ожидаемый результат. + +```text +без замечаний +``` + +Проверка подтверждает отсутствие: + +- trailing whitespace; +- конфликтов merge; +- нарушений форматирования; +- ошибок окончания строк. + +--- + +# Scope Build 060.17 + +Настоящий Build ограничен исключительно реализацией инфраструктуры **Trades Feed Core**. + +В область реализации входят: + +- Protocol; +- Registry; +- Handler; +- Feed; +- специализированное исключение; +- unit-тесты новой инфраструктуры. + +Build **не включает**: + +- интеграцию с Runtime; +- непрерывное получение Trade; +- подписки WebSocket; +- дедупликацию; +- сортировку Trade; +- восстановление истории; +- REST Backfill; +- агрегирование сделок; +- хранение состояния Feed. + +Все перечисленные возможности относятся к следующим этапам дорожной карты. + +--- + +# Архитектурный результат + +После завершения Build 060.17 подсистема Market Data Acquisition получила полностью сформированный базовый уровень **Trades Feed**. + +Архитектура обработки сделок приобрела следующий законченный вид. + +```text +TradeDocumentSource + │ + ▼ +TradeDocumentHandler + │ + ▼ +Schema Validation + │ + ▼ +Trade Adapter + │ + ▼ +Parser + │ + ▼ +Value Validation + │ + ▼ +Mapper + │ + ▼ +Trade +``` + +Над данной цепочкой располагается компонент + +```text +TradesFeed +``` + +который организует взаимодействие между Source и Handler. + +При этом Feed не вмешивается в процесс преобразования данных и не содержит собственной бизнес-логики. + +Каждый компонент отвечает исключительно за собственную область ответственности. + +Подобная архитектура полностью соответствует принципу строгого разделения ответственности, принятому в проекте Dzentra. + +--- + +# Соблюдение архитектурных принципов + +В рамках настоящего Build полностью сохранены все архитектурные инварианты проекта. + +--- + +## Локальность изменений + +Все изменения ограничены исключительно инфраструктурой Trade Feed. + +Не изменялись: + +- Runtime; +- Subscription Layer; +- Parser; +- Mapper; +- Value Validation; +- транспортные модели; +- каноническая модель `Trade`. + +Это позволило свести риск возникновения регрессий к минимуму. + +--- + +## Повторное использование существующей инфраструктуры + +Одним из ключевых требований Build являлось максимально возможное повторное использование уже реализованных компонентов. + +Новый Feed не создаёт собственных механизмов обработки. + +Он использует уже существующие: + +- Schema Validation; +- Trade Adapter; +- Parser; +- Value Validation; +- Mapper. + +Подобное решение обеспечивает единый конвейер обработки Trade независимо от точки входа данных. + +--- + +## Разделение ответственности + +Каждый новый компонент обладает единственной областью ответственности. + +TradeDocumentSource отвечает исключительно за получение транспортного документа. + +TradeDocumentHandler отвечает исключительно за подготовку и преобразование документа. + +TradesFeed отвечает исключительно за координацию работы предыдущих компонентов. + +TradesFeedRegistry отвечает исключительно за регистрацию реализаций Feed. + +Ни один компонент не выполняет обязанности другого. + +Это полностью соответствует принципу Single Responsibility. + +--- + +## Использование Protocol + +Все зависимости нового Feed построены через Protocol. + +Ни один компонент не зависит от конкретной реализации другого. + +Подобный подход обеспечивает: + +- слабую связанность компонентов; +- простоту unit-тестирования; +- возможность замены реализаций без изменения архитектуры. + +--- + +## Отсутствие состояния + +Новый Feed не хранит собственного состояния. + +Он не содержит: + +- кеш; +- историю; +- активные подписки; +- очередь сообщений; +- накопленные сделки. + +Каждый вызов метода + +```text +load_trade() +``` + +является полностью независимым. + +Подобное решение соответствует текущему этапу развития архитектуры. + +Stateful-механизмы будут реализованы отдельными Build. + +--- + +## Подготовка к масштабированию + +Несмотря на минимальный объём собственной логики, настоящий Build закладывает архитектурную основу для дальнейшего развития Trade Pipeline. + +Поверх текущего Feed могут быть реализованы: + +- непрерывный WebSocket Feed; +- буферизация сделок; +- сортировка по времени; +- дедупликация сообщений; +- REST Backfill; +- синхронизация после reconnect; +- агрегирование Trade. + +При этом публичный контракт Feed останется неизменным. + +--- + +# Архитектурные решения Build (ADR) + +## ADR-060.17-001 + +**TradesFeed является исключительно Orchestration Layer.** + +Feed не содержит бизнес-логики. + +Его единственная обязанность заключается в организации взаимодействия между Source и Handler. + +Все преобразования данных выполняются нижележащими компонентами. + +--- + +## ADR-060.17-002 + +**Schema Validation выполняется внутри Handler.** + +Проверка структуры транспортного документа относится к уровню обработки документа. + +Feed не должен знать устройство WebSocket-сообщения. + +Adapter не должен заниматься проверкой транспортного контракта. + +Поэтому ответственность за Schema Validation закреплена за Handler. + +--- + +## ADR-060.17-003 + +**Trade Adapter остаётся неизменным.** + +Существующий Adapter уже реализует: + +- Parser; +- Value Validation; +- Mapper. + +Расширение его обязанностей привело бы к нарушению принципа единственной ответственности. + +В рамках Build Adapter используется повторно без каких-либо изменений. + +--- + +## ADR-060.17-004 + +**Все зависимости Feed определяются через Protocol.** + +Feed взаимодействует исключительно с контрактами: + +```text +TradeDocumentSource + +TradeDocumentHandler +``` + +Конкретные реализации остаются полностью взаимозаменяемыми. + +Это обеспечивает слабую связанность архитектуры и упрощает тестирование. + +--- + +## ADR-060.17-005 + +**TradesFeedRegistry повторяет существующий архитектурный шаблон Registry.** + +Для нового Feed не создавался отдельный механизм регистрации. + +Использован уже принятый архитектурный шаблон проекта. + +Благодаря этому все Registry подсистемы Market Data Acquisition обладают единым поведением и единым жизненным циклом. + +--- + +## ADR-060.17-006 + +**Новый Feed не зависит от Runtime.** + +Несмотря на то что в дальнейшем Trade будет поступать через WebSocket Runtime, настоящий Build сознательно не связывает эти компоненты. + +Feed остаётся независимым от источника транспортных документов. + +Это позволяет использовать его: + +- с WebSocket; +- с REST; +- с Mock-реализациями; +- в unit-тестах. + +Без изменения собственного кода. + +--- + +# Критерии завершения Build + +Build 060.17 считается завершённым, поскольку выполнены все поставленные задачи. + +- ✔ добавлен `TradeDocumentSource`; +- ✔ добавлен `TradeDocumentHandler`; +- ✔ добавлен `TradesFeedProtocol`; +- ✔ реализован `TradesFeedRegistry`; +- ✔ реализован `TradeFeedRegistryError`; +- ✔ реализован `DzengiTradeDocumentHandler`; +- ✔ реализован `TradesFeed`; +- ✔ выполнен экспорт `adapt_websocket_trade_document`; +- ✔ реализованы unit-тесты Handler; +- ✔ реализованы unit-тесты Feed; +- ✔ расширены тесты Protocol; +- ✔ расширены тесты Registry; +- ✔ успешно пройден полный набор тестов (`950 passed`); +- ✔ существующая архитектура не нарушена; +- ✔ изменения полностью соответствуют согласованному scope Build. + +--- + +# Следующий этап + +Следующим этапом дорожной карты является + +```text +Build 060.18 — Trades Feed Runtime Integration +``` + +Основной целью следующего Build станет интеграция нового Trades Feed с инфраструктурой WebSocket Runtime. + +На данном этапе предстоит реализовать: + +- подключение Feed к Runtime; +- использование Subscription Layer; +- получение непрерывного потока сделок; +- взаимодействие с Runtime Commands; +- подготовку инфраструктуры для обработки непрерывного потока Trade. + +Build 060.17 создаёт необходимый фундамент для данной интеграции. + +--- + +# Итог + +Build 060.17 завершил формирование базовой инфраструктуры **Trades Feed Core** внутри подсистемы Market Data Acquisition. + +В рамках Build был реализован полный набор контрактов, необходимых для организации потока сделок: + +- `TradeDocumentSource`; +- `TradeDocumentHandler`; +- `TradesFeedProtocol`; +- `TradesFeedRegistry`; +- `TradeFeedRegistryError`; +- `DzengiTradeDocumentHandler`; +- `TradesFeed`. + +Все новые компоненты интегрированы в существующую архитектуру без изменения ранее реализованной инфраструктуры. + +Особое внимание было уделено сохранению архитектурных принципов проекта. + +Feed реализован как лёгкий Orchestration Layer. + +Handler инкапсулирует подготовку транспортного документа. + +Adapter продолжает выполнять исключительно преобразование транспортной модели в каноническую модель `Trade`. + +В результате Build сформировал завершённый базовый слой получения сделок, который полностью соответствует архитектуре Market Data Acquisition, успешно прошёл полное регрессионное тестирование (`950 passed`) и стал фундаментом для последующей интеграции с WebSocket Runtime. \ No newline at end of file