From 14f698d48ad14e064e8b86e068bf6ae3b24a894a Mon Sep 17 00:00:00 2001 From: Sergey Date: Fri, 24 Jul 2026 18:34:33 +0300 Subject: [PATCH] Build 060.20.2: integrate Consistency with Trade Runtime --- .../trade_stream_consistency_controller.py | 39 ++++++++-- ...est_trade_stream_consistency_controller.py | 78 +++++++++++++++---- 2 files changed, 92 insertions(+), 25 deletions(-) diff --git a/app/src/market_data/acquisition/consistency/trade_stream_consistency_controller.py b/app/src/market_data/acquisition/consistency/trade_stream_consistency_controller.py index ee32a09..1a04ce4 100644 --- a/app/src/market_data/acquisition/consistency/trade_stream_consistency_controller.py +++ b/app/src/market_data/acquisition/consistency/trade_stream_consistency_controller.py @@ -9,6 +9,9 @@ from src.market_data.acquisition.consistency.trade_stream_state import ( TradeStreamState, ) from src.market_data.acquisition.models.trade import Trade +from src.market_data.acquisition.runtime.trade.trade_runtime_protocol import ( + TradeRuntimeProtocol, +) class TradeStreamConsistencyController( @@ -18,11 +21,17 @@ class TradeStreamConsistencyController( Контроллер проверки согласованности Canonical Trade Stream. Для каждого торгового инструмента поддерживается - независимое состояние проверки. + независимое состояние проверки, которое хранится + в Trade Runtime Registry. """ - def __init__(self) -> None: - self._states: dict[str, TradeStreamState] = {} + _RUNTIME_NAMESPACE = "consistency" + + def __init__( + self, + runtime: TradeRuntimeProtocol, + ) -> None: + self._runtime = runtime def accept( self, @@ -36,10 +45,24 @@ class TradeStreamConsistencyController( self, symbol: str, ) -> TradeStreamState: - state = self._states.get(symbol) + key = self._runtime_key(symbol) - if state is None: - state = TradeStreamState(symbol=symbol) - self._states[symbol] = state + if self._runtime.is_registered(key): + return self._runtime.get(key) - return state \ No newline at end of file + state = TradeStreamState(symbol=symbol) + self._runtime.register(key, state) + + return state + + @classmethod + def _runtime_key( + cls, + symbol: str, + ) -> str: + """ + Возвращает ключ Runtime Registry + для состояния проверки Trade Stream. + """ + + return f"{cls._RUNTIME_NAMESPACE}:{symbol}" diff --git a/app/tests/unit/market_data/acquisition/consistency/test_trade_stream_consistency_controller.py b/app/tests/unit/market_data/acquisition/consistency/test_trade_stream_consistency_controller.py index d350be5..372ee5c 100644 --- a/app/tests/unit/market_data/acquisition/consistency/test_trade_stream_consistency_controller.py +++ b/app/tests/unit/market_data/acquisition/consistency/test_trade_stream_consistency_controller.py @@ -14,10 +14,30 @@ from src.market_data.acquisition.consistency.trade_stream_exceptions import ( TradeConsistencyError, TradeOrderingError, ) +from src.market_data.acquisition.consistency.trade_stream_state import ( + TradeStreamState, +) from src.market_data.acquisition.models.trade import ( Trade, TradeAggressorSide, ) +from src.market_data.acquisition.runtime.trade.trade_runtime_registry import ( + TradeRuntimeRegistry, +) + + +@pytest.fixture +def runtime() -> TradeRuntimeRegistry: + return TradeRuntimeRegistry() + + +@pytest.fixture +def controller( + runtime: TradeRuntimeRegistry, +) -> TradeStreamConsistencyController: + return TradeStreamConsistencyController( + runtime=runtime, + ) def _trade( @@ -46,19 +66,35 @@ def _trade( ) -def test_creates_state_for_first_symbol() -> None: - controller = TradeStreamConsistencyController() - +def test_creates_state_for_first_symbol( + controller: TradeStreamConsistencyController, + runtime: TradeRuntimeRegistry, +) -> None: trade = _trade(symbol="BTCUSD") result = controller.accept(trade) assert result == trade + assert runtime.is_registered("consistency:BTCUSD") -def test_reuses_state_for_same_symbol() -> None: - controller = TradeStreamConsistencyController() +def test_registers_trade_stream_state( + controller: TradeStreamConsistencyController, + runtime: TradeRuntimeRegistry, +) -> None: + controller.accept( + _trade(symbol="BTCUSD"), + ) + state = runtime.get("consistency:BTCUSD") + + assert isinstance(state, TradeStreamState) + + +def test_reuses_state_for_same_symbol( + controller: TradeStreamConsistencyController, + runtime: TradeRuntimeRegistry, +) -> None: first_trade = _trade( symbol="BTCUSD", trade_id=100, @@ -69,15 +105,19 @@ def test_reuses_state_for_same_symbol() -> None: ) controller.accept(first_trade) + first_state = runtime.get("consistency:BTCUSD") result = controller.accept(second_trade) + second_state = runtime.get("consistency:BTCUSD") assert result == second_trade + assert second_state is first_state -def test_keeps_independent_state_per_symbol() -> None: - controller = TradeStreamConsistencyController() - +def test_keeps_independent_state_per_symbol( + controller: TradeStreamConsistencyController, + runtime: TradeRuntimeRegistry, +) -> None: btc_trade = _trade( symbol="BTCUSD", trade_id=100, @@ -90,13 +130,17 @@ def test_keeps_independent_state_per_symbol() -> None: btc_result = controller.accept(btc_trade) eth_result = controller.accept(eth_trade) + btc_state = runtime.get("consistency:BTCUSD") + eth_state = runtime.get("consistency:ETHUSD") + assert btc_result == btc_trade assert eth_result == eth_trade + assert btc_state is not eth_state -def test_returns_none_for_duplicate_trade() -> None: - controller = TradeStreamConsistencyController() - +def test_returns_none_for_duplicate_trade( + controller: TradeStreamConsistencyController, +) -> None: trade = _trade() controller.accept(trade) @@ -106,9 +150,9 @@ def test_returns_none_for_duplicate_trade() -> None: assert result is None -def test_propagates_ordering_error() -> None: - controller = TradeStreamConsistencyController() - +def test_propagates_ordering_error( + controller: TradeStreamConsistencyController, +) -> None: controller.accept( _trade(trade_id=100), ) @@ -119,9 +163,9 @@ def test_propagates_ordering_error() -> None: ) -def test_propagates_consistency_error() -> None: - controller = TradeStreamConsistencyController() - +def test_propagates_consistency_error( + controller: TradeStreamConsistencyController, +) -> None: controller.accept( _trade( trade_id=100,