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
@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.ts
— CREATE 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.ts ↔
apps/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 }):
processedIds→processed_at = NOW().retryIds→attempts += 1; еслиattempts < BOT_OUTBOX_MAX_ATTEMPTS (10)—next_retry_at = NOW() + BOT_OUTBOX_BACKOFF_BASE_SECONDS(60) * 2^attempts(экспоненциальный backoff, секунды); достигнут предел —registered→ dead-letter: остаётсяprocessed_at = NULLнавсегда, толькоERROR-лог[bot-outbox][DEAD-LETTER]для ручного разбора (событие регистрации никогда не роняется молча);- любой другой тип → принудительно помечается
processed_at = NOW()(не блокирует очередь бесконечно).
Consumer: поллинг и per-user ordering
OutboxConsumerService.poll() — @Cron(EVERY_10_SECONDS):
- Single-flight через
RedlockService(ключonboarding:outbox, TTL 25s) — если лок не взят, цикл просто выходит (другой цикл/инстанс уже работает). processOnce(): тянет пачку (BATCH_LIMIT = 100), группирует по идентичности юзера (tg:<telegramId>илиuid:<magnetUid>), внутри группы сортирует по возрастаниюid.- Для каждого события группы:
tryMark()—INSERTвbot_processed_eventsпоoutboxEventId(дедуп-леджер,UNIQUE); еслиER_DUP_ENTRY/23505— событие уже обработано раньше, всё равно идёт вprocessedIds(чтобы курсор ack'а двигался), обработчик не вызывается повторно. - Ес ли обработчик (
EventHandlersService.handle) бросает — событие уходит вretryIds, и обработка этого юзера останавливается (break) — следующие события того же пользователя в этой пачке не трогаются, чтобы не нарушить порядок (например, нельзя применитьregisteredраньше не доехавшегоwallet_connected). Другие пользователи в пачке продолжают обрабатываться независимо. - В конце — один вызов
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/:id (в AuthController, тоже
под 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'а).