Build 060.23: integrate Trade Stream Acquisition

This commit is contained in:
2026-07-28 07:20:09 +03:00
parent 9953bab993
commit ee8765b716
7 changed files with 2989 additions and 0 deletions

View File

@@ -0,0 +1,74 @@
# app/src/market_data/acquisition/trade_stream_acquisition_protocol.py
from __future__ import annotations
"""
Публичный контракт интеграции Trade Stream в Acquisition Layer.
Build 060.23 вводит единую сервисную границу между:
- Trade Subscription Layer;
- Acquisition Runtime Service;
- Unified WebSocket Adapter;
- Trade Stream Consistency.
Контракт не определяет production WebSocket transport, receive-loop,
reconnect, recovery, хранение или публикацию канонических Trade.
"""
from typing import Protocol, runtime_checkable
from src.market_data.acquisition.models.trade import Trade
@runtime_checkable
class TradeStreamAcquisitionServiceProtocol(Protocol):
"""
Контракт сервиса интеграции Trade Stream.
Сервис отвечает за:
- передачу команды подписки в Acquisition Runtime;
- преобразование одного входящего WebSocket-документа;
- передачу канонического Trade в Consistency Layer;
- возврат согласованного Trade либо None для дубликата.
Сервис не владеет WebSocket lifecycle и не запускает receive-loop.
"""
async def subscribe(
self,
symbols: tuple[str, ...],
*,
correlation_id: str | None = None,
) -> None:
"""
Передать в Acquisition Runtime команду подписки Trade Stream.
Args:
symbols:
Символы торговых инструментов для подписки.
correlation_id:
Необязательный идентификатор транспортного запроса.
При отсутствии значение создаётся Subscription Builder.
"""
...
def handle_message(
self,
document: object,
) -> Trade | None:
"""
Обработать один входящий WebSocket-документ.
Args:
document:
Сырой документ, полученный из WebSocket Runtime.
Returns:
Канонический Trade после проверки согласованности;
None, если сообщение не относится к Trade либо является
корректным дубликатом.
"""
...

View File

@@ -0,0 +1,111 @@
# app/src/market_data/acquisition/trade_stream_acquisition_service.py
from __future__ import annotations
"""
Сервис интеграции Trade Stream с инфраструктурой Acquisition Runtime.
Build 060.23 вводит сервисную границу между:
- Trade Subscription Layer;
- Acquisition Runtime Service;
- Unified WebSocket Adapter;
- Trade Stream Consistency.
На текущем этапе реализована передача команды подписки.
Обработка входящих сообщений будет добавлена следующим шагом Build.
"""
from src.market_data.acquisition.trade_stream_message_adapter_protocol import (
TradeStreamMessageAdapterProtocol,
)
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.runtime.acquisition_runtime_service_protocol import (
AcquisitionRuntimeServiceProtocol,
)
from src.market_data.acquisition.subscriptions.trades import (
build_trade_subscribe_command,
)
from src.market_data.acquisition.trade_stream_acquisition_protocol import (
TradeStreamAcquisitionServiceProtocol,
)
class TradeStreamAcquisitionService(
TradeStreamAcquisitionServiceProtocol,
):
"""
Координатор инфраструктуры получения Trade Stream.
Сервис объединяет:
- Acquisition Runtime;
- Unified WebSocket Adapter;
- Trade Stream Consistency.
При этом сам сервис не реализует транспорт,
жизненный цикл WebSocket либо Recovery.
"""
def __init__(
self,
runtime_service: AcquisitionRuntimeServiceProtocol,
adapter: TradeStreamMessageAdapterProtocol,
consistency_controller: TradeStreamConsistencyProtocol,
) -> None:
self._runtime_service = runtime_service
self._adapter = adapter
self._consistency_controller = consistency_controller
async def subscribe(
self,
symbols: tuple[str, ...],
*,
correlation_id: str | None = None,
) -> None:
"""
Передать в Acquisition Runtime команду подписки Trade Stream.
Args:
symbols:
Символы торговых инструментов для подписки.
correlation_id:
Необязательный идентификатор транспортного запроса.
При отсутствии значение создаётся Subscription Builder.
"""
command = build_trade_subscribe_command(
symbols,
correlation_id=correlation_id,
)
await self._runtime_service.dispatch(command)
def handle_message(
self,
document: object,
) -> Trade | None:
"""
Обработать один входящий WebSocket-документ.
Документ преобразуется через Unified Adapter. Только канонический
Trade передаётся в Trade Stream Consistency.
Args:
document:
Сырой WebSocket-документ.
Returns:
Trade после проверки согласованности;
None, если сообщение не относится к Trade либо является
корректным дубликатом.
"""
result = self._adapter.map_message(document)
if not isinstance(result, Trade):
return None
return self._consistency_controller.accept(result)

View File

@@ -0,0 +1,40 @@
# app/src/market_data/acquisition/trade_stream_message_adapter_protocol.py
from __future__ import annotations
"""
Контракт адаптера входящих сообщений Trade Stream.
Build 060.23 отделяет сервис координации Acquisition
от конкретной реализации WebSocket-протокола биржи.
"""
from typing import Protocol, runtime_checkable
from src.market_data.acquisition.models.candle_close import (
CandleCloseEvent,
)
from src.market_data.acquisition.models.quote import Quote
from src.market_data.acquisition.models.trade import Trade
TradeStreamMappedMessage = Quote | CandleCloseEvent | Trade
@runtime_checkable
class TradeStreamMessageAdapterProtocol(Protocol):
"""
Контракт преобразования одного входящего WebSocket-документа.
Реализация может быть exchange-specific, но вызывающий сервис
зависит только от результата преобразования.
"""
def map_message(
self,
document: object,
) -> TradeStreamMappedMessage:
"""
Преобразовать один WebSocket-документ в каноническую модель.
"""
...

View File

@@ -0,0 +1,117 @@
# app/tests/unit/market_data/acquisition/test_dzengi_trade_stream_acquisition_integration.py
from __future__ import annotations
from datetime import datetime, timezone
from decimal import Decimal
from src.market_data.acquisition.adapters.dzengi.websocket import (
DzengiUnifiedWebSocketAdapter,
)
from src.market_data.acquisition.models.trade import (
Trade,
TradeAggressorSide,
)
from src.market_data.acquisition.runtime.websocket_protocol import (
AcquisitionRuntimeCommand,
)
from src.market_data.acquisition.trade_stream_acquisition_service import (
TradeStreamAcquisitionService,
)
class FakeRuntimeService:
def __init__(self) -> None:
self.commands: list[AcquisitionRuntimeCommand] = []
async def dispatch(
self,
command: AcquisitionRuntimeCommand,
) -> None:
self.commands.append(command)
class RecordingConsistencyController:
def __init__(self) -> None:
self.accepted_trades: list[Trade] = []
def accept(
self,
trade: Trade,
) -> Trade | None:
self.accepted_trades.append(trade)
return trade
def _dzengi_trade_document() -> object:
return {
"status": "OK",
"destination": "internal.trade",
"payload": {
"id": 2134857062,
"price": "64497.25",
"size": "0.005",
"ts": 1784218066823,
"symbol": "BTC/USD_LEVERAGE",
"buyer": True,
"orderId": "order-123",
},
}
def test_dzengi_trade_document_passes_complete_acquisition_pipeline() -> None:
runtime = FakeRuntimeService()
consistency = RecordingConsistencyController()
service = TradeStreamAcquisitionService(
runtime_service=runtime,
adapter=DzengiUnifiedWebSocketAdapter(),
consistency_controller=consistency,
)
result = service.handle_message(
_dzengi_trade_document(),
)
assert isinstance(result, Trade)
assert consistency.accepted_trades == [result]
assert result.symbol == "BTC/USD_LEVERAGE"
assert result.trade_id == 2134857062
assert result.price == Decimal("64497.25")
assert result.quantity == Decimal("0.005")
assert result.executed_at == datetime.fromtimestamp(
1784218066823 / 1000,
tz=timezone.utc,
)
assert result.aggressor_side is TradeAggressorSide.BUY
assert result.source == "dzengi_websocket_trade"
def test_dzengi_duplicate_trade_is_rejected_by_real_consistency_pipeline() -> None:
from src.market_data.acquisition.consistency.trade_stream_consistency_controller import (
TradeStreamConsistencyController,
)
from src.market_data.acquisition.consistency.trade_stream_state_store import (
TradeStreamStateStore,
)
runtime = FakeRuntimeService()
consistency = TradeStreamConsistencyController(
state_store=TradeStreamStateStore(),
)
service = TradeStreamAcquisitionService(
runtime_service=runtime,
adapter=DzengiUnifiedWebSocketAdapter(),
consistency_controller=consistency,
)
document = _dzengi_trade_document()
first_result = service.handle_message(document)
duplicate_result = service.handle_message(document)
assert isinstance(first_result, Trade)
assert duplicate_result is None

View File

@@ -0,0 +1,383 @@
# app/tests/unit/market_data/acquisition/test_trade_stream_acquisition_service.py
from __future__ import annotations
import asyncio
import json
from datetime import datetime, timezone
from decimal import Decimal
from typing import cast
import pytest
from src.market_data.acquisition.models.candle_close import (
CandleCloseEvent,
)
from src.market_data.acquisition.models.quote import Quote
from src.market_data.acquisition.models.trade import (
Trade,
TradeAggressorSide,
)
from src.market_data.acquisition.runtime.runtime_commands import (
SubscribeCommand,
)
from src.market_data.acquisition.runtime.transport_messages import (
TransportTextMessage,
)
from src.market_data.acquisition.runtime.websocket_protocol import (
AcquisitionRuntimeCommand,
)
from src.market_data.acquisition.trade_stream_acquisition_protocol import (
TradeStreamAcquisitionServiceProtocol,
)
from src.market_data.acquisition.trade_stream_acquisition_service import (
TradeStreamAcquisitionService,
)
from src.market_data.acquisition.trade_stream_message_adapter_protocol import (
TradeStreamMappedMessage,
)
class FakeRuntimeService:
def __init__(self) -> None:
self.commands: list[AcquisitionRuntimeCommand] = []
async def dispatch(
self,
command: AcquisitionRuntimeCommand,
) -> None:
self.commands.append(command)
class FakeMessageAdapter:
def __init__(
self,
result: TradeStreamMappedMessage,
*,
error: Exception | None = None,
) -> None:
self._result = result
self._error = error
self.documents: list[object] = []
def map_message(
self,
document: object,
) -> TradeStreamMappedMessage:
self.documents.append(document)
if self._error is not None:
raise self._error
return self._result
class FakeConsistencyController:
def __init__(
self,
result: Trade | None = None,
*,
use_input_trade: bool = True,
error: Exception | None = None,
) -> None:
self._result = result
self._use_input_trade = use_input_trade
self._error = error
self.accepted_trades: list[Trade] = []
def accept(
self,
trade: Trade,
) -> Trade | None:
self.accepted_trades.append(trade)
if self._error is not None:
raise self._error
if self._use_input_trade:
return trade
return self._result
def _trade(
*,
trade_id: int = 2134857062,
) -> Trade:
return Trade(
symbol="BTC/USD_LEVERAGE",
trade_id=trade_id,
price=Decimal("64497.25"),
quantity=Decimal("0.005"),
executed_at=datetime(
2026,
7,
16,
11,
27,
46,
823000,
tzinfo=timezone.utc,
),
aggressor_side=TradeAggressorSide.BUY,
source="dzengi",
)
def _non_trade_quote() -> Quote:
return cast(
Quote,
object.__new__(Quote),
)
def _non_trade_candle() -> CandleCloseEvent:
return cast(
CandleCloseEvent,
object.__new__(CandleCloseEvent),
)
def create_service(
*,
adapter_result: TradeStreamMappedMessage | None = None,
adapter_error: Exception | None = None,
consistency_result: Trade | None = None,
consistency_uses_input_trade: bool = True,
consistency_error: Exception | None = None,
) -> tuple[
TradeStreamAcquisitionService,
FakeRuntimeService,
FakeMessageAdapter,
FakeConsistencyController,
]:
runtime = FakeRuntimeService()
adapter = FakeMessageAdapter(
adapter_result if adapter_result is not None else _trade(),
error=adapter_error,
)
consistency = FakeConsistencyController(
consistency_result,
use_input_trade=consistency_uses_input_trade,
error=consistency_error,
)
service = TradeStreamAcquisitionService(
runtime_service=runtime,
adapter=adapter,
consistency_controller=consistency,
)
return (
service,
runtime,
adapter,
consistency,
)
def test_service_implements_protocol() -> None:
service, _, _, _ = create_service()
assert isinstance(
service,
TradeStreamAcquisitionServiceProtocol,
)
def test_subscribe_dispatches_subscribe_command() -> None:
service, runtime, _, _ = create_service()
asyncio.run(
service.subscribe(
(
"BTCUSDT",
"ETHUSDT",
)
)
)
assert len(runtime.commands) == 1
command = runtime.commands[0]
assert isinstance(
command,
SubscribeCommand,
)
def test_subscribe_preserves_symbols() -> None:
service, runtime, _, _ = create_service()
asyncio.run(
service.subscribe(
(
"BTCUSDT",
"ETHUSDT",
)
)
)
command = runtime.commands[0]
assert isinstance(
command,
SubscribeCommand,
)
assert "BTCUSDT" in command.subscription_key
assert "ETHUSDT" in command.subscription_key
def test_subscribe_accepts_correlation_id() -> None:
service, runtime, _, _ = create_service()
asyncio.run(
service.subscribe(
(
"BTCUSDT",
),
correlation_id="corr-123",
)
)
command = runtime.commands[0]
assert isinstance(command, SubscribeCommand)
assert isinstance(command.message, TransportTextMessage)
document = json.loads(command.message.payload)
assert document["correlationId"] == "corr-123"
def test_handle_message_passes_trade_to_consistency() -> None:
trade = _trade()
document = {"destination": "internal.trade"}
service, _, adapter, consistency = create_service(
adapter_result=trade,
)
result = service.handle_message(document)
assert adapter.documents == [document]
assert consistency.accepted_trades == [trade]
assert result is trade
def test_handle_message_preserves_consistency_result_identity() -> None:
adapted_trade = _trade(trade_id=100)
consistency_result = _trade(trade_id=101)
service, _, _, consistency = create_service(
adapter_result=adapted_trade,
consistency_result=consistency_result,
consistency_uses_input_trade=False,
)
result = service.handle_message(
{"destination": "internal.trade"},
)
assert consistency.accepted_trades == [adapted_trade]
assert result is consistency_result
def test_handle_message_returns_none_for_duplicate() -> None:
trade = _trade()
service, _, _, consistency = create_service(
adapter_result=trade,
consistency_result=None,
consistency_uses_input_trade=False,
)
result = service.handle_message(
{"destination": "internal.trade"},
)
assert consistency.accepted_trades == [trade]
assert result is None
def test_handle_message_ignores_quote() -> None:
quote = _non_trade_quote()
document = {"destination": "quote"}
service, _, adapter, consistency = create_service(
adapter_result=quote,
)
result = service.handle_message(document)
assert adapter.documents == [document]
assert consistency.accepted_trades == []
assert result is None
def test_handle_message_ignores_candle_close_event() -> None:
candle = _non_trade_candle()
document = {"destination": "ohlc"}
service, _, adapter, consistency = create_service(
adapter_result=candle,
)
result = service.handle_message(document)
assert adapter.documents == [document]
assert consistency.accepted_trades == []
assert result is None
def test_handle_message_calls_adapter_once() -> None:
document = {"destination": "internal.trade"}
service, _, adapter, _ = create_service()
service.handle_message(document)
assert adapter.documents == [document]
def test_handle_message_propagates_adapter_error() -> None:
original = RuntimeError("adapter failed")
service, _, adapter, consistency = create_service(
adapter_error=original,
)
with pytest.raises(RuntimeError, match="adapter failed") as exc_info:
service.handle_message(
{"destination": "internal.trade"},
)
assert exc_info.value is original
assert len(adapter.documents) == 1
assert consistency.accepted_trades == []
def test_handle_message_propagates_consistency_error() -> None:
original = RuntimeError("consistency failed")
trade = _trade()
service, _, adapter, consistency = create_service(
adapter_result=trade,
consistency_error=original,
)
with pytest.raises(
RuntimeError,
match="consistency failed",
) as exc_info:
service.handle_message(
{"destination": "internal.trade"},
)
assert exc_info.value is original
assert len(adapter.documents) == 1
assert consistency.accepted_trades == [trade]