Files
dzentra_bot/docs/architecture/trades_feed.md

27 KiB
Raw Permalink Blame History

Trades Feed (Time & Sales) — текущая архитектура

Статус: Current; accepted in Build 060.30.1

Область: Production vertical slice Canonical Trades

Версия документа: 1.0


Связанные документы


1. Назначение и фактический scope

Trades Feed получает сделки Dzengi, преобразует их в Canonical Trade, проверяет порядок и дубликаты, восстанавливает разрывы через REST и при включённом Storage атомарно сохраняет сделку вместе с persistent checkpoint.

После перезапуска Runtime восстанавливает Consistency state из PostgreSQL, выполняет Startup Recovery и только затем пропускает накопленные live-сообщения. Сохранённые Canonical Market Data доступны через отдельные Historical Access и deterministic Replay API.

Краткая итоговая цепочка:

Dzengi WebSocket → Transport → validation/mapping → Canonical Trade ─┐
                                                                     ├→ Consistency
Dzengi REST Recovery → validation/mapping → Canonical Trade ─────────┘
                                                                     ↓
                                   optional PostgreSQL Trade + checkpoint
                                             ├→ next persistent startup
                                             │  → Hydration
                                             │  → Consistency state
                                             │  → Startup Recovery
                                             ├→ explicit caller
                                             │  → Historical Access
                                             └→ explicit caller
                                                → Replay Plan
                                                → Replay Session → consumer

Production ingestion и persistence в этой цепочке реализованы только для Trades. Наличие Quote/Candle моделей и repositories не означает, что Quotes и Candles уже подключены как Production feeds.


2. Матрица готовности Market Data

Возможность Trades Quotes Candle revisions
Canonical model и Dzengi adapters Да Да Да
Production WebSocket Runtime consumer Да Нет Нет
Live/Recovery Consistency Да Нет Нет
Production persistence wiring Да Нет Нет
PostgreSQL repository и schema Да Да Да
Persistent operational checkpoint Да Нет Нет
Historical reader Да Да Да
Bounded deterministic Replay Да Да Да

Другие acquisition feeds проекта не входят в завершённую production вертикаль Trades Feed.


3. Режимы по feature flags

Trade Stream Market Data Storage Фактическое поведение
выключен выключен Trade Runtime и Market Data pool не создаются
включён выключен Live, in-memory Consistency и reconnect Recovery без restart persistence
включён включён Durable Trades, checkpoint, Hydration и Startup Recovery
выключен включён Недопустимая конфигурация; settings завершаются ошибкой

Оба флага по умолчанию выключены. При включённом Trade Stream требуются явные WebSocket URL, REST base URL и symbols. Скрытых URL или symbol fallback нет.

Без Storage первоначальный startup подписывается без persistent Hydration, ожидания startup ACK и Startup REST Recovery. Reconnect Recovery во время уже работающего процесса остаётся доступным.


4. Карта компонентов

Граница Ответственность Основные компоненты
Transport WebSocket lifecycle, send/receive, Ping/Pong probe DzengiWebSocketTransport, WebSocketSession
Protocol Команды, события, subscription и control routing AcquisitionRuntimeService, WebSocketSubscriptionManager
Adapter JSON/schema/value validation и Canonical mapping DzengiUnifiedWebSocketAdapter, REST adapter
Runtime Startup, receive loop, buffering, tasks и shutdown TradeStreamProductionRuntime
Consistency Порядок, deduplication и per-symbol state TradeStreamConsistencyController, TradeStreamStateStore
Recovery REST windows и повторная подача в Consistency RuntimeRecoveryCoordinator, TradeRecoveryController
Persistence Canonical write, provenance и checkpoint TradeStorageObservationSink, PostgresTradeRepository
Startup Hydration durable tail и Startup Recovery TradeStreamStateHydrator, RuntimeStartupRecoveryCoordinator
Read-side Historical pages MarketDataHistoricalAccess, PostgreSQL readers
Replay Snapshot, virtual clock и one-shot playback PostgresReplayPlanBuilder, ReplaySessionFactory, ReplaySession
Composition Concrete dependency graph и root lifecycle bootstrap, ApplicationComposition

Transport не знает Canonical models, Recovery или PostgreSQL. Bootstrap является внешним composition root и создаёт concrete adapters.


5. Canonical Trade и Consistency

Один TradeStreamState принадлежит одному нормализованному symbol. Состояние хранит:

  • последнюю принятую сделку;
  • in-memory checkpoint;
  • ограниченное окно Trade ID для deduplication;
  • Canonical payload уже наблюдавшихся сделок.

Размер deduplication tail по умолчанию равен 10 000 Trades.

Trade ID соответствует signed 32-bit контракту. Сравнение учитывает rollover:

  • INT32_MAX → INT32_MIN является продвижением;
  • -1 → 0 является продвижением;
  • расстояние ровно в половину 32-битного цикла неоднозначно и отклоняется.

Две доставки считаются одним биржевым фактом, если совпадают symbol, Trade ID, цена, количество, execution time и aggressor side. Поле source описывает путь доставки WebSocket/REST и не участвует в этом сравнении. При этом source сохраняется как часть Canonical payload и provenance.

Конфликтующий дубликат и нарушение rollover-aware порядка являются ошибками, а не молча отбрасываемыми данными.

Исходники:


6. Live processing

Точная последовательность одного live-документа:

WebSocket receive
→ зафиксировать transport activity
→ опубликовать Runtime Event
→ JSON decode
→ войти в LiveProcessingGate
→ разделить control и market document
→ Dzengi adapter
→ Canonical Trade
→ общий Consistency Controller
→ optional durable sink
→ продвинуть in-memory state

Blocking adapter/Consistency/Persistence path выполняется через asyncio.to_thread. Runtime владеет worker-задачей и при cancellation дожидается уже начатой обработки, чтобы не оставить неизвестный результат записи.

При выключенном Storage Consistency продвигает только память процесса. При включённом Storage durable callback завершается до изменения in-memory checkpoint.

Runtime events доставляются in-process последовательно. Publisher не имеет фоновой очереди или отдельной root task; production composition регистрирует logging consumer.

Исходники:


7. Durable write и persistent checkpoint

Для новой принятой сделки гарантия имеет следующий порядок:

BEGIN
→ insert либо validate identical Trade
→ compare-and-set persistent checkpoint
→ COMMIT
→ advance in-memory state

Trade и checkpoint записываются одним PostgresTradeRepository и одной транзакцией. Ошибка вставки, payload conflict или checkpoint conflict откатывает операцию и не продвигает in-memory state.

Полный дубликат обрабатывается отдельно:

validate stored Trade
→ optional provenance update
→ persistent checkpoint не двигается
→ in-memory checkpoint не двигается

Идентичность строки Trade:

venue + symbol + trade_id + executed_at

Persistent checkpoint имеет один ключ venue + symbol и ссылается на точную durable Trade identity через DEFERRABLE INITIALLY DEFERRED foreign key с NO ACTION.

Parent table Trades partitioned по диапазону executed_at; Quotes — по received_at, Candle revisions — по open_time. Миграции всегда создают default partition. Конкретная месячная UTC-partition появляется только после явного вызова Partition Manager. Symbol и venue являются identity/query scope, но не отдельными физическими уровнями partitioning.

Исходники:


8. Startup с persistent state

Application сначала открывает Market Data pool и применяет migrations. Только после успешного Storage startup запускаются Telegram polling и Trade Runtime tasks.

Persistent Trade Runtime выполняет:

load checkpoint и bounded durable tail
→ либо атомарно принять последнюю durable Trade при первом запуске
→ опубликовать восстановленный StateStore целиком
→ WebSocket connect
→ отправить subscriptions
→ дождаться соответствующего ACK
→ сохранить ранние market documents в bounded FIFO
→ REST Recovery до одной зафиксированной границы времени
→ обработать FIFO в исходном порядке
→ запустить Supervisor, receive loop и Scheduler

Если для symbol нет checkpoint и durable Trades, state создаётся пустым; исторический backfill с биржи до начала доступного REST-окна не выполняется.

Если durable Trades есть, а checkpoint отсутствует, Hydrator может атомарно принять последнюю durable Trade как начальную точку. Corrupt, orphan или несовместимый checkpoint является фатальной ошибкой. Тихого сброса checkpoint нет.

Startup buffer по умолчанию ограничен 10 000 documents. Negative ACK, ACK timeout, overflow, Hydration error и Recovery error не открывают Live gate.

Исходники:


9. Reconnect и Recovery

Transport receive error и Heartbeat timeout используют одну generation-aware single-flight операцию:

закрыть общий Live gate
→ Disconnect
→ Connect
→ переотправить desired subscriptions
→ зафиксировать recovery_end_time
→ последовательно выполнить REST Recovery для symbols
→ тот же Consistency Controller и durable sink
→ открыть Live gate

Вторая причина reconnect для того же connection generation присоединяется к уже выполняющейся операции. Старое transport error не запускает новый reconnect после смены generation.

Важные границы:

  • Reconnect выполняет одну попытку; retry loop и backoff отсутствуют.
  • Restore subscriptions означает успешную повторную отправку; отдельное ожидание ACK перед reconnect Recovery не реализовано.
  • Recovery использует REST /api/v1/aggTrades через настроенный base URL.
  • Окна ограничены TRADE_STREAM_RECOVERY_WINDOW_MS, по умолчанию 3 599 999 ms.
  • Возможный повтор на границе окон удаляет общий Consistency layer.
  • Если state/checkpoint отсутствует, Recovery не придумывает начальную историю и возвращает пустой результат.
  • Ошибка reconnect или Recovery переводит gate в terminal failure и распространяется в Application.

Recovery гарантирует упорядоченную обработку фактически полученных REST данных, но не может гарантировать полноту данных, которых нет в ответе биржи.

Исходники:


10. Liveness, Heartbeat и Scheduler

Встроенный WebSocket keepalive библиотеки отключён. Источником transport liveness является явный Ping/Pong probe.

Scheduler периодически:

  1. вызывает transport probe;
  2. при положительном ответе обновляет Heartbeat activity;
  3. проверяет Heartbeat timeout;
  4. передаёт подтверждённый timeout Supervisor.

Любое успешно полученное WebSocket-сообщение также считается activity. Отрицательный probe сам по себе не запускает reconnect немедленно: решение принимается по Heartbeat timeout. Supervisor не допускает параллельный второй timeout reconnect.

Scheduler не владеет своей asyncio task. Task принадлежит Production Runtime, который предварительно claim()-ит Scheduler и освобождает его после остановки.

Исходники:


11. Historical Access

Historical Access является отдельной synchronous read-side границей и не расширяет write-only Storage facade.

Поддерживаются:

  • Trade History;
  • Quote History;
  • Candle Revision History.

Общие свойства:

  • timezone-aware полуоткрытый диапазон [start_time, end_time);
  • forward-only keyset pagination;
  • строгий cursor scope;
  • Canonical validation прочитанных PostgreSQL rows;
  • limit по умолчанию 500, максимум 1 000.

Порядок страниц:

Тип Порядок
Trade executed_at, затем replay_sequence
Quote received_at, затем replay_sequence
Candle revision open_time, затем replay_sequence

Каждая страница выполняет отдельный запрос. Между страницами не удерживается одна REPEATABLE READ транзакция, поэтому concurrent writes или Retention могут изменить dataset следующего запроса.

Исходники:


12. Deterministic Replay

Replay отделён от live Runtime:

caller
→ ReplaySessionFactory.prepare_session()
→ PostgresReplayPlanBuilder
→ READ ONLY REPEATABLE READ snapshot
→ bounded immutable ReplayPlan
→ закрыть transaction и connection
→ fresh Clock + fresh consumer + fresh ReplaySession
→ caller await session.run()

ReplayPlan ограничен максимум 100 000 events. Превышение limit приводит к явной ошибке без частичного плана.

Глобальный порядок задаётся парой:

replay_at + replay_sequence
Тип replay_at
Trade executed_at
Quote received_at
Candle revision observed_at

replay_sequence выдаётся одним PostgreSQL sequence для всех трёх семейств и остаётся неизменяемым.

Clock начинается в request.time_range.start_time, разрешает равное время и запрещает движение назад. Session является one-shot:

CREATED → RUNNING → COMPLETED | FAILED | CANCELLED

Повторный run() запрещён. Отдельного ReplayEngine, default consumer, background task, automatic startup и wall-clock pacing нет. Вызвавший код предоставляет ReplayConsumerFactoryProtocol. Затем ReplaySessionFactory.prepare_session() создаёт fresh consumer, Clock и Session; caller владеет полученной Session и её run().

Исходники:


13. Ownership и lifecycle

Ресурс Владелец
PostgreSQL pool и migrations Application через MarketDataStorageLifecycle
Telegram polling и root Runtime task run_application()
startup, receive, scheduler и market-processing tasks TradeStreamProductionRuntime
Heartbeat state RuntimeSupervisor
reconnect/recovery single-flight RuntimeReconnectRecoveryCoordinator
Hydration/Startup Recovery worker Runtime через RuntimeStartupRecoveryCoordinator
WebSocket connection DzengiWebSocketTransport через Session
Runtime Event delivery caller publish(); фонового owner нет
Historical connection/cursor один repository call
Replay snapshot transaction PostgresReplayPlanBuilder.create_plan()
Replay execution caller session.run()

Application shutdown:

отменить Telegram polling
→ runtime.stop()
    → завершить startup worker, если он выполняется
    → остановить и дождаться Scheduler
    → остановить Supervisor
    → отменить и дождаться receive loop
    → остановить WebSocket Session
    → очистить subscriptions
→ дождаться root Runtime task
→ закрыть PostgreSQL pool
→ закрыть bot HTTP session

Ошибка включённого Runtime или persistent write считается фатальной для всего приложения. Degraded fallback к in-memory режиму не выполняется.


14. Направление зависимостей

bootstrap
   ↓
concrete Dzengi adapters + Runtime composition + PostgreSQL adapters

acquisition → Canonical models + own Protocol boundaries
            → narrow Storage checkpoint contracts/exceptions
Dzengi REST adapter → integrations.exchange REST client
storage     → Canonical models
access      → Canonical models + read contracts
replay      → Canonical models + access row materialization

Подтверждённые границы:

  • market_data не импортирует Telegram, Trading или Bootstrap;
  • Bootstrap импортирует concrete adapters как composition root;
  • Acquisition не зависит от Historical Access или Replay;
  • Production Runtime не импортирует concrete PostgreSQL repository;
  • Replay не зависит от Live Runtime и не вызывает Consistency;
  • Dzengi-specific код преимущественно находится в acquisition/adapters/dzengi, но provider subscription document и часть validation/handlers физически размещены в соседних acquisition пакетах.

Canonical models сейчас физически находятся в market_data/acquisition/models. Это фактическое расположение, а не утверждение, что общая Domain package уже выделена.


15. Failure policy

Ошибка Поведение
Неверные settings Fail fast до запуска Runtime
Storage pool или migration Application startup завершается ошибкой
Corrupt/orphan checkpoint Startup завершается ошибкой
Subscription ACK timeout/negative ACK при persistent startup Live gate остаётся закрытым, Runtime завершается
Startup buffer overflow Runtime завершается
WebSocket receive failure Одна generation-aware reconnect/recovery попытка
Reconnect или REST Recovery failure Terminal gate failure и завершение Application
Canonical conflict/order error Не принимается и распространяется как ошибка
Persistent write/checkpoint conflict Transaction rollback, in-memory state не двигается
Replay consumer failure Session переходит в FAILED
Replay cancellation Session переходит в CANCELLED

Ошибки не превращаются в молчаливое продолжение с потенциально несогласованным состоянием.


16. Доступный период и ограничения

Доступная история начинается не с момента существования инструмента на бирже, а с наиболее ранней Canonical записи, которая фактически находится в PostgreSQL. Обычно это момент успешного включения Storage, но история может также включать ранее импортированные durable rows.

Граница может сдвигаться из-за явно применённой Retention Policy.

Система не предоставляет:

  • initial exchange historical backfill;
  • completeness metadata и доказательство полной биржевой истории;
  • storage raw exchange documents;
  • automatic monthly partition creation;
  • automatic Retention Scheduler;
  • Production persistence Quotes/Candles;
  • общий snapshot между Historical pages;
  • восстановление исходного порядка сетевых пакетов;
  • Replay HTTP/CLI/UI endpoint;
  • automatic Replay startup;
  • infinite reconnect retry/backoff;
  • HA, leader election, distributed lease или multi-instance ownership;
  • автоматические backup, metrics и alerting.

Default PostgreSQL partitions принимают строки без заранее созданной месячной partition. Явный Partition Manager позднее может атомарно перенести соответствующие строки.

Migration 9 выполняет блокирующий backfill replay_sequence. Для крупной базы нужны backup, измерение на сопоставимом объёме и отдельное maintenance window.


17. Безопасная итоговая формулировка

Trades Feed реализует единый Canonical Trade pipeline с live-получением, in-process consistency, generation-aware reconnect, REST gap recovery, опциональной durable-записью и persistent startup recovery. Historical Access и deterministic Replay работают над фактически сохранёнными Canonical Market Data. Полнота истории ограничена моментом включения persistence, доступностью REST Recovery и применяемой Retention Policy.