Build 060.18: implement Trade Stream Consistency Controller
This commit is contained in:
1
app/src/market_data/acquisition/consistency/__init__.py
Normal file
1
app/src/market_data/acquisition/consistency/__init__.py
Normal file
@@ -0,0 +1 @@
|
||||
# app/src/market_data/acquisition/consistency/__init__.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
|
||||
@@ -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
|
||||
@@ -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 при обнаружении
|
||||
конфликтующего дубликата.
|
||||
"""
|
||||
...
|
||||
@@ -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
|
||||
@@ -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"),
|
||||
),
|
||||
)
|
||||
@@ -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,
|
||||
)
|
||||
Reference in New Issue
Block a user