3 Commits
Author SHA1 Message Date
asepharyana c57ee12da1 docs(gateway): consolidate README/ARCHITECTURE, drop stale MODULE_STRUCTURE
README.md was the extraction-era document (referenced winston, mock-crc.ts,
llmModerationClient.ts, indonesianTextNormalizer.ts — all long gone) and
duplicated ARCHITECTURE.md. Rewritten as a short run-the-service guide;
layout/design lives only in ARCHITECTURE.md.

MODULE_STRUCTURE.md deleted: it was a stale duplicate of ARCHITECTURE.md,
referenced by nothing but itself.

ARCHITECTURE.md updated to the post-refactor reality: app/ lifecycle split
(bootstrap/lifecycle/process-guards/metrics-collector), ai-moderation
recovery-worker + cache-prune, per-module index.ts facades, one-way
dependency rule, corrected init/shutdown/observability sections.
2026-09-24 15:09:13 +07:00
asepharyana 494e16b3b3 refactor(gateway): split bootstrap + aiAnalyzer, add module barrels
app/:
- bootstrap.ts 277 -> 145 lines: config guard, DB connect, client debug
  logging and startup order are now named steps with a comment header
- lifecycle.ts (new): everything wired on the Discord 'ready' hook, in
  explicit order (inject broadcaster -> register listeners -> start workers)
- process-guards.ts (new): SIGINT/SIGTERM/uncaughtException/unhandledRejection
  in ONE place, using isTransientStreamError() instead of two duplicated
  inline code lists
- metrics-collector.ts (new): AI pipeline Prometheus gauges

modules/:
- ai-moderation/index.ts + message-capture/index.ts (new): public facades so
  app/ never reaches into internal files
- aiAnalyzer.ts 317 -> 146 lines: pure entry API; skip-verdict recording
  extracted into recordSkip()
- recovery-worker.ts (new): stranded-message recovery + stale lane/CB pruning
- cache-prune.ts (new): 6h expired-verdict sweep, throttled + resettable
- drop 3 dead re-exports (pickBatchWithinBudget/onCircuitBreakerAlert/
  getConversationKey) whose consumers import the origin files directly

No behavior change. typecheck + lint + 138 tests green; nix build OK.
2026-09-24 15:04:53 +07:00
asepharyana 6bf3b40cc7 refactor(gateway): drop shared barrel, migrate to granular imports + shared error helpers
- delete src/shared/index.ts fat barrel; point 10 importers at the exact
  module they use (redis-channels, moderation-types, utils/pagination)
- message-capture/types.ts re-exports from shared/moderation-types directly
- shared/errors: add errorMessage() + isTransientStreamError() helpers,
  replacing the repeated err-message and transient-code checks
- drop unused imports flagged by biome

No behavior change. typecheck + lint + 138 tests green.
2026-09-24 14:58:19 +07:00
25 changed files with 772 additions and 833 deletions
+65 -28
View File
@@ -5,10 +5,9 @@ messages/attachments/reactions/threads/presence, runs LLM-based AI
moderation, and publishes everything to Redis pub/sub for the backend to
consume. The backend serves the HTTP/WS API to the frontend.
> NOTE: this doc is the source of truth for the module layout. The older
> `MODULE_STRUCTURE.md` was stale (referenced `winston`, `mock-crc.ts`,
> `indonesianTextNormalizer.ts`, and `aiAnalysisWorker.ts`/`llmModerationClient.ts`
> which were renamed/merged). If they disagree, this file wins.
> 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.
## Top-level layout
@@ -16,33 +15,41 @@ consume. The backend serves the HTTP/WS API to the frontend.
services/discord-gateway/
├── src/
│ ├── index.ts # Entry point → initializeDiscordGateway()
│ ├── app/
│ │ ├── bootstrap.ts # Wires client, DB, Redis, workers, schedulers
│ │ ├── shutdown.ts # Graceful shutdown (SIGINT/SIGTERM + transient errors)
│ ├── 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
│ │ └── retention.ts # Expired-record cleanup scheduler
│ ├── shared/
│ ├── shared/ # Infrastructure — never imports from modules/
│ │ ├── config/ # Zod-validated env (index.ts = schema+loader)
│ │ ├── database/ # Drizzle ORM + pg Pool + migrations
│ │ │ ├── init.ts drizzle.ts pool.ts migrate.ts migrateCli.ts
│ │ │ └── schema/ # messages, cache, meta, analytics
│ │ ├── logger/ # pino wrapper + createChildLogger()
│ │ ├── errors/ # AppError / ConfigError ...
│ │ ├── errors/ # AppError / ConfigError ... + errorMessage()
│ │ │ # + isTransientStreamError()
│ │ ├── utils/ # retry, pagination
│ │ ├── discord/clientOptions.ts # discord.js-selfbot-v13 client options
│ │ ├── uploader.ts # Shared attachment upload helper
│ │ ├── redis-channels.ts # Redis channel-name constants
│ │ ├── redis-channels.ts # Redis channel + command constants
│ │ └── moderation-types.ts # Shared AI analysis domain types
│ └── modules/
│ └── modules/ # Feature modules, each with an index.ts facade
│ ├── message-capture/ # Discord event listeners + DB store
│ ├── ai-moderation/ # LLM moderation pipeline (see below)
│ ├── attachment-upload/ # Download + (sharp) resize + upload
│ ├── event-broadcaster/ # RedisEventPublisher + EventBroadcaster
│ ├── event-broadcaster/ # RedisEventPublisher + EventBroadcaster
│ ├── command-handler/ # Redis-subscribed backend→gateway commands
│ ├── reaction-tracking/ thread-tracking/ user-presence/
│ ├── channel-topic/ guild-member-events/
│ ├── channel-topic/ guild-member-events/ monitor/
│ └── gateway-metrics/ # Prometheus /metrics endpoint (port 4016)
```
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.
## AI moderation pipeline (`ai-moderation/`)
LLM-only judge — no regex/heuristic classification. One orchestrator call
@@ -54,8 +61,15 @@ 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.
- `aiAnalyzer.ts` — public API: `queueMessageAnalysis`, `getAnalysisQueueStatus`,
`startPendingAIAnalysisWorker` (recovery worker + cache-prune).
- `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.
- `batchScheduler.ts` — per-conversation per-LANE debounce → `processBatch`
(lane-aware). `splitMessagesByLane` / `laneOfMessage` live in
`analysisLanes.ts` (pure, unit-testable).
@@ -123,28 +137,51 @@ See `src/shared/redis-channels.ts` for the canonical names.
## Initialization flow
`bootstrap.ts` runs these steps in order (each is a named function):
1. Validate env (Zod). Refuse to start if `AI_ANALYSIS_ENABLED` but no key.
2. `AUTO_MIGRATE_ON_STARTUP` → run pending Drizzle migrations.
3. `initializeDatabase()` (pg Pool, min 0).
4. Create discord.js-selfbot-v13 client; register listeners on `ready`.
5. Start `gmw-discord-gateway` metrics server (port `METRICS_PORT`, default 4016).
6. `client.login(token)`.
→ `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.
## Graceful shutdown
`SIGINT`/`SIGTERM` (and uncaught transient stream errors: EPIPE / ECONNRESET /
ERR_STREAM_DESTROYED / ERR_STREAM_WRITE_AFTER_END are treated as non-fatal):
stop metrics → close event broadcaster (Redis) → close command handler →
close DB → destroy client → exit.
`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.
## Observability
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
pipeline gauges — `ai_analysis_queued_conversations`,
`ai_analysis_active_batch_requests`, `ai_analysis_active_individual_requests`,
`ai_analysis_individual_in_flight`, `ai_analysis_individual_circuit_breaker_active`,
`ai_analysis_worker_threads`, `ai_analysis_worker_threads_active`.
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}`.
## Key invariants (do not break)
@@ -1,73 +0,0 @@
# Discord Gateway Service — Module Structure
> Kept as a compact module map. For the authoritative layout, design
> decisions, and invariants, see `ARCHITECTURE.md`. This file was rewritten
> on 2026-08-16 to fix stale references (`winston` → pino,
> `mock-crc.ts`/`indonesianTextNormalizer.ts` removed,
> `aiAnalysisWorker.ts` → `ai-analysis-worker.ts`,
> `llmModerationClient.ts` → `llmClient.ts`).
## Top-level
```
services/discord-gateway/
├── src/
│ ├── index.ts # Entry point
│ ├── app/ # bootstrap, shutdown, retention
│ ├── shared/ # config, database, logger, errors, utils, discord, uploader
│ └── modules/
│ ├── message-capture/ # Discord listeners + DB store + metadata
│ ├── ai-moderation/ # LLM moderation pipeline (largest module)
│ ├── attachment-upload/ # Download + sharp resize + upload
│ ├── event-broadcaster/ # RedisEventPublisher + EventBroadcaster
│ ├── command-handler/ # Backend→gateway Redis commands
│ ├── reaction-tracking/ thread-tracking/ user-presence/
│ ├── channel-topic/ guild-member-events/
│ └── gateway-metrics/ # Prometheus /metrics (port 4016)
├── tests/ # Vitest suites (129 tests)
├── drizzle/ # Drizzle migration SQL + journal
├── ARCHITECTURE.md README.md package.json tsconfig.json vitest.config.ts
```
## Module responsibilities (summary)
### message-capture
Captures `messageCreate`/`messageUpdate`/`messageDelete`, extracts metadata,
stores to Postgres, publishes to Redis. Controller–Service–Repository split:
`messageCapture.ts` (listener) → `messageStore.ts` (DB) + `messageMetadata.ts`
(service).
### ai-moderation
LLM-only moderation. Entry: `aiAnalyzer.ts` (`queueMessageAnalysis`,
`startPendingAIAnalysisWorker`, `getAnalysisQueueStatus`). Scheduling:
`batchScheduler.ts` → `batchProcessor.ts` (batch lock + circuit breaker) →
`individualFallbackProcessor.ts` (per-message retry). Heavy work runs in the
Piscina pool via `ai-analysis-worker.ts` (jobs `batch` / `individual`).
Orchestration/caching: `moderationOrchestrator.ts` (exact hash → batched
semantic Qdrant → LLM), `textBatchProcessor.ts` / `mediaBatchProcessor.ts`
(one LLM call per sub-batch), `llmClient.ts` (central streaming client),
`embeddingClient.ts` + `qdrantClient.ts` (semantic cache), plus
`channelCultureStore.ts` / `userProfileStore.ts`.
### attachment-upload
`attachmentUploader.ts` (download → upload to storage) + `imageResizer.ts`
(sharp resize). Emits `discord:attachment:*`.
### event-broadcaster
`RedisEventPublisher` (ioredis publish) + `EventBroadcaster` (typed methods).
Channel names in `src/shared/redis-channels.ts`.
### gateway-metrics
`metrics.ts` Prometheus HTTP server on `METRICS_PORT` (4016). Collectors run
per scrape; live pipeline gauges registered in `bootstrap.ts`.
## Shared infrastructure
- **config** — Zod schema in `shared/config/index.ts` (single source of truth).
- **database** — Drizzle ORM over `pg`; pool `min:0` (`shared/config`).
- **logger** — `pino` wrapper, `createChildLogger()` for context loggers.
- **errors** — `AppError` hierarchy (`ConfigError`, …).
## Notes
- No HTTP server (other than the metrics endpoint). Pure event-driven.
- `MODULE_STRUCTURE.md` is intentionally a sketch; `ARCHITECTURE.md` is the
detailed reference. When they diverge, `ARCHITECTURE.md` wins.
+50 -306
View File
@@ -1,319 +1,63 @@
# Discord Gateway Service - Extraction Complete
# Discord Gateway
## Overview
Event-driven selfbot service: captures Discord events, runs LLM moderation,
publishes everything to Redis for the backend to consume.
Successfully extracted Discord Gateway service with **Modular MVC + Event-Driven Architecture** using Redis pub/sub for inter-service communication.
> Architecture, invariants and the AI pipeline are documented in
> **`ARCHITECTURE.md`** — that file is the source of truth. This README only
> covers how to run it.
## Directory Structure
## Commands
```
services/discord-gateway/
├── src/
│ ├── app/
│ │ ├── bootstrap.ts # Service initialization (Discord client, DB, Redis)
│ │ └── shutdown.ts # Graceful shutdown handler
│ ├── shared/ # Shared infrastructure layer
│ │ ├── config/
│ │ │ └── config.ts # Zod-validated environment config
│ │ ├── database/
│ │ │ ├── schema.ts # Drizzle ORM schema
│ │ │ ├── drizzle.ts # PostgreSQL connection
│ │ │ ├── migrate.ts # Migration runner
│ │ ├── errors/
│ │ │ └── errors.ts # Custom error classes
│ │ ├── logger/
│ │ │ ├── logger.ts # Winston logger wrapper
│ │ │ └── serialization.ts # Log serialization
│ │ ├── utils/
│ │ │ └── retry.ts # Retry with exponential backoff
│ │ └── discord/
│ │ └── clientOptions.ts # Discord.js client config
│ ├── modules/ # Feature modules (Modular MVC)
│ │ ├── message-capture/ # Controller-Service-Repository
│ │ │ ├── messageCapture.ts # Controller: Discord event listeners
│ │ │ ├── messageStore.ts # Repository: DB operations
│ │ │ ├── messageMetadata.ts # Service: Metadata extraction
│ │ │ ├── types.ts # Domain types
│ │ │ └── index.ts # Module exports
│ │ ├── ai-moderation/ # Controller-Service-Repository
│ │ │ ├── aiAnalyzer.ts # Controller: Analysis orchestration
│ │ │ ├── llmModerationClient.ts # Service: LLM API client
│ │ │ ├── aiAnalysisWorker.ts # Service: Worker pool
│ │ │ ├── indonesianTextNormalizer.ts # Service: Text normalization
│ │ │ ├── moderationPrompt.ts # Service: Prompt generation
│ │ │ └── index.ts # Module exports
│ │ ├── attachment-upload/ # Controller-Service-Repository
│ │ │ ├── attachmentUploader.ts # Service: Upload orchestration
│ │ │ ├── imageResizer.ts # Service: Image resizing
│ │ │ └── index.ts # Module exports
│ │ └── event-broadcaster/ # Event-driven layer
│ │ ├── eventBroadcaster.ts # Service: Redis pub/sub publisher
│ │ ├── eventTypes.ts # Domain: Event type definitions
│ │ └── index.ts # Module exports
│ ├── mock-crc.ts # CRC polyfill for discord.js
│ └── index.ts # Service entry point
├── ARCHITECTURE.md # Detailed architecture documentation
├── package.json # Service dependencies
└── tsconfig.json # TypeScript configuration (inherited)
```bash
pnpm install
pnpm typecheck # tsc --noEmit
pnpm lint # biome check --diagnostic-level=error .
pnpm test # vitest run (138 tests)
pnpm build # tsc — CI/prod builds run this inside nix, which also
# runs scripts/fix-imports.mjs to rewrite @/ aliases and
# extensionless imports for Node ESM
pnpm dev # tsx watch src/index.ts
pnpm start # node dist/index.js
```
## Architecture Patterns
Deployment is CI-only: `nix build .#discord-gateway` → Attic cache → systemd
restart on the VPS. Do not build/hand-copy the artifact.
### 1. Modular MVC Structure
Each feature module follows **Controller-Service-Repository** pattern:
**Message Capture Module**:
- **Controller** (`messageCapture.ts`): Listens to Discord events (messageCreate, messageUpdate, messageDelete)
- **Service** (`messageMetadata.ts`): Extracts and normalizes message metadata
- **Repository** (`messageStore.ts`): Database CRUD operations
**AI Moderation Module**:
- **Controller** (`aiAnalyzer.ts`): Orchestrates analysis workflow
- **Service** (`llmModerationClient.ts`): LLM API integration
- **Service** (`aiAnalysisWorker.ts`): Worker pool management
- **Service** (`indonesianTextNormalizer.ts`): Text preprocessing
**Attachment Upload Module**:
- **Service** (`attachmentUploader.ts`): Upload orchestration
- **Service** (`imageResizer.ts`): Image processing
### 2. Event-Driven Architecture
**Redis Pub/Sub** replaces WebSocket broadcaster:
## Layout
```
Discord Events → Discord Gateway Service → Redis Pub/Sub → Backend Service
↓
Event Channels:
- discord:message:created
- discord:message:updated
- discord:message:deleted
- discord:message:analyzed
- discord:attachment:created
- discord:attachment:uploaded
- discord:analysis:queue_status
src/
├── index.ts # entry → initializeDiscordGateway()
├── 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
│ └── retention.ts # expired-record cleanup scheduler
├── shared/ # infrastructure — never imports from modules/
│ ├── config/ database/ logger/ errors/ utils/
│ ├── discord/clientOptions.ts
│ ├── redis-channels.ts # canonical Redis channel + command constants
│ └── moderation-types.ts # domain types shared across services
└── modules/ # feature modules (each exposes an index.ts facade)
├── ai-moderation/ # LLM moderation pipeline (largest module)
├── message-capture/ # Discord listeners + message/attachment DB
├── attachment-upload/ # download → resize → upload
├── event-broadcaster/ # Redis pub/sub publisher
├── command-handler/ # backend → gateway commands over Redis
├── gateway-metrics/ # Prometheus /metrics (port 4016)
├── monitor/ # weekly digest scheduler
└── reaction-tracking/ thread-tracking/ user-presence/
channel-topic/ guild-member-events/
```
### 3. Shared Infrastructure Layer
Centralized, reusable components:
- **Config**: Zod-validated environment variables
- **Logger**: Winston logger with context support
- **Database**: Drizzle ORM with PostgreSQL
- **Errors**: Custom error classes with codes and HTTP status codes
- **Utils**: Retry logic with exponential backoff
- **Discord**: Client configuration and options
Dependency direction is one-way: `index.ts` → `app/` → `modules/` → `shared/`.
Callers outside a module import its `index.ts` facade, never an internal file.
### 4. No HTTP Server
- **Event-driven only**: No Express, WebSocket, or HTTP routes
- **Redis pub/sub**: All inter-service communication via Redis
- **Backend service**: Consumes events and serves HTTP API
- **Frontend**: Continues to use Backend HTTP API
## Testing
## Key Features
### Message Capture
1. Discord emits `messageCreate`, `messageUpdate`, `messageDelete` events
2. `messageCapture.ts` listener receives and validates event
3. Extract metadata: user, channel, content, timestamp, attachments
4. `messageStore.ts` inserts into PostgreSQL
5. `eventBroadcaster.messageCreated()` publishes to Redis
6. Backend service subscribes and processes
### AI Moderation
1. `aiAnalyzer.ts` queues messages for analysis
2. `llmModerationClient.ts` calls LLM API with context
3. `indonesianTextNormalizer.ts` preprocesses text
4. Results stored in database
5. `eventBroadcaster.messageAnalyzed()` publishes results
6. Backend service receives and updates UI
### Attachment Upload
1. `messageCapture.ts` detects attachments
2. `attachmentUploader.ts` downloads from Discord
3. `imageResizer.ts` resizes images if needed
4. Upload to external storage with retry logic
5. `eventBroadcaster.attachmentUploaded()` publishes
6. Backend service stores metadata
## Initialization Flow
```
1. Load environment config (Zod validation)
↓
2. Initialize PostgreSQL connection
↓
3. Run pending database migrations
↓
4. Create Discord client with optimized cache
↓
5. Initialize Redis event broadcaster
↓
6. Register Discord event listeners
- messageCapture (message events)
- aiAnalyzer (analysis worker)
↓
7. Login to Discord
↓
8. Listen for graceful shutdown signals
```
## Graceful Shutdown
On SIGINT/SIGTERM/uncaughtException/unhandledRejection:
1. Close PostgreSQL connection
2. Close Redis connection
3. Destroy Discord client
4. Exit process (code 0 for clean, 1 for error)
## Dependencies
**Core Discord**:
- `discord.js-selfbot-v13` — Discord client (selfbot variant)
**Media Processing**:
- `sharp` — Image resizing
**Data & Config**:
- `drizzle-orm` — Type-safe ORM
- `pg` — PostgreSQL driver
- `zod` — Config validation
- `ioredis` — Redis client
**Logging & Utilities**:
- `winston` — Structured logging
- `p-retry` — Retry with backoff
- `p-limit` — Concurrency limiting
- `piscina` — Worker pool
## No Breaking Changes
- Original `src/` remains untouched
- Discord Gateway is a **new service** in `services/discord-gateway/`
- Can run alongside existing monolith during transition
- Backend service will consume Redis events
- Frontend continues to use Backend HTTP API
## Next Steps
1. **Create Backend service** (`services/backend/`)
- HTTP API endpoints
- Redis event subscribers
- Database models
- WebSocket broadcaster
2. **Update Frontend** (`frontend/`)
- Connect to Backend HTTP API
- Subscribe to WebSocket events
3. **Nix & CI/CD**
- flake.nix package for Discord Gateway
- systemd services (gmw-backend, gmw-discord-gateway)
- GitHub Actions for build/deploy (nix copy → systemctl restart)
4. **Documentation**
- API documentation
- Event schema documentation
- Deployment guide
## Files Created
**Total: 43 files**
### Shared Infrastructure (9 files)
- `src/shared/config/config.ts`
- `src/shared/database/` (5 files)
- `@bete/shared/errors` (shared package)
- `src/shared/logger/logger.ts`
- `src/shared/logger/serialization.ts`
- `src/shared/utils/retry.ts`
- `src/shared/discord/clientOptions.ts`
### Modules (28 files)
- `src/modules/message-capture/` (5 files)
- `src/modules/ai-moderation/` (6 files)
- `src/modules/attachment-upload/` (3 files)
- `src/modules/event-broadcaster/` (3 files)
### App & Entry (4 files)
- `src/app/bootstrap.ts`
- `src/app/shutdown.ts`
- `src/index.ts`
- `src/mock-crc.ts`
### Configuration (2 files)
- `package.json`
- `ARCHITECTURE.md`
## Verification Checklist
✅ Directory structure created
✅ Shared infrastructure migrated
✅ Message capture module migrated
✅ AI moderation module migrated
✅ Attachment upload module migrated
✅ Event broadcaster module created (Redis pub/sub)
✅ Bootstrap and entry point created
✅ Package.json with dependencies
✅ No HTTP server code (Express, WebSocket removed)
✅ Event-driven architecture implemented
✅ Graceful shutdown handler
✅ Module index files for clean exports
✅ Architecture documentation
## Event Flow Diagram
```
┌─────────────────────────────────────────────────────────────────┐
│ Discord Gateway Service │
├─────────────────────────────────────────────────────────────────┤
│ │
│ ┌────────────────────────────┐ ┌────────────────────────────┐ │
│ │ Message Capture │ │ AI Moderation │ │
│ │ (Controller) │ │ (Controller) │ │
│ └──────────────┬─────────────┘ └──────────────┬─────────────┘ │
│ │ │ │
│ ├───────────────────────────────┤ │
│ │ │ │
│ ▼ ▼ │
│ ┌───────────────────────────────────────────────────────────┐ │
│ │ Event Broadcaster (Redis Pub/Sub) │ │
│ │ - discord:message:created │ │
│ │ - discord:message:updated │ │
│ │ - discord:message:deleted │ │
│ │ - discord:message:analyzed │ │
│ │ - discord:attachment:created │ │
│ │ - discord:attachment:uploaded │ │
│ │ - discord:analysis:queue_status │ │
│ └───────────────────────────────────────────────────────────┘ │
│ │ │
└────────────────────────────────┼────────────────────────────────┘
│
│ Redis Pub/Sub
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ Backend Service │
│ (Subscribes to events, serves HTTP API, manages WebSocket) │
└─────────────────────────────────────────────────────────────────┘
│
│ HTTP API
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ Frontend Application │
│ (React SPA, real-time updates via WebSocket) │
└─────────────────────────────────────────────────────────────────┘
```
## Summary
The Discord Gateway service has been successfully extracted with:
- **Modular MVC architecture** for clean separation of concerns
- **Event-driven design** using Redis pub/sub for inter-service communication
- **Shared infrastructure layer** for reusable components
- **No HTTP server** — pure event-driven service
- **Graceful shutdown** handling
- **Type-safe configuration** with Zod validation
- **Structured logging** with Winston
- **PostgreSQL integration** with Drizzle ORM
The service is ready for integration with the Backend service, which will consume Redis events and serve the HTTP API to the Frontend.
Vitest, tests in `tests/`. Config supplies dummy env vars so the suite runs
without live Postgres/Redis/Qdrant; external services are mocked. `llmE2e.test.ts`
is skipped by default and needs real credentials (`pnpm test:e2e:live`).
+71 -203
View File
@@ -1,36 +1,19 @@
import { Client } from "discord.js-selfbot-v13";
import { ConfigError, DatabaseError } from "@/shared/errors/index";
import { createChildLogger } from "@/shared/logger/index";
import {
getAnalysisQueueStatus,
startPendingAIAnalysisWorker,
} from "../modules/ai-moderation/aiAnalyzer.js";
import {
mediaWorkerPool,
textWorkerPool,
} from "../modules/ai-moderation/circuitBreaker.js";
import { registerChannelTopicCapture } from "../modules/channel-topic/index.js";
ConfigError,
DatabaseError,
errorMessage,
} from "@/shared/errors/index.js";
import { createChildLogger } from "@/shared/logger/index.js";
import { CommandHandler } from "../modules/command-handler/commandHandler.js";
import {
EventBroadcaster,
RedisEventPublisher,
} from "../modules/event-broadcaster/index.js";
import {
registerCollector,
setGauge,
startMetricsServer,
stopMetricsServer,
} from "../modules/gateway-metrics/index.js";
import { registerGuildMemberEvents } from "../modules/guild-member-events/index.js";
import {
registerMessageCapture,
setEventBroadcaster as setMessageCaptureEventBroadcaster,
} from "../modules/message-capture/messageCapture.js";
import { setModerationEventBroadcaster } from "../modules/message-capture/moderationActionsDb.js";
import { startDigestScheduler } from "../modules/monitor/digestScheduler.js";
import { registerReactionCapture } from "../modules/reaction-tracking/index.js";
import { registerThreadCapture } from "../modules/thread-tracking/index.js";
import { registerPresenceCapture } from "../modules/user-presence/index.js";
import { config } from "../shared/config/index.js";
import {
closeDatabase,
@@ -38,45 +21,34 @@ import {
} from "../shared/database/drizzle.js";
import { runMigrations } from "../shared/database/migrate.js";
import { createDiscordClientOptions } from "../shared/discord/clientOptions.js";
import { startRetentionCleanup } from "./retention.js";
import { startGatewayLifecycle } from "./lifecycle.js";
import { registerPipelineMetrics } from "./metrics-collector.js";
import { registerProcessGuards } from "./process-guards.js";
import { createGracefulShutdown } from "./shutdown.js";
const logger = createChildLogger("discord-gateway");
// ─── Bootstrap ─────────────────────────────────────────────────────────────
//
// Startup order:
// 1. validate config (fail fast on missing AI credentials)
// 2. connect infrastructure (migrations → DB pool)
// 3. build long-lived services (Discord client, Redis publisher, command
// handler) + install shutdown/process guards
// 4. start observability (pipeline gauges → metrics server)
// 5. log in (ready-hook wires listeners via lifecycle.ts)
export async function initializeDiscordGateway() {
/** Refuse to start when AI analysis is on but no LLM credentials exist. */
function assertConfigIsUsable(): void {
if (config.AI_ANALYSIS_ENABLED && !config.AI_LLM_API_KEY) {
throw new ConfigError(
"AI_ANALYSIS_ENABLED=true but AI_LLM_API_KEY is missing from environment. AI analysis cannot run without credentials.",
);
}
}
const token = config.DISCORD_TOKEN;
logger.info(
{ hasToken: token.length > 0, tokenLength: token.length },
"Config loaded",
);
logger.info("Creating Discord client");
const client = new Client(createDiscordClientOptions());
// Initialize Redis event broadcaster
const redisPublisher = new RedisEventPublisher(config.REDIS_URL, logger);
const eventBroadcaster = new EventBroadcaster(redisPublisher);
// Initialize Redis command handler for backend→gateway commands
const commandHandler = new CommandHandler();
const gracefulShutdown = createGracefulShutdown({
logger,
closeDatabase,
client,
eventBroadcaster,
commandHandler,
stopMetricsServer,
});
/** Run migrations (when enabled) then open the PostgreSQL pool. */
async function connectDatabase(): Promise<void> {
try {
if (config.AUTO_MIGRATE_ON_STARTUP) {
logger.info(
@@ -90,183 +62,79 @@ export async function initializeDiscordGateway() {
logger.info("PostgreSQL database initialized");
} catch (err) {
logger.error(
{ err, errorMsg: err instanceof Error ? err.message : String(err) },
{ err, errorMsg: errorMessage(err) },
"Failed to initialize database",
);
throw new DatabaseError(
`Database initialization failed: ${err instanceof Error ? err.message : String(err)}`,
`Database initialization failed: ${errorMessage(err)}`,
);
}
}
/** Log only client debug lines that carry signal (errors/streams, or VERBOSE). */
function registerClientDebugLogging(client: Client): void {
client.on("debug", (msg) => {
if (
msg.toLowerCase().includes("error") ||
msg.toLowerCase().includes("stream")
) {
const lower = msg.toLowerCase();
if (lower.includes("error") || lower.includes("stream")) {
logger.info({ debugMsg: msg }, "Discord Client Debug");
} else if (config.VERBOSE) {
logger.debug({ debugMsg: msg }, "Discord Client Debug");
}
});
}
client.on("ready", async () => {
export async function initializeDiscordGateway() {
assertConfigIsUsable();
const token = config.DISCORD_TOKEN;
logger.info(
{ hasToken: token.length > 0, tokenLength: token.length },
"Config loaded",
);
logger.info("Creating Discord client");
const client = new Client(createDiscordClientOptions());
// Long-lived services: Redis event broadcaster (gateway → backend) and the
// Redis command handler (backend → gateway).
const redisPublisher = new RedisEventPublisher(config.REDIS_URL, logger);
const eventBroadcaster = new EventBroadcaster(redisPublisher);
const commandHandler = new CommandHandler();
const gracefulShutdown = createGracefulShutdown({
logger,
closeDatabase,
client,
eventBroadcaster,
commandHandler,
stopMetricsServer,
});
await connectDatabase();
registerClientDebugLogging(client);
client.on("ready", () => {
logger.info({ user: client.user?.tag }, "Bot logged in");
setMessageCaptureEventBroadcaster(eventBroadcaster);
setModerationEventBroadcaster(eventBroadcaster);
registerMessageCapture(client);
startPendingAIAnalysisWorker(client, eventBroadcaster);
// Register new event captures
registerReactionCapture(client, eventBroadcaster);
registerThreadCapture(client, eventBroadcaster);
registerPresenceCapture(client, eventBroadcaster);
registerChannelTopicCapture(client, eventBroadcaster);
registerGuildMemberEvents(client, eventBroadcaster);
// Start command handler after Discord is ready
commandHandler.start(client);
logger.info("Command handler started");
// Start retention cleanup scheduler
startRetentionCleanup();
// Start weekly moderation digest (public, automated)
startDigestScheduler();
startGatewayLifecycle({
client,
eventBroadcaster,
commandHandler,
logger,
});
});
client.on("error", (err) => {
logger.error(
{ err, errorMsg: err instanceof Error ? err.message : String(err) },
"Client error",
);
logger.error({ err, errorMsg: errorMessage(err) }, "Client error");
});
process.on("SIGINT", () => {
gracefulShutdown("SIGINT");
});
registerProcessGuards(logger, gracefulShutdown);
process.on("SIGTERM", () => {
gracefulShutdown("SIGTERM");
});
process.on("uncaughtException", (err) => {
const code =
typeof (err as NodeJS.ErrnoException).code === "string"
? (err as NodeJS.ErrnoException).code
: "";
// Transient stream-teardown errors (voice stop/disconnect races, child
// process stdin closed while we still write) are NOT fatal — crashing the
// gateway on EPIPE takes the whole bot offline mid-music. Log + continue.
if (
code === "EPIPE" ||
code === "ERR_STREAM_DESTROYED" ||
code === "ERR_STREAM_WRITE_AFTER_END" ||
code === "ECONNRESET"
) {
logger.warn(
{ error: err },
"Uncaught transient stream error — continuing",
);
return;
}
logger.error(
{
err,
errorMsg: err instanceof Error ? err.message : String(err),
stack: err?.stack,
},
"Uncaught exception",
);
gracefulShutdown("uncaughtException");
});
process.on("unhandledRejection", (reason) => {
const err =
reason instanceof Error ? reason : new Error(String(reason ?? "unknown"));
const code = (err as NodeJS.ErrnoException).code ?? "";
// Same transient-teardown policy as uncaughtException: a rejection that
// fires while a stream is being torn down (EPIPE after ffmpeg stdin
// closes, write-after-destroy, socket reset) must NOT take the whole
// gateway offline. Log detail + continue. Everything else still shuts
// down so real bugs surface.
if (
code === "EPIPE" ||
code === "ERR_STREAM_DESTROYED" ||
code === "ERR_STREAM_WRITE_AFTER_END" ||
code === "ECONNRESET"
) {
logger.warn(
{ error: err },
"Unhandled rejection transient stream error — continuing",
);
return;
}
logger.error({ error: err, reason: String(reason) }, "Unhandled rejection");
gracefulShutdown("unhandledRejection");
});
// ── Metrics: register live pipeline collectors before starting server ──
// These refresh on every scrape so Prometheus sees real AI-analysis
// queue depth, concurrency, and DB pool state instead of an empty stub.
registerCollector(() => {
if (!config.AI_ANALYSIS_ENABLED) return;
try {
const status = getAnalysisQueueStatus();
setGauge("ai_analysis_queued_conversations", status.queuedConversations);
setGauge("ai_analysis_active_batch_requests", status.activeRequests);
setGauge(
"ai_analysis_active_text_requests",
status.activeTextRequests ?? status.activeRequests,
);
setGauge(
"ai_analysis_active_media_requests",
status.activeMediaRequests ?? 0,
);
setGauge(
"ai_analysis_active_individual_requests",
status.activeIndividualRequests,
);
setGauge(
"ai_analysis_individual_in_flight",
status.individualInFlightCount,
);
setGauge(
"ai_analysis_individual_circuit_breaker_active",
status.individualCircuitBreakerActive ? 1 : 0,
);
if (typeof status.lastError === "string") {
setGauge("ai_analysis_last_error_present", status.lastError ? 1 : 0);
}
type PoolState = { _poolState?: { size: number; active: number } };
const textPool = textWorkerPool as unknown as PoolState;
const mediaPool = mediaWorkerPool as unknown as PoolState;
// Reported per queue (2026-08-31 text/media pool split) so the text
// and media backlogs are distinguishable in dashboards/alerts instead
// of one combined "worker threads" number.
if (textPool._poolState) {
setGauge("ai_analysis_worker_threads_text", textPool._poolState.size);
setGauge(
"ai_analysis_worker_threads_active_text",
textPool._poolState.active,
);
}
if (mediaPool._poolState) {
setGauge("ai_analysis_worker_threads_media", mediaPool._poolState.size);
setGauge(
"ai_analysis_worker_threads_active_media",
mediaPool._poolState.active,
);
}
} catch (err) {
logger.warn({ error: String(err) }, "AI metrics collector failed");
}
});
// Start metrics server
// Metrics: register live pipeline collectors before starting the server.
registerPipelineMetrics(logger);
startMetricsServer();
logger.info("Calling Discord client.login");
// Fix: use await + try/catch instead of .then().catch()
try {
await client.login(token);
logger.info("Discord client logged in successfully");
@@ -0,0 +1,61 @@
import type { Client } from "discord.js-selfbot-v13";
import type { Logger } from "@/shared/logger/index.js";
import { startPendingAIAnalysisWorker } from "../modules/ai-moderation/index.js";
import { registerChannelTopicCapture } from "../modules/channel-topic/index.js";
import type { CommandHandler } from "../modules/command-handler/commandHandler.js";
import type { EventBroadcaster } from "../modules/event-broadcaster/index.js";
import { registerGuildMemberEvents } from "../modules/guild-member-events/index.js";
import {
registerMessageCapture,
setEventBroadcaster as setMessageCaptureEventBroadcaster,
setModerationEventBroadcaster,
} from "../modules/message-capture/index.js";
import { startDigestScheduler } from "../modules/monitor/digestScheduler.js";
import { registerReactionCapture } from "../modules/reaction-tracking/index.js";
import { registerThreadCapture } from "../modules/thread-tracking/index.js";
import { registerPresenceCapture } from "../modules/user-presence/index.js";
import { startRetentionCleanup } from "./retention.js";
export interface GatewayLifecycleOptions {
client: Client;
eventBroadcaster: EventBroadcaster;
commandHandler: CommandHandler;
logger: Logger;
}
/**
* Wires everything that must start once Discord is connected.
*
* Ordering matters:
* 1. Inject the event broadcaster into the modules that publish events —
* they must be able to publish before their listeners are registered.
* 2. Register the Discord event listeners (capture modules).
* 3. Start the background workers/schedulers.
*/
export function startGatewayLifecycle({
client,
eventBroadcaster,
commandHandler,
logger,
}: GatewayLifecycleOptions): void {
// 1. Inject broadcaster first so no captured event is dropped.
setMessageCaptureEventBroadcaster(eventBroadcaster);
setModerationEventBroadcaster(eventBroadcaster);
// 2. Discord event listeners.
registerMessageCapture(client);
registerReactionCapture(client, eventBroadcaster);
registerThreadCapture(client, eventBroadcaster);
registerPresenceCapture(client, eventBroadcaster);
registerChannelTopicCapture(client, eventBroadcaster);
registerGuildMemberEvents(client, eventBroadcaster);
// 3. Background workers + schedulers.
startPendingAIAnalysisWorker(client, eventBroadcaster);
commandHandler.start(client);
logger.info("Command handler started");
startRetentionCleanup();
// Weekly moderation digest (public, automated)
startDigestScheduler();
}
@@ -0,0 +1,77 @@
import type { Logger } from "@/shared/logger/index.js";
import {
getAnalysisQueueStatus,
mediaWorkerPool,
textWorkerPool,
} from "../modules/ai-moderation/index.js";
import {
registerCollector,
setGauge,
} from "../modules/gateway-metrics/index.js";
import { config } from "../shared/config/index.js";
/** Piscina exposes its live thread counters on `_poolState`. */
type PoolState = { _poolState?: { size: number; active: number } };
/**
* Registers the AI-pipeline Prometheus gauges.
*
* The collector refreshes on every scrape, so Prometheus sees real queue
* depth / concurrency / worker-thread state instead of an empty stub.
* Registered before the metrics server starts.
*/
export function registerPipelineMetrics(logger: Logger): void {
registerCollector(() => {
if (!config.AI_ANALYSIS_ENABLED) return;
try {
const status = getAnalysisQueueStatus();
setGauge("ai_analysis_queued_conversations", status.queuedConversations);
setGauge("ai_analysis_active_batch_requests", status.activeRequests);
setGauge(
"ai_analysis_active_text_requests",
status.activeTextRequests ?? status.activeRequests,
);
setGauge(
"ai_analysis_active_media_requests",
status.activeMediaRequests ?? 0,
);
setGauge(
"ai_analysis_active_individual_requests",
status.activeIndividualRequests,
);
setGauge(
"ai_analysis_individual_in_flight",
status.individualInFlightCount,
);
setGauge(
"ai_analysis_individual_circuit_breaker_active",
status.individualCircuitBreakerActive ? 1 : 0,
);
if (typeof status.lastError === "string") {
setGauge("ai_analysis_last_error_present", status.lastError ? 1 : 0);
}
// Reported per queue (2026-08-31 text/media pool split) so the text
// and media backlogs are distinguishable in dashboards/alerts instead
// of one combined "worker threads" number.
const textPool = textWorkerPool as unknown as PoolState;
const mediaPool = mediaWorkerPool as unknown as PoolState;
if (textPool._poolState) {
setGauge("ai_analysis_worker_threads_text", textPool._poolState.size);
setGauge(
"ai_analysis_worker_threads_active_text",
textPool._poolState.active,
);
}
if (mediaPool._poolState) {
setGauge("ai_analysis_worker_threads_media", mediaPool._poolState.size);
setGauge(
"ai_analysis_worker_threads_active_media",
mediaPool._poolState.active,
);
}
} catch (err) {
logger.warn({ error: String(err) }, "AI metrics collector failed");
}
});
}
@@ -0,0 +1,61 @@
import { errorMessage, isTransientStreamError } from "@/shared/errors/index.js";
import type { Logger } from "@/shared/logger/index.js";
import type { GracefulShutdown } from "./shutdown.js";
/**
* Process-level signal + error guards.
*
* Extracted from bootstrap so the "what keeps the gateway alive vs what
* shuts it down" policy lives in exactly one place.
*
* Policy: transient stream-teardown failures (EPIPE / ERR_STREAM_DESTROYED /
* ERR_STREAM_WRITE_AFTER_END / ECONNRESET) are logged and IGNORED — crashing
* the gateway on them (voice stop/disconnect races, a child process stdin
* closed while we still write) takes the whole bot offline mid-operation.
* Anything else is a real bug: log with stack and shut down cleanly.
*/
export function registerProcessGuards(
logger: Logger,
gracefulShutdown: GracefulShutdown,
): void {
process.on("SIGINT", () => {
gracefulShutdown("SIGINT");
});
process.on("SIGTERM", () => {
gracefulShutdown("SIGTERM");
});
process.on("uncaughtException", (err) => {
if (isTransientStreamError(err)) {
logger.warn(
{ error: err },
"Uncaught transient stream error — continuing",
);
return;
}
logger.error(
{
err,
errorMsg: errorMessage(err),
stack: err?.stack,
},
"Uncaught exception",
);
gracefulShutdown("uncaughtException");
});
process.on("unhandledRejection", (reason) => {
const err =
reason instanceof Error ? reason : new Error(String(reason ?? "unknown"));
if (isTransientStreamError(err)) {
logger.warn(
{ error: err },
"Unhandled rejection transient stream error — continuing",
);
return;
}
logger.error({ error: err, reason: String(reason) }, "Unhandled rejection");
gracefulShutdown("unhandledRejection");
});
}
@@ -1,8 +1,5 @@
import { inArray, lt } from "drizzle-orm";
import type {
NodePgDatabase,
NodePgQueryResultHKT,
} from "drizzle-orm/node-postgres";
import type { NodePgDatabase } from "drizzle-orm/node-postgres";
import { createChildLogger } from "@/shared/logger/index";
import { config } from "../shared/config/index.js";
import { getDatabase } from "../shared/database/drizzle.js";
@@ -12,28 +12,14 @@ import {
buildSkipAnalysisUserResult,
isAgeRestrictedMessage,
isSkipAnalysisUser,
skipAgeRestrictedMessages,
skipAnalysisUserMessages,
} from "./batchProcessor.js";
import { scheduleConversationAnalysis } from "./batchScheduler.js";
import { getConversationKey } from "./circuitBreaker.js";
import {
ANALYSIS_LANES,
type AnalysisLane,
clearConversationProcessing,
conversationConsecutiveErrors,
conversationDebounceTimers,
conversationErrorCooldown,
conversationProcessing,
isConversationProcessingLocked,
} from "./conversationState.js";
import { conversationDebounceTimers } from "./conversationState.js";
import {
activeIndividualRequests,
enqueueIndividualFallbacks,
individualCooldownUntil,
individualInFlight,
individualInFlightByConversation,
individualInFlightLastTouched,
} from "./individualFallbackProcessor.js";
import {
broadcastAnalysisCompleted,
@@ -41,31 +27,20 @@ import {
setModerationClient,
setSharedEventBroadcaster,
} from "./moderationState.js";
import { deleteExpiredQdrantPoints } from "./qdrantClient.js";
import { pruneExpiredTexts } from "./textCacheStore.js";
import { startRecoveryWorker } from "./recovery-worker.js";
const logger = createChildLogger("ai-analyzer");
// ---------------------------------------------------------------------------
// Cache hygiene (expired verdict sweep)
// ---------------------------------------------------------------------------
const CACHE_PRUNE_INTERVAL_MS = 6 * 60 * 60 * 1000; // every 6 hours
let lastCachePruneAt = 0;
// ---------------------------------------------------------------------------
// Re-exports from sub-modules (preserving original public API)
// ---------------------------------------------------------------------------
export { pickBatchWithinBudget } from "./batchProcessor.js";
export { getConversationKey } from "./circuitBreaker.js";
export { onCircuitBreakerAlert } from "./conversationState.js";
// ---------------------------------------------------------------------------
// Public API
// Public API — queueing, status, worker startup
// ---------------------------------------------------------------------------
/**
* Queues a message for analysis (debounced by conversation).
*
* Messages that never need an LLM call are short-circuited here and recorded
* with their skip verdict: age-restricted messages and configured skip-list
* users.
*/
export async function queueMessageAnalysis(messageId: string): Promise<void> {
if (!config.AI_ANALYSIS_ENABLED) return;
@@ -78,13 +53,7 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
}
if (isAgeRestrictedMessage(message)) {
const updated = await messageStore.updateMessageAIAnalysis(
message.id,
buildAgeRestrictedSkipResult(),
);
if (updated) {
broadcastAnalysisCompleted(updated);
}
await recordSkip(message.id, buildAgeRestrictedSkipResult());
logger.debug(
{ messageId },
"Skipped AI analysis for age-restricted message",
@@ -93,13 +62,7 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
}
if (isSkipAnalysisUser(message)) {
const updated = await messageStore.updateMessageAIAnalysis(
message.id,
buildSkipAnalysisUserResult(),
);
if (updated) {
broadcastAnalysisCompleted(updated);
}
await recordSkip(message.id, buildSkipAnalysisUserResult());
logger.debug(
{ messageId, userId: message.user_id },
"Skipped AI analysis for configured skip-list user",
@@ -119,6 +82,17 @@ export async function queueMessageAnalysis(messageId: string): Promise<void> {
}
}
/** Persist a skip verdict and broadcast it so the dashboard reflects it. */
async function recordSkip(
messageId: string,
result: Parameters<typeof messageStore.updateMessageAIAnalysis>[1],
): Promise<void> {
const updated = await messageStore.updateMessageAIAnalysis(messageId, result);
if (updated) {
broadcastAnalysisCompleted(updated);
}
}
/**
* Queues a conversation for analysis (debounced).
*/
@@ -144,11 +118,12 @@ export function getAnalysisQueueStatus(): AnalysisQueueStatus {
}
/**
* Starts the periodic recovery worker.
* Starts the background workers behind the analysis pipeline:
* - the recovery worker (stranded pending / incomplete messages + cache prune)
* - the optional culture and user-profile learners.
*
* Now also recovers messages stuck in `error/analysis_incomplete`
* state (not just `pending`), and skips conversations that already have
* individual fallback work in progress to avoid DB last-write-wins races.
* Also injects the Discord client and event broadcaster into the pipeline
* state so downstream modules can act and publish.
*/
export function startPendingAIAnalysisWorker(
client?: Client,
@@ -167,151 +142,5 @@ export function startPendingAIAnalysisWorker(
.catch(console.error);
}
setInterval(() => {
// [D] Periodic cache hygiene: purge expired moderation verdicts from
// Postgres and Qdrant. Expired entries are never reused (filters check
// expires_at) but accumulate forever without this sweep.
const now = Date.now();
if (now - lastCachePruneAt >= CACHE_PRUNE_INTERVAL_MS) {
lastCachePruneAt = now;
Promise.all([pruneExpiredTexts(), deleteExpiredQdrantPoints()])
.then(([pgDeleted, qdDeleted]) => {
if (pgDeleted > 0 || qdDeleted > 0) {
logger.info(
{ pgDeleted, qdDeleted },
"Expired moderation cache pruned",
);
}
})
.catch((err: unknown) => {
logger.warn({ error: String(err) }, "Moderation cache prune failed");
});
}
// Only revert stuck processing messages if there's active processing.
// Avoids a DB query every recovery interval when the pipeline is idle.
if (conversationProcessing.size > 0) {
messageStore
.revertStuckProcessingMessages(300000)
.catch((err: unknown) => {
logger.error(
{ error: String(err) },
"Failed to run stuck processing recovery",
);
});
}
Promise.all([
messageStore.getPendingConversationKeys(500),
messageStore.getConversationKeysWithIncompleteAnalysis(200),
])
.then(([pendingKeys, incompleteKeys]) => {
const now = Date.now();
for (const [key, expiry] of conversationErrorCooldown) {
if (now >= expiry) conversationErrorCooldown.delete(key);
}
// conversationProcessing is now a Partial<Record<lane, startedAt>>.
// Prune stale lane slots individually so one stale lane never clears
// the other lane's healthy lock.
for (const [key, record] of conversationProcessing) {
for (const lane of ANALYSIS_LANES as readonly AnalysisLane[]) {
const startedAt = record?.[lane];
if (
startedAt &&
now - startedAt >= config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS
) {
clearConversationProcessing(key, lane);
}
}
}
const staleThreshold = config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS * 2;
for (const [key, lastTouched] of individualInFlightLastTouched) {
if (now - lastTouched >= staleThreshold) {
individualInFlightLastTouched.delete(key);
individualInFlightByConversation.delete(key);
logger.warn(
{ key },
"Pruned stale individualInFlightByConversation entry",
);
}
}
// Also prune stale per-conversation CB error counts that have cooled
// down so old conversations can be retried.
for (const [key] of conversationConsecutiveErrors) {
const cbExpire = conversationErrorCooldown.get(key) ?? 0;
if (cbExpire && now >= cbExpire) {
conversationConsecutiveErrors.delete(key);
}
}
const incompleteKeySet = new Set(incompleteKeys);
// --- Batch recovery for pending messages ---
for (const key of pendingKeys) {
if (
ANALYSIS_LANES.some((lane) =>
conversationDebounceTimers.has(`${key}::${lane}`),
)
) {
continue;
}
// Batch recovery must not race ANY in-flight batch lane, so the
// lock check is lane-agnostic here (individual fallback handles
// error rows separately).
if (isConversationProcessingLocked(key)) continue;
if (individualInFlightByConversation.has(key)) continue;
if (incompleteKeySet.has(key)) continue;
const cooldownUntil = conversationErrorCooldown.get(key);
if (cooldownUntil && now < cooldownUntil) continue;
// No lane specified → schedule BOTH lanes; each fetches its own
// pending subset from the DB.
scheduleConversationAnalysis(key);
}
// --- Individual recovery for error/analysis_incomplete messages ---
// Circuit breaker check: no point iterating if individual CB is active.
if (now >= individualCooldownUntil) {
const promises: Promise<void>[] = [];
for (const key of incompleteKeys) {
// Skip if individual work is already running for this conversation.
if (individualInFlightByConversation.has(key)) continue;
// Skip if batch processing is running.
if (isConversationProcessingLocked(key)) continue;
promises.push(
messageStore
.getIncompleteMessagesByConversation(key, 500)
.then(async (msgs) => {
const processableMessages = await skipAnalysisUserMessages(
await skipAgeRestrictedMessages(msgs),
);
return processableMessages;
})
.then((msgs) => {
if (msgs.length > 0) {
enqueueIndividualFallbacks(msgs);
}
})
.catch((err: unknown) => {
logger.error(
{ key, error: String(err) },
"Failed to fetch incomplete messages for recovery",
);
}),
);
}
// Errors are handled per-key; return the combined promise for observability.
return Promise.all(promises);
}
})
.catch((err: unknown) => {
logger.error(
{ error: err instanceof Error ? err.message : String(err) },
"Pending AI analysis recovery worker failed",
);
});
}, config.AI_ANALYSIS_RECOVERY_INTERVAL_MS);
startRecoveryWorker();
}
@@ -0,0 +1,41 @@
import { createChildLogger } from "@/shared/logger/index";
import { config } from "../../shared/config/index.js";
import { deleteExpiredQdrantPoints } from "./qdrantClient.js";
import { pruneExpiredTexts } from "./textCacheStore.js";
const logger = createChildLogger("cache-prune");
/** Expired-verdict sweep cadence. */
const CACHE_PRUNE_INTERVAL_MS = 6 * 60 * 60 * 1000; // every 6 hours
let lastCachePruneAt = 0;
/**
* Cache hygiene: purge expired moderation verdicts from Postgres and Qdrant.
*
* Expired entries are never reused (read filters check `expires_at`) but they
* accumulate forever without a sweep. Called from the recovery interval; the
* 6-hour throttle keeps it to one sweep per window.
*/
export function runCachePruneIfDue(now: number = Date.now()): void {
if (now - lastCachePruneAt < CACHE_PRUNE_INTERVAL_MS) return;
lastCachePruneAt = now;
Promise.all([pruneExpiredTexts(), deleteExpiredQdrantPoints()])
.then(([pgDeleted, qdDeleted]) => {
if (pgDeleted > 0 || qdDeleted > 0) {
logger.info(
{ pgDeleted, qdDeleted },
"Expired moderation cache pruned",
);
}
})
.catch((err: unknown) => {
logger.warn({ error: String(err) }, "Moderation cache prune failed");
});
}
/** Reset the throttle window (tests). */
export function resetCachePruneState(): void {
lastCachePruneAt = 0;
}
@@ -0,0 +1,35 @@
/**
* Public surface of the AI-moderation module.
*
* The module has ~50 internal files; callers outside it (app/, tests, other
* modules) should import from THIS barrel so internal files can be moved
* without touching call sites.
*
* Deep imports remain valid inside the module itself.
*/
export type {
AIRecommendedAction,
AISeverity,
AIStatus,
AnalysisQueueStatus,
AnalysisResult,
} from "../../shared/moderation-types.js";
// ── Entry API: queueing, status, recovery worker ──────────────────────────
export {
getAnalysisQueueStatus,
queueConversationAnalysis,
queueMessageAnalysis,
startPendingAIAnalysisWorker,
} from "./aiAnalyzer.js";
// ── Worker pools (app/metrics-collector reads their live thread counters) ──
export {
getConversationKey,
mediaWorkerPool,
textWorkerPool,
} from "./circuitBreaker.js";
// ── Pipeline state hooks the bootstrap injects into ───────────────────────
export {
setModerationClient,
setSharedEventBroadcaster,
} from "./moderationState.js";
@@ -0,0 +1,184 @@
import { createChildLogger } from "@/shared/logger/index";
import { config } from "../../shared/config/index.js";
import { messageStore } from "../message-capture/messageStore.js";
import {
skipAgeRestrictedMessages,
skipAnalysisUserMessages,
} from "./batchProcessor.js";
import { scheduleConversationAnalysis } from "./batchScheduler.js";
import { runCachePruneIfDue } from "./cache-prune.js";
import {
ANALYSIS_LANES,
type AnalysisLane,
clearConversationProcessing,
conversationConsecutiveErrors,
conversationDebounceTimers,
conversationErrorCooldown,
conversationProcessing,
isConversationProcessingLocked,
} from "./conversationState.js";
import {
enqueueIndividualFallbacks,
individualCooldownUntil,
individualInFlightByConversation,
individualInFlightLastTouched,
} from "./individualFallbackProcessor.js";
const logger = createChildLogger("ai-recovery");
/** Revert messages stuck in `processing` for longer than this. */
const STUCK_PROCESSING_AGE_MS = 300_000;
/**
* Starts the periodic recovery worker.
*
* Recovers two classes of stranded work:
* - `pending` messages → re-scheduled through the normal per-lane debounce.
* - `error` / `analysis_incomplete` messages → individual fallback queue.
*
* Also prunes stale in-memory bookkeeping (lane locks, per-conversation
* circuit-breaker counters, individual-fallback in-flight markers) so a
* crashed batch cannot wedge a conversation forever, and triggers the
* throttled cache-prune sweep.
*
* Skips conversations that already have individual fallback work in progress
* to avoid DB last-write-wins races.
*/
export function startRecoveryWorker(): void {
setInterval(() => {
runCachePruneIfDue();
// Only revert stuck processing messages if there's active processing.
// Avoids a DB query every recovery interval when the pipeline is idle.
if (conversationProcessing.size > 0) {
messageStore
.revertStuckProcessingMessages(STUCK_PROCESSING_AGE_MS)
.catch((err: unknown) => {
logger.error(
{ error: String(err) },
"Failed to run stuck processing recovery",
);
});
}
Promise.all([
messageStore.getPendingConversationKeys(500),
messageStore.getConversationKeysWithIncompleteAnalysis(200),
])
.then(([pendingKeys, incompleteKeys]) => {
const now = Date.now();
pruneStaleConversationState(now);
const incompleteKeySet = new Set(incompleteKeys);
// --- Batch recovery for pending messages ---
for (const key of pendingKeys) {
if (
ANALYSIS_LANES.some((lane) =>
conversationDebounceTimers.has(`${key}::${lane}`),
)
) {
continue;
}
// Batch recovery must not race ANY in-flight batch lane, so the
// lock check is lane-agnostic here (individual fallback handles
// error rows separately).
if (isConversationProcessingLocked(key)) continue;
if (individualInFlightByConversation.has(key)) continue;
if (incompleteKeySet.has(key)) continue;
const cooldownUntil = conversationErrorCooldown.get(key);
if (cooldownUntil && now < cooldownUntil) continue;
// No lane specified → schedule BOTH lanes; each fetches its own
// pending subset from the DB.
scheduleConversationAnalysis(key);
}
// --- Individual recovery for error/analysis_incomplete messages ---
// Circuit breaker check: no point iterating if individual CB is active.
if (now >= individualCooldownUntil) {
const promises: Promise<void>[] = [];
for (const key of incompleteKeys) {
// Skip if individual work is already running for this conversation.
if (individualInFlightByConversation.has(key)) continue;
// Skip if batch processing is running.
if (isConversationProcessingLocked(key)) continue;
promises.push(
recoverIncompleteConversation(key).catch((err: unknown) => {
logger.error(
{ key, error: String(err) },
"Failed to fetch incomplete messages for recovery",
);
}),
);
}
// Errors are handled per-key; return the combined promise for observability.
return Promise.all(promises);
}
})
.catch((err: unknown) => {
logger.error(
{ error: err instanceof Error ? err.message : String(err) },
"Pending AI analysis recovery worker failed",
);
});
}, config.AI_ANALYSIS_RECOVERY_INTERVAL_MS);
}
/** Fetch one conversation's incomplete messages and queue them individually. */
async function recoverIncompleteConversation(key: string): Promise<void> {
const msgs = await messageStore.getIncompleteMessagesByConversation(key, 500);
const processable = await skipAnalysisUserMessages(
await skipAgeRestrictedMessages(msgs),
);
if (processable.length > 0) {
enqueueIndividualFallbacks(processable);
}
}
/**
* Drop stale in-memory bookkeeping:
* - per-lane processing locks past the timeout (pruned PER LANE so one stale
* lane never clears the other lane's healthy lock),
* - individual-fallback in-flight markers that stopped being touched,
* - per-conversation circuit-breaker error counts whose cooldown has lapsed.
*/
function pruneStaleConversationState(now: number): void {
for (const [key, expiry] of conversationErrorCooldown) {
if (now >= expiry) conversationErrorCooldown.delete(key);
}
for (const [key, record] of conversationProcessing) {
for (const lane of ANALYSIS_LANES as readonly AnalysisLane[]) {
const startedAt = record?.[lane];
if (
startedAt &&
now - startedAt >= config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS
) {
clearConversationProcessing(key, lane);
}
}
}
const staleThreshold = config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS * 2;
for (const [key, lastTouched] of individualInFlightLastTouched) {
if (now - lastTouched >= staleThreshold) {
individualInFlightLastTouched.delete(key);
individualInFlightByConversation.delete(key);
logger.warn(
{ key },
"Pruned stale individualInFlightByConversation entry",
);
}
}
// Also prune stale per-conversation CB error counts that have cooled
// down so old conversations can be retried.
for (const [key] of conversationConsecutiveErrors) {
const cbExpire = conversationErrorCooldown.get(key) ?? 0;
if (cbExpire && now >= cbExpire) {
conversationConsecutiveErrors.delete(key);
}
}
}
@@ -1,12 +1,12 @@
import type { Client } from "discord.js-selfbot-v13";
import Redis from "ioredis";
import { config } from "../../shared/config/index.js";
import { createChildLogger } from "../../shared/logger/index.js";
import {
BACKEND_COMMAND,
type CommandMessage,
type CommandReply,
} from "../../shared/index.js";
import { createChildLogger } from "../../shared/logger/index.js";
} from "../../shared/redis-channels.js";
import { GuildHandler } from "./guild.handler.js";
import {
type CommandHandlerFn,
@@ -1,6 +1,9 @@
import type { Client } from "discord.js-selfbot-v13";
import type { CommandMessage, CommandReply } from "../../shared/index.js";
import { createChildLogger } from "../../shared/logger/index.js";
import type {
CommandMessage,
CommandReply,
} from "../../shared/redis-channels.js";
// ---------------------------------------------------------------------------
// GuildHandler
@@ -4,7 +4,7 @@ import {
COMMAND_MODERATION_ACTION,
type CommandMessage,
type CommandReply,
} from "../../shared/index.js";
} from "../../shared/redis-channels.js";
import type { GuildHandler } from "./guild.handler.js";
import type { ModerationHandler } from "./moderation.handler.js";
@@ -1,6 +1,9 @@
import type { Client } from "discord.js-selfbot-v13";
import type { CommandMessage, CommandReply } from "../../shared/index.js";
import { createChildLogger } from "../../shared/logger/index.js";
import type {
CommandMessage,
CommandReply,
} from "../../shared/redis-channels.js";
import { messageStore } from "../message-capture/messageStore.js";
// ---------------------------------------------------------------------------
@@ -1,9 +1,12 @@
import Redis from "ioredis";
import type { AttachmentRecord, MessageRecord } from "../../shared/index.js";
import {
type CustomLogger,
createChildLogger,
} from "../../shared/logger/index.js";
import type {
AttachmentRecord,
MessageRecord,
} from "../../shared/moderation-types.js";
import { type DiscordGatewayEvent, EventChannels } from "./eventTypes.js";
export class RedisEventPublisher {
@@ -17,7 +17,7 @@ import {
DISCORD_THREAD_DELETED,
DISCORD_THREAD_UPDATED,
type DiscordGatewayEvent,
} from "../../shared/index.js";
} from "../../shared/redis-channels.js";
export type { DiscordGatewayEvent };
@@ -0,0 +1,41 @@
/**
* Public surface of the message-capture module.
*
* Callers outside this module import from here instead of reaching into
* messageCapture.ts / moderationActionsDb.ts / messageStore.ts directly, so
* the internal file layout can change without touching call sites.
*
* Deep imports remain valid inside the module itself.
*/
// ── Discord listener registration (called from app/lifecycle.ts) ───────────
export {
registerMessageCapture,
setEventBroadcaster,
} from "./messageCapture.js";
// ── Message persistence facade ────────────────────────────────────────────
export { messageStore } from "./messageStore.js";
// ── Live moderation-action publishing ─────────────────────────────────────
export { setModerationEventBroadcaster } from "./moderationActionsDb.js";
// ── Domain types ──────────────────────────────────────────────────────────
export type {
AIRecommendedAction,
AISeverity,
AIStatus,
AnalysisQueueStatus,
AnalysisResult,
AttachmentRecord,
DashboardMessage,
MessageQuery,
MessageRecord,
MessageReview,
ModerationAction,
ModerationActionType,
ModerationWsEvent,
PageResult,
RetentionPolicy,
ReviewStatus,
RoleMetadata,
UserMetadata,
} from "./types.js";
@@ -2,8 +2,11 @@ import { and, desc, eq, inArray, type SQL, sql } from "drizzle-orm";
import type { NodePgDatabase } from "drizzle-orm/node-postgres";
import type * as schema from "../../shared/database/schema.js";
import { messagesTable } from "../../shared/database/schema.js";
import { buildCursorCondition, pageResult } from "../../shared/index.js";
import { createChildLogger, type Logger } from "../../shared/logger/index.js";
import {
buildCursorCondition,
pageResult,
} from "../../shared/utils/pagination.js";
import type {
MessageQuery,
MessageRecord,
@@ -2,8 +2,11 @@ import { and, desc, eq, inArray, type SQL } from "drizzle-orm";
import type { NodePgDatabase } from "drizzle-orm/node-postgres";
import type * as schema from "../../shared/database/schema.js";
import { moderationActionsTable } from "../../shared/database/schema.js";
import { buildCursorCondition, pageResult } from "../../shared/index.js";
import { createChildLogger, type Logger } from "../../shared/logger/index.js";
import {
buildCursorCondition,
pageResult,
} from "../../shared/utils/pagination.js";
import type { EventBroadcaster } from "../event-broadcaster/eventBroadcaster.js";
import type { ModerationAction, PageResult } from "../message-capture/types.js";
@@ -2,8 +2,11 @@ import { and, desc, eq, inArray, type SQL } from "drizzle-orm";
import type { NodePgDatabase } from "drizzle-orm/node-postgres";
import type * as schema from "../../shared/database/schema.js";
import { messageReviewsTable } from "../../shared/database/schema.js";
import { buildCursorCondition, pageResult } from "../../shared/index.js";
import { createChildLogger, type Logger } from "../../shared/logger/index.js";
import {
buildCursorCondition,
pageResult,
} from "../../shared/utils/pagination.js";
import type { MessageReview, PageResult } from "../message-capture/types.js";
// ─── ReviewsDb Class ────────────────────────────────────────────────────────
@@ -2,8 +2,7 @@ import type {
AnalysisQueueStatus,
AttachmentRecord,
MessageRecord,
UserMetadata,
} from "../../shared/index.js";
} from "../../shared/moderation-types.js";
// Re-export all shared types for backward compatibility
export type {
@@ -26,7 +25,7 @@ export type {
ReviewStatus,
RoleMetadata,
UserMetadata,
} from "../../shared/index.js";
} from "../../shared/moderation-types.js";
// Local-only types (not shared across services)
export type ModerationWsEvent =
@@ -46,3 +46,32 @@ export class ConfigError extends AppError {
this.name = "ConfigError";
}
}
// ---------------------------------------------------------------------------
// Generic error helpers (shared by every module — avoids the repeated
// `err instanceof Error ? err.message : String(err)` pattern, 18+ sites)
// ---------------------------------------------------------------------------
/** Normalize an unknown thrown value to a readable message. */
export function errorMessage(err: unknown): string {
return err instanceof Error ? err.message : String(err);
}
/**
* True when an error code is a transient stream-teardown failure
* (EPIPE / stream destroyed / write-after-end / socket reset).
*
* These are NOT fatal: crashing the gateway on them (e.g. voice stop races,
* ffmpeg stdin closed while we still write) takes the whole bot offline
* mid-operation. Callers that install process-level handlers use this to
* log-and-continue instead of shutting down.
*/
export function isTransientStreamError(err: unknown): boolean {
const code = (err as NodeJS.ErrnoException)?.code ?? "";
return (
code === "EPIPE" ||
code === "ERR_STREAM_DESTROYED" ||
code === "ERR_STREAM_WRITE_AFTER_END" ||
code === "ECONNRESET"
);
}
@@ -1,9 +0,0 @@
export * from "./config/index.js";
export * from "./database/init.js";
export * from "./database/pool.js";
export * from "./database/schema.js";
export * from "./errors/index.js";
export * from "./logger/index.js";
export * from "./moderation-types.js";
export * from "./redis-channels.js";
export * from "./utils/index.js";