Перейти к основному содержимому

Outbox-контракт (backend ↔ бот)

Описание

Бот не подписан ни на какие события backend'а напрямую (ни очередь, ни webhook) — вся доменная интеграция идёт через транзакционный outbox: таблицу bot_event_outbox, которой владеет apps/backend, и HTTP-контракт /bot/events/pending + /bot/events/ack, который бот поллит.

Модуль на стороне backend: apps/backend/src/bot/ (BotOutboxModule@Global(), экспортирует BotOutboxService; BotModule — контроллер /bot/*). Модуль на стороне бота: apps/onboarding-bot/src/outbox/ (OutboxConsumerService + EventHandlersService).

Таблица bot_event_outbox

apps/backend/src/bot/models/bot-event-outbox.entity.ts
@Entity({ name: 'bot_event_outbox' })
@Index('idx_bot_event_outbox_processed_id', ['processedAt', 'id'])
export class BotEventOutbox {
id: number;
eventType: BotEventType; // varchar(64)
telegramId: string | null; // bigint, хранится строкой (precision)
magnetUid: number | null;
payload: Record<string, unknown> | null; // json
createdAt: Date;
processedAt: Date | null;
attempts: number;
claimedAt: Date | null;
lockedUntil: Date | null;
lastError: string | null;
lastAttemptAt: Date | null;
nextRetryAt: Date | null;
}

Миграция apps/backend/src/migrations/1782769264312-create-bot-event-outbox.tsCREATE TABLE IF NOT EXISTS с инлайн-индексом, безопасна при повторном запуске migrationsRun.

Типы событий (BotEventType)

ЗначениеЭмитируется из (backend)Обрабатывается ботом (EventHandlersService)
wallet_connectedПодключение кошелька к пользователюПродвигает этап 7 (WalletConnected), инициирует confirm-промпт привязки
registeredОнчейн-синк регистрации (sync.service)Продвигает этап 8 (Registered), детект «ушёл к другому лидеру», confirm-промпт привязки
new_partnerРегистрация партнёра под реферомDM рефереру «🎉 Новый партнёр…» (или пассивный фоллбэк-нотиф, если DM недоступен)
status_nft_activatedАктивация статусного NFT (LP1/LP2)Продвигает этап 9 (NftActivated)
reactivation(зарезервировано, backend-эмиссия не включена — см. known limitations)DM «Рады видеть вас снова…»
struct_overtake(зарезервировано, backend-эмиссия не включена — см. known limitations)DM «Важно: изменения в вашей структуре…»

BotEventType продублирован на обеих сторонах и обязан оставаться в синхроне — значения персистятся как строки и не версионируются: apps/backend/src/bot/bot.constants.tsapps/onboarding-bot/src/models/funnel.enums.ts (BotEventType).

Эмиссия — методы BotOutboxService (apps/backend/src/bot/bot-outbox.service.ts): emitWalletConnected, emitRegistered, emitNewPartner, emitStatusNftActivated. Каждый метод делает ровно один INSERT в bot_event_outbox — сам по себе неатомарен с бизнес-транзакцией вызывающего кода (не обёрнут в общий QueryRunner), поэтому доставка — at-least-once, не exactly-once (см. ниже).

Producer: GET /bot/events/pending

BotService.getPendingEvents(limit) возвращает пачку необработанных событий:

WHERE processed_at IS NULL
AND (next_retry_at IS NULL OR next_retry_at <= NOW())
ORDER BY COALESCE(magnet_uid, telegram_id) ASC, id ASC
LIMIT <limit>

Сортировка по идентичности пользователя группирует все события одного юзера подряд в одной странице — это важно для per-user ordering на стороне консьюмера (ниже). limit по умолчанию BOT_OUTBOX_DEFAULT_PENDING_LIMIT = 100, жёсткий кап BOT_OUTBOX_MAX_PENDING_LIMIT = 500.

Consumer: POST /bot/events/ack

BotService.acknowledge({ processedIds, retryIds }):

  • processedIdsprocessed_at = NOW().
  • retryIdsattempts += 1; если attempts < BOT_OUTBOX_MAX_ATTEMPTS (10)next_retry_at = NOW() + BOT_OUTBOX_BACKOFF_BASE_SECONDS(60) * 2^attempts (экспоненциальный backoff, секунды); достигнут предел —
    • registereddead-letter: остаётся processed_at = NULL навсегда, только ERROR-лог [bot-outbox][DEAD-LETTER] для ручного разбора (событие регистрации никогда не роняется молча);
    • любой другой тип → принудительно помечается processed_at = NOW() (не блокирует очередь бесконечно).

Consumer: поллинг и per-user ordering

OutboxConsumerService.poll()@Cron(EVERY_10_SECONDS):

  1. Single-flight через RedlockService (ключ onboarding:outbox, TTL 25s) — если лок не взят, цикл просто выходит (другой цикл/инстанс уже работает).
  2. processOnce(): тянет пачку (BATCH_LIMIT = 100), группирует по идентичности юзера (tg:<telegramId> или uid:<magnetUid>), внутри группы сортирует по возрастанию id.
  3. Для каждого события группы: tryMark()INSERT в bot_processed_events по outboxEventId (дедуп-леджер, UNIQUE); если ER_DUP_ENTRY/23505 — событие уже обработано раньше, всё равно идёт в processedIds (чтобы курсор ack'а двигался), обработчик не вызывается повторно.
  4. Если обработчик (EventHandlersService.handle) бросает — событие уходит в retryIds, и обработка этого юзера останавливается (break) — следующие события того же пользователя в этой пачке не трогаются, чтобы не нарушить порядок (например, нельзя применить registered раньше не доехавшего wallet_connected). Другие пользователи в пачке продолжают обрабатываться независимо.
  5. В конце — один вызов ack(processedIds, retryIds).

Гарантии доставки

  • At-least-once, не exactly-once. Между применением побочного эффекта (handlers.handle(event)) и записью в дедуп-леджер (tryMark вставляется до обработки, что на самом деле защищает от повторной обработки при повторном поллинге того же события, но не от гонки двух процессов — здесь подстраховывает single-flight redlock) возможен краш процесса; ack идемпотентен, а обработчики (onWalletConnected, onRegistered, …) сами идемпотентны по данным (GREATEST при продвижении этапа, findOne + save вместо insert), поэтому повторное применение события неопасно.
  • Dead-letter только для registered — единственный тип события, потеря которого означает необратимо пропущенную привязку/продвижение воронки; остальные типы (уведомления, апдейт NFT-уровня) при исчерпании попыток считаются менее критичными и принудительно закрываются, чтобы не блокировать очередь.
  • Дедуп-леджер бота растёт бесконечно без чисткиProcessedPruneService (@Cron(EVERY_DAY_AT_MIDNIGHT)) ежедневно удаляет строки bot_processed_events старше 30 дней (RETENTION_DAYS). Безопасно: backend сам не переигрывает уже подтверждённые (processed_at != NULL) события, так что старые записи дедупа не нужны для корректности.

Прочие сервис-эндпоинты /bot/*

Не часть outbox-очереди, но под тем же BotAuthGuard и swagger tag bot:

EndpointНазначение
GET /bot/leader/resolve?code=Резолв реф-кода (users.link) в { uid, link, name } лидера
GET /bot/user/state?telegramId=Состояние юзера для гейта консоли (isRegistrationConfirmed, lp1, lp2, addr, timeActive, ownLink)
POST /bot/telegram/bindПромоут previewTelegramId → telegram_id после in-bot подтверждения
POST /bot/notificationПассивный фоллбэк-нотиф в веб-кабинет, когда DM недоступен (403)
POST /bot/telegram/connect-linkМинт (или переиспользование) одноразового telegramConnect=<uuid> токена

А также GET /auth/check-account-by-telegram/:idAuthController, тоже под BotAuthGuard) — используется ViewerService/OnboardingUpdate для проверки, привязан ли telegram_id к Magnet-аккаунту.

Аутентификация

BotAuthGuard (apps/backend/src/shared/guards/bot.guard.ts) требует Authorization: Bearer <BOT_TO_BACKEND_SECRET>, сравнение — crypto.timingSafeEqual (с предварительной проверкой длины буферов, т.к. timingSafeEqual бросает при несовпадении длины) — защита от утечки секрета через time-based side channel. Секрет общий для обеих сторон (BOT_TO_BACKEND_SECRET в .env бота и backend'а).