Compare commits
3
Commits
8743fcc0b5
...
c57ee12da1
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c57ee12da1 | ||
|
|
494e16b3b3 | ||
|
|
6bf3b40cc7 |
@@ -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.
|
||||
@@ -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`).
|
||||
|
||||
@@ -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";
|
||||
Reference in New Issue
Block a user