build 059.9: add websocket runtime protocols
This commit is contained in:
@@ -0,0 +1,79 @@
|
|||||||
|
# app/src/market_data/acquisition/runtime/websocket_protocol.py
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Protocol, runtime_checkable
|
||||||
|
|
||||||
|
|
||||||
|
@runtime_checkable
|
||||||
|
class WebSocketTransportProtocol(Protocol):
|
||||||
|
"""
|
||||||
|
Низкоуровневый транспорт WebSocket.
|
||||||
|
|
||||||
|
Контракт отвечает только за открытие и закрытие соединения,
|
||||||
|
отправку исходящих сообщений и получение входящих сообщений.
|
||||||
|
|
||||||
|
Реализация не должна выполнять:
|
||||||
|
- parsing;
|
||||||
|
- schema validation;
|
||||||
|
- value validation;
|
||||||
|
- mapping;
|
||||||
|
- routing бизнес-событий.
|
||||||
|
"""
|
||||||
|
|
||||||
|
async def connect(self) -> None:
|
||||||
|
"""Открыть WebSocket-соединение."""
|
||||||
|
...
|
||||||
|
|
||||||
|
async def disconnect(self) -> None:
|
||||||
|
"""Закрыть WebSocket-соединение."""
|
||||||
|
...
|
||||||
|
|
||||||
|
async def send(self, message: str | bytes) -> None:
|
||||||
|
"""Отправить сериализованное транспортное сообщение."""
|
||||||
|
...
|
||||||
|
|
||||||
|
async def receive(self) -> str | bytes:
|
||||||
|
"""Получить одно сырое транспортное сообщение."""
|
||||||
|
...
|
||||||
|
|
||||||
|
|
||||||
|
@runtime_checkable
|
||||||
|
class WebSocketSessionProtocol(Protocol):
|
||||||
|
"""
|
||||||
|
Контракт жизненного цикла WebSocket-сессии.
|
||||||
|
|
||||||
|
Session координирует транспорт и последующие runtime-компоненты,
|
||||||
|
но не содержит exchange-specific parsing или mapping.
|
||||||
|
"""
|
||||||
|
|
||||||
|
@property
|
||||||
|
def is_connected(self) -> bool:
|
||||||
|
"""Показывает, открыта ли активная WebSocket-сессия."""
|
||||||
|
...
|
||||||
|
|
||||||
|
async def start(self) -> None:
|
||||||
|
"""Запустить WebSocket-сессию."""
|
||||||
|
...
|
||||||
|
|
||||||
|
async def stop(self) -> None:
|
||||||
|
"""Остановить WebSocket-сессию."""
|
||||||
|
...
|
||||||
|
|
||||||
|
|
||||||
|
@runtime_checkable
|
||||||
|
class WebSocketSubscriptionManagerProtocol(Protocol):
|
||||||
|
"""
|
||||||
|
Контракт управления активными WebSocket-подписками.
|
||||||
|
|
||||||
|
Конкретные модели подписок и формат сообщений будут добавлены
|
||||||
|
в последующих Build'ах.
|
||||||
|
"""
|
||||||
|
|
||||||
|
async def restore_subscriptions(self) -> None:
|
||||||
|
"""Восстановить активные подписки после подключения."""
|
||||||
|
...
|
||||||
|
|
||||||
|
async def clear_subscriptions(self) -> None:
|
||||||
|
"""Очистить runtime-состояние активных подписок."""
|
||||||
|
...
|
||||||
@@ -0,0 +1,66 @@
|
|||||||
|
# app/tests/unit/market_data/acquisition/runtime/test_websocket_protocol.py
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from src.market_data.acquisition.runtime.websocket_protocol import (
|
||||||
|
WebSocketSessionProtocol,
|
||||||
|
WebSocketSubscriptionManagerProtocol,
|
||||||
|
WebSocketTransportProtocol,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class FakeWebSocketTransport:
|
||||||
|
async def connect(self) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
async def disconnect(self) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
async def send(self, message: str | bytes) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
async def receive(self) -> str | bytes:
|
||||||
|
return "{}"
|
||||||
|
|
||||||
|
|
||||||
|
class FakeWebSocketSession:
|
||||||
|
@property
|
||||||
|
def is_connected(self) -> bool:
|
||||||
|
return True
|
||||||
|
|
||||||
|
async def start(self) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
async def stop(self) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
class FakeWebSocketSubscriptionManager:
|
||||||
|
async def restore_subscriptions(self) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
async def clear_subscriptions(self) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def test_transport_implementation_satisfies_protocol() -> None:
|
||||||
|
assert isinstance(FakeWebSocketTransport(), WebSocketTransportProtocol)
|
||||||
|
|
||||||
|
|
||||||
|
def test_session_implementation_satisfies_protocol() -> None:
|
||||||
|
assert isinstance(FakeWebSocketSession(), WebSocketSessionProtocol)
|
||||||
|
|
||||||
|
|
||||||
|
def test_subscription_manager_implementation_satisfies_protocol() -> None:
|
||||||
|
assert isinstance(
|
||||||
|
FakeWebSocketSubscriptionManager(),
|
||||||
|
WebSocketSubscriptionManagerProtocol,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_incomplete_transport_does_not_satisfy_protocol() -> None:
|
||||||
|
class IncompleteTransport:
|
||||||
|
async def connect(self) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
assert not isinstance(IncompleteTransport(), WebSocketTransportProtocol)
|
||||||
297
docs/migrations/build_059_9.md
Normal file
297
docs/migrations/build_059_9.md
Normal file
@@ -0,0 +1,297 @@
|
|||||||
|
# Build 059.9 — WebSocket Runtime Protocol
|
||||||
|
|
||||||
|
**Migration Build**
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Цель Build
|
||||||
|
|
||||||
|
Build 059.9 открывает новый этап миграции подсистемы получения рыночных данных.
|
||||||
|
|
||||||
|
До настоящего момента была полностью построена независимая цепочка обработки входящих сообщений:
|
||||||
|
|
||||||
|
```text
|
||||||
|
Raw Message
|
||||||
|
│
|
||||||
|
Schema Validation
|
||||||
|
│
|
||||||
|
Parser
|
||||||
|
│
|
||||||
|
Value Validation
|
||||||
|
│
|
||||||
|
Mapper
|
||||||
|
│
|
||||||
|
Canonical Model
|
||||||
|
│
|
||||||
|
Adapter
|
||||||
|
```
|
||||||
|
|
||||||
|
Данная архитектура уже позволяет корректно обрабатывать сообщения независимо от конкретного транспорта доставки.
|
||||||
|
|
||||||
|
Следующим этапом становится построение собственного Runtime уровня, который будет отвечать исключительно за транспортировку сообщений.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Причина появления Runtime
|
||||||
|
|
||||||
|
Исторически получение данных осуществлялось через класс ExchangeWebSocketClient.
|
||||||
|
|
||||||
|
Этот класс одновременно выполнял множество различных обязанностей:
|
||||||
|
|
||||||
|
- открытие соединения;
|
||||||
|
- закрытие соединения;
|
||||||
|
- reconnect;
|
||||||
|
- heartbeat;
|
||||||
|
- отправку запросов;
|
||||||
|
- получение сообщений;
|
||||||
|
- обработку ошибок;
|
||||||
|
- маршрутизацию сообщений;
|
||||||
|
- работу с конкретной биржей.
|
||||||
|
|
||||||
|
Подобная архитектура нарушает принцип единственной ответственности (Single Responsibility Principle) и значительно усложняет расширение системы.
|
||||||
|
|
||||||
|
В новой архитектуре эти обязанности разделяются между несколькими независимыми компонентами Runtime.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Основная идея новой архитектуры
|
||||||
|
|
||||||
|
Runtime не должен знать ничего о бизнес-логике.
|
||||||
|
|
||||||
|
Он отвечает исключительно за транспортный уровень.
|
||||||
|
|
||||||
|
После завершения всей серии Build'ов структура будет выглядеть следующим образом:
|
||||||
|
|
||||||
|
```text
|
||||||
|
Exchange Runtime
|
||||||
|
│
|
||||||
|
Connection / Session
|
||||||
|
│
|
||||||
|
Heartbeat / Reconnect
|
||||||
|
│
|
||||||
|
Subscription Manager
|
||||||
|
│
|
||||||
|
WebSocket Transport
|
||||||
|
│
|
||||||
|
▼
|
||||||
|
Unified WebSocket Adapter
|
||||||
|
│
|
||||||
|
┌──────────────────┼──────────────────┐
|
||||||
|
▼ ▼ ▼
|
||||||
|
Quote Candle Trades
|
||||||
|
│
|
||||||
|
▼
|
||||||
|
Feed → Handler → Registry
|
||||||
|
```
|
||||||
|
|
||||||
|
Таким образом Runtime становится полностью независимым от конкретного вида сообщений.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Почему Runtime создаётся отдельно
|
||||||
|
|
||||||
|
В подсистеме Acquisition уже существует файл
|
||||||
|
|
||||||
|
```
|
||||||
|
acquisition/protocol.py
|
||||||
|
```
|
||||||
|
|
||||||
|
Он содержит контракты уровня получения данных:
|
||||||
|
|
||||||
|
- InstrumentDocumentSource
|
||||||
|
- QuoteDocumentSource
|
||||||
|
- CandlesDocumentSource
|
||||||
|
- Feed Protocol
|
||||||
|
- Handler Protocol
|
||||||
|
|
||||||
|
Эти контракты ничего не знают о WebSocket.
|
||||||
|
|
||||||
|
Они работают уже после доставки сообщения.
|
||||||
|
|
||||||
|
Поэтому смешивать Acquisition Protocol и Runtime Protocol архитектурно неправильно.
|
||||||
|
|
||||||
|
В результате Runtime получает собственный уровень контрактов.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Новые контракты Runtime
|
||||||
|
|
||||||
|
В Build 059.9 создаётся новый файл:
|
||||||
|
|
||||||
|
```text
|
||||||
|
src/market_data/acquisition/runtime/websocket_protocol.py
|
||||||
|
```
|
||||||
|
|
||||||
|
В нём определяются исключительно транспортные контракты:
|
||||||
|
|
||||||
|
- WebSocketTransportProtocol
|
||||||
|
- WebSocketSessionProtocol
|
||||||
|
- WebSocketSubscriptionManagerProtocol
|
||||||
|
|
||||||
|
Все контракты являются Protocol без реализации.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Что входит в ответственность Runtime
|
||||||
|
|
||||||
|
Runtime отвечает только за транспортный уровень.
|
||||||
|
|
||||||
|
Например:
|
||||||
|
|
||||||
|
- открытие соединения;
|
||||||
|
- закрытие соединения;
|
||||||
|
- получение транспортных сообщений;
|
||||||
|
- отправку транспортных сообщений;
|
||||||
|
- жизненный цикл WebSocket Session;
|
||||||
|
- управление подписками.
|
||||||
|
|
||||||
|
Runtime **не выполняет**:
|
||||||
|
|
||||||
|
- schema validation;
|
||||||
|
- parsing;
|
||||||
|
- value validation;
|
||||||
|
- mapping;
|
||||||
|
- business routing;
|
||||||
|
- преобразование моделей;
|
||||||
|
- обработку рыночной логики.
|
||||||
|
|
||||||
|
Все эти обязанности уже реализованы в Acquisition.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Что намеренно НЕ реализуется в Build 059.9
|
||||||
|
|
||||||
|
В рамках данного Build отсутствуют:
|
||||||
|
|
||||||
|
- WebSocket клиент;
|
||||||
|
- asyncio логика;
|
||||||
|
- reconnect;
|
||||||
|
- heartbeat;
|
||||||
|
- ping/pong;
|
||||||
|
- сериализация сообщений;
|
||||||
|
- JSON;
|
||||||
|
- модели транспортных сообщений;
|
||||||
|
- реальные подписки.
|
||||||
|
|
||||||
|
Build создаёт исключительно архитектурный фундамент.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# План дальнейшего развития Runtime
|
||||||
|
|
||||||
|
После Build 059.9 развитие Runtime планируется небольшими независимыми этапами.
|
||||||
|
|
||||||
|
## Build 059.10
|
||||||
|
|
||||||
|
Transport Message Models
|
||||||
|
|
||||||
|
Добавление immutable моделей транспортных сообщений.
|
||||||
|
|
||||||
|
Например:
|
||||||
|
|
||||||
|
- Connect
|
||||||
|
- Disconnect
|
||||||
|
- Subscribe
|
||||||
|
- Ping
|
||||||
|
- Pong
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Build 059.11
|
||||||
|
|
||||||
|
Subscription Manager
|
||||||
|
|
||||||
|
Появится полноценное управление активными подписками.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Build 059.12
|
||||||
|
|
||||||
|
WebSocket Session
|
||||||
|
|
||||||
|
Жизненный цикл соединения.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Build 059.13
|
||||||
|
|
||||||
|
Reconnect Manager
|
||||||
|
|
||||||
|
Автоматическое восстановление соединения.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Build 059.14
|
||||||
|
|
||||||
|
Heartbeat
|
||||||
|
|
||||||
|
Поддержание активности соединения.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Архитектурный принцип миграции
|
||||||
|
|
||||||
|
Все существующие Parser, Validator, Mapper и Adapter остаются неизменными.
|
||||||
|
|
||||||
|
Меняется исключительно транспорт доставки сообщений.
|
||||||
|
|
||||||
|
В результате один и тот же Adapter сможет работать:
|
||||||
|
|
||||||
|
- со старым ExchangeWebSocketClient;
|
||||||
|
- с новым Runtime;
|
||||||
|
- с тестовыми сообщениями;
|
||||||
|
- с любым будущим источником данных.
|
||||||
|
|
||||||
|
Это позволяет выполнять миграцию постепенно без остановки работы существующего торгового бота.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Новое архитектурное соглашение проекта
|
||||||
|
|
||||||
|
Во время выполнения Build принято дополнительное архитектурное решение.
|
||||||
|
|
||||||
|
Для крупных подсистем проекта рекомендуется использовать описательные имена файлов вместо слишком общих.
|
||||||
|
|
||||||
|
Например:
|
||||||
|
|
||||||
|
Хорошо:
|
||||||
|
|
||||||
|
```text
|
||||||
|
websocket_protocol.py
|
||||||
|
market_cache.py
|
||||||
|
exchange_service.py
|
||||||
|
```
|
||||||
|
|
||||||
|
Предпочтительно избегать новых файлов с названиями:
|
||||||
|
|
||||||
|
```text
|
||||||
|
protocol.py
|
||||||
|
service.py
|
||||||
|
utils.py
|
||||||
|
helpers.py
|
||||||
|
```
|
||||||
|
|
||||||
|
если их назначение можно описать более конкретно.
|
||||||
|
|
||||||
|
Это правило принято с расчётом на дальнейший рост проекта и будет постепенно применяться при плановой архитектурной ревизии существующих модулей.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Проверка Build
|
||||||
|
|
||||||
|
Выполнены проверки:
|
||||||
|
|
||||||
|
- compileall
|
||||||
|
- unit tests
|
||||||
|
- runtime regression
|
||||||
|
- git diff --check
|
||||||
|
|
||||||
|
Все проверки успешно пройдены.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
# Итог
|
||||||
|
|
||||||
|
Build 059.9 не добавляет новую функциональность получения рыночных данных.
|
||||||
|
|
||||||
|
Его задача — создать фундамент Runtime уровня, который позволит в последующих Build'ах полностью заменить существующий транспорт WebSocket без изменения уже реализованной цепочки обработки сообщений.
|
||||||
Reference in New Issue
Block a user