From 618b716e71d5d22820fd4a7016ad6d594fd82f2e Mon Sep 17 00:00:00 2001 From: Sergey Date: Mon, 20 Jul 2026 19:16:19 +0300 Subject: [PATCH] Build 060.18: implement Trade Stream Consistency Controller --- .../acquisition/consistency/__init__.py | 1 + .../trade_stream_consistency_controller.py | 45 ++++ .../consistency/trade_stream_exceptions.py | 17 ++ .../consistency/trade_stream_protocol.py | 29 +++ .../consistency/trade_stream_state.py | 103 ++++++++++ ...est_trade_stream_consistency_controller.py | 138 +++++++++++++ .../consistency/test_trade_stream_state.py | 192 ++++++++++++++++++ 7 files changed, 525 insertions(+) create mode 100644 app/src/market_data/acquisition/consistency/__init__.py create mode 100644 app/src/market_data/acquisition/consistency/trade_stream_consistency_controller.py create mode 100644 app/src/market_data/acquisition/consistency/trade_stream_exceptions.py create mode 100644 app/src/market_data/acquisition/consistency/trade_stream_protocol.py create mode 100644 app/src/market_data/acquisition/consistency/trade_stream_state.py create mode 100644 app/tests/unit/market_data/acquisition/consistency/test_trade_stream_consistency_controller.py create mode 100644 app/tests/unit/market_data/acquisition/consistency/test_trade_stream_state.py diff --git a/app/src/market_data/acquisition/consistency/__init__.py b/app/src/market_data/acquisition/consistency/__init__.py new file mode 100644 index 0000000..0c602ef --- /dev/null +++ b/app/src/market_data/acquisition/consistency/__init__.py @@ -0,0 +1 @@ +# app/src/market_data/acquisition/consistency/__init__.py \ No newline at end of file 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 new file mode 100644 index 0000000..ee32a09 --- /dev/null +++ b/app/src/market_data/acquisition/consistency/trade_stream_consistency_controller.py @@ -0,0 +1,45 @@ +# app/src/market_data/acquisition/consistency/trade_stream_consistency_controller.py + +from __future__ import annotations + +from src.market_data.acquisition.consistency.trade_stream_protocol import ( + TradeStreamConsistencyProtocol, +) +from src.market_data.acquisition.consistency.trade_stream_state import ( + TradeStreamState, +) +from src.market_data.acquisition.models.trade import Trade + + +class TradeStreamConsistencyController( + TradeStreamConsistencyProtocol, +): + """ + Контроллер проверки согласованности Canonical Trade Stream. + + Для каждого торгового инструмента поддерживается + независимое состояние проверки. + """ + + def __init__(self) -> None: + self._states: dict[str, TradeStreamState] = {} + + def accept( + self, + trade: Trade, + ) -> Trade | None: + state = self._get_state(trade.symbol) + + return state.accept(trade) + + def _get_state( + self, + symbol: str, + ) -> TradeStreamState: + state = self._states.get(symbol) + + if state is None: + state = TradeStreamState(symbol=symbol) + self._states[symbol] = state + + return state \ No newline at end of file diff --git a/app/src/market_data/acquisition/consistency/trade_stream_exceptions.py b/app/src/market_data/acquisition/consistency/trade_stream_exceptions.py new file mode 100644 index 0000000..ffb4dbf --- /dev/null +++ b/app/src/market_data/acquisition/consistency/trade_stream_exceptions.py @@ -0,0 +1,17 @@ +# app/src/market_data/acquisition/consistency/trade_stream_exceptions.py + +from __future__ import annotations + +from src.market_data.acquisition.exceptions import ( + MarketDataAcquisitionError, +) + + +# Нарушение монотонности Canonical Trade Stream. +class TradeOrderingError(MarketDataAcquisitionError): + pass + + +# Обнаружен конфликтующий дубликат сделки. +class TradeConsistencyError(MarketDataAcquisitionError): + pass \ No newline at end of file diff --git a/app/src/market_data/acquisition/consistency/trade_stream_protocol.py b/app/src/market_data/acquisition/consistency/trade_stream_protocol.py new file mode 100644 index 0000000..076f823 --- /dev/null +++ b/app/src/market_data/acquisition/consistency/trade_stream_protocol.py @@ -0,0 +1,29 @@ +# app/src/market_data/acquisition/consistency/trade_stream_protocol.py + +from __future__ import annotations + +from typing import Protocol, runtime_checkable + +from src.market_data.acquisition.models.trade import Trade + + +# Контракт проверки согласованности Canonical Trade Stream. +@runtime_checkable +class TradeStreamConsistencyProtocol(Protocol): + def accept( + self, + trade: Trade, + ) -> Trade | None: + """ + Проверить согласованность поступившей сделки. + + Возвращает исходную модель Trade, если сделка принята. + + Возвращает None, если обнаружен полный дубликат. + + Генерирует TradeOrderingError при нарушении порядка. + + Генерирует TradeConsistencyError при обнаружении + конфликтующего дубликата. + """ + ... \ No newline at end of file diff --git a/app/src/market_data/acquisition/consistency/trade_stream_state.py b/app/src/market_data/acquisition/consistency/trade_stream_state.py new file mode 100644 index 0000000..cdca0a4 --- /dev/null +++ b/app/src/market_data/acquisition/consistency/trade_stream_state.py @@ -0,0 +1,103 @@ +# app/src/market_data/acquisition/consistency/trade_stream_state.py + +from __future__ import annotations + +from collections import deque +from dataclasses import dataclass, field + +from src.market_data.acquisition.models.trade import Trade +from src.market_data.acquisition.consistency.trade_stream_exceptions import ( + TradeConsistencyError, + TradeOrderingError, +) + +DEFAULT_DEDUPLICATION_WINDOW_SIZE = 10_000 + + +@dataclass(slots=True) +class TradeStreamState: + """ + Состояние проверки согласованности Canonical Trade Stream + для одного торгового инструмента. + """ + + symbol: str + deduplication_window_size: int = DEFAULT_DEDUPLICATION_WINDOW_SIZE + + last_trade_id: int | None = None + + _trade_window: deque[int] = field(init=False, repr=False) + _trades: dict[int, Trade] = field(init=False, repr=False) + + def __post_init__(self) -> None: + if not self.symbol: + raise ValueError("symbol must not be empty") + + if self.deduplication_window_size <= 0: + raise ValueError( + "deduplication_window_size must be positive" + ) + + self._trade_window = deque( + maxlen=self.deduplication_window_size, + ) + self._trades = {} + + def accept( + self, + trade: Trade, + ) -> Trade | None: + """ + Проверить сделку на согласованность. + + Возвращает принятую сделку. + + Возвращает None при полном дубликате. + + Генерирует исключение при нарушении порядка + либо конфликтующем дубликате. + """ + + if trade.symbol != self.symbol: + raise ValueError( + f"Unexpected symbol: {trade.symbol!r}" + ) + + trade_id = trade.trade_id + + if self.last_trade_id is not None: + if trade_id < self.last_trade_id: + previous = self._trades.get(trade_id) + + if previous is None: + raise TradeOrderingError() + + if previous == trade: + return None + + raise TradeConsistencyError() + + previous = self._trades.get(trade_id) + + if previous is not None: + if previous == trade: + return None + + raise TradeConsistencyError() + + self._append(trade) + + self.last_trade_id = trade_id + + return trade + + def _append( + self, + trade: Trade, + ) -> None: + if len(self._trade_window) == self.deduplication_window_size: + oldest_trade_id = self._trade_window[0] + self._trades.pop(oldest_trade_id, None) + + self._trade_window.append(trade.trade_id) + self._trades[trade.trade_id] = trade \ No newline at end of file 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 new file mode 100644 index 0000000..d350be5 --- /dev/null +++ b/app/tests/unit/market_data/acquisition/consistency/test_trade_stream_consistency_controller.py @@ -0,0 +1,138 @@ +# app/tests/unit/market_data/acquisition/consistency/test_trade_stream_consistency_controller.py + +from __future__ import annotations + +from datetime import datetime, timezone +from decimal import Decimal + +import pytest + +from src.market_data.acquisition.consistency.trade_stream_consistency_controller import ( + TradeStreamConsistencyController, +) +from src.market_data.acquisition.consistency.trade_stream_exceptions import ( + TradeConsistencyError, + TradeOrderingError, +) +from src.market_data.acquisition.models.trade import ( + Trade, + TradeAggressorSide, +) + + +def _trade( + *, + symbol: str = "BTCUSD", + trade_id: int = 100, + price: Decimal = Decimal("50000.00"), + quantity: Decimal = Decimal("0.25"), + aggressor_side: TradeAggressorSide = TradeAggressorSide.BUY, +) -> Trade: + return Trade( + symbol=symbol, + trade_id=trade_id, + price=price, + quantity=quantity, + executed_at=datetime( + 2026, + 1, + 1, + 12, + 0, + tzinfo=timezone.utc, + ), + aggressor_side=aggressor_side, + source="dzengi", + ) + + +def test_creates_state_for_first_symbol() -> None: + controller = TradeStreamConsistencyController() + + trade = _trade(symbol="BTCUSD") + + result = controller.accept(trade) + + assert result == trade + + +def test_reuses_state_for_same_symbol() -> None: + controller = TradeStreamConsistencyController() + + first_trade = _trade( + symbol="BTCUSD", + trade_id=100, + ) + second_trade = _trade( + symbol="BTCUSD", + trade_id=101, + ) + + controller.accept(first_trade) + + result = controller.accept(second_trade) + + assert result == second_trade + + +def test_keeps_independent_state_per_symbol() -> None: + controller = TradeStreamConsistencyController() + + btc_trade = _trade( + symbol="BTCUSD", + trade_id=100, + ) + eth_trade = _trade( + symbol="ETHUSD", + trade_id=100, + ) + + btc_result = controller.accept(btc_trade) + eth_result = controller.accept(eth_trade) + + assert btc_result == btc_trade + assert eth_result == eth_trade + + +def test_returns_none_for_duplicate_trade() -> None: + controller = TradeStreamConsistencyController() + + trade = _trade() + + controller.accept(trade) + + result = controller.accept(trade) + + assert result is None + + +def test_propagates_ordering_error() -> None: + controller = TradeStreamConsistencyController() + + controller.accept( + _trade(trade_id=100), + ) + + with pytest.raises(TradeOrderingError): + controller.accept( + _trade(trade_id=99), + ) + + +def test_propagates_consistency_error() -> None: + controller = TradeStreamConsistencyController() + + controller.accept( + _trade( + trade_id=100, + price=Decimal("50000.00"), + ), + ) + + with pytest.raises(TradeConsistencyError): + controller.accept( + _trade( + trade_id=100, + price=Decimal("50001.00"), + ), + ) \ No newline at end of file diff --git a/app/tests/unit/market_data/acquisition/consistency/test_trade_stream_state.py b/app/tests/unit/market_data/acquisition/consistency/test_trade_stream_state.py new file mode 100644 index 0000000..7a8e35a --- /dev/null +++ b/app/tests/unit/market_data/acquisition/consistency/test_trade_stream_state.py @@ -0,0 +1,192 @@ +from __future__ import annotations + +from datetime import datetime, timezone +from decimal import Decimal + +import pytest + +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, +) + + +def _trade( + *, + symbol: str = "BTCUSD", + trade_id: int = 100, + price: Decimal = Decimal("50000.00"), + quantity: Decimal = Decimal("0.25"), + executed_at: datetime | None = None, + aggressor_side: TradeAggressorSide = TradeAggressorSide.BUY, + source: str = "dzengi", +) -> Trade: + return Trade( + symbol=symbol, + trade_id=trade_id, + price=price, + quantity=quantity, + executed_at=executed_at + or datetime( + 2026, + 1, + 1, + 12, + 0, + tzinfo=timezone.utc, + ), + aggressor_side=aggressor_side, + source=source, + ) + + +def test_trade_stream_state_uses_slots() -> None: + state = TradeStreamState(symbol="BTCUSD") + + assert not hasattr(state, "__dict__") + + +def test_accepts_first_trade() -> None: + state = TradeStreamState(symbol="BTCUSD") + trade = _trade() + + result = state.accept(trade) + + assert result == trade + assert state.last_trade_id == trade.trade_id + + +def test_accepts_trade_with_greater_trade_id() -> None: + state = TradeStreamState(symbol="BTCUSD") + first_trade = _trade(trade_id=100) + second_trade = _trade(trade_id=101) + + state.accept(first_trade) + result = state.accept(second_trade) + + assert result == second_trade + assert state.last_trade_id == second_trade.trade_id + + +def test_accepts_trade_with_gap() -> None: + state = TradeStreamState(symbol="BTCUSD") + first_trade = _trade(trade_id=100) + trade_after_gap = _trade(trade_id=105) + + state.accept(first_trade) + result = state.accept(trade_after_gap) + + assert result == trade_after_gap + assert state.last_trade_id == trade_after_gap.trade_id + + +def test_returns_none_for_identical_duplicate() -> None: + state = TradeStreamState(symbol="BTCUSD") + trade = _trade() + + state.accept(trade) + result = state.accept(trade) + + assert result is None + assert state.last_trade_id == trade.trade_id + + +def test_raises_consistency_error_for_conflicting_duplicate() -> None: + state = TradeStreamState(symbol="BTCUSD") + original_trade = _trade( + trade_id=100, + price=Decimal("50000.00"), + ) + conflicting_trade = _trade( + trade_id=100, + price=Decimal("50001.00"), + ) + + state.accept(original_trade) + + with pytest.raises(TradeConsistencyError): + state.accept(conflicting_trade) + + +def test_raises_ordering_error_for_older_trade() -> None: + state = TradeStreamState(symbol="BTCUSD") + current_trade = _trade(trade_id=100) + older_trade = _trade(trade_id=99) + + state.accept(current_trade) + + with pytest.raises(TradeOrderingError): + state.accept(older_trade) + + +def test_returns_none_for_duplicate_still_inside_window() -> None: + state = TradeStreamState( + symbol="BTCUSD", + deduplication_window_size=3, + ) + first_trade = _trade(trade_id=100) + second_trade = _trade(trade_id=101) + third_trade = _trade(trade_id=102) + + state.accept(first_trade) + state.accept(second_trade) + state.accept(third_trade) + + result = state.accept(first_trade) + + assert result is None + assert state.last_trade_id == third_trade.trade_id + + +def test_raises_ordering_error_after_trade_leaves_window() -> None: + state = TradeStreamState( + symbol="BTCUSD", + deduplication_window_size=2, + ) + first_trade = _trade(trade_id=100) + second_trade = _trade(trade_id=101) + third_trade = _trade(trade_id=102) + + state.accept(first_trade) + state.accept(second_trade) + state.accept(third_trade) + + with pytest.raises(TradeOrderingError): + state.accept(first_trade) + + +def test_rejects_unexpected_symbol() -> None: + state = TradeStreamState(symbol="BTCUSD") + trade = _trade(symbol="ETHUSD") + + with pytest.raises(ValueError): + state.accept(trade) + + +def test_rejects_empty_symbol() -> None: + with pytest.raises(ValueError): + TradeStreamState(symbol="") + + +@pytest.mark.parametrize( + "window_size", + [ + 0, + -1, + ], +) +def test_rejects_non_positive_window_size( + window_size: int, +) -> None: + with pytest.raises(ValueError): + TradeStreamState( + symbol="BTCUSD", + deduplication_window_size=window_size, + ) \ No newline at end of file