Build 060.20.2: integrate Consistency with Trade Runtime
This commit is contained in:
@@ -9,6 +9,9 @@ from src.market_data.acquisition.consistency.trade_stream_state import (
|
|||||||
TradeStreamState,
|
TradeStreamState,
|
||||||
)
|
)
|
||||||
from src.market_data.acquisition.models.trade import Trade
|
from src.market_data.acquisition.models.trade import Trade
|
||||||
|
from src.market_data.acquisition.runtime.trade.trade_runtime_protocol import (
|
||||||
|
TradeRuntimeProtocol,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
class TradeStreamConsistencyController(
|
class TradeStreamConsistencyController(
|
||||||
@@ -18,11 +21,17 @@ class TradeStreamConsistencyController(
|
|||||||
Контроллер проверки согласованности Canonical Trade Stream.
|
Контроллер проверки согласованности Canonical Trade Stream.
|
||||||
|
|
||||||
Для каждого торгового инструмента поддерживается
|
Для каждого торгового инструмента поддерживается
|
||||||
независимое состояние проверки.
|
независимое состояние проверки, которое хранится
|
||||||
|
в Trade Runtime Registry.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
def __init__(self) -> None:
|
_RUNTIME_NAMESPACE = "consistency"
|
||||||
self._states: dict[str, TradeStreamState] = {}
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
runtime: TradeRuntimeProtocol,
|
||||||
|
) -> None:
|
||||||
|
self._runtime = runtime
|
||||||
|
|
||||||
def accept(
|
def accept(
|
||||||
self,
|
self,
|
||||||
@@ -36,10 +45,24 @@ class TradeStreamConsistencyController(
|
|||||||
self,
|
self,
|
||||||
symbol: str,
|
symbol: str,
|
||||||
) -> TradeStreamState:
|
) -> TradeStreamState:
|
||||||
state = self._states.get(symbol)
|
key = self._runtime_key(symbol)
|
||||||
|
|
||||||
if state is None:
|
if self._runtime.is_registered(key):
|
||||||
state = TradeStreamState(symbol=symbol)
|
return self._runtime.get(key)
|
||||||
self._states[symbol] = state
|
|
||||||
|
state = TradeStreamState(symbol=symbol)
|
||||||
|
self._runtime.register(key, state)
|
||||||
|
|
||||||
return state
|
return state
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def _runtime_key(
|
||||||
|
cls,
|
||||||
|
symbol: str,
|
||||||
|
) -> str:
|
||||||
|
"""
|
||||||
|
Возвращает ключ Runtime Registry
|
||||||
|
для состояния проверки Trade Stream.
|
||||||
|
"""
|
||||||
|
|
||||||
|
return f"{cls._RUNTIME_NAMESPACE}:{symbol}"
|
||||||
|
|||||||
@@ -14,10 +14,30 @@ from src.market_data.acquisition.consistency.trade_stream_exceptions import (
|
|||||||
TradeConsistencyError,
|
TradeConsistencyError,
|
||||||
TradeOrderingError,
|
TradeOrderingError,
|
||||||
)
|
)
|
||||||
|
from src.market_data.acquisition.consistency.trade_stream_state import (
|
||||||
|
TradeStreamState,
|
||||||
|
)
|
||||||
from src.market_data.acquisition.models.trade import (
|
from src.market_data.acquisition.models.trade import (
|
||||||
Trade,
|
Trade,
|
||||||
TradeAggressorSide,
|
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(
|
def _trade(
|
||||||
@@ -46,19 +66,35 @@ def _trade(
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def test_creates_state_for_first_symbol() -> None:
|
def test_creates_state_for_first_symbol(
|
||||||
controller = TradeStreamConsistencyController()
|
controller: TradeStreamConsistencyController,
|
||||||
|
runtime: TradeRuntimeRegistry,
|
||||||
|
) -> None:
|
||||||
trade = _trade(symbol="BTCUSD")
|
trade = _trade(symbol="BTCUSD")
|
||||||
|
|
||||||
result = controller.accept(trade)
|
result = controller.accept(trade)
|
||||||
|
|
||||||
assert result == trade
|
assert result == trade
|
||||||
|
assert runtime.is_registered("consistency:BTCUSD")
|
||||||
|
|
||||||
|
|
||||||
def test_reuses_state_for_same_symbol() -> None:
|
def test_registers_trade_stream_state(
|
||||||
controller = TradeStreamConsistencyController()
|
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(
|
first_trade = _trade(
|
||||||
symbol="BTCUSD",
|
symbol="BTCUSD",
|
||||||
trade_id=100,
|
trade_id=100,
|
||||||
@@ -69,15 +105,19 @@ def test_reuses_state_for_same_symbol() -> None:
|
|||||||
)
|
)
|
||||||
|
|
||||||
controller.accept(first_trade)
|
controller.accept(first_trade)
|
||||||
|
first_state = runtime.get("consistency:BTCUSD")
|
||||||
|
|
||||||
result = controller.accept(second_trade)
|
result = controller.accept(second_trade)
|
||||||
|
second_state = runtime.get("consistency:BTCUSD")
|
||||||
|
|
||||||
assert result == second_trade
|
assert result == second_trade
|
||||||
|
assert second_state is first_state
|
||||||
|
|
||||||
|
|
||||||
def test_keeps_independent_state_per_symbol() -> None:
|
def test_keeps_independent_state_per_symbol(
|
||||||
controller = TradeStreamConsistencyController()
|
controller: TradeStreamConsistencyController,
|
||||||
|
runtime: TradeRuntimeRegistry,
|
||||||
|
) -> None:
|
||||||
btc_trade = _trade(
|
btc_trade = _trade(
|
||||||
symbol="BTCUSD",
|
symbol="BTCUSD",
|
||||||
trade_id=100,
|
trade_id=100,
|
||||||
@@ -90,13 +130,17 @@ def test_keeps_independent_state_per_symbol() -> None:
|
|||||||
btc_result = controller.accept(btc_trade)
|
btc_result = controller.accept(btc_trade)
|
||||||
eth_result = controller.accept(eth_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 btc_result == btc_trade
|
||||||
assert eth_result == eth_trade
|
assert eth_result == eth_trade
|
||||||
|
assert btc_state is not eth_state
|
||||||
|
|
||||||
|
|
||||||
def test_returns_none_for_duplicate_trade() -> None:
|
def test_returns_none_for_duplicate_trade(
|
||||||
controller = TradeStreamConsistencyController()
|
controller: TradeStreamConsistencyController,
|
||||||
|
) -> None:
|
||||||
trade = _trade()
|
trade = _trade()
|
||||||
|
|
||||||
controller.accept(trade)
|
controller.accept(trade)
|
||||||
@@ -106,9 +150,9 @@ def test_returns_none_for_duplicate_trade() -> None:
|
|||||||
assert result is None
|
assert result is None
|
||||||
|
|
||||||
|
|
||||||
def test_propagates_ordering_error() -> None:
|
def test_propagates_ordering_error(
|
||||||
controller = TradeStreamConsistencyController()
|
controller: TradeStreamConsistencyController,
|
||||||
|
) -> None:
|
||||||
controller.accept(
|
controller.accept(
|
||||||
_trade(trade_id=100),
|
_trade(trade_id=100),
|
||||||
)
|
)
|
||||||
@@ -119,9 +163,9 @@ def test_propagates_ordering_error() -> None:
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def test_propagates_consistency_error() -> None:
|
def test_propagates_consistency_error(
|
||||||
controller = TradeStreamConsistencyController()
|
controller: TradeStreamConsistencyController,
|
||||||
|
) -> None:
|
||||||
controller.accept(
|
controller.accept(
|
||||||
_trade(
|
_trade(
|
||||||
trade_id=100,
|
trade_id=100,
|
||||||
|
|||||||
Reference in New Issue
Block a user