2026-08-16 08:41:54 +07:00
|
|
|
|
# Discord Gateway — Architecture
|
|
|
|
|
|
|
|
|
|
|
|
Pure event-driven microservice (no HTTP server). Captures Discord
|
2026-09-23 21:25:03 +07:00
|
|
|
|
messages/attachments/reactions/threads/presence, runs LLM-based AI
|
2026-08-16 08:41:54 +07:00
|
|
|
|
moderation, and publishes everything to Redis pub/sub for the backend to
|
|
|
|
|
|
consume. The backend serves the HTTP/WS API to the frontend.
|
|
|
|
|
|
|
2026-09-24 15:09:13 +07:00
|
|
|
|
> NOTE: this doc is the source of truth for the module layout. The old
|
|
|
|
|
|
> `MODULE_STRUCTURE.md` was a stale duplicate and has been removed. `README.md`
|
|
|
|
|
|
> only covers how to run the service.
|
2026-08-16 08:41:54 +07:00
|
|
|
|
|
|
|
|
|
|
## Top-level layout
|
|
|
|
|
|
|
|
|
|
|
|
```
|
2026-06-01 21:44:29 +07:00
|
|
|
|
services/discord-gateway/
|
|
|
|
|
|
├── src/
|
2026-08-16 08:41:54 +07:00
|
|
|
|
│ ├── index.ts # Entry point → initializeDiscordGateway()
|
2026-09-24 15:09:13 +07:00
|
|
|
|
│ ├── app/ # Process lifecycle
|
|
|
|
|
|
│ │ ├── bootstrap.ts # Startup order: config → DB → services → metrics → login
|
|
|
|
|
|
│ │ ├── lifecycle.ts # Everything wired on the Discord 'ready' hook
|
|
|
|
|
|
│ │ ├── process-guards.ts # SIGINT/SIGTERM + uncaught-error policy
|
|
|
|
|
|
│ │ ├── metrics-collector.ts # AI pipeline Prometheus gauges
|
|
|
|
|
|
│ │ ├── shutdown.ts # Graceful shutdown sequence
|
2026-08-16 08:41:54 +07:00
|
|
|
|
│ │ └── retention.ts # Expired-record cleanup scheduler
|
2026-09-24 15:09:13 +07:00
|
|
|
|
│ ├── shared/ # Infrastructure — never imports from modules/
|
2026-08-16 08:41:54 +07:00
|
|
|
|
│ │ ├── config/ # Zod-validated env (index.ts = schema+loader)
|
|
|
|
|
|
│ │ ├── database/ # Drizzle ORM + pg Pool + migrations
|
|
|
|
|
|
│ │ │ ├── init.ts drizzle.ts pool.ts migrate.ts migrateCli.ts
|
2026-09-23 21:25:03 +07:00
|
|
|
|
│ │ │ └── schema/ # messages, cache, meta, analytics
|
2026-08-16 08:41:54 +07:00
|
|
|
|
│ │ ├── logger/ # pino wrapper + createChildLogger()
|
2026-09-24 15:09:13 +07:00
|
|
|
|
│ │ ├── errors/ # AppError / ConfigError ... + errorMessage()
|
|
|
|
|
|
│ │ │ # + isTransientStreamError()
|
2026-08-16 08:41:54 +07:00
|
|
|
|
│ │ ├── utils/ # retry, pagination
|
|
|
|
|
|
│ │ ├── discord/clientOptions.ts # discord.js-selfbot-v13 client options
|
|
|
|
|
|
│ │ ├── uploader.ts # Shared attachment upload helper
|
2026-09-24 15:09:13 +07:00
|
|
|
|
│ │ ├── redis-channels.ts # Redis channel + command constants
|
2026-08-16 08:41:54 +07:00
|
|
|
|
│ │ └── moderation-types.ts # Shared AI analysis domain types
|
2026-09-24 15:09:13 +07:00
|
|
|
|
│ └── modules/ # Feature modules, each with an index.ts facade
|
2026-08-16 08:41:54 +07:00
|
|
|
|
│ ├── message-capture/ # Discord event listeners + DB store
|
|
|
|
|
|
│ ├── ai-moderation/ # LLM moderation pipeline (see below)
|
|
|
|
|
|
│ ├── attachment-upload/ # Download + (sharp) resize + upload
|
2026-09-24 15:09:13 +07:00
|
|
|
|
│ ├── event-broadcaster/ # RedisEventPublisher + EventBroadcaster
|
2026-08-16 08:41:54 +07:00
|
|
|
|
│ ├── command-handler/ # Redis-subscribed backend→gateway commands
|
|
|
|
|
|
│ ├── reaction-tracking/ thread-tracking/ user-presence/
|
2026-09-24 15:09:13 +07:00
|
|
|
|
│ ├── channel-topic/ guild-member-events/ monitor/
|
2026-08-16 08:41:54 +07:00
|
|
|
|
│ └── gateway-metrics/ # Prometheus /metrics endpoint (port 4016)
|
|
|
|
|
|
```
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-09-24 15:09:13 +07:00
|
|
|
|
Dependency direction is one-way: `index.ts` → `app/` → `modules/` → `shared/`.
|
|
|
|
|
|
Code outside a module imports its `index.ts` facade, never an internal file;
|
|
|
|
|
|
deep imports stay valid inside the module itself.
|
|
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
## AI moderation pipeline (`ai-moderation/`)
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
LLM-only judge — no regex/heuristic classification. One orchestrator call
|
2026-09-24 14:09:53 +07:00
|
|
|
|
handles a whole batch. **Independent text/media lanes** (2026-09-24): a
|
|
|
|
|
|
conversation batch is split into a text lane (messages with no media) and a
|
|
|
|
|
|
media lane (attachments/stickers/embeds) that are dispatched to separate
|
|
|
|
|
|
pools, hold SEPARATE per-lane processing locks, and run under SEPARATE LLM
|
|
|
|
|
|
concurrency semaphores. The text lane frees its lock and saves+broadcasts the
|
|
|
|
|
|
moment text analysis finishes — it never waits on a slow vision/media batch
|
|
|
|
|
|
of the same conversation, and vice versa.
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-09-24 15:09:13 +07:00
|
|
|
|
- `aiAnalyzer.ts` — public API: `queueMessageAnalysis`, `queueConversationAnalysis`,
|
|
|
|
|
|
`getAnalysisQueueStatus`, `startPendingAIAnalysisWorker`. Short-circuits
|
|
|
|
|
|
age-restricted and skip-list messages before any LLM work.
|
|
|
|
|
|
- `recovery-worker.ts` — periodic sweep for stranded `pending` messages
|
|
|
|
|
|
(re-scheduled per lane) and `error`/`analysis_incomplete` messages
|
|
|
|
|
|
(individual fallback queue); prunes stale lane locks, per-conversation CB
|
|
|
|
|
|
counters and individual in-flight markers.
|
|
|
|
|
|
- `cache-prune.ts` — throttled (6h) expired-verdict sweep across Postgres and
|
|
|
|
|
|
Qdrant, driven from the recovery interval.
|
2026-09-24 14:09:53 +07:00
|
|
|
|
- `batchScheduler.ts` — per-conversation per-LANE debounce → `processBatch`
|
|
|
|
|
|
(lane-aware). `splitMessagesByLane` / `laneOfMessage` live in
|
|
|
|
|
|
`analysisLanes.ts` (pure, unit-testable).
|
|
|
|
|
|
- `batchProcessor.ts` — per-lane batch lock/circuit-breaker, fans failed
|
|
|
|
|
|
targets to individual fallback. `processBatch` releases ITS lane's lock the
|
|
|
|
|
|
moment that lane's worker job finishes; the other lane owns its own lock.
|
2026-08-16 08:41:54 +07:00
|
|
|
|
- `individualFallbackProcessor.ts` — one-message-at-a-time retry path, own CB.
|
2026-09-24 14:09:53 +07:00
|
|
|
|
- `conversationState.ts` / `circuitBreaker.ts` — per-conversation PER-LANE
|
|
|
|
|
|
state (`conversationProcessing` holds a lane → startedAt map per key),
|
|
|
|
|
|
Piscina `textWorkerPool`/`mediaWorkerPool`, `getConversationKey`.
|
|
|
|
|
|
- `ai-analysis-worker.ts` — Piscina entry point (`batch` (lane) /
|
|
|
|
|
|
`individual` jobs). Runs `runModerationAnalysis` off the main thread.
|
2026-08-16 08:41:54 +07:00
|
|
|
|
- `moderationOrchestrator.ts` — exact-hash cache → batched semantic (Qdrant)
|
|
|
|
|
|
cache → LLM. Text and media paths run in parallel.
|
|
|
|
|
|
- `textBatchProcessor.ts` / `mediaBatchProcessor.ts` — actual LLM calls
|
2026-09-24 14:09:53 +07:00
|
|
|
|
(one call per sub-batch, not per message). `mediaBatchProcessor` routes its
|
|
|
|
|
|
moderation LLM call through the MEDIA semaphore.
|
2026-08-16 08:41:54 +07:00
|
|
|
|
- `llmClient.ts` — central OpenAI-compatible chat client (streaming, retries,
|
2026-09-24 14:09:53 +07:00
|
|
|
|
thinking-disable injection). TWO concurrency semaphores:
|
|
|
|
|
|
`AI_LLM_MAX_CONCURRENT` (text lane, default 8) and
|
|
|
|
|
|
`AI_LLM_MEDIA_MAX_CONCURRENT` (media lane, default 4) — a vision backlog
|
|
|
|
|
|
can never consume text slots. `visionAnalyzer.ts` / `mediaAnalysisClient.ts`
|
2026-08-16 08:41:54 +07:00
|
|
|
|
share the same router/base URL (different model alias for vision).
|
|
|
|
|
|
- `embeddingClient.ts` + `qdrantClient.ts` — semantic cache (one embed call +
|
|
|
|
|
|
one batched Qdrant search for all uncached targets).
|
|
|
|
|
|
- `textCacheStore.ts` / `channelCultureStore.ts` / `userProfileStore.ts` /
|
2026-08-18 18:27:15 +07:00
|
|
|
|
`userProfileStore.ts` — caches learned user profile summaries (optional).
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
### Concurrency model
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-09-24 14:09:53 +07:00
|
|
|
|
- Main thread owns TWO per-lane LLM semaphores (2026-09-24):
|
|
|
|
|
|
`AI_LLM_MAX_CONCURRENT` (text, default 8) and `AI_LLM_MEDIA_MAX_CONCURRENT`
|
|
|
|
|
|
(media, default 4) via `llmClient.withLlmConcurrency(fn, { lane })`.
|
2026-08-31 22:37:28 +07:00
|
|
|
|
- Two Piscina pools run the heavy LLM work off the event loop: a text pool
|
|
|
|
|
|
(`PISCINA_MAX_THREADS`, default 4) and a dedicated media pool
|
2026-09-24 14:09:53 +07:00
|
|
|
|
(`PISCINA_MEDIA_MAX_THREADS`, default 2). A batch is routed by lane to the
|
|
|
|
|
|
matching pool — this keeps a slow image/vision batch from occupying every
|
|
|
|
|
|
thread and blocking unrelated text-only batches behind it. **Each worker
|
|
|
|
|
|
thread (in either pool) initializes its own pg Pool** (min 0, grows to
|
|
|
|
|
|
`POSTGRES_POOL_MAX`). See "Memory & connections" below.
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
## Memory & DB connections
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
`MemoryMax=1G` (raised from 512M — live RSS sits at ~500 MiB, peak 508 MiB,
|
|
|
|
|
|
so 512M left ~2% headroom and risked an OOM-kill restart). Host has 8 GB free.
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-31 22:37:28 +07:00
|
|
|
|
`POSTGRES_POOL_MIN=0` (default). The gateway = main process + up to 4 text
|
|
|
|
|
|
Piscina worker threads + up to 2 media Piscina worker threads, each with its
|
|
|
|
|
|
own pg Pool. With min:0 the pools stay empty until a query runs and drop
|
|
|
|
|
|
idle clients afterward, instead of holding `(1 main + 4 text + 2 media) × 2
|
|
|
|
|
|
= 14` permanently-open idle connections against PgBouncer. The pool still
|
|
|
|
|
|
grows on demand up to `POSTGRES_POOL_MAX`.
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
## Event channels (Redis pub/sub)
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
`discord:message:{created,updated,deleted,analyzed}`,
|
|
|
|
|
|
`discord:attachment:{created,uploaded}`,
|
|
|
|
|
|
`discord:analysis:queue_status`,
|
|
|
|
|
|
`discord:reaction:{added,removed}`,
|
|
|
|
|
|
`discord:thread:{created,deleted,updated}`,
|
|
|
|
|
|
`discord:channel_topic:updated`,
|
|
|
|
|
|
`discord:presence:updated`,
|
|
|
|
|
|
`discord:guild_member:{added,removed}`.
|
|
|
|
|
|
See `src/shared/redis-channels.ts` for the canonical names.
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
## Initialization flow
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-09-24 15:09:13 +07:00
|
|
|
|
`bootstrap.ts` runs these steps in order (each is a named function):
|
|
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
1. Validate env (Zod). Refuse to start if `AI_ANALYSIS_ENABLED` but no key.
|
2026-09-24 15:09:13 +07:00
|
|
|
|
→ `assertConfigIsUsable()`
|
|
|
|
|
|
2. Build long-lived services: Discord client, `RedisEventPublisher` +
|
|
|
|
|
|
`EventBroadcaster`, `CommandHandler`; install the shutdown handler.
|
|
|
|
|
|
3. Connect infrastructure → `connectDatabase()`:
|
|
|
|
|
|
`AUTO_MIGRATE_ON_STARTUP` runs pending Drizzle migrations, then
|
|
|
|
|
|
`initializeDatabase()` (pg Pool, min 0).
|
|
|
|
|
|
4. `registerClientDebugLogging()` — only client debug lines carrying signal.
|
|
|
|
|
|
5. Install process guards (`registerProcessGuards`).
|
|
|
|
|
|
6. Register pipeline gauges + start the metrics server (port `METRICS_PORT`,
|
|
|
|
|
|
default 4016).
|
|
|
|
|
|
7. `client.login(token)`.
|
|
|
|
|
|
|
|
|
|
|
|
On the Discord `ready` event, `lifecycle.ts` runs `startGatewayLifecycle()`:
|
|
|
|
|
|
|
|
|
|
|
|
1. Inject the event broadcaster into message-capture and moderation-actions
|
|
|
|
|
|
(before any listener can fire).
|
|
|
|
|
|
2. Register Discord listeners: message-capture, reaction, thread, presence,
|
|
|
|
|
|
channel-topic, guild-member.
|
|
|
|
|
|
3. Start background work: AI analysis worker + recovery worker, command
|
|
|
|
|
|
handler, retention cleanup, weekly digest.
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
## Graceful shutdown
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-09-24 15:09:13 +07:00
|
|
|
|
`process-guards.ts` owns the policy. `SIGINT`/`SIGTERM` and non-transient
|
|
|
|
|
|
uncaught exceptions/rejections run `shutdown.ts`; transient stream errors
|
|
|
|
|
|
(EPIPE / ECONNRESET / ERR_STREAM_DESTROYED / ERR_STREAM_WRITE_AFTER_END, see
|
|
|
|
|
|
`isTransientStreamError()`) are logged and IGNORED so the bot stays online.
|
|
|
|
|
|
|
|
|
|
|
|
Shutdown order: stop metrics → close event broadcaster (Redis) → close command
|
|
|
|
|
|
handler → close DB → destroy client → exit.
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
## Observability
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
Prometheus scrapes `127.0.0.1:4016/metrics` (`bete_*` prefix). Collectors run
|
|
|
|
|
|
per-scrape and expose: process memory/uptime, and (when AI analysis is on) live
|
2026-09-24 15:09:13 +07:00
|
|
|
|
pipeline gauges registered by `app/metrics-collector.ts` —
|
|
|
|
|
|
`ai_analysis_queued_conversations`, `ai_analysis_active_batch_requests`,
|
|
|
|
|
|
`ai_analysis_active_text_requests`, `ai_analysis_active_media_requests`,
|
|
|
|
|
|
`ai_analysis_active_individual_requests`, `ai_analysis_individual_in_flight`,
|
|
|
|
|
|
`ai_analysis_individual_circuit_breaker_active`,
|
|
|
|
|
|
`ai_analysis_worker_threads_{text,media}`,
|
|
|
|
|
|
`ai_analysis_worker_threads_active_{text,media}`.
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
## Key invariants (do not break)
|
2026-06-01 21:44:29 +07:00
|
|
|
|
|
2026-08-16 08:41:54 +07:00
|
|
|
|
- **LLM is the only judge.** Failed LLM → `status:"error"` + recovery retry.
|
|
|
|
|
|
Never reintroduce regex/heuristic content classification.
|
|
|
|
|
|
- **Discord tokens are sanitized** (`discordTokens.ts`: `<:emoji:id>` →
|
|
|
|
|
|
`[emoji:name]`, `<@id>` → `@user`, etc.) before content reaches the LLM, so
|
|
|
|
|
|
numeric snowflake IDs never trigger false positives.
|
|
|
|
|
|
- **Semantic cache is batched** (one embed call + one Qdrant batch search),
|
|
|
|
|
|
not N sequential round-trips. `ensureQdrantCollection` is memoized.
|
2026-08-28 20:18:48 +07:00
|
|
|
|
- **Streaming is mandatory** against the omniroute base URL (non-stream waits for
|
2026-08-16 08:41:54 +07:00
|
|
|
|
the full body and times out). `llmClient` aggregates SSE chunks.
|