build 059.12: add runtime event models
This commit is contained in:
114
app/src/market_data/acquisition/runtime/runtime_events.py
Normal file
114
app/src/market_data/acquisition/runtime/runtime_events.py
Normal file
@@ -0,0 +1,114 @@
|
||||
# app/src/market_data/acquisition/runtime/runtime_events.py
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
|
||||
from src.market_data.acquisition.runtime.transport_messages import (
|
||||
TransportBinaryMessage,
|
||||
TransportTextMessage,
|
||||
)
|
||||
|
||||
|
||||
TransportMessage = TransportTextMessage | TransportBinaryMessage
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class ConnectedEvent:
|
||||
"""
|
||||
Событие успешного открытия транспортного соединения.
|
||||
|
||||
Событие фиксирует уже произошедший факт подключения и не содержит
|
||||
логики управления транспортом.
|
||||
"""
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class DisconnectedEvent:
|
||||
"""
|
||||
Событие завершения транспортного соединения.
|
||||
|
||||
Событие не определяет причину отключения и не запускает reconnect
|
||||
самостоятельно.
|
||||
"""
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class ConnectFailedEvent:
|
||||
"""
|
||||
Событие неуспешной попытки открытия транспортного соединения.
|
||||
|
||||
reason содержит безопасное текстовое описание причины ошибки.
|
||||
"""
|
||||
|
||||
reason: str
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class MessageReceivedEvent:
|
||||
"""
|
||||
Событие получения транспортного сообщения.
|
||||
|
||||
message содержит непрозрачный текстовый или бинарный payload.
|
||||
Runtime Event не интерпретирует содержимое сообщения.
|
||||
"""
|
||||
|
||||
message: TransportMessage
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class MessageSentEvent:
|
||||
"""
|
||||
Событие успешной отправки транспортного сообщения.
|
||||
|
||||
message содержит сообщение, переданное транспортному уровню.
|
||||
"""
|
||||
|
||||
message: TransportMessage
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class ReconnectStartedEvent:
|
||||
"""
|
||||
Событие начала попытки восстановления соединения.
|
||||
|
||||
attempt содержит номер текущей попытки reconnect.
|
||||
"""
|
||||
|
||||
attempt: int
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class ReconnectCompletedEvent:
|
||||
"""
|
||||
Событие успешного завершения восстановления соединения.
|
||||
|
||||
attempt содержит номер попытки, на которой соединение было
|
||||
восстановлено.
|
||||
"""
|
||||
|
||||
attempt: int
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class ReconnectFailedEvent:
|
||||
"""
|
||||
Событие неуспешной попытки восстановления соединения.
|
||||
|
||||
attempt содержит номер попытки.
|
||||
reason содержит безопасное текстовое описание причины ошибки.
|
||||
"""
|
||||
|
||||
attempt: int
|
||||
reason: str
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class HeartbeatTimeoutEvent:
|
||||
"""
|
||||
Событие превышения допустимого интервала heartbeat.
|
||||
|
||||
timeout_seconds содержит настроенный порог ожидания ответа.
|
||||
"""
|
||||
|
||||
timeout_seconds: float
|
||||
@@ -0,0 +1,116 @@
|
||||
# app/tests/unit/market_data/acquisition/runtime/test_runtime_events.py
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import FrozenInstanceError
|
||||
|
||||
import pytest
|
||||
|
||||
from src.market_data.acquisition.runtime.runtime_events import (
|
||||
ConnectedEvent,
|
||||
ConnectFailedEvent,
|
||||
DisconnectedEvent,
|
||||
HeartbeatTimeoutEvent,
|
||||
MessageReceivedEvent,
|
||||
MessageSentEvent,
|
||||
ReconnectCompletedEvent,
|
||||
ReconnectFailedEvent,
|
||||
ReconnectStartedEvent,
|
||||
)
|
||||
from src.market_data.acquisition.runtime.transport_messages import (
|
||||
TransportBinaryMessage,
|
||||
TransportTextMessage,
|
||||
)
|
||||
|
||||
|
||||
def test_connected_event_is_marker_event() -> None:
|
||||
assert ConnectedEvent() == ConnectedEvent()
|
||||
|
||||
|
||||
def test_disconnected_event_is_marker_event() -> None:
|
||||
assert DisconnectedEvent() == DisconnectedEvent()
|
||||
|
||||
|
||||
def test_connect_failed_event_contains_reason() -> None:
|
||||
event = ConnectFailedEvent(reason="connection refused")
|
||||
|
||||
assert event.reason == "connection refused"
|
||||
|
||||
|
||||
def test_message_received_event_supports_text_message() -> None:
|
||||
message = TransportTextMessage(payload="market-data")
|
||||
event = MessageReceivedEvent(message=message)
|
||||
|
||||
assert event.message is message
|
||||
|
||||
|
||||
def test_message_received_event_supports_binary_message() -> None:
|
||||
message = TransportBinaryMessage(payload=b"\x01\x02")
|
||||
event = MessageReceivedEvent(message=message)
|
||||
|
||||
assert event.message is message
|
||||
|
||||
|
||||
def test_message_sent_event_contains_transport_message() -> None:
|
||||
message = TransportTextMessage(payload="ping")
|
||||
event = MessageSentEvent(message=message)
|
||||
|
||||
assert event.message is message
|
||||
|
||||
|
||||
def test_reconnect_started_event_contains_attempt() -> None:
|
||||
event = ReconnectStartedEvent(attempt=1)
|
||||
|
||||
assert event.attempt == 1
|
||||
|
||||
|
||||
def test_reconnect_completed_event_contains_attempt() -> None:
|
||||
event = ReconnectCompletedEvent(attempt=2)
|
||||
|
||||
assert event.attempt == 2
|
||||
|
||||
|
||||
def test_reconnect_failed_event_contains_attempt_and_reason() -> None:
|
||||
event = ReconnectFailedEvent(
|
||||
attempt=3,
|
||||
reason="timeout",
|
||||
)
|
||||
|
||||
assert event.attempt == 3
|
||||
assert event.reason == "timeout"
|
||||
|
||||
|
||||
def test_heartbeat_timeout_event_contains_timeout_seconds() -> None:
|
||||
event = HeartbeatTimeoutEvent(timeout_seconds=30.0)
|
||||
|
||||
assert event.timeout_seconds == 30.0
|
||||
|
||||
|
||||
def test_connect_failed_event_is_immutable() -> None:
|
||||
event = ConnectFailedEvent(reason="timeout")
|
||||
|
||||
with pytest.raises(FrozenInstanceError):
|
||||
setattr(event, "reason", "connection refused")
|
||||
|
||||
|
||||
def test_message_received_event_is_immutable() -> None:
|
||||
event = MessageReceivedEvent(
|
||||
message=TransportTextMessage(payload="payload"),
|
||||
)
|
||||
|
||||
with pytest.raises(FrozenInstanceError):
|
||||
setattr(
|
||||
event,
|
||||
"message",
|
||||
TransportTextMessage(payload="other"),
|
||||
)
|
||||
|
||||
|
||||
def test_reconnect_failed_event_is_immutable() -> None:
|
||||
event = ReconnectFailedEvent(
|
||||
attempt=1,
|
||||
reason="timeout",
|
||||
)
|
||||
|
||||
with pytest.raises(FrozenInstanceError):
|
||||
setattr(event, "attempt", 2)
|
||||
588
docs/migrations/build_059_12.md
Normal file
588
docs/migrations/build_059_12.md
Normal file
@@ -0,0 +1,588 @@
|
||||
# Build 059.12 — Runtime Event Models
|
||||
|
||||
**Migration Build**
|
||||
|
||||
---
|
||||
|
||||
# Цель Build
|
||||
|
||||
После завершения Build 059.11 Runtime уже содержит три независимых архитектурных уровня:
|
||||
|
||||
- Runtime Protocols;
|
||||
- Runtime Transport Messages;
|
||||
- Runtime Commands.
|
||||
|
||||
Следующим шагом является создание слоя **Runtime Events**.
|
||||
|
||||
Данный Build вводит immutable-модели событий Runtime.
|
||||
|
||||
События описывают **уже произошедшие факты** и полностью отделены от команд Runtime.
|
||||
|
||||
---
|
||||
|
||||
# Причина появления Runtime Events
|
||||
|
||||
До настоящего момента Runtime умел описывать:
|
||||
|
||||
- транспортные интерфейсы;
|
||||
- транспортные сообщения;
|
||||
- намерения системы.
|
||||
|
||||
Однако отсутствовал механизм описания результатов выполнения этих намерений.
|
||||
|
||||
Например.
|
||||
|
||||
После выполнения
|
||||
|
||||
```
|
||||
ConnectCommand
|
||||
```
|
||||
|
||||
может произойти:
|
||||
|
||||
```
|
||||
ConnectedEvent
|
||||
```
|
||||
|
||||
или
|
||||
|
||||
```
|
||||
ConnectFailedEvent
|
||||
```
|
||||
|
||||
Это два разных факта.
|
||||
|
||||
Команда выражает желание выполнить действие.
|
||||
|
||||
Событие фиксирует результат выполнения этого действия.
|
||||
|
||||
---
|
||||
|
||||
# Архитектурная идея
|
||||
|
||||
Build 059.12 завершает построение базовой event-driven модели Runtime.
|
||||
|
||||
Получается следующая архитектура.
|
||||
|
||||
```text
|
||||
Runtime
|
||||
|
||||
Commands
|
||||
│
|
||||
▼
|
||||
Session
|
||||
│
|
||||
▼
|
||||
Transport
|
||||
│
|
||||
▼
|
||||
Events
|
||||
```
|
||||
|
||||
Commands направляются вниз.
|
||||
|
||||
Events распространяются вверх.
|
||||
|
||||
Именно такое разделение используется в большинстве современных event-driven систем.
|
||||
|
||||
---
|
||||
|
||||
# Новый модуль
|
||||
|
||||
Создан файл
|
||||
|
||||
```
|
||||
src/market_data/acquisition/runtime/runtime_events.py
|
||||
```
|
||||
|
||||
В модуле определены immutable-модели Runtime Events.
|
||||
|
||||
---
|
||||
|
||||
# Использование Transport Messages
|
||||
|
||||
Транспортные события используют модели Build 059.10.
|
||||
|
||||
```
|
||||
TransportTextMessage
|
||||
TransportBinaryMessage
|
||||
```
|
||||
|
||||
Runtime Event не знает содержимого сообщения.
|
||||
|
||||
Он лишь фиксирует факт передачи или получения транспортного payload.
|
||||
|
||||
---
|
||||
|
||||
# TransportMessage
|
||||
|
||||
Для удобства типизации используется локальный alias.
|
||||
|
||||
```python
|
||||
TransportMessage =
|
||||
TransportTextMessage
|
||||
| TransportBinaryMessage
|
||||
```
|
||||
|
||||
Он используется только внутри Runtime Events.
|
||||
|
||||
Никакой новой модели данных не создаётся.
|
||||
|
||||
---
|
||||
|
||||
# Новые события
|
||||
|
||||
Build вводит девять Runtime Events.
|
||||
|
||||
---
|
||||
|
||||
# ConnectedEvent
|
||||
|
||||
Фиксирует успешное открытие транспортного соединения.
|
||||
|
||||
Событие является маркерным.
|
||||
|
||||
Не содержит дополнительных данных.
|
||||
|
||||
Причина.
|
||||
|
||||
На текущем этапе Runtime отсутствует модель идентификатора соединения.
|
||||
|
||||
Добавлять подобные поля преждевременно.
|
||||
|
||||
---
|
||||
|
||||
# DisconnectedEvent
|
||||
|
||||
Фиксирует завершение транспортного соединения.
|
||||
|
||||
Событие также является маркерным.
|
||||
|
||||
Причина отключения намеренно отсутствует.
|
||||
|
||||
Определение причин разрыва соединения относится к будущему уровню Reliability.
|
||||
|
||||
---
|
||||
|
||||
# ConnectFailedEvent
|
||||
|
||||
Фиксирует неудачную попытку подключения.
|
||||
|
||||
Содержит:
|
||||
|
||||
```
|
||||
reason
|
||||
```
|
||||
|
||||
Причина хранится в виде строки.
|
||||
|
||||
---
|
||||
|
||||
## Почему используется строка
|
||||
|
||||
Во время проектирования рассматривались варианты хранения:
|
||||
|
||||
- Exception;
|
||||
- traceback;
|
||||
- transport-specific error;
|
||||
- websocket exception.
|
||||
|
||||
От данных вариантов было принято решение отказаться.
|
||||
|
||||
Runtime Event не должен зависеть от конкретной реализации транспорта.
|
||||
|
||||
Строковое описание полностью соответствует принципу transport-agnostic Runtime.
|
||||
|
||||
---
|
||||
|
||||
# MessageReceivedEvent
|
||||
|
||||
Фиксирует получение транспортного сообщения.
|
||||
|
||||
Содержит:
|
||||
|
||||
```
|
||||
TransportMessage
|
||||
```
|
||||
|
||||
Runtime Event не анализирует payload.
|
||||
|
||||
Его обработкой занимаются последующие уровни Acquisition.
|
||||
|
||||
---
|
||||
|
||||
# MessageSentEvent
|
||||
|
||||
Фиксирует успешную передачу транспортного сообщения.
|
||||
|
||||
Содержит:
|
||||
|
||||
```
|
||||
TransportMessage
|
||||
```
|
||||
|
||||
Это позволяет журналировать транспортный обмен, не анализируя содержимое сообщений.
|
||||
|
||||
---
|
||||
|
||||
# ReconnectStartedEvent
|
||||
|
||||
Фиксирует начало новой попытки восстановления соединения.
|
||||
|
||||
Содержит:
|
||||
|
||||
```
|
||||
attempt
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Почему хранится номер попытки
|
||||
|
||||
Номер попытки является устойчивым Runtime-фактом.
|
||||
|
||||
Он потребуется:
|
||||
|
||||
- Supervisor;
|
||||
- журналированию;
|
||||
- диагностике;
|
||||
- мониторингу Runtime.
|
||||
|
||||
---
|
||||
|
||||
# ReconnectCompletedEvent
|
||||
|
||||
Фиксирует успешное восстановление соединения.
|
||||
|
||||
Также содержит:
|
||||
|
||||
```
|
||||
attempt
|
||||
```
|
||||
|
||||
что позволяет определить, на какой попытке произошло восстановление.
|
||||
|
||||
---
|
||||
|
||||
# ReconnectFailedEvent
|
||||
|
||||
Фиксирует завершение очередной попытки восстановления соединения с ошибкой.
|
||||
|
||||
Содержит:
|
||||
|
||||
```
|
||||
attempt
|
||||
reason
|
||||
```
|
||||
|
||||
Таким образом Runtime может описывать процесс восстановления без зависимости от конкретной реализации транспорта.
|
||||
|
||||
---
|
||||
|
||||
# HeartbeatTimeoutEvent
|
||||
|
||||
Фиксирует превышение допустимого интервала heartbeat.
|
||||
|
||||
Содержит:
|
||||
|
||||
```
|
||||
timeout_seconds
|
||||
```
|
||||
|
||||
Хранится только настроенный порог ожидания.
|
||||
|
||||
---
|
||||
|
||||
## Почему отсутствуют timestamp
|
||||
|
||||
Во время проектирования обсуждалось хранение:
|
||||
|
||||
```
|
||||
timestamp
|
||||
occurred_at
|
||||
created_at
|
||||
```
|
||||
|
||||
Было принято решение отказаться.
|
||||
|
||||
Причины.
|
||||
|
||||
---
|
||||
|
||||
### Причина №1
|
||||
|
||||
В Runtime ещё отсутствует единая модель времени.
|
||||
|
||||
---
|
||||
|
||||
### Причина №2
|
||||
|
||||
В разных окружениях время может поступать из различных источников.
|
||||
|
||||
Например:
|
||||
|
||||
- системные часы;
|
||||
- монотонные часы;
|
||||
- серверное время биржи.
|
||||
|
||||
До появления общей модели времени вводить timestamp преждевременно.
|
||||
|
||||
---
|
||||
|
||||
# Почему отсутствует RuntimeEvent
|
||||
|
||||
Также обсуждалось создание общего базового класса.
|
||||
|
||||
Например.
|
||||
|
||||
```
|
||||
RuntimeEvent
|
||||
```
|
||||
|
||||
или
|
||||
|
||||
```
|
||||
BaseEvent
|
||||
```
|
||||
|
||||
От идеи отказались.
|
||||
|
||||
Причины.
|
||||
|
||||
---
|
||||
|
||||
## Причина №1
|
||||
|
||||
Все события являются immutable dataclass.
|
||||
|
||||
Общего поведения между ними нет.
|
||||
|
||||
---
|
||||
|
||||
## Причина №2
|
||||
|
||||
Преждевременная иерархия только усложняет архитектуру.
|
||||
|
||||
---
|
||||
|
||||
## Причина №3
|
||||
|
||||
Общий базовый тип можно добавить позднее без нарушения обратной совместимости.
|
||||
|
||||
---
|
||||
|
||||
# Почему отсутствует RuntimeEventType
|
||||
|
||||
Рассматривалось использование enum.
|
||||
|
||||
Например.
|
||||
|
||||
```
|
||||
CONNECTED
|
||||
DISCONNECTED
|
||||
MESSAGE_RECEIVED
|
||||
```
|
||||
|
||||
Было принято решение отказаться.
|
||||
|
||||
Тип события полностью определяется его классом.
|
||||
|
||||
Дополнительный enum создавал бы дублирование информации.
|
||||
|
||||
---
|
||||
|
||||
# Почему отсутствуют Subscription Events
|
||||
|
||||
Изначально предполагалось добавить:
|
||||
|
||||
- SubscriptionRegisteredEvent;
|
||||
- SubscriptionRemovedEvent;
|
||||
- SubscriptionRestoredEvent.
|
||||
|
||||
После анализа архитектуры было принято решение перенести их в следующий Build.
|
||||
|
||||
Причина.
|
||||
|
||||
В Build 059.12 Subscription Manager ещё отсутствует.
|
||||
|
||||
Следовательно, пока отсутствует компонент, способный генерировать подобные события.
|
||||
|
||||
События должны появляться одновременно с соответствующим уровнем архитектуры.
|
||||
|
||||
Поэтому Subscription Events перенесены в Build 059.13.
|
||||
|
||||
---
|
||||
|
||||
# Что НЕ входит в Build
|
||||
|
||||
Сознательно не реализованы:
|
||||
|
||||
- Event Bus;
|
||||
- Dispatcher;
|
||||
- Observer;
|
||||
- callbacks;
|
||||
- asyncio;
|
||||
- Publisher;
|
||||
- Subscriber;
|
||||
- Session;
|
||||
- Subscription Manager;
|
||||
- обработчики событий;
|
||||
- очередь событий.
|
||||
|
||||
Build содержит исключительно immutable-модели Runtime Events.
|
||||
|
||||
---
|
||||
|
||||
# Обновлённая архитектура Runtime
|
||||
|
||||
После завершения Build Runtime приобретает следующий вид.
|
||||
|
||||
```text
|
||||
Runtime
|
||||
|
||||
Protocols
|
||||
│
|
||||
Transport Messages
|
||||
│
|
||||
Runtime Commands
|
||||
│
|
||||
Runtime Events
|
||||
│
|
||||
Transport
|
||||
│
|
||||
Network
|
||||
```
|
||||
|
||||
Каждый уровень отвечает только за собственную область ответственности.
|
||||
|
||||
---
|
||||
|
||||
# Следующие Build
|
||||
|
||||
## Build 059.13
|
||||
|
||||
Subscription Manager.
|
||||
|
||||
Появятся:
|
||||
|
||||
- управление подписками;
|
||||
- хранение Runtime-состояния;
|
||||
- восстановление подписок;
|
||||
- SubscriptionRegisteredEvent;
|
||||
- SubscriptionRemovedEvent;
|
||||
- SubscriptionRestoredEvent.
|
||||
|
||||
---
|
||||
|
||||
## Build 059.14
|
||||
|
||||
WebSocket Session.
|
||||
|
||||
Появится координация:
|
||||
|
||||
- Commands;
|
||||
- Events;
|
||||
- Transport;
|
||||
- Subscription Manager.
|
||||
|
||||
---
|
||||
|
||||
## Build 059.15
|
||||
|
||||
Runtime Reliability.
|
||||
|
||||
Будут реализованы:
|
||||
|
||||
- Heartbeat;
|
||||
- Reconnect;
|
||||
- Supervisor;
|
||||
- Scheduler;
|
||||
- Ping/Pong.
|
||||
|
||||
---
|
||||
|
||||
# Проверка Build
|
||||
|
||||
Выполнена полная проверка.
|
||||
|
||||
---
|
||||
|
||||
## Компиляция
|
||||
|
||||
```
|
||||
python -m compileall
|
||||
```
|
||||
|
||||
Успешно.
|
||||
|
||||
---
|
||||
|
||||
## Unit Tests
|
||||
|
||||
Созданы тесты:
|
||||
|
||||
```
|
||||
test_runtime_events.py
|
||||
```
|
||||
|
||||
Проверяется:
|
||||
|
||||
- создание каждого события;
|
||||
- корректность хранения данных;
|
||||
- поддержка TransportMessage;
|
||||
- immutable-поведение dataclass.
|
||||
|
||||
Все тесты успешно пройдены.
|
||||
|
||||
---
|
||||
|
||||
## Runtime Regression
|
||||
|
||||
Совместно проверены:
|
||||
|
||||
- Runtime Protocols;
|
||||
- Runtime Transport Messages;
|
||||
- Runtime Commands;
|
||||
- Runtime Events.
|
||||
|
||||
Все Runtime-тесты успешно завершены.
|
||||
|
||||
---
|
||||
|
||||
## Проверка репозитория
|
||||
|
||||
Выполнен
|
||||
|
||||
```
|
||||
git diff --check
|
||||
```
|
||||
|
||||
Ошибок форматирования не обнаружено.
|
||||
|
||||
---
|
||||
|
||||
# Итог
|
||||
|
||||
Build 059.12 завершает формирование слоя Runtime Events.
|
||||
|
||||
Теперь Runtime имеет четыре полностью независимых уровня.
|
||||
|
||||
```text
|
||||
Protocols
|
||||
│
|
||||
Transport Messages
|
||||
│
|
||||
Runtime Commands
|
||||
│
|
||||
Runtime Events
|
||||
```
|
||||
|
||||
Каждый уровень описывает собственную область ответственности.
|
||||
|
||||
Commands выражают намерения системы.
|
||||
|
||||
Events фиксируют уже произошедшие факты.
|
||||
|
||||
Подобное разделение делает архитектуру Runtime предсказуемой, расширяемой и соответствует общепринятым принципам построения event-driven систем.
|
||||
|
||||
Следующий Build посвящён созданию Subscription Manager, который станет первым Runtime-компонентом, использующим одновременно Commands и Events, сохраняя при этом независимость транспортного уровня от бизнес-логики Acquisition.
|
||||
Reference in New Issue
Block a user