Eventing And Queue Baseline — Событийная шина, очереди и асинхронная дисциплина платформы
Версия: 1.0
Дата: 24.04.2026
Статус: Готов к обсуждению
Назначение документа
Этот документ фиксирует рекомендуемый baseline для eventing, queues, async processing и delivery semantics платформы vitrip.store.
Его задача — определить:
- где платформа обязана мыслить событиями, а где синхронным запросом;
- какие классы очередей и событий нужны разным execution contours;
- какие delivery guarantees достаточно сильны для платформы, а где нужна более строгая дисциплина;
- как eventing связан с replay, observability, release safety и domain truth.
Опорные документы
- Архитектурная основа платформы vitrip.store
- Business Services — Сервисная декомпозиция платформы
- Ingestion Layer — Приём, нормализация, маппинг и governance
- Offer Pricing Booking Semantics — Семантика предложения, цены и бронирования
- Partner Finance And Clearing — Балансы, лимиты, взаиморасчёты и финансовая дисциплина партнёров
- API Metering And Usage Governance — Учёт потребления, квоты и дисциплина использования
- Implementation Technology Baseline — Рекомендуемый технологический фундамент реализации
- Initial Event Taxonomy — Первичная таксономия событий платформы
- Deployment And Operating Model — Развёртывание и эксплуатационная модель
- Observability And Incident Response — Наблюдаемость и реагирование на инциденты
- Release Engineering And Migrations — Релизы, совместимость и эволюция схем
Почему Этот Документ Нужен Отдельно
Платформа уже не может мыслиться как чисто synchronous request/response система.
У неё уже есть и будут усиливаться:
- supplier ingestion;
- replay;
- invalidation waves;
- repricing cascades;
- quote expiry and refresh;
- booking follow-up;
- settlement event generation;
- reconciliation queues;
- metering aggregation;
- governance and review cases.
Без отдельной фиксации eventing and queue baseline система рискует скатиться либо в:
- хаотичный background jobs zoo;
- либо в forced synchronous architecture там, где она operationally вредна.
Главный Принцип
Async processing должен использоваться не потому, что “так модно масштабировать”, а потому, что разные доменные контуры имеют разные требования к latency, durability, retry semantics and recoverability.
Из этого следует:
очередь или событие не являются заменой truth. Они являются transport and coordination mechanism вокруг truth-bearing domains.
Где Eventing Обязателен
1. Supplier Intake And Ingestion
Нужны:
- intake jobs;
- parsing jobs;
- normalization jobs;
- review triggers;
- replay tasks;
- downstream invalidation events.
2. Offer And Quote Downstream Effects
Нужны:
- repricing triggers;
- quote invalidation signals;
- cache invalidation events;
- publication gating notifications.
3. Booking Follow-Up
Нужны:
- supplier confirmation follow-up;
- unknown-state recovery jobs;
- cancellation/amendment processing;
- post-booking case routing.
4. Settlement And Clearing
Нужны:
- settlement event generation;
- reconciliation queue items;
- adjustment and correction jobs;
- clearing recalculation triggers.
5. Governance And Operations
Нужны:
- review queue creation;
- anomaly case distribution;
- operator intervention tasks;
- SLA escalation jobs.
Где Нельзя Прятать Смысл За Eventual Consistency
Платформа не должна прятать за asynchronous optimism следующие критические переходы:
- final quote promise creation;
- booking commit decision;
- critical access/tenant authorization decision;
- final partner financial admissibility check;
- externally visible “booking confirmed” promise.
Эти переходы могут порождать downstream events, но их доменная фиксация должна происходить в controlled transactional boundary.
Классы Очередей
1. Durable Work Queues
Для задач, которые нельзя терять:
- supplier ingestion stages;
- replay jobs;
- booking follow-up;
- financial correction work;
- governance review tasks.
2. Event Streams
Для широкого fan-out of domain-significant signals:
- offer updated;
- quote invalidated;
- booking state changed;
- settlement event posted;
- tenant throttled;
- publication blocked.
3. Delay / Scheduled Queues
Для:
- quote expiry checks;
- retry with backoff;
- SLA timers;
- deferred reconciliation checks;
- delayed supplier follow-up.
4. Dead-Letter / Quarantine Queues
Для:
- poison messages;
- non-repeatable errors;
- suspicious supplier payloads;
- manual triage candidates.
Delivery Guarantees
At-Least-Once — Базовый Практический Выбор
Для большинства industrial use cases платформы baseline-дисциплина должна исходить из:
- at-least-once delivery;
- idempotent consumers;
- explicit deduplication where needed;
- replay safety.
Exactly-Once
Не должен становиться blanket-requirement для всего eventing слоя.
Он слишком дорог и усложняет систему.
Вместо этого платформа должна предпочитать:
- strong transactional write of truth;
- deterministic event emission strategy;
- idempotent downstream handlers;
- correction events where needed.
Ordering
Ordering нужен не глобально, а там, где он доменно важен.
Особенно это касается:
- booking state transitions;
- settlement and clearing events;
- quote invalidation after repricing;
- review-case lifecycle.
Глобальный total ordering для всей платформы не нужен и вреден.
Replay As First-Class Capability
Платформа должна изначально проектироваться так, чтобы replay был нормальной операцией, а не emergency hack.
Replay важен для:
- ingestion improvements;
- supplier incident recovery;
- corrected normalization logic;
- rebuilding projections;
- re-evaluating publication integrity.
Long-Running Idempotency для саг и многоэтапных операций
Раздел добавлен 30.04.2026 как самостоятельное наблюдение архитектора (gap, не указанный во внешнем ревью).
Раздел фиксирует каноничную семантику idempotency для долгих многошаговых операций, типичными примерами которых являются Tour Builder саги (booking commit нескольких segment'ов через нескольких supplier'ов), сложные refund-flows, multi-supplier orchestration.
Принцип
Стандартная idempotency (один ключ — один результат, повторы безопасны) хорошо работает для синхронных операций, завершающихся за секунды. Но для длинных операций (минуты, иногда часы) эта семантика недостаточна: партнёр может сделать повтор запроса с тем же idempotency-key до того как операция завершилась, и платформа должна обработать это корректно.
Каноничный принцип: платформа различает три фазы по idempotency-key:
- Operation in_progress — операция начата, ещё не завершена. Повтор запроса — не запускает повторное выполнение, а возвращает текущее состояние операции с partial result если применимо.
- Operation completed — операция завершилась успехом или явной ошибкой. Повтор запроса — возвращает финальный результат операции.
- Operation expired — операция вышла за timeout без явного результата. Повтор запроса — может новую операцию с тем же ключом запустить, но только после явного
force_retryпараметра.
Каноничный жизненный цикл idempotency-key для long-running операций
Запрос с idempotency-key K приходит первый раз
↓
Создаётся IdempotencyRecord(K):
state = in_progress
started_at = now
timeout_at = now + operation_timeout
last_progress_event_id = null
final_result = null
↓
Операция стартует асинхронно через Saga
↓
Партнёр повторяет запрос с тем же K до завершения
↓
Платформа смотрит IdempotencyRecord(K), state = in_progress
↓
Возвращает 202 Accepted с структурированным payload:
{
"operation_id": "K_internal",
"state": "in_progress",
"started_at": "...",
"estimated_completion_at": "...",
"progress": {
"completed_steps": [...],
"current_step": "...",
"remaining_steps": [...]
},
"partial_result": {... доступные на текущий момент данные ...}
}
↓
Партнёр периодически polls с тем же idempotency-key
↓
Операция завершилась
↓
IdempotencyRecord(K).state = completed (or failed)
IdempotencyRecord(K).final_result = {...}
↓
Любой следующий повтор с K возвращает 200 OK + final_result
↓
IdempotencyRecord(K) хранится N дней (по умолчанию 30) для возможности retrieve финального результата
Каноничные timeout'ы для разных типов операций
| Тип операции | Operation timeout | Idempotency record retention |
|---|---|---|
| Простой booking commit (один supplier) | 5 минут | 30 дней |
| Tour booking saga (несколько supplier'ов) | 30 минут | 30 дней |
| Multi-segment refund saga | 60 минут | 90 дней |
| Async ticketing (post-booking) | 24 часа | 90 дней |
| Long-running data export | 7 дней | 30 дней |
При истечении operation timeout без завершения:
- Saga сама завершается через compensation flow (rollback всех уже сделанных шагов).
- IdempotencyRecord переходит в state =
expired. - Партнёр получает финальный response с
state: expiredи details ошибки.
Что НЕ делает long-running idempotency
Чтобы избежать неправильных интерпретаций:
- Не гарантирует exactly-once execution в распределённой системе — гарантирует, что наблюдаемый партнёром результат соответствует одному завершению операции (либо его атомарной отмены).
- Не позволяет партнёру «узнать ответ заранее» — пока операция в in_progress, payload не содержит финального результата, только partial state.
- Не освобождает партнёра от обязанности обрабатывать
state: in_progressкорректно — клиент должен polls с reasonable cadence (не более 1 запроса в секунду на operation), не делать infinite retry без backoff'а.
Как партнёр должен polls in_progress операции
Каноничный паттерн polling'а:
- Initial request возвращает
state: in_progressсRetry-Afterheader'ом (рекомендуемое окно polling'а, например 10 секунд). - Subsequent polls с тем же idempotency-key выполняются не чаще чем
Retry-After. Слишком частые polls могут rate-limit'ed. - Exponential backoff при затягивании операции: первый poll через 5s, второй через 10s, третий через 20s, далее по 60s.
- При получении
state: in_progressсpartial_result— клиент может отображать прогресс end-customer'у (например, «Бронируем перелёт... [✓] Отель... [⌛] Трансфер...»).
Связь с partial booking states
Long-running idempotency тесно интегрирована с booking state machine (см. reference/booking-state-machine.md):
partially_confirmedсостояние booking'а видно через partial_result в idempotency response.- Каждый supplier confirmation/failure обновляет обе структуры: booking state machine и idempotency record (через event handlers).
- Partner может из partial_result определить, какие segment'ы уже подтверждены — что критично для UX (например, показать пользователю, какие части тура уже зафиксированы).
Force retry — каноничный mechanism
Иногда партнёр сознательно хочет retry completed-but-failed операции без создания нового idempotency-key (например, transient failure при booking, и партнёр хочет повторить с теми же параметрами).
Каноничный паттерн:
- Партнёр посылает запрос с тем же idempotency-key + явным header
X-Idempotency-Force-Retry: true. - Платформа проверяет, что existing IdempotencyRecord(K) находится в state =
failedилиexpired. - Создаётся new attempt с modified key
K@retry-N, но партнёр продолжает видеть его как операцию K. - Original record archived; new record создаётся в state =
in_progress.
Безопасность: force retry работает только если original failed/expired. Для completed (успешных) операций force retry отвергается с 409 Conflict (нельзя повторно execute уже успешно завершённую операцию через тот же ключ).
Метрики Long-Running Idempotency
long_running_in_progress_polls_per_operation— среднее число polls пока партнёр ждёт результат (показатель UX, target < 5).idempotent_response_rate— доля повторов, удачно обработанных через idempotency (target ≈ 100%).force_retry_rate— частота явных retry attempts (показатель надёжности supplier-cети, цель ≤ 1%).expired_operation_rate— частота операций, перешедших в state = expired (target ≤ 0.5%).polling_violations— нарушения cadence polling'а партнёрами (для rate-limit policies).
Хранение и cleanup
- IdempotencyRecord хранится в dedicated хранилище (не в основной БД доменных сущностей), партиционированное по дате.
- Retention в зависимости от типа (см. таблицу выше).
- После expiry — архивирование в холодное хранение для audit, удаление через 1 год.
Обоснование (тезисы)
Тезис 1. Idempotency-семантика для длинных операций отличается от синхронных операций.
Альтернативы: (а) единая семантика для всех операций; (б) расширенная семантика для длинных.
Trade-off: вариант (а) приводит к ситуациям где партнёр получает или ошибки при повторе, или начинает новую операцию параллельно (race condition); вариант (б) — каноничный, обеспечивает predictable UX для длинных операций.
Тезис 2. Polling — нормальный pattern для длинных операций, но требует discipline.
Альтернативы: (а) только webhooks для длинных операций; (б) только polling; (в) и то, и то.
Trade-off: вариант (а) требует партнёру держать публичный endpoint для callbacks, что не всегда возможно (особенно на early stages); вариант (б) — каноничный baseline, всегда работает; (в) — extended вариант, добавляется на Stage 2+ для удобства.
Тезис 3. Force retry — отдельный явный механизм, не side-effect повтора.
Альтернативы: (а) автоматический retry при повторе (idempotency-key с тем же значением запускает new attempt); (б) явный force retry header.
Trade-off: вариант (а) маскирует ошибки и приводит к неожиданным retries; вариант (б) — каноничный, явный intent, audit trail.
Technology Baseline
На текущем этапе наиболее разумный baseline выглядит так:
- durable queue / stream backbone such as Kafka or a comparable durable event log for high-value domain fan-out and replay-sensitive flows;
- simpler broker-backed work queues where full log semantics are unnecessary;
- scheduled/delay processing either through broker capabilities or dedicated scheduler workers;
- explicit dead-letter and quarantine handling as mandatory.
Важная Оговорка
Платформе не нужно насильно тащить один и тот же messaging product во все async scenarios.
Но ей нужно удерживать единый operational model:
- traceability;
- replayability;
- queue visibility;
- failure handling;
- backpressure discipline.
Queue Governance
Для каждой queue or stream family должно быть ясно:
- кто producer;
- кто consumer;
- что является payload contract;
- что является idempotency key;
- сколько хранится сообщение;
- как выглядит retry policy;
- когда сообщение считается poisoned;
- куда уходит operator action.
Связь С Observability
Очереди и события должны быть fully observable:
- backlog;
- lag;
- retry rate;
- age;
- poison rate;
- handler latency;
- correlation to domain ids.
Queue invisibility для этой платформы равна operational blindness.
Связь С Release Engineering
Любая серьёзная эволюция event contracts должна учитывать:
- compatibility window;
- old and new consumers;
- replay impact;
- schema evolution;
- correction strategy for already emitted events.
Что Должно Быть Сделано Дальше
После фиксации этого baseline нужно:
- синхронизировать
ingestion,deployment,release-engineering,observability,implementation-technology-baseline; - опираться на Initial Event Taxonomy — Первичная таксономия событий платформы как на первый управляемый словарь event families;
- определить queue families per execution contour;
- зафиксировать preferred broker strategy на implementation-ready уровне.
Уточнение под Фазы 4–6 (28.04.2026) — расширение классов событий и связи
Документ опубликован 24.04.2026 в Фазе 3 как baseline для eventing. После Фаз 4–6 опубликованы документы, которые расширяют классификацию событий и используют этот baseline. Эта секция фиксирует расширения и обязательные связи.
Расширение классификации событий
В этом документе зафиксированы 2 главных класса (domain events + analytical events). После Фаз 4–6 — 6 классов:
| Класс | Источник истины | Каноничный документ |
|---|---|---|
| Domain events (state transitions, audit-significant) | сервисы платформы | eventing-and-queue-baseline.md (этот документ); каноничные имена в booking-state-machine.md, payment-domain.md, и т.д. |
| Analytical events (поведение subjects на surfaces) | поверхности взаимодействия | data-platform-and-events-tracking.md |
| Saga events (Tour Builder transactions) | tour-builder-service | tour-builder-operational-model.md |
| Audit / security events (immutable WORM) | все services через audit pipeline | security-architecture.md |
| Operational events (DR, capacity, incidents, SLA breach) | runbooks, capacity, DR services | disaster-recovery-and-capacity.md, sla-and-on-call-model.md, runbooks-incident-playbooks.md |
| Compliance events (DSR, consent, breach notification) | compliance pipeline | compliance-and-legal.md |
Каждый класс имеет разные guarantees:
| Класс | Delivery guarantee | Retention | Order | Encryption |
|---|---|---|---|---|
| Domain | at-least-once + ordered per entity | долгий (audit) | per-entity | at-rest |
| Analytical | at-least-once с допустимой потерей малой доли | class-зависимый | best-effort | at-rest |
| Saga | exactly-once-effective (idempotent) + ordered | долгий | strict | at-rest |
| Audit / security | at-least-once + tamper-detection | 7 лет (Tier 1) | timestamp-ordered | at-rest + tamper-evident |
| Operational | at-least-once | долгий (regulated) | best-effort | at-rest |
| Compliance | at-least-once + immutable | 7+ лет | timestamp-ordered | at-rest + audit-grade |
Связь со state machine Booking
reference/booking-state-machine.md (Фаза 5) — каноничные 14 состояний Booking. Каждое state transition генерирует domain event через event bus этого baseline. Имена событий — каноничные (см. data-platform-and-events-tracking.md, секция «Связь со статусной машиной бронирования»).
Связь с partner webhooks
reference/notification-and-communication.md (Фаза 4) — webhook delivery как отдельный класс external delivery:
- domain events trigger webhook delivery;
- webhook signing через HMAC (см. security-architecture.md);
- retry policy с exponential backoff;
- DLQ для unreachable endpoints;
- webhook events visible через partner-facing API as Product (см. api-as-product.md).
Связь с replay capability
Этот документ упоминает replay для recovery. Реализация — через append-only event log в DWH (см. data-platform-and-events-tracking.md). Каноничные replay scenarios:
- Изменение нормализации — replay supplier ingestion events с новой логикой;
- Bug fix downstream — replay events для пересчёта affected aggregations;
- DR recovery — replay events с last verified point;
- ML feature recompute — replay events для нового feature engineering.
Replay events generates replay.started, replay.progress, replay.completed — отдельный класс operational events.
Связь с метеринг и квотами
reference/api-metering-and-usage-governance.md использует events этого baseline для:
- usage metering (по факту получения events);
- quota enforcement (real-time через event bus);
- billing aggregation (batch через DWH).
Каноничный итог уточнения
Eventing baseline остаётся техническим фундаментом (event bus choice, delivery semantics, queue patterns). Расширения:
- 6 классов событий вместо 2 (с разными guarantees per class);
- Каноничные имена событий — в специализированных документах (booking-state-machine, payment-domain, tour-builder-operational и т.д.);
- Replay capability реализуется через append-only event log в DWH;
- Webhook delivery — отдельный класс external delivery с дополнительными гарантиями.
Уточнение выполнено через no-destruction.