Build 060.19: implement Trade Recovery subsystem
This commit is contained in:
@@ -0,0 +1 @@
|
||||
# app/src/market_data/acquisition/adapters/dzengi/auth.py
|
||||
1
app/src/market_data/acquisition/recovery/__init__.py
Normal file
1
app/src/market_data/acquisition/recovery/__init__.py
Normal file
@@ -0,0 +1 @@
|
||||
# app/src/market_data/acquisition/recovery/__init__.py
|
||||
@@ -0,0 +1,118 @@
|
||||
# app/src/market_data/acquisition/recovery/trade_recovery_controller.py
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from src.market_data.acquisition.adapters.dzengi.rest import (
|
||||
DzengiTradesDocumentSource,
|
||||
)
|
||||
from src.market_data.acquisition.adapters.dzengi.rest_trade_adapter import (
|
||||
adapt_rest_agg_trades_document,
|
||||
)
|
||||
from src.market_data.acquisition.consistency.trade_stream_protocol import (
|
||||
TradeStreamConsistencyProtocol,
|
||||
)
|
||||
from src.market_data.acquisition.models.trade import Trade
|
||||
from src.market_data.acquisition.recovery.trade_recovery_normalizer import (
|
||||
normalize_recovered_trades,
|
||||
)
|
||||
from src.market_data.acquisition.recovery.trade_recovery_protocol import (
|
||||
TradeRecoveryProtocol,
|
||||
)
|
||||
from src.market_data.acquisition.recovery.trade_recovery_request import (
|
||||
TradeRecoveryRequest,
|
||||
)
|
||||
from src.market_data.acquisition.recovery.trade_recovery_result import (
|
||||
TradeRecoveryResult,
|
||||
)
|
||||
from src.market_data.acquisition.validation.schema import (
|
||||
validate_rest_agg_trades_schema,
|
||||
)
|
||||
|
||||
|
||||
class TradeRecoveryController(TradeRecoveryProtocol):
|
||||
"""
|
||||
Контроллер восстановления пропущенных сделок через Dzengi REST API.
|
||||
|
||||
Контроллер выполняет одну stateless-операцию восстановления:
|
||||
|
||||
1. получает сырой документ aggTrades;
|
||||
2. выполняет schema validation;
|
||||
3. преобразует документ в канонические Trade;
|
||||
4. нормализует порядок сделок;
|
||||
5. пропускает сделки через общий consistency-контроллер;
|
||||
6. возвращает только принятые сделки.
|
||||
|
||||
Контроллер не владеет состоянием согласованности потока. Для Recovery
|
||||
должен передаваться тот же экземпляр TradeStreamConsistencyProtocol,
|
||||
который используется основным потоком сделок.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
document_source: DzengiTradesDocumentSource,
|
||||
consistency_controller: TradeStreamConsistencyProtocol,
|
||||
) -> None:
|
||||
self._document_source = document_source
|
||||
self._consistency_controller = consistency_controller
|
||||
|
||||
def recover(
|
||||
self,
|
||||
request: TradeRecoveryRequest,
|
||||
) -> TradeRecoveryResult:
|
||||
"""
|
||||
Выполнить одну операцию восстановления сделок.
|
||||
|
||||
Идентичные дубликаты, возвращённые consistency-контроллером
|
||||
как None, не включаются в результат.
|
||||
|
||||
Исключения transport, schema, parsing, value validation, mapping
|
||||
и consistency не перехватываются и не оборачиваются.
|
||||
"""
|
||||
|
||||
document = self._document_source.fetch_trades_document(
|
||||
request.symbol,
|
||||
start_time=request.start_time,
|
||||
end_time=request.end_time,
|
||||
limit=request.limit,
|
||||
)
|
||||
|
||||
validated_document = validate_rest_agg_trades_schema(
|
||||
document,
|
||||
)
|
||||
|
||||
trades = adapt_rest_agg_trades_document(
|
||||
validated_document,
|
||||
symbol=request.symbol,
|
||||
)
|
||||
|
||||
normalized_trades = normalize_recovered_trades(
|
||||
trades,
|
||||
)
|
||||
|
||||
recovered_trades = self._accept_trades(
|
||||
normalized_trades,
|
||||
)
|
||||
|
||||
return TradeRecoveryResult(
|
||||
symbol=request.symbol,
|
||||
requested_start_time=request.start_time,
|
||||
requested_end_time=request.end_time,
|
||||
recovered_trades=recovered_trades,
|
||||
)
|
||||
|
||||
def _accept_trades(
|
||||
self,
|
||||
trades: tuple[Trade, ...],
|
||||
) -> tuple[Trade, ...]:
|
||||
accepted_trades: list[Trade] = []
|
||||
|
||||
for trade in trades:
|
||||
accepted_trade = self._consistency_controller.accept(
|
||||
trade,
|
||||
)
|
||||
|
||||
if accepted_trade is not None:
|
||||
accepted_trades.append(accepted_trade)
|
||||
|
||||
return tuple(accepted_trades)
|
||||
@@ -0,0 +1,37 @@
|
||||
# app/src/market_data/acquisition/recovery/trade_recovery_exceptions.py
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from src.market_data.acquisition.exceptions import (
|
||||
MarketDataAcquisitionError,
|
||||
)
|
||||
|
||||
|
||||
class TradeRecoveryError(MarketDataAcquisitionError):
|
||||
"""
|
||||
Базовое исключение подсистемы восстановления сделок.
|
||||
"""
|
||||
|
||||
|
||||
class TradeRecoveryWindowError(TradeRecoveryError):
|
||||
"""
|
||||
Некорректный диапазон восстановления сделок.
|
||||
"""
|
||||
|
||||
|
||||
class TradeRecoveryLimitError(TradeRecoveryError):
|
||||
"""
|
||||
Некорректное значение параметра limit.
|
||||
"""
|
||||
|
||||
|
||||
class TradeRecoveryNormalizationError(TradeRecoveryError):
|
||||
"""
|
||||
Ошибка нормализации восстановленных сделок.
|
||||
"""
|
||||
|
||||
|
||||
class TradeRecoveryControllerError(TradeRecoveryError):
|
||||
"""
|
||||
Ошибка контроллера восстановления сделок.
|
||||
"""
|
||||
@@ -0,0 +1,40 @@
|
||||
# app/src/market_data/acquisition/recovery/trade_recovery_normalizer.py
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Iterable
|
||||
|
||||
from src.market_data.acquisition.models.trade import Trade
|
||||
|
||||
|
||||
def normalize_recovered_trades(
|
||||
trades: Iterable[Trade],
|
||||
) -> tuple[Trade, ...]:
|
||||
"""
|
||||
Нормализовать порядок восстановленных сделок.
|
||||
|
||||
Сделки возвращаются в возрастающем порядке по ``trade_id``.
|
||||
Исходная последовательность не изменяется.
|
||||
|
||||
Нормализатор намеренно не выполняет дедупликацию и не проверяет
|
||||
согласованность сделок. Идентичные и конфликтующие дубликаты должны
|
||||
обрабатываться экземпляром ``TradeStreamConsistencyController``.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
trades
|
||||
Последовательность канонических сделок.
|
||||
|
||||
Returns
|
||||
-------
|
||||
tuple[Trade, ...]
|
||||
Неизменяемая последовательность сделок, отсортированная
|
||||
по возрастанию ``trade_id``.
|
||||
"""
|
||||
|
||||
return tuple(
|
||||
sorted(
|
||||
trades,
|
||||
key=lambda trade: trade.trade_id,
|
||||
)
|
||||
)
|
||||
@@ -0,0 +1,55 @@
|
||||
# app/src/market_data/acquisition/recovery/trade_recovery_protocol.py
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from abc import abstractmethod
|
||||
from typing import Protocol
|
||||
|
||||
from src.market_data.acquisition.recovery.trade_recovery_request import (
|
||||
TradeRecoveryRequest,
|
||||
)
|
||||
from src.market_data.acquisition.recovery.trade_recovery_result import (
|
||||
TradeRecoveryResult,
|
||||
)
|
||||
|
||||
|
||||
class TradeRecoveryProtocol(Protocol):
|
||||
"""
|
||||
Контракт подсистемы восстановления сделок.
|
||||
|
||||
Реализация должна:
|
||||
|
||||
- получать сделки из внешнего источника;
|
||||
- нормализовать их порядок;
|
||||
- выполнять согласование через
|
||||
TradeStreamConsistencyController;
|
||||
- возвращать канонический результат восстановления.
|
||||
|
||||
Реализация не должна:
|
||||
|
||||
- выполнять повторные попытки;
|
||||
- управлять WebSocket;
|
||||
- управлять Runtime;
|
||||
- выполнять кэширование;
|
||||
- принимать решения о реконнекте.
|
||||
"""
|
||||
|
||||
@abstractmethod
|
||||
def recover(
|
||||
self,
|
||||
request: TradeRecoveryRequest,
|
||||
) -> TradeRecoveryResult:
|
||||
"""
|
||||
Выполнить одну операцию восстановления сделок.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
request
|
||||
Параметры операции восстановления.
|
||||
|
||||
Returns
|
||||
-------
|
||||
TradeRecoveryResult
|
||||
Канонический результат восстановления.
|
||||
"""
|
||||
...
|
||||
@@ -0,0 +1,78 @@
|
||||
# app/src/market_data/acquisition/recovery/trade_recovery_request.py
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
|
||||
|
||||
_MAX_RECOVERY_WINDOW_MS = 60 * 60 * 1000
|
||||
_MIN_RECOVERY_LIMIT = 1
|
||||
_MAX_RECOVERY_LIMIT = 1000
|
||||
|
||||
|
||||
# Неизменяемое описание одного REST-запроса восстановления сделок.
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class TradeRecoveryRequest:
|
||||
symbol: str
|
||||
|
||||
start_time: int
|
||||
end_time: int
|
||||
|
||||
limit: int | None = None
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
"""
|
||||
Проверить локальные инварианты запроса восстановления.
|
||||
|
||||
Временные границы задаются в миллисекундах Unix time и передаются
|
||||
в Dzengi REST API как параметры startTime и endTime.
|
||||
"""
|
||||
|
||||
if not isinstance(self.symbol, str):
|
||||
raise TypeError("symbol должен иметь тип str.")
|
||||
|
||||
if not self.symbol.strip():
|
||||
raise ValueError("symbol не должен быть пустым.")
|
||||
|
||||
if isinstance(self.start_time, bool) or not isinstance(
|
||||
self.start_time,
|
||||
int,
|
||||
):
|
||||
raise TypeError("start_time должен иметь тип int.")
|
||||
|
||||
if isinstance(self.end_time, bool) or not isinstance(
|
||||
self.end_time,
|
||||
int,
|
||||
):
|
||||
raise TypeError("end_time должен иметь тип int.")
|
||||
|
||||
if self.start_time < 0:
|
||||
raise ValueError(
|
||||
"start_time не должен быть отрицательным."
|
||||
)
|
||||
|
||||
if self.end_time < 0:
|
||||
raise ValueError(
|
||||
"end_time не должен быть отрицательным."
|
||||
)
|
||||
|
||||
if self.start_time > self.end_time:
|
||||
raise ValueError(
|
||||
"start_time не должен быть больше end_time."
|
||||
)
|
||||
|
||||
if self.end_time - self.start_time >= _MAX_RECOVERY_WINDOW_MS:
|
||||
raise ValueError(
|
||||
"Диапазон восстановления должен быть меньше одного часа."
|
||||
)
|
||||
|
||||
if self.limit is None:
|
||||
return
|
||||
|
||||
if isinstance(self.limit, bool) or not isinstance(self.limit, int):
|
||||
raise TypeError("limit должен иметь тип int или None.")
|
||||
|
||||
if not _MIN_RECOVERY_LIMIT <= self.limit <= _MAX_RECOVERY_LIMIT:
|
||||
raise ValueError(
|
||||
"limit должен находиться в диапазоне от 1 до 1000."
|
||||
)
|
||||
@@ -0,0 +1,56 @@
|
||||
# app/src/market_data/acquisition/recovery/trade_recovery_result.py
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
|
||||
from src.market_data.acquisition.models.trade import Trade
|
||||
|
||||
|
||||
# Неизменяемый результат одной операции восстановления сделок.
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class TradeRecoveryResult:
|
||||
symbol: str
|
||||
|
||||
requested_start_time: int
|
||||
requested_end_time: int
|
||||
|
||||
recovered_trades: tuple[Trade, ...]
|
||||
|
||||
@property
|
||||
def recovered_count(self) -> int:
|
||||
"""
|
||||
Количество успешно восстановленных сделок.
|
||||
"""
|
||||
|
||||
return len(self.recovered_trades)
|
||||
|
||||
@property
|
||||
def is_empty(self) -> bool:
|
||||
"""
|
||||
Признак отсутствия восстановленных сделок.
|
||||
"""
|
||||
|
||||
return not self.recovered_trades
|
||||
|
||||
@property
|
||||
def first_trade(self) -> Trade | None:
|
||||
"""
|
||||
Первая сделка после нормализации.
|
||||
"""
|
||||
|
||||
if not self.recovered_trades:
|
||||
return None
|
||||
|
||||
return self.recovered_trades[0]
|
||||
|
||||
@property
|
||||
def last_trade(self) -> Trade | None:
|
||||
"""
|
||||
Последняя сделка после нормализации.
|
||||
"""
|
||||
|
||||
if not self.recovered_trades:
|
||||
return None
|
||||
|
||||
return self.recovered_trades[-1]
|
||||
@@ -0,0 +1 @@
|
||||
# app/src/market_data/acquisition/runtime/reconnect.py
|
||||
@@ -0,0 +1 @@
|
||||
# app/src/market_data/acquisition/runtime/supervisor.py
|
||||
Reference in New Issue
Block a user