49 KiB
Build 060.25 — Production Runtime Integration Architecture
Статус: Accepted
Build: 060.25
Подсистема: Market Data Acquisition / Trade Stream Runtime
Дата начала: 2026-07-30
Дата завершения: 2026-07-31
Версия документа: 1.0
1. Назначение
Документ фиксирует архитектуру Production Runtime Integration для Trade Stream.
Build 060.24 создал внутренний граф Runtime Recovery, но намеренно:
- не создавал реальное WebSocket-соединение;
- не запускал receive loop;
- не управлял asyncio-задачами;
- не связывал reconnect с Trade Recovery;
- не подключал Runtime к bootstrap приложения.
Build 060.25 превратил этот граф в управляемый production lifecycle без нарушения границ Transport, Runtime, Acquisition, Consistency и Recovery.
Документ является живой архитектурной спецификацией.
Он разделяет:
- принятые и реализованные решения 060.25.0–060.25.8;
- end-to-end, reconnect и stress validation Build 060.26;
- persistent storage и startup recovery Build 060.27–060.28.
Итоговый migration report зафиксирован в build_060_25.md.
2. Статус подэтапов
| Подэтап | Название | Статус |
|---|---|---|
| 060.25.0 | Static Contract Cleanup | Accepted |
| 060.25.1 | Dzengi WebSocket Transport | Accepted |
| 060.25.2 | WebSocket Session and Subscription Manager | Accepted |
| 060.25.3 | Async Runtime Event Publisher | Accepted |
| 060.25.4 | Trade Stream Production Runtime | Accepted |
| 060.25.5 | Reconnect and Recovery Integration | Accepted |
| 060.25.6 | Heartbeat and Scheduler Integration | Accepted |
| 060.25.7 | Settings, Bootstrap and Graceful Shutdown | Accepted |
| 060.25.8 | Targeted Runtime Verification | Accepted |
Статус Accepted означает:
- реализация прошла целевые тесты;
- архитектурные findings устранены;
- пользователь принял результат подэтапа;
- файлы ещё не зафиксированы отдельным Build commit.
3. Исходная архитектура
Build 060.24 предоставил:
TradeStreamRuntimeComposition
├── TradeStreamStateStore
├── TradeStreamConsistencyController
├── TradeRecoveryController
├── TradeRecoveryWindowPlanner
├── RuntimeRecoveryCoordinator
├── AcquisitionRuntimeService
├── TradeStreamAcquisitionService
├── ReconnectCoordinator
├── HeartbeatMonitor
├── RuntimeSupervisor
└── RuntimeScheduler
Composition создаёт единые экземпляры stateful-компонентов, но не запускает их.
Внешними зависимостями оставались:
WebSocketTransportProtocol
WebSocketSessionProtocol
WebSocketSubscriptionManagerProtocol
AcquisitionRuntimeEventPublisherProtocol
TradeStreamMessageAdapterProtocol
DzengiTradesDocumentSource
До Build 060.25 production-реализации первых четырёх WebSocket Runtime контрактов отсутствовали.
4. Целевая модель
APPLICATION BOOTSTRAP
│
▼
Trade Stream Production Runtime
│
┌───────────────────────┼───────────────────────┐
│ │ │
▼ ▼ ▼
WebSocket lifecycle Runtime supervision Recovery gate
│ │ │
▼ ▼ ▼
Dzengi Transport Heartbeat / Scheduler Runtime Recovery
│ │ │
▼ ▼ ▼
WebSocket Session Runtime Supervisor REST Recovery
│ │ │
▼ ▼ │
Subscription Manager ───► Reconnect Coordinator │
│ │
▼ │
receive loop │
│ │
▼ ▼
TradeStreamAcquisitionService ─► TradeStreamConsistencyController
│
▼
TradeStreamStateStore
Production Runtime является владельцем выполнения, но не владельцем бизнес-состояния сделок.
5. Архитектурные границы
5.1. Dzengi WebSocket Transport
Transport отвечает только за:
- нормализацию WebSocket URL;
- открытие и закрытие соединения;
- отправку
str | bytes; - получение
str | bytes; - transport timeouts и ping/pong параметры библиотеки;
- преобразование инфраструктурных ошибок в WebSocket transport errors.
Transport не:
- декодирует JSON;
- различает ACK и market events;
- строит subscription payload;
- изменяет Trade Stream state;
- выполняет reconnect policy;
- выполняет Recovery.
5.2. WebSocket Session
Session отвечает за:
- идемпотентные
start()иstop(); - сериализацию lifecycle через
asyncio.Lock; - согласование локального состояния с Transport;
- повторный
start()после фактического закрытия Transport.
Session не должна превращать повторный start() в неявный force
reconnect. Принудительная замена соединения является явной операцией
ReconnectCoordinator.
5.3. Subscription Manager
Subscription Manager отвечает за:
- desired subscription registry;
- active state текущего соединения;
- идемпотентную подписку;
- последовательное восстановление подписок;
- очистку runtime-состояния;
- явную политику unsupported unsubscribe.
Manager работает с:
TransportTextMessage
TransportBinaryMessage
Он не анализирует содержимое payload и не знает о Trade, Quote, Candle или конкретной JSON destination.
5.4. Production Runtime
Production Runtime должен отвечать за:
- запуск и остановку полного графа;
- владение asyncio-задачами;
- единственный receive loop;
- JSON decoding;
- отделение control/ACK сообщений от market events;
- передачу market documents в Acquisition Service;
- координацию reconnect и recovery;
- graceful cancellation и shutdown.
Production Runtime не должен дублировать:
- Trade consistency;
- subscription registry;
- heartbeat calculations;
- recovery window planning;
- WebSocket transport logic.
6. Build 060.25.0 — Static Contract Cleanup
6.1. Исходная проблема
Тестовый helper create_coordinator() принимал:
RecordingWindowPlanner | None
Негативный тест передавал:
BrokenPlanner(TradeRecoveryWindowPlanner)
Runtime-поведение было корректным, но Pylance обнаруживал несовместимый тип аргумента.
6.2. Принятое решение
Тестовая иерархия согласована с реальной:
TradeRecoveryWindowPlanner
│
▼
RecordingWindowPlanner
│
▼
BrokenPlanner
Удалён локальный type: ignore[arg-type].
Production-контракт RuntimeRecoveryCoordinator не изменён.
7. Build 060.25.1 — Dzengi WebSocket Transport
7.1. Реализация
Добавлен:
DzengiWebSocketTransport
Transport структурно соответствует:
WebSocketTransportProtocol
Прямое наследование от Protocol не используется, чтобы сохранить
рабочий __slots__.
7.2. URL policy
Поддерживаемые исходные схемы:
http → ws
https → wss
ws → ws
wss → wss
URL нормализуется до endpoint:
/connect
Примеры:
https://api-adapter.dzengi.com
↓
wss://api-adapter.dzengi.com/connect
http://localhost:8080/
↓
ws://localhost:8080/connect
Пустые, относительные, неподдерживаемые и fragment URL отклоняются до сетевого вызова.
7.3. Connection options
Через constructor внедряются:
headers
open_timeout
ping_interval
ping_timeout
close_timeout
connector
Инъекция connector позволяет тестировать Transport без реальной сети.
Transport использует WebSocket subprotocol:
json
7.4. Lifecycle
DISCONNECTED
│ connect()
▼
OPEN
│ disconnect()
▼
DISCONNECTED
Повторные connect() и disconnect() идемпотентны.
Lifecycle сериализован отдельным asyncio.Lock.
Transport считает соединение активным только если:
connection is not None
and connection.state is State.OPEN
7.5. Error policy
Введены:
WebSocketTransportError
WebSocketTransportNotConnectedError
Connect, disconnect, send и receive errors оборачиваются с сохранением
исходной ошибки в __cause__.
Попытка send или receive без открытого соединения завершается
WebSocketTransportNotConnectedError.
asyncio.CancelledError не поглощается обработчиками Exception.
8. Build 060.25.2 — WebSocket Session
8.1. Реализация
Добавлен:
WebSocketSession
Session структурно соответствует:
WebSocketSessionProtocol
8.2. Идемпотентность
start()
├── Session и Transport connected → no-op
└── disconnected → transport.connect()
stop()
├── Session и Transport disconnected → no-op
└── active → transport.disconnect()
Параллельные lifecycle-вызовы сериализуются Session lock.
После ошибки start() Session остаётся disconnected.
После ошибки stop() локальный Session state гарантированно
сбрасывается в finally.
8.3. Remote disconnect
Session не доверяет только локальному _started.
Состояние:
is_connected =
_started
and transport.is_connected
Если Transport обнаружил закрытие соединения, следующий start()
действительно открывает новое соединение.
9. Build 060.25.2 — Subscription State
9.1. Два уровня состояния
Subscription Manager разделяет:
desired subscriptions
и:
active subscription keys
Desired registry описывает требуемое состояние Runtime и переживает ошибку текущей сетевой попытки.
Active keys относятся только к текущему WebSocket connection.
9.2. Subscribe
validate key and message
│
▼
store desired subscription
│
▼
transport.send()
│
├── success → mark active
└── failure → keep pending desired subscription
Повторный вызов для active key является no-op.
Повторный вызов для pending key выполняет новую отправку и может обновить сохранённое transport message.
9.3. Restore
После нового соединения:
clear active keys
│
▼
for each desired subscription in registration order:
transport.send()
│
▼
mark key active
Ошибка восстановления не удаляет desired registry.
9.4. Unsubscribe
Публичный Dzengi endpoint не поддерживает надёжный
trades.unsubscribe.
Поэтому production-default:
supports_unsubscribe = False
Попытка unsubscribe завершается:
WebSocketUnsubscribeNotSupportedError
и не изменяет registry.
Generic supported branch сохраняется для будущих провайдеров.
9.5. Clear
clear_subscriptions() очищает desired и active state без отправки
сетевого сообщения.
Это локальная lifecycle-операция, а не имитация exchange unsubscribe.
10. Корректирующее решение Reconnect
10.1. Обнаруженный риск
Первоначальная реализация Build 060.24 выполняла:
ConnectCommand
│
▼
restore_subscriptions()
Новая идемпотентная Session показала скрытую проблему.
При heartbeat timeout WebSocket мог формально оставаться OPEN.
В таком состоянии Session.start() являлся no-op, и подписки повторно
отправлялись в старое соединение.
10.2. Принятое решение
Одна reconnect-попытка теперь выполняет:
ReconnectStartedEvent
│
▼
DisconnectCommand
│
▼
ConnectCommand
│
▼
restore_subscriptions()
│
▼
ReconnectCompletedEvent
Это гарантирует замену старого connection даже при формальном
State.OPEN.
Registry подписок не очищается при disconnect, поэтому desired state доступен для восстановления.
10.3. Failure policy
Ошибка disconnect, connect или restore:
state = FAILED
│
▼
ReconnectFailedEvent
│
▼
raise original exception
Retry loop и backoff по-прежнему находятся вне ReconnectCoordinator.
11. Concurrency invariants
11.1. Локальные locks
DzengiWebSocketTransport._lifecycle_lock
защищает connect/disconnect.
WebSocketSession._lifecycle_lock
защищает Session start/stop.
WebSocketSubscriptionManager._lock
защищает subscribe, unsubscribe, restore и clear.
11.2. Lock ordering
Допустимый порядок:
Session lock
│
▼
Transport lifecycle lock
Subscription Manager lock
│
▼
Transport send
Transport не вызывает Session или Subscription Manager обратно. Обратного lock ordering нет.
11.3. Cross-component coordination
Локальные locks не заменяют общий lifecycle gate.
Production Runtime обязан исключить гонку:
live message processing
X
reconnect / recovery state mutation
Этот gate относится к 060.25.4–060.25.5.
12. Build 060.25.3 — Async Runtime Event Publisher
Статус: Accepted
Реализованы:
AcquisitionRuntimeEventConsumerProtocol
AcquisitionRuntimeEventPublisher
AcquisitionRuntimeEventLoggingConsumer
AcquisitionRuntimeEventPublisher является production-реализацией:
AcquisitionRuntimeEventPublisherProtocol
12.1. Delivery model
Publisher:
- принимает immutable-набор Consumer через constructor injection;
- последовательно await-ит Consumer в порядке регистрации;
- сериализует конкурентные
publish()однимasyncio.Lock; - возвращается только после завершения fan-out;
- не создаёт
asyncio.Task; - не хранит event queue;
- не имеет собственного
start()илиstop().
Рекурсивная публикация через тот же Publisher отклоняется до входа
в lock. Для определения прямой и child-task reentrancy используется
ContextVar.
12.2. Consumer failure policy
Обычная ошибка одного Consumer:
- не прерывает доставку остальным Consumer;
- не прерывает Heartbeat или Reconnect lifecycle;
- фиксируется только безопасными метаданными:
event type,consumer type,error type; - не добавляет exception message, traceback или
repr(event)в лог.
Ошибка logging handler также не выходит из Publisher.
CancelledError и другие BaseException не перехватываются.
Publisher освобождает lock и reentrancy context в finally, поэтому
может использоваться после отменённого вызова.
12.3. Logging Consumer
AcquisitionRuntimeEventLoggingConsumer использует стандартный
Python logging.
Lifecycle и failure events получают соответствующие уровни
INFO, WARNING и ERROR.
Для MessageReceivedEvent и MessageSentEvent на уровне DEBUG
фиксируются только:
- тип transport message;
- размер payload.
Содержимое payload не журналируется.
12.4. Legacy Runtime Events
Существующий src/runtime_events/publisher.py не изменён.
Новый Publisher не зависит от:
NotificationService;- Telegram;
JournalService;- legacy
RuntimeEvent; - market models.
Преобразование Acquisition Runtime Event в legacy Runtime Event в настоящий Build не добавлено.
13. Build 060.25.4 — Trade Stream Production Runtime
Статус: Accepted
Реализован один высокоуровневый владелец lifecycle:
TradeStreamProductionRuntime
13.1. Задачи
Production Runtime владеет:
- startup sequence;
- receive task;
- lifecycle gate;
- cancellation;
- graceful shutdown;
- terminal error propagation.
Startup, Scheduler и receive выполняются как явно сохранённые owned
tasks.
Внешний bootstrap остаётся владельцем корневой coroutine run().
Production Runtime синхронно закрепляет Scheduler за собой до первого await сетевого startup. Поэтому другой lifecycle owner не может запустить тот же Scheduler во время подключения или подписки.
13.2. Receive loop
Существует ровно один consumer:
transport.receive()
│
▼
decode JSON
│
├── ACK / control → runtime control handling
├── market document → TradeStreamAcquisitionService
└── invalid document → explicit error policy
ACK и control messages нельзя передавать в Unified Market Adapter как неизвестный market event.
Production Runtime создаёт correlation ID Trade subscription и передаёт его через Acquisition Service. Provider-specific handler:
- принимает только совпадающий correlation ID;
- подтверждает успешный ACK;
- идемпотентно принимает повторный успешный ACK;
- делает отрицательный ACK terminal error;
- отклоняет неизвестный correlation ID и malformed control message.
13.3. Task ownership
Production Runtime создаёт, хранит, отменяет и await-ит собственные startup, Scheduler и receive tasks.
stop() может отменить незавершившиеся:
- Session startup;
- публикацию
ConnectedEvent; - отправку Trade subscription.
Ни Transport, ни Session, ни Heartbeat не создают скрытые background tasks на уровне приложения.
14. Build 060.25.5 — Reconnect and Recovery
Статус: Accepted
14.1. Основная последовательность
connection failure or confirmed liveness timeout
│
▼
acquire single reconnect/recovery gate
│
▼
pause live message processing
│
▼
disconnect old connection
│
▼
connect new connection
│
▼
restore desired subscriptions
│
▼
capture recovery_end_time
│
▼
run RuntimeRecoveryCoordinator
│
▼
resume buffered live messages
14.2. Recovery execution
RuntimeRecoveryCoordinator.recover() является синхронным и использует
блокирующий REST.
В asyncio Runtime он должен выполняться:
await asyncio.to_thread(...)
либо через отдельный async adapter с эквивалентной семантикой.
14.3. Shared state
Live и Recovery используют один:
TradeStreamStateStore
TradeStreamConsistencyController
Поэтому live processing должен быть приостановлен на время Recovery.
Одновременное изменение state из event loop и recovery thread запрещено.
14.4. Boundary policy
После восстановления подписки новые WebSocket frames могут накапливаться в transport buffer.
Recovery обрабатывается первым, затем buffered live frames.
Повтор на общей временной границе устраняется существующим Consistency Layer.
14.5. Single-flight и поколения connection
Добавлен:
RuntimeReconnectRecoveryCoordinator
Он объединяет базовый ReconnectCoordinator и
RuntimeRecoveryCoordinator в одну single-flight операцию.
Одновременные запросы reconnect присоединяются к уже выполняющейся операции. Ошибка старого connection generation не запускает повторный reconnect после того, как её поколение уже было обработано.
14.6. Cancellation и ошибки Recovery
Recovery запускается через asyncio.to_thread() при закрытом live gate.
Если владеющая asyncio-задача отменена, Runtime продолжает ждать фактического завершения Recovery worker. Повторные запросы отмены не могут преждевременно открыть gate и оставить worker изменяющим Consistency state в фоне.
После ошибки reconnect или Recovery gate переходит в failed state. Накопленные live messages не передаются в Consistency Layer. Failed state очищается только при начале нового полного Production Runtime lifecycle.
14.7. Composition identity
Composition создаёт один экземпляр live gate и передаёт один
RuntimeReconnectRecoveryCoordinator в RuntimeSupervisor.
Production Runtime проверяет:
- identity общего live gate;
- равенство нормализованного набора symbols для Live и Recovery.
Несогласованный dependency graph отклоняется до запуска WebSocket lifecycle.
15. Build 060.25.6 — Heartbeat and Scheduler
Статус: Accepted
15.1. Liveness
Отсутствие сделок не означает потерю соединения.
Для тихого symbol допустима длительная пауза между market events.
Поэтому reconnect нельзя запускать только по условию:
no Trade messages
Liveness подтверждается двумя transport-level источниками:
- успешным WebSocket Ping/Pong;
- успешно полученным WebSocket message независимо от его market/control назначения.
15.2. Реализованная граница
Реализован отдельный контракт:
RuntimeLivenessProbeProtocol
│
▼
WebSocket ping/pong implementation
│
├── success → supervisor.notify_activity()
└── failure → reconnect path
DzengiWebSocketTransport.probe():
- возвращает
Trueтолько после соответствующего Pong; - возвращает
False, если connection отсутствует, закрыт или Pong не получен вовремя; - не перехватывает cancellation;
- распространяет неожиданную ошибку открытого connection как
WebSocketTransportError.
Ожидание Pong всегда имеет конечную границу probe_timeout.
Если отдельное значение не передано, оно берётся из конечного
ping_timeout. При отключённом ping_timeout применяется безопасное
значение 20 секунд. Неположительное и бесконечное значение запрещено.
15.3. Scheduler
Scheduler продолжает владеть только временем вызовов.
Каждая итерация:
- выполняет transport probe;
- при успешном Pong уведомляет Supervisor об активности;
- проверяет Heartbeat timeout;
- передаёт подтверждённый timeout Supervisor.
Он не должен:
- отправлять WebSocket payload самостоятельно;
- изменять subscription registry;
- выполнять Trade Recovery;
- владеть receive loop.
15.4. Reconnect generation
Supervisor связывает каждый Heartbeat monitoring period с generation текущего connection.
Heartbeat timeout передаётся в
reconnect_after_transport_failure(observed_generation=...).
Поэтому timeout и receive error одного поколения:
- присоединяются к одной активной reconnect/recovery operation;
- используют результат уже завершённой operation;
- не запускают повторный reconnect после восстановления того же connection generation.
15.5. Task ownership и shutdown
Production Runtime выполняет Scheduler.claim(owner) до первого
сетевого await и освобождает claim только после завершения owned task.
Это исключает запуск Scheduler вторым владельцем во время startup.
Порядок остановки:
Scheduler stop / cancel / await
↓
Supervisor stop
↓
receive task cancel / await
↓
Session stop
Scheduler-led reconnect/recovery достигает безопасной границы до остановки Supervisor. Receive-led operation безопасно завершается при последующем cancel/await receive task и до остановки Session.
16. Build 060.25.7 — Settings and Bootstrap
Статус: Accepted
16.1. Settings
Добавлена отдельная неизменяемая группа:
TradeStreamSettings
├── enabled
├── websocket_url
├── symbols
├── open_timeout_seconds
├── probe_timeout_seconds
├── close_timeout_seconds
├── heartbeat_timeout_seconds
├── scheduler_interval_seconds
└── recovery_window_ms
Production Trade Stream должен иметь отдельный feature flag с безопасным значением по умолчанию.
Общий EXCHANGE_ENABLED не должен неявно запускать новый Runtime без
явной настройки.
Используются переменные:
TRADE_STREAM_ENABLED=false
TRADE_STREAM_WS_URL
TRADE_STREAM_SYMBOLS
TRADE_STREAM_OPEN_TIMEOUT_SECONDS=10
TRADE_STREAM_PROBE_TIMEOUT_SECONDS=20
TRADE_STREAM_CLOSE_TIMEOUT_SECONDS=10
TRADE_STREAM_HEARTBEAT_TIMEOUT_SECONDS=30
TRADE_STREAM_SCHEDULER_INTERVAL_SECONDS=5
TRADE_STREAM_RECOVERY_WINDOW_MS=3599999
При выключенном flag зависимые Trade Stream значения не проверяются и
не создают Runtime. При включённом flag обязательны отдельные
TRADE_STREAM_WS_URL, TRADE_STREAM_SYMBOLS и EXCHANGE_BASE_URL для
REST Recovery.
Fallback с EXCHANGE_WS_URL и DEFAULT_SYMBOL запрещён.
EXCHANGE_BASE_URL и EXCHANGE_TIMEOUT_SEC используются только
явным REST-клиентом Recovery.
Все интервалы должны быть положительными и конечными, recovery window — положительным целым числом. Символы очищаются от пробелов, дедуплицируются и сортируются.
16.2. Bootstrap
Добавлен отдельный production composition root:
build_trade_stream_production_runtime(settings)
│
├── feature flag disabled → None
│
└── enabled
├── DzengiWebSocketTransport
├── WebSocketSession
├── WebSocketSubscriptionManager
├── AcquisitionRuntimeEventPublisher
├── DzengiTradesDocumentSource
├── TradeStreamRuntimeComposition
└── TradeStreamProductionRuntime
Factory использует один снимок Settings. Создание графа:
- не открывает WebSocket;
- не выполняет REST;
- не создаёт asyncio-задачи;
- сохраняет identity общих stateful-зависимостей.
Встроенный keepalive библиотеки websockets отключён:
ping_interval=None
ping_timeout=None
Единственным владельцем активного Ping/Pong является Runtime Scheduler
через конечный probe_timeout.
Корневой ApplicationComposition содержит Telegram Bot, Dispatcher и
опциональный Trade Stream Runtime. run_application() создаёт и
контролирует обе root tasks.
Ошибка или неожиданное нормальное завершение включённого Trade Stream фатальны для всего процесса. Ошибка Telegram polling симметрично останавливает Trade Stream. Исходная ошибка сохраняет приоритет над ошибками cleanup.
main.py не должен содержать детали WebSocket protocol или recovery
алгоритма.
16.3. Shutdown
Внешний порядок приложения:
stop accepting Telegram work
│
▼
TradeStreamProductionRuntime.stop()
│
▼
await both root tasks
│
▼
close bot session exactly once
Внутренний порядок Trade Stream Runtime:
Scheduler stop / cancel / await
│
▼
Supervisor stop
│
▼
receive task cancel / await
│
▼
Session stop
│
▼
clear subscriptions
Application Runner передаёт Aiogram параметр
close_bot_session=False, поэтому сессию закрывает ровно один корневой
владелец.
Cancellation не прерывает начатый cleanup: cleanup выполняется в отдельной owned task и ожидается через shield. Все созданные root tasks отменяются либо завершаются и обязательно await-ятся.
17. Build 060.25.8 — Verification
Статус: Accepted
060.25.8 не добавляет production-поведение. Этап подтверждает детерминированными in-process тестами, что принятые компоненты работают как один production graph.
Добавлен управляемый test harness:
Production Trade Stream Factory
│
├── real DzengiWebSocketTransport
├── real WebSocketSession
├── real WebSocketSubscriptionManager
├── real Runtime Composition
└── real TradeStreamProductionRuntime
│
├── controlled fake WebSocket connector
├── controlled fake REST client
└── fake Telegram boundary
Harness не использует реальную сеть, реальные credentials или недетерминированные внешние задержки.
Подтверждены следующие vertical-slice сценарии:
- выключенный feature flag не создаёт Runtime graph;
- production factory выполняет subscribe → ACK → Trade → checkpoint;
- ошибка WebSocket connect фатальна и не оставляет задач;
- ошибка первой subscription send откатывает partial startup;
- reconnect открывает новое connection и восстанавливает subscription до запуска REST Recovery;
- recovered Trade обрабатывается раньше buffered Live Trade;
- ошибка Recovery оставляет gate в failed state и отклоняет buffered Trade;
- одновременные ошибки Telegram и Trade Stream обрабатываются детерминированно;
- повторная cancellation не прерывает начатый cleanup;
- Bot session закрывается один раз, owned tasks всегда await-ятся.
Production-код в 060.25.8 не изменён.
В Build 060.26 остаются:
- live exchange integration;
- длительные reconnect scenarios;
- recovery scenarios с реальными задержками;
- stress testing;
- network fault injection;
- финальная Runtime documentation verification.
Unit tests по умолчанию не должны обращаться к внешней сети.
18. Реализованные файлы
18.1. Production code
Добавлены:
app/src/market_data/acquisition/adapters/dzengi/
websocket_control_message_handler.py
websocket_inbound_message_classifier.py
websocket_transport.py
app/src/market_data/acquisition/runtime/
acquisition_runtime_event_logging_consumer.py
acquisition_runtime_event_publisher.py
live_processing_gate.py
runtime_liveness_probe.py
runtime_reconnect_recovery_coordinator.py
trade_stream_production_runtime.py
websocket_inbound_message.py
websocket_session.py
websocket_subscription_manager.py
app/src/bootstrap/
application.py
trade_stream_runtime.py
Изменены:
app/.env.example
app/src/bootstrap/app_factory.py
app/src/core/config.py
app/src/integrations/exchange/rest_client.py
app/src/main.py
app/src/market_data/acquisition/exceptions.py
app/src/market_data/acquisition/runtime/reconnect.py
app/src/market_data/acquisition/runtime/scheduler.py
app/src/market_data/acquisition/runtime/supervisor.py
app/src/market_data/acquisition/trade_stream_runtime_composition.py
18.2. Tests
Добавлены:
app/tests/unit/market_data/acquisition/adapters/dzengi/
test_websocket_control_message_handler.py
test_websocket_inbound_message_classifier.py
test_websocket_transport.py
app/tests/unit/market_data/acquisition/runtime/
test_acquisition_runtime_event_logging_consumer.py
test_acquisition_runtime_event_publisher.py
test_live_processing_gate.py
test_runtime_reconnect_recovery_coordinator.py
test_trade_stream_production_runtime.py
test_websocket_session.py
test_websocket_subscription_manager.py
test_websocket_runtime_reconnect_integration.py
app/tests/unit/bootstrap/
test_app_factory.py
test_application.py
test_trade_stream_production_integration.py
test_trade_stream_runtime.py
app/tests/unit/core/
test_config.py
app/tests/unit/test_main.py
Изменены:
app/tests/unit/market_data/acquisition/runtime/
test_reconnect_coordinator.py
test_runtime_recovery_coordinator.py
test_runtime_scheduler.py
test_runtime_supervisor.py
app/tests/unit/market_data/acquisition/
test_trade_stream_runtime_composition.py
Состав файлов соответствует итоговому состоянию Build 060.25.
19. Test evidence
После реализации 060.25.8 выполнены:
Settings / Bootstrap tests: 30 passed
New vertical / lifecycle scenarios: 8 passed
Expanded Production Runtime target: 445 passed
Full project regression: 1869 passed
Расширенный целевой набор включает:
- Trade Stream settings;
- Application Composition, factory и
main.py; - production factory vertical-slice integration;
- весь
market_data/acquisition/runtime; - Dzengi WebSocket Transport, classifier и control handler;
- Trade Stream Runtime Composition.
Отдельно проверены:
OPEN connection
│ reconnect()
▼
old connection closed
│
▼
new connection opened
│
▼
subscriptions restored
и:
initial subscribe send failure
│
▼
desired subscription remains pending
│ reconnect()
▼
subscription sent through new connection
Для Async Runtime Event Publisher отдельно подтверждены:
- последовательная доставка;
- сериализация конкурентных публикаций;
- отсутствие fire-and-forget;
- consumer error isolation;
- отсутствие payload в error log;
- устойчивость к ошибке logging handler;
- распространение cancellation;
- повторное использование после cancellation;
- защита от прямой и child-task reentrancy;
- сохранение Heartbeat и Reconnect lifecycle при ошибке Consumer.
После read-only review 060.25.8 целевой набор и полная регрессия повторены 2026-07-31. Новых findings не обнаружено.
20. Принятые архитектурные решения
ADR-060.25-001 — Transport остаётся raw boundary
Статус: Accepted
DzengiWebSocketTransport работает только с str | bytes и не содержит
JSON или market routing.
ADR-060.25-002 — Session start является идемпотентным
Статус: Accepted
Повторный start не заменяет исправное открытое соединение.
ADR-060.25-003 — Reconnect всегда выполняет disconnect
Статус: Accepted
Force reconnect выражен последовательностью Disconnect → Connect, а не скрытым поведением Session.start().
ADR-060.25-004 — Subscription registry хранит desired state
Статус: Accepted
Сетевой send failure не удаляет намерение подписаться. Desired и active state разделены.
ADR-060.25-005 — Unsupported unsubscribe является явной ошибкой
Статус: Accepted
Dzengi unsubscribe не имитируется локальным успешным результатом.
ADR-060.25-006 — Production Runtime владеет задачами
Статус: Accepted
Внешний bootstrap владеет корневой task run(). Production Runtime
до сетевого startup закрепляет Scheduler за собой, а затем создаёт,
хранит, отменяет и await-ит собственные startup, Scheduler и receive
tasks. Claim освобождается после завершения Scheduler task.
ADR-060.25-007 — Recovery выполняется при закрытом live gate
Статус: Accepted
REST Recovery и live message processing не изменяют общий Consistency state параллельно. Отмена ожидает фактического завершения Recovery worker, а ошибка Recovery блокирует buffered live processing до нового Runtime lifecycle.
ADR-060.25-008 — Heartbeat основан на transport liveness
Статус: Accepted
Отсутствие Trade events не является достаточным признаком разрыва. Liveness подтверждается конечным Ping/Pong probe либо успешно полученным WebSocket message. Timeout и receive error используют одно connection generation и одну single-flight reconnect/recovery operation.
ADR-060.25-009 — Runtime Events доставляются inline
Статус: Accepted
Publisher последовательно await-ит фиксированный набор Consumer и не создаёт скрытую очередь или background task.
ADR-060.25-010 — Consumer errors не управляют Runtime lifecycle
Статус: Accepted
Обычная ошибка Consumer изолируется, безопасно диагностируется и не прерывает доставку остальным Consumer. Cancellation распространяется.
ADR-060.25-011 — Trade Stream включается только отдельным flag
Статус: Accepted
TRADE_STREAM_ENABLED по умолчанию выключен. EXCHANGE_ENABLED не
запускает новый Runtime, а legacy WebSocket URL и default symbol не
используются как fallback.
ADR-060.25-012 — Ошибка включённого Runtime фатальна
Статус: Accepted
Telegram polling и Trade Stream принадлежат одному Application Runner. Terminal error любой включённой root task останавливает вторую root task, после чего обе задачи ожидаются и сессия Bot закрывается ровно один раз.
ADR-060.25-013 — Scheduler является единственным keepalive owner
Статус: Accepted
Автоматический keepalive websockets отключён. Ping/Pong выполняется
только Scheduler через явный transport liveness probe с конечным
timeout.
21. Инварианты Build
К завершению 060.25 должны выполняться:
exactly one production receive loop
exactly one scheduler task
at most one reconnect/recovery sequence
reconnect =
disconnect
→ connect
→ restore subscriptions
recovery and live state mutation are not concurrent
desired subscriptions survive transient send failure
quiet market does not cause false reconnect
startup and shutdown are deterministic and awaitable
Composition construction has no network side effects
disabled Trade Stream creates no Runtime graph
application owns and awaits both root tasks
bot session is closed exactly once
runtime event delivery is ordered and awaitable
consumer failure does not interrupt runtime lifecycle
22. Не входит в Build
Build 060.25 не реализует:
- persistent market data storage;
- persistent checkpoint;
- startup recovery после перезапуска процесса;
- historical query API;
- replay API;
- analytics API;
- общий retry framework приложения;
- production stress certification.
Эти задачи относятся к Build 060.26–060.29.
23. Критерии завершения
Build 060.25 получил статус Completed, поскольку:
- все этапы 060.25.0–060.25.8 приняты;
- concrete Runtime dependencies собраны через composition root;
- Trade Stream запускается и останавливается из bootstrap;
- receive loop передаёт Trade messages в общий Consistency Layer;
- ACK/control messages не разрушают market pipeline;
- reconnect действительно заменяет connection;
- desired subscriptions восстанавливаются;
- Recovery выполняется до возобновления live processing;
- heartbeat не зависит только от частоты сделок;
- cancellation не оставляет фоновые задачи;
- целевые тесты и полная регрессия проходят;
git diff --checkне содержит новых ошибок;- создан итоговый
build_060_25.md; - настоящий документ обновлён до финального фактического состояния.
24. Текущий итог
На момент версии 1.0:
060.25.0 Accepted
060.25.1 Accepted
060.25.2 Accepted
060.25.3 Accepted
060.25.4 Accepted
060.25.5 Accepted
060.25.6 Accepted
060.25.7 Accepted
060.25.8 Accepted
Реализован production-ready фундамент WebSocket lifecycle:
- concrete Dzengi Transport;
- идемпотентная Session;
- desired/active Subscription Manager;
- принудительная замена connection во время reconnect;
- интеграционные unit-сценарии двух критических recovery cases;
- последовательный Async Runtime Event Publisher;
- безопасный Logging Consumer;
- изоляция Consumer errors от Heartbeat и Reconnect lifecycle;
- единый Trade Stream Production Runtime;
- отменяемые и awaitable startup/receive tasks;
- явная маршрутизация и проверка Trade subscription ACK;
- single-flight reconnect → subscription restore → Recovery;
- общий live/recovery gate с блокировкой buffered сообщений после Recovery failure;
- безопасное ожидание Recovery worker при повторной cancellation;
- проверка единого набора Live и Recovery symbols;
- конечный WebSocket Ping/Pong liveness probe;
- generation-aware объединение heartbeat и receive reconnect;
- единственная owned Scheduler task;
- детерминированная остановка Scheduler → Supervisor → receive loop;
- отдельный безопасный Trade Stream feature flag;
- production composition root без сетевых side effects;
- один Application Runner для Telegram и Trade Stream;
- фатальность terminal error включённого Trade Stream;
- явные WebSocket URL и symbols без legacy fallback;
- отключённый встроенный WebSocket keepalive;
- детерминированное ожидание root tasks и однократное закрытие Bot session.
Runtime подключён к bootstrap приложения и прошёл финальную детерминированную in-process verification. Отдельный read-only review 060.25.8 не выявил findings; целевой набор из 445 тестов и полная регрессия из 1869 тестов прошли. Build 060.25 завершён.