Compare commits
25
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
54d02098c8 | ||
|
|
5fc0e8b3cb | ||
|
|
045cdf1f75 | ||
|
|
deed7bdfb0 | ||
|
|
c4d9ade85e | ||
|
|
80daa9f045 | ||
|
|
ef7708bf7d | ||
|
|
750f3aa598 | ||
|
|
c57ee12da1 | ||
|
|
494e16b3b3 | ||
|
|
6bf3b40cc7 | ||
|
|
8743fcc0b5 | ||
|
|
e34dcd6bc6 | ||
|
|
f9fecfc144 | ||
|
|
a59f3132ee | ||
|
|
c7f53e4f7e | ||
|
|
8583bcdf17 | ||
|
|
b12eb0a038 | ||
|
|
b72423c64d | ||
|
|
6b34fc97ec | ||
|
|
8ce8978755 | ||
|
|
869cad88e4 | ||
|
|
f136cf3f6b | ||
|
|
cd3ee5b5d8 | ||
|
|
93288afa0b |
@@ -74,12 +74,6 @@ AI_LLM_IMAGE_MAX_DIMENSION=1024 # Max image dimension in pixels before r
|
||||
AI_LLM_TEXT_BATCH_SIZE=20 # Max messages per text-only moderation batch (default: 20)
|
||||
AI_LLM_MEDIA_ANALYSIS_TIMEOUT_MS=60000 # Timeout in ms for media analysis calls (default: 60000)
|
||||
AI_LLM_TEXT_ANALYSIS_TIMEOUT_MS=30000 # Timeout in ms for text-only analysis calls (default: 30000)
|
||||
AI_LLM_JEV_ENABLED=false # Use TypeSafe Jev (System One) as PRIMARY text analyzer; LLM is fallback
|
||||
AI_LLM_JEV_API_KEY= # REQUIRED if AI_LLM_JEV_ENABLED=true. 9router/TypeSafe API key for /v1/systemone
|
||||
AI_LLM_JEV_BASE_URL=http://127.0.0.1:4014 # 9router base URL (default: local 9router; prod: https://9router.asepharyana.my.id)
|
||||
AI_LLM_JEV_MODEL=oc/jev-1.13-free # Jev model id (default: oc/jev-1.13-free)
|
||||
AI_LLM_JEV_TIMEOUT_MS=45000 # Timeout in ms for a Jev systemone batch call (default: 45000)
|
||||
AI_LLM_JEV_MIN_CONFIDENCE=0.9 # Min status-choice confidence to accept a Jev verdict (default: 0.9)
|
||||
|
||||
# === AI Analysis Tuning ===
|
||||
AI_ANALYSIS_DEBOUNCE_MS=500 # Debounce window for batching messages in ms (default: 500)
|
||||
|
||||
@@ -151,9 +151,17 @@ jobs:
|
||||
|
||||
attic_push_vps_hop() {
|
||||
echo "Fallback: VPS-hop attic push"
|
||||
# Recover the client binary BEFORE the fallback can use it: the
|
||||
# bootstrap cascade below resets ATTIC_BIN="" and never restores it
|
||||
# in the fallback branch, so `sudo $ATTIC_BIN push` used to run as
|
||||
# `sudo push` -> "sudo: 'push': command not found". On the VPS the
|
||||
# closure lives at the canonical ATTIC_DIR path.
|
||||
VPS_ATTIC="/nix/store/fygyy3yk4rqdknxkiwkqambpnhyax0k4-attic-0.1.0/bin/attic"
|
||||
ssh "$VPS_USER@$VPS_HOST" "test -x '$VPS_ATTIC'" \
|
||||
|| ssh "$VPS_USER@$VPS_HOST" "sudo /nix/var/nix/profiles/default/bin/nix-store --realise '$ATTIC_DIR'"
|
||||
# Copy closure to VPS (fast if attic already has it via substitute)
|
||||
ssh "$VPS_USER@$VPS_HOST" "sudo /nix/var/nix/profiles/default/bin/nix-store --realise '$STORE_PATH'" 2>/dev/null \
|
||||
|| nix copy --to "ssh://$VPS_USER@$VPS_HOST" "$STORE_PATH"
|
||||
|| nix copy --to "ssh://***@$VPS_HOST" "$STORE_PATH"
|
||||
# Push from VPS → Attic over Tailscale.
|
||||
# --ignore-upstream-cache-filter is REQUIRED: without it, attic skips
|
||||
# writing the narinfo to gmw when chunks exist in the upstream
|
||||
@@ -162,7 +170,7 @@ jobs:
|
||||
# sudo: attic must read root's config (~/.config/attic), which has
|
||||
# the imrnes-ts server → Tailscale. Non-root users' configs only
|
||||
# have the public `pub` server → "Server imrnes-ts does not exist".
|
||||
ssh "$VPS_USER@$VPS_HOST" "sudo $ATTIC_BIN push imrnes-ts:gmw '$STORE_PATH' --jobs 4 --ignore-upstream-cache-filter" \
|
||||
ssh "$VPS_USER@$VPS_HOST" "sudo $VPS_ATTIC push imrnes-ts:gmw '$STORE_PATH' --jobs 4 --ignore-upstream-cache-filter" \
|
||||
|| echo "attic push failed (non-fatal; ssh copy fallback below)"
|
||||
}
|
||||
|
||||
|
||||
@@ -23,3 +23,7 @@ result
|
||||
|
||||
# Playwright MCP artifacts
|
||||
.playwright-mcp/
|
||||
findings.md
|
||||
findings.md
|
||||
progress.md
|
||||
task_plan.md
|
||||
|
||||
@@ -1,10 +0,0 @@
|
||||
import { z } from "zod";
|
||||
|
||||
export const searchQuerySchema = z.object({
|
||||
q: z.string().default(""),
|
||||
channelId: z.string().optional(),
|
||||
guildId: z.string().optional(),
|
||||
limit: z.coerce.number().int().positive().max(100).default(20),
|
||||
});
|
||||
|
||||
export type SearchQuery = z.infer<typeof searchQuerySchema>;
|
||||
@@ -1,5 +0,0 @@
|
||||
import { z } from "zod";
|
||||
|
||||
export const healthCheckSchema = z.object({
|
||||
verbose: z.coerce.boolean().optional().default(false),
|
||||
});
|
||||
@@ -1,80 +0,0 @@
|
||||
/**
|
||||
* moderationMetrics.ts
|
||||
*
|
||||
* Prometheus metrics for AI moderation pipeline.
|
||||
* Defined in backend (where prom-client is installed + /api/metrics endpoint).
|
||||
*/
|
||||
import { Counter, Histogram } from "prom-client";
|
||||
|
||||
// ── LLM Call Metrics ──
|
||||
export const llmCallsTotal = new Counter({
|
||||
name: "moderation_llm_calls_total",
|
||||
help: "Total LLM moderation calls",
|
||||
labelNames: ["path", "model"] as const,
|
||||
});
|
||||
|
||||
export const llmCallDuration = new Histogram({
|
||||
name: "moderation_llm_call_duration_ms",
|
||||
help: "LLM moderation call duration (ms)",
|
||||
labelNames: ["path", "status"] as const,
|
||||
buckets: [500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 120000],
|
||||
});
|
||||
|
||||
export const llmTokensTotal = new Counter({
|
||||
name: "moderation_llm_tokens_total",
|
||||
help: "Total tokens consumed by LLM moderation",
|
||||
labelNames: ["type"] as const,
|
||||
});
|
||||
|
||||
// ── Cache Metrics ──
|
||||
export const moderationCacheHits = new Counter({
|
||||
name: "moderation_cache_hits_total",
|
||||
help: "Moderation cache hits",
|
||||
labelNames: ["layer"] as const,
|
||||
});
|
||||
|
||||
export const moderationCacheMisses = new Counter({
|
||||
name: "moderation_cache_misses_total",
|
||||
help: "Moderation cache misses",
|
||||
labelNames: ["layer"] as const,
|
||||
});
|
||||
|
||||
// ── Media Analysis Metrics ──
|
||||
export const mediaAnalysesTotal = new Counter({
|
||||
name: "moderation_media_analyses_total",
|
||||
help: "Media analyses performed",
|
||||
labelNames: ["type"] as const,
|
||||
});
|
||||
|
||||
export const mediaDownloadDuration = new Histogram({
|
||||
name: "moderation_media_download_duration_ms",
|
||||
help: "Media download duration (ms)",
|
||||
labelNames: ["source"] as const,
|
||||
buckets: [100, 500, 1000, 2000, 5000, 10000, 30000],
|
||||
});
|
||||
|
||||
// ── Batch & Error Metrics ──
|
||||
export const moderationBatchSize = new Histogram({
|
||||
name: "moderation_batch_size",
|
||||
help: "Messages per batch",
|
||||
labelNames: ["path"] as const,
|
||||
buckets: [1, 5, 10, 20, 50, 100],
|
||||
});
|
||||
|
||||
export const moderationErrors = new Counter({
|
||||
name: "moderation_errors_total",
|
||||
help: "Moderation errors",
|
||||
labelNames: ["type"] as const,
|
||||
});
|
||||
|
||||
export const webSearchCalls = new Counter({
|
||||
name: "moderation_websearch_calls_total",
|
||||
help: "Wikipedia web-search calls",
|
||||
labelNames: ["status"] as const,
|
||||
});
|
||||
|
||||
export const autoDeleteActions = new Counter({
|
||||
name: "moderation_auto_delete_total",
|
||||
help: "Auto-delete actions",
|
||||
labelNames: ["action"] as const,
|
||||
});
|
||||
@@ -1,25 +0,0 @@
|
||||
import type { CommandReply } from "./index.js";
|
||||
import { createChildLogger } from "./logger/index.js";
|
||||
|
||||
export { createChildLogger };
|
||||
|
||||
/**
|
||||
* Attempt a Redis command first; if it fails or times out, fall back.
|
||||
*
|
||||
* @param commandFn - Function that issues the publishCommand and returns the reply.
|
||||
* @param fallbackFn - Async fallback, typically reads from Redis status key.
|
||||
* @param commandLabel - Label used for logging (e.g. "voice:connect").
|
||||
*/
|
||||
export async function tryCommandThenFallback<T>(
|
||||
commandFn: () => Promise<CommandReply<T> | null>,
|
||||
fallbackFn: () => Promise<T>,
|
||||
commandLabel: string,
|
||||
): Promise<T> {
|
||||
const logger = createChildLogger(`command-helper:${commandLabel}`);
|
||||
const reply = await commandFn();
|
||||
if (reply?.success && reply.data !== undefined && reply.data !== null) {
|
||||
return reply.data;
|
||||
}
|
||||
logger.warn("discord-gateway unreachable, falling back");
|
||||
return fallbackFn();
|
||||
}
|
||||
@@ -96,10 +96,10 @@ export const configSchema = z
|
||||
.transform((v) => v === "true")
|
||||
.default(false),
|
||||
AI_LLM_API_KEY: z.string().optional(),
|
||||
AI_LLM_BASE_URL: z
|
||||
.string()
|
||||
.url()
|
||||
.default("http://100.121.180.82:20128/api/v1"),
|
||||
// 9router — OpenAI-compatible router on this host (127.0.0.1:4014).
|
||||
// Loopback on purpose: backend runs on the same machine as 9router, so no
|
||||
// TLS/proxy hop is needed.
|
||||
AI_LLM_BASE_URL: z.string().url().default("http://127.0.0.1:4014/v1"),
|
||||
AI_LLM_MODEL: z.string().default("text"),
|
||||
AI_LLM_VISION_MODEL: z.string().optional(),
|
||||
AI_LLM_EMBEDDING_MODEL: z.string().optional(),
|
||||
|
||||
@@ -1,8 +0,0 @@
|
||||
export {
|
||||
broadcastBinary,
|
||||
broadcastEvent,
|
||||
clearBroadcastFunctions,
|
||||
setBroadcastFunctions,
|
||||
} from "./broadcast.js";
|
||||
export { startRedisBridge, stopRedisBridge } from "./redis-bridge.js";
|
||||
export { closeWebSocketServer, createWebSocketServer } from "./server.js";
|
||||
@@ -54,7 +54,7 @@ src/
|
||||
**Never** reintroduce regex/heuristic content classification.
|
||||
2. **Discord tokens sanitized** before reaching LLM (`discordTokens.ts`).
|
||||
3. **Semantic cache is batched** — one embed call + one Qdrant batch search.
|
||||
4. **Streaming is mandatory** against the omniroute base URL.
|
||||
4. **Streaming is mandatory** against the router base URL.
|
||||
|
||||
## AI moderation pipeline
|
||||
|
||||
@@ -69,8 +69,14 @@ textBatchProcessor.ts mediaBatchProcessor.ts llmClient.ts
|
||||
```
|
||||
|
||||
- Entry: `aiAnalyzer.ts` (`queueMessageAnalysis`, `startPendingAIAnalysisWorker`)
|
||||
- Concurrency: LLM semaphore (`AI_LLM_MAX_CONCURRENT`, default 5)
|
||||
- Concurrency: **two per-lane LLM semaphores** (2026-09-24) — text
|
||||
(`AI_LLM_MAX_CONCURRENT`, default 8) and media/vision
|
||||
(`AI_LLM_MEDIA_MAX_CONCURRENT`, default 4); a media backlog can never
|
||||
consume text slots
|
||||
- Piscina: text pool (4 threads) + media pool (2 threads)
|
||||
- Locks are **per conversation per lane** (`conversationProcessing` maps key →
|
||||
lane → startedAt): the text lane of a conversation never waits on that
|
||||
conversation's media lane (this was the "image blocks the queue" bug)
|
||||
- **Each worker thread has its own pg Pool** (min 0, grows to `POSTGRES_POOL_MAX`)
|
||||
|
||||
## Module: message-capture
|
||||
|
||||
@@ -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,54 +15,83 @@ 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/
|
||||
│ └── gateway-metrics/ # Prometheus /metrics endpoint (port 4016)
|
||||
│ ├── channel-topic/ guild-member-events/ monitor/
|
||||
│ └── gateway-metrics/ # Prometheus /metrics endpoint (METRICS_PORT)
|
||||
```
|
||||
|
||||
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
|
||||
handles a whole batch (text + media split internally, parallel paths).
|
||||
handles a whole batch. **Independent text/media lanes** (2026-09-24): a
|
||||
conversation batch is split into a text lane (messages with no media) and a
|
||||
media lane (attachments/stickers/embeds) that are dispatched to separate
|
||||
pools, hold SEPARATE per-lane processing locks, and run under SEPARATE LLM
|
||||
concurrency semaphores. The text lane frees its lock and saves+broadcasts the
|
||||
moment text analysis finishes — it never waits on a slow vision/media batch
|
||||
of the same conversation, and vice versa.
|
||||
|
||||
- `aiAnalyzer.ts` — public API: `queueMessageAnalysis`, `getAnalysisQueueStatus`,
|
||||
`startPendingAIAnalysisWorker` (recovery worker + cache-prune).
|
||||
- `batchScheduler.ts` — per-conversation debounce → `processBatch`.
|
||||
- `batchProcessor.ts` — batch lock/circuit-breaker, fans failed targets to
|
||||
individual fallback.
|
||||
- `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).
|
||||
- `batchProcessor.ts` — per-lane batch lock/circuit-breaker, fans failed
|
||||
targets to individual fallback. `processBatch` releases ITS lane's lock the
|
||||
moment that lane's worker job finishes; the other lane owns its own lock.
|
||||
- `individualFallbackProcessor.ts` — one-message-at-a-time retry path, own CB.
|
||||
- `conversationState.ts` / `circuitBreaker.ts` — per-conversation state,
|
||||
Piscina `workerPool`, `getConversationKey`.
|
||||
- `ai-analysis-worker.ts` — Piscina entry point (`batch` / `individual` jobs).
|
||||
Runs `runModerationAnalysis` off the main thread.
|
||||
- `conversationState.ts` / `circuitBreaker.ts` — per-conversation PER-LANE
|
||||
state (`conversationProcessing` holds a lane → startedAt map per key),
|
||||
Piscina `textWorkerPool`/`mediaWorkerPool`, `getConversationKey`.
|
||||
- `ai-analysis-worker.ts` — Piscina entry point (`batch` (lane) /
|
||||
`individual` jobs). Runs `runModerationAnalysis` off the main thread.
|
||||
- `moderationOrchestrator.ts` — exact-hash cache → batched semantic (Qdrant)
|
||||
cache → LLM. Text and media paths run in parallel.
|
||||
- `textBatchProcessor.ts` / `mediaBatchProcessor.ts` — actual LLM calls
|
||||
(one call per sub-batch, not per message).
|
||||
(one call per sub-batch, not per message). `mediaBatchProcessor` routes its
|
||||
moderation LLM call through the MEDIA semaphore.
|
||||
- `llmClient.ts` — central OpenAI-compatible chat client (streaming, retries,
|
||||
thinking-disable injection). `visionAnalyzer.ts` / `mediaAnalysisClient.ts`
|
||||
thinking-disable injection). TWO concurrency semaphores:
|
||||
`AI_LLM_MAX_CONCURRENT` (text lane, default 8) and
|
||||
`AI_LLM_MEDIA_MAX_CONCURRENT` (media lane, default 4) — a vision backlog
|
||||
can never consume text slots. `visionAnalyzer.ts` / `mediaAnalysisClient.ts`
|
||||
share the same router/base URL (different model alias for vision).
|
||||
- `embeddingClient.ts` + `qdrantClient.ts` — semantic cache (one embed call +
|
||||
one batched Qdrant search for all uncached targets).
|
||||
@@ -72,16 +100,16 @@ handles a whole batch (text + media split internally, parallel paths).
|
||||
|
||||
### Concurrency model
|
||||
|
||||
- Main thread owns the LLM semaphore (`AI_LLM_MAX_CONCURRENT`, default 5) via
|
||||
`llmClient.withLlmConcurrency`.
|
||||
- Main thread owns TWO per-lane LLM semaphores (2026-09-24):
|
||||
`AI_LLM_MAX_CONCURRENT` (text, default 8) and `AI_LLM_MEDIA_MAX_CONCURRENT`
|
||||
(media, default 4) via `llmClient.withLlmConcurrency(fn, { lane })`.
|
||||
- Two Piscina pools run the heavy LLM work off the event loop: a text pool
|
||||
(`PISCINA_MAX_THREADS`, default 4) and a dedicated media pool
|
||||
(`PISCINA_MEDIA_MAX_THREADS`, default 2). A batch is routed to the media
|
||||
pool if ANY of its messages carries an attachment/sticker/embed — this
|
||||
keeps a slow image/vision batch from occupying every thread and blocking
|
||||
unrelated text-only batches behind it. **Each worker thread (in either
|
||||
pool) initializes its own pg Pool** (min 0, grows to `POSTGRES_POOL_MAX`).
|
||||
See "Memory & connections" below.
|
||||
(`PISCINA_MEDIA_MAX_THREADS`, default 2). A batch is routed by lane to the
|
||||
matching pool — this keeps a slow image/vision batch from occupying every
|
||||
thread and blocking unrelated text-only batches behind it. **Each worker
|
||||
thread (in either pool) initializes its own pg Pool** (min 0, grows to
|
||||
`POSTGRES_POOL_MAX`). See "Memory & connections" below.
|
||||
|
||||
## Memory & DB connections
|
||||
|
||||
@@ -109,28 +137,53 @@ 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 (`METRICS_PORT`, code
|
||||
default 9090, set per deployment).
|
||||
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
|
||||
Prometheus scrapes the metrics server at `127.0.0.1:$METRICS_PORT/metrics`
|
||||
(`bete_*` prefix; the code default is 9090 — deployments set it explicitly,
|
||||
this host uses 4018). 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)
|
||||
|
||||
@@ -141,5 +194,5 @@ pipeline gauges — `ai_analysis_queued_conversations`,
|
||||
numeric snowflake IDs never trigger false positives.
|
||||
- **Semantic cache is batched** (one embed call + one Qdrant batch search),
|
||||
not N sequential round-trips. `ensureQdrantCollection` is memoized.
|
||||
- **Streaming is mandatory** against the omniroute base URL (non-stream waits for
|
||||
- **Streaming is mandatory** against the router base URL (non-stream waits for
|
||||
the full body and times out). `llmClient` aggregates SSE chunks.
|
||||
|
||||
@@ -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 (METRICS_PORT)
|
||||
├── 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`).
|
||||
|
||||
@@ -23,7 +23,6 @@
|
||||
"test:e2e:live": "bash scripts/run-llm-e2e.sh"
|
||||
},
|
||||
"dependencies": {
|
||||
"@typesafe-ai/sdk": "^0.6.0",
|
||||
"axios": "^1.20.0",
|
||||
"discord.js-selfbot-v13": "^3.7.1",
|
||||
"dotenv": "^18.0.0",
|
||||
@@ -48,7 +47,7 @@
|
||||
"@types/pg": "^8.23.1",
|
||||
"@types/ws": "^8.18.1",
|
||||
"drizzle-kit": "^0.31.10",
|
||||
"tsx": "^4.23.13",
|
||||
"tsx": "^4.23.15",
|
||||
"typescript": "^7.0.2",
|
||||
"vitest": "latest"
|
||||
}
|
||||
|
||||
Generated
+14
-23
@@ -8,9 +8,6 @@ importers:
|
||||
|
||||
.:
|
||||
dependencies:
|
||||
'@typesafe-ai/sdk':
|
||||
specifier: ^0.6.0
|
||||
version: 0.6.0
|
||||
axios:
|
||||
specifier: ^1.20.0
|
||||
version: 1.20.0(debug@4.4.3(supports-color@7.2.0))(supports-color@7.2.0)
|
||||
@@ -79,14 +76,14 @@ importers:
|
||||
specifier: ^0.31.10
|
||||
version: 0.31.10
|
||||
tsx:
|
||||
specifier: ^4.23.13
|
||||
version: 4.23.13
|
||||
specifier: ^4.23.15
|
||||
version: 4.23.15
|
||||
typescript:
|
||||
specifier: ^7.0.2
|
||||
version: 7.0.2
|
||||
vitest:
|
||||
specifier: latest
|
||||
version: 5.0.1(@types/node@26.4.0)(vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.13))
|
||||
version: 5.0.1(@types/node@26.4.0)(vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.15))
|
||||
|
||||
packages:
|
||||
|
||||
@@ -1090,10 +1087,6 @@ packages:
|
||||
'@types/ws@8.18.1':
|
||||
resolution: {integrity: sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==}
|
||||
|
||||
'@typesafe-ai/sdk@0.6.0':
|
||||
resolution: {integrity: sha512-IddX+Q0XM+VagOUZFeP7wZjaO4SHMdvnh2zEBdrZZnXedWI3BNK1lKhMx3ayrkFWvVLbVcUHJy6AVZlY+e6Jaw==}
|
||||
engines: {node: '>=20'}
|
||||
|
||||
'@typescript/typescript-aix-ppc64@7.0.2':
|
||||
resolution: {integrity: sha512-MTKKkWB7p/0E9xi1d1tHtZ5PiLkGEMIq88pK2CubZjOsLtYTLqhgIgi6zepFa+9GHZ6h05NMCkQxGKiPXMxXtQ==}
|
||||
engines: {node: '>=16.20.0'}
|
||||
@@ -2172,8 +2165,8 @@ packages:
|
||||
tslib@2.8.1:
|
||||
resolution: {integrity: sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==}
|
||||
|
||||
tsx@4.23.13:
|
||||
resolution: {integrity: sha512-BL5MGkRln6aDYhb0xbQlEAGw743BaZYWdbWtdJOBriYJboKgUUYCadFp2/FpBBZquBC/ezNBn7wMMPx7FDZUDw==}
|
||||
tsx@4.23.15:
|
||||
resolution: {integrity: sha512-Yiex1Ovn8z2xPpOWckIiysV1SSyRMY9BkLF++q0yKiDxCqRhosKfMg3janKkiLBwZ5c/YryloKwGZcrEmtwxKw==}
|
||||
engines: {node: '>=18.0.0'}
|
||||
hasBin: true
|
||||
|
||||
@@ -2993,8 +2986,6 @@ snapshots:
|
||||
dependencies:
|
||||
'@types/node': 26.4.0
|
||||
|
||||
'@typesafe-ai/sdk@0.6.0': {}
|
||||
|
||||
'@typescript/typescript-aix-ppc64@7.0.2':
|
||||
optional: true
|
||||
|
||||
@@ -3055,14 +3046,14 @@ snapshots:
|
||||
'@typescript/typescript-win32-x64@7.0.2':
|
||||
optional: true
|
||||
|
||||
'@vitest/mocker@5.0.1(vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.13))':
|
||||
'@vitest/mocker@5.0.1(vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.15))':
|
||||
dependencies:
|
||||
'@jridgewell/trace-mapping': 0.3.31
|
||||
'@vitest/spy': 5.0.1
|
||||
estree-walker: 3.0.3
|
||||
magic-string: 1.2.3
|
||||
optionalDependencies:
|
||||
vite: 8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.13)
|
||||
vite: 8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.15)
|
||||
|
||||
'@vitest/spy@5.0.1': {}
|
||||
|
||||
@@ -3230,7 +3221,7 @@ snapshots:
|
||||
'@drizzle-team/brocli': 0.10.2
|
||||
'@esbuild-kit/esm-loader': 2.6.5
|
||||
esbuild: 0.25.12
|
||||
tsx: 4.23.13
|
||||
tsx: 4.23.15
|
||||
|
||||
drizzle-orm@0.45.2(@types/pg@8.23.1)(pg@8.23.0):
|
||||
optionalDependencies:
|
||||
@@ -3960,7 +3951,7 @@ snapshots:
|
||||
|
||||
tslib@2.8.1: {}
|
||||
|
||||
tsx@4.23.13:
|
||||
tsx@4.23.15:
|
||||
dependencies:
|
||||
esbuild: 0.28.2
|
||||
optionalDependencies:
|
||||
@@ -3996,7 +3987,7 @@ snapshots:
|
||||
util-deprecate@1.0.2:
|
||||
optional: true
|
||||
|
||||
vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.13):
|
||||
vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.15):
|
||||
dependencies:
|
||||
lightningcss: 1.33.0
|
||||
picomatch: 4.0.7
|
||||
@@ -4007,12 +3998,12 @@ snapshots:
|
||||
'@types/node': 26.4.0
|
||||
esbuild: 0.28.2
|
||||
fsevents: 2.3.3
|
||||
tsx: 4.23.13
|
||||
tsx: 4.23.15
|
||||
|
||||
vitest@5.0.1(@types/node@26.4.0)(vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.13)):
|
||||
vitest@5.0.1(@types/node@26.4.0)(vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.15)):
|
||||
dependencies:
|
||||
'@types/chai': 5.2.3
|
||||
'@vitest/mocker': 5.0.1(vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.13))
|
||||
'@vitest/mocker': 5.0.1(vite@8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.15))
|
||||
chai: 6.2.2
|
||||
es-module-lexer: 2.3.2
|
||||
expect-type: 1.4.0
|
||||
@@ -4023,7 +4014,7 @@ snapshots:
|
||||
tinybench: 6.1.4
|
||||
tinyexec: 1.3.0
|
||||
tinyglobby: 0.2.17
|
||||
vite: 8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.13)
|
||||
vite: 8.1.5(@types/node@26.4.0)(esbuild@0.28.2)(tsx@4.23.15)
|
||||
why-is-node-running: 2.3.0
|
||||
optionalDependencies:
|
||||
'@types/node': 26.4.0
|
||||
|
||||
@@ -1,82 +1,54 @@
|
||||
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/config.js";
|
||||
import { config } from "../shared/config/index.js";
|
||||
import {
|
||||
closeDatabase,
|
||||
initializeDatabase,
|
||||
} 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,175 +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_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,10 +1,7 @@
|
||||
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/config.js";
|
||||
import { config } from "../shared/config/index.js";
|
||||
import { getDatabase } from "../shared/database/drizzle.js";
|
||||
import type * as schema from "../shared/database/schema.js";
|
||||
import { attachmentsTable, messagesTable } from "../shared/database/schema.js";
|
||||
|
||||
@@ -1,5 +1,9 @@
|
||||
import type { Client } from "discord.js-selfbot-v13";
|
||||
import type { createChildLogger } from "@/shared/logger/index";
|
||||
import {
|
||||
mediaWorkerPool,
|
||||
textWorkerPool,
|
||||
} from "../modules/ai-moderation/circuitBreaker.js";
|
||||
import type { CommandHandler } from "../modules/command-handler/commandHandler.js";
|
||||
import type { EventBroadcaster } from "../modules/event-broadcaster/index.js";
|
||||
import type { stopMetricsServer } from "../modules/gateway-metrics/index.js";
|
||||
@@ -45,6 +49,30 @@ export function createGracefulShutdown(
|
||||
options.logger.info("Closing command handler...");
|
||||
await options.commandHandler.close();
|
||||
|
||||
// ½. Tear down AI-analysis worker pools BEFORE closing the DB.
|
||||
// Piscina worker threads survive process.exit() as orphans otherwise —
|
||||
// they keep holding DB connections/locks after the main process is gone.
|
||||
// (Two live gateways fighting over the same rows was the root cause of
|
||||
// messages stuck in ai_status='processing'.)
|
||||
options.logger.info("Destroying AI worker pools...");
|
||||
const destroyPool = (pool: { destroy: () => Promise<void> }) =>
|
||||
Promise.race([
|
||||
pool.destroy(),
|
||||
new Promise<void>((resolve) =>
|
||||
setTimeout(() => {
|
||||
options.logger.warn(
|
||||
"Timed out destroying worker pool; exiting anyway",
|
||||
);
|
||||
resolve();
|
||||
}, 5000),
|
||||
),
|
||||
]);
|
||||
await Promise.allSettled([
|
||||
destroyPool(textWorkerPool),
|
||||
destroyPool(mediaWorkerPool),
|
||||
]);
|
||||
options.logger.info("AI worker pools destroyed");
|
||||
|
||||
// 2. DB pool
|
||||
options.logger.info("Closing database...");
|
||||
await options.closeDatabase();
|
||||
|
||||
@@ -17,7 +17,7 @@
|
||||
*/
|
||||
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { initializeDatabase } from "../../shared/database/drizzle.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
@@ -86,7 +86,12 @@ export interface MessageBatch {
|
||||
|
||||
// Worker job types (Piscina entry point)
|
||||
type WorkerJob =
|
||||
| { type: "batch"; conversationKey: string; messages: MessageRecord[] }
|
||||
| {
|
||||
type: "batch";
|
||||
conversationKey: string;
|
||||
lane: "text" | "media";
|
||||
messages: MessageRecord[];
|
||||
}
|
||||
| { type: "individual"; message: MessageRecord; skipNormalAnalysis: boolean };
|
||||
|
||||
type BatchOkResponse = {
|
||||
@@ -263,6 +268,7 @@ function normalizeResult(
|
||||
async function processBatch(job: {
|
||||
type: "batch";
|
||||
conversationKey: string;
|
||||
lane: "text" | "media";
|
||||
messages: MessageRecord[];
|
||||
}): Promise<BatchOkResponse | BatchErrorResponse> {
|
||||
const { conversationKey, messages } = job;
|
||||
|
||||
@@ -1,34 +1,25 @@
|
||||
import type { Client } from "discord.js-selfbot-v13";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type { EventBroadcaster } from "../event-broadcaster/index.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type { AnalysisQueueStatus } from "../message-capture/types.js";
|
||||
import {
|
||||
activeMediaRequests,
|
||||
activeRequests,
|
||||
activeTextRequests,
|
||||
buildAgeRestrictedSkipResult,
|
||||
buildSkipAnalysisUserResult,
|
||||
isAgeRestrictedMessage,
|
||||
isSkipAnalysisUser,
|
||||
skipAgeRestrictedMessages,
|
||||
skipAnalysisUserMessages,
|
||||
} from "./batchProcessor.js";
|
||||
import { scheduleConversationAnalysis } from "./batchScheduler.js";
|
||||
import { getConversationKey } from "./circuitBreaker.js";
|
||||
import {
|
||||
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,
|
||||
@@ -36,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;
|
||||
@@ -73,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",
|
||||
@@ -88,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",
|
||||
@@ -114,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).
|
||||
*/
|
||||
@@ -129,6 +108,8 @@ export function getAnalysisQueueStatus(): AnalysisQueueStatus {
|
||||
return {
|
||||
queuedConversations: conversationDebounceTimers.size,
|
||||
activeRequests,
|
||||
activeTextRequests,
|
||||
activeMediaRequests,
|
||||
activeIndividualRequests,
|
||||
individualInFlightCount: individualInFlight.size,
|
||||
individualCircuitBreakerActive: Date.now() < individualCooldownUntil,
|
||||
@@ -137,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,
|
||||
@@ -160,131 +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);
|
||||
}
|
||||
for (const [key, startedAt] of conversationProcessing) {
|
||||
if (now - startedAt >= config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS) {
|
||||
conversationProcessing.delete(key);
|
||||
}
|
||||
}
|
||||
|
||||
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 (conversationDebounceTimers.has(key)) continue;
|
||||
if (isConversationProcessingLocked(key)) continue;
|
||||
if (individualInFlightByConversation.has(key)) continue;
|
||||
if (incompleteKeySet.has(key)) continue;
|
||||
const cooldownUntil = conversationErrorCooldown.get(key);
|
||||
if (cooldownUntil && now < cooldownUntil) continue;
|
||||
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,35 @@
|
||||
/**
|
||||
* analysisLanes.ts
|
||||
*
|
||||
* Pure lane helpers for the AI-analysis queue. Kept free of any import chain
|
||||
* that pulls Piscina/worker/DB so they can be unit-tested in isolation (the
|
||||
* scheduler's `splitMessagesByLane` used to live in batchScheduler.ts, which
|
||||
* transitively imports the worker pool).
|
||||
*/
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
import type { AnalysisLane } from "./conversationState.js";
|
||||
import { hasMediaContent } from "./mediaAnalysisClient.js";
|
||||
|
||||
export type { AnalysisLane } from "./conversationState.js";
|
||||
|
||||
/** True when this message belongs to the media lane (has attachment/sticker/embed). */
|
||||
export function laneOfMessage(message: MessageRecord): AnalysisLane {
|
||||
return hasMediaContent(message) ? "media" : "text";
|
||||
}
|
||||
|
||||
/**
|
||||
* Splits an arbitrary message array into per-lane lists. Used when the
|
||||
* scheduler runs a conversation-wide pass (lane omitted): each lane gets its
|
||||
* own subset so text and media never share a worker job.
|
||||
*/
|
||||
export function splitMessagesByLane(messages: MessageRecord[]): {
|
||||
text: MessageRecord[];
|
||||
media: MessageRecord[];
|
||||
} {
|
||||
const text: MessageRecord[] = [];
|
||||
const media: MessageRecord[] = [];
|
||||
for (const m of messages) {
|
||||
(laneOfMessage(m) === "media" ? media : text).push(m);
|
||||
}
|
||||
return { text, media };
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type {
|
||||
AnalysisResult,
|
||||
MessageRecord,
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { Guild } from "discord.js-selfbot-v13";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
|
||||
interface ChannelWithSend {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import type { Client, PermissionString } from "discord.js-selfbot-v13";
|
||||
import { LRUCache } from "lru-cache";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { parseRichMessageMetadata } from "../message-capture/messageMetadata.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { Client } from "discord.js-selfbot-v13";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
|
||||
const logger = createChildLogger("auto-delete-notify");
|
||||
|
||||
@@ -43,3 +43,22 @@ export function pickBatchWithinBudget(
|
||||
|
||||
return batch;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the messages that were fetched/claimed but did NOT make it into the
|
||||
* trimmed batch (i.e. the tail past the token budget).
|
||||
*
|
||||
* The DB claim step flips every fetched pending row to `processing`; the batch
|
||||
* trim may then stop early on the token budget. Those tail rows would stay
|
||||
* stuck in `processing` forever unless the caller explicitly un-claims them —
|
||||
* this helper identifies exactly which rows that is, so the caller can write
|
||||
* them back to `pending` for the next wave.
|
||||
*/
|
||||
export function computeBudgetOverflowMessages(
|
||||
claimed: MessageRecord[],
|
||||
trimmed: MessageRecord[],
|
||||
): MessageRecord[] {
|
||||
if (trimmed.length === 0) return claimed;
|
||||
const trimmedIds = new Set(trimmed.map((m) => m.id));
|
||||
return claimed.filter((m) => !trimmedIds.has(m.id));
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { isAgeRestrictedMetadata } from "../message-capture/messageMetadata.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
@@ -8,13 +8,14 @@ import { partitionBatchOutcome } from "./batchOutcomeClassifier.js";
|
||||
import { mediaWorkerPool, textWorkerPool } from "./circuitBreaker.js";
|
||||
import { estimateTokens } from "./conversationContext.js";
|
||||
import {
|
||||
type AnalysisLane,
|
||||
clearConversationProcessing,
|
||||
conversationErrorCooldown,
|
||||
conversationProcessing,
|
||||
getConversationProcessingStartedAt,
|
||||
recordConversationBatchFailure,
|
||||
resetConversationBatchFailures,
|
||||
} from "./conversationState.js";
|
||||
import { enqueueIndividualFallbacks } from "./individualFallbackProcessor.js";
|
||||
import { hasMediaContent } from "./mediaAnalysisClient.js";
|
||||
import {
|
||||
broadcastAnalysisCompleted,
|
||||
LAST_ERROR,
|
||||
@@ -34,10 +35,12 @@ export interface AnalysisWorkerResponse {
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Observability
|
||||
// Observability (per-lane counters live alongside the aggregate)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export let activeRequests = 0;
|
||||
export let activeTextRequests = 0;
|
||||
export let activeMediaRequests = 0;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Exported helpers
|
||||
@@ -182,32 +185,36 @@ export async function skipAnalysisUserMessages(
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Runs one worker job (either the text-only or the media sub-batch of a
|
||||
* Runs ONE worker job for a single lane (text-only or media sub-batch of a
|
||||
* conversation) end-to-end: dispatch → broadcast/save → fallback routing.
|
||||
* Returns whether the *caller* should schedule the next debounce pass for
|
||||
* this conversation (mirrors the old single-job semantics, now evaluated
|
||||
* per queue).
|
||||
*
|
||||
* Broadcasting happens here, inside each queue's own call — NOT after
|
||||
* waiting on the other queue. That's the actual fix for "text menunggu
|
||||
* image": previously one mixed conversation batch made ONE worker call
|
||||
* with both text and media targets, and `runModerationAnalysis` only
|
||||
* resolves (so results only get saved/broadcast) once BOTH finish — so a
|
||||
* fast text verdict sat unused until the slow vision/image verdict was
|
||||
* ready too. Splitting into two independent jobs means the text queue
|
||||
* saves+broadcasts its rows the moment IT finishes, regardless of how long
|
||||
* the media queue takes.
|
||||
* The lock for this conversation+lane is RELEASED here as soon as THIS lane's
|
||||
* worker job resolves — never after waiting on the other lane. That's the
|
||||
* core fix for "text menunggu image": previously one conversation batch made
|
||||
* ONE worker call with both text and media targets, and processing finished
|
||||
* only once BOTH lanes completed, so a fast text verdict sat unused until the
|
||||
* slow vision/image verdict was ready. Now each lane's results save+broadcast
|
||||
* the moment ITS job finishes, and the conversation lock for that lane is
|
||||
* freed independently.
|
||||
*
|
||||
* Returns whether the *caller* should schedule the next debounce pass for
|
||||
* this conversation's LANE.
|
||||
*/
|
||||
async function runQueueBatch(
|
||||
pool: typeof textWorkerPool,
|
||||
conversationKey: string,
|
||||
lane: AnalysisLane,
|
||||
messages: MessageRecord[],
|
||||
): Promise<boolean> {
|
||||
activeRequests++;
|
||||
if (lane === "media") activeMediaRequests++;
|
||||
else activeTextRequests++;
|
||||
|
||||
try {
|
||||
const result = (await pool.run({
|
||||
type: "batch",
|
||||
conversationKey,
|
||||
lane,
|
||||
messages,
|
||||
})) as AnalysisWorkerResponse;
|
||||
|
||||
@@ -222,12 +229,13 @@ async function runQueueBatch(
|
||||
}
|
||||
|
||||
if (!result.ok) {
|
||||
recordConversationBatchFailure(conversationKey);
|
||||
recordConversationBatchFailure(conversationKey, lane);
|
||||
|
||||
// Batch failed entirely -- fall back all messages to individual queue
|
||||
logger.warn(
|
||||
{
|
||||
conversationKey,
|
||||
lane,
|
||||
messageCount: messages.length,
|
||||
error: result.error,
|
||||
},
|
||||
@@ -243,6 +251,7 @@ async function runQueueBatch(
|
||||
logger.error(
|
||||
{
|
||||
conversationKey,
|
||||
lane,
|
||||
error: LAST_ERROR.value,
|
||||
messageCount: messages.length,
|
||||
messageIds: messages.map((m) => m.id),
|
||||
@@ -284,6 +293,7 @@ async function runQueueBatch(
|
||||
logger.warn(
|
||||
{
|
||||
conversationKey,
|
||||
lane,
|
||||
count: messagesForIndividualQueue.length,
|
||||
ids: messagesForIndividualQueue.map((m) => m.id),
|
||||
totalBatchSize: messages.length,
|
||||
@@ -297,6 +307,7 @@ async function runQueueBatch(
|
||||
logger.warn(
|
||||
{
|
||||
conversationKey,
|
||||
lane,
|
||||
count: apiFailedMessages.length,
|
||||
ids: apiFailedMessages.map((m) => m.id),
|
||||
},
|
||||
@@ -335,7 +346,7 @@ async function runQueueBatch(
|
||||
}
|
||||
|
||||
// Trigger conversation cooldown
|
||||
recordConversationBatchFailure(conversationKey);
|
||||
recordConversationBatchFailure(conversationKey, lane);
|
||||
const existingCooldown =
|
||||
conversationErrorCooldown.get(conversationKey) ?? 0;
|
||||
const newCooldown = Date.now() + config.AI_ANALYSIS_ERROR_COOLDOWN_MS;
|
||||
@@ -347,14 +358,14 @@ async function runQueueBatch(
|
||||
return false;
|
||||
}
|
||||
|
||||
resetConversationBatchFailures(conversationKey);
|
||||
resetConversationBatchFailures(conversationKey, lane);
|
||||
conversationErrorCooldown.delete(conversationKey);
|
||||
return true;
|
||||
} catch (error) {
|
||||
recordConversationBatchFailure(conversationKey);
|
||||
recordConversationBatchFailure(conversationKey, lane);
|
||||
|
||||
logger.warn(
|
||||
{ conversationKey, messageCount: messages.length },
|
||||
{ conversationKey, lane, messageCount: messages.length },
|
||||
"Batch threw exception -- routing all messages to individual fallback queue",
|
||||
);
|
||||
enqueueIndividualFallbacks(messages);
|
||||
@@ -370,6 +381,7 @@ async function runQueueBatch(
|
||||
logger.error(
|
||||
{
|
||||
conversationKey,
|
||||
lane,
|
||||
error: LAST_ERROR.value,
|
||||
stack: errorStack,
|
||||
messageCount: messages.length,
|
||||
@@ -384,59 +396,67 @@ async function runQueueBatch(
|
||||
return false;
|
||||
} finally {
|
||||
activeRequests--;
|
||||
if (lane === "media") activeMediaRequests--;
|
||||
else activeTextRequests--;
|
||||
}
|
||||
}
|
||||
|
||||
export async function processBatch(
|
||||
conversationKey: string,
|
||||
lane: AnalysisLane,
|
||||
messages: MessageRecord[],
|
||||
processingStartedAt: number,
|
||||
): Promise<void> {
|
||||
// Release this lane's lock immediately when there's nothing to do. The
|
||||
// messages array was already labelled with the lane it belongs to by the
|
||||
// scheduler (which fetched them from the DB), so an empty array means this
|
||||
// lane has no work — free it so the debounce can re-arm right away.
|
||||
if (messages.length === 0) {
|
||||
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
if (
|
||||
getConversationProcessingStartedAt(conversationKey, lane) ===
|
||||
processingStartedAt
|
||||
) {
|
||||
clearConversationProcessing(conversationKey, lane);
|
||||
}
|
||||
return;
|
||||
}
|
||||
const cooldownUntil = conversationErrorCooldown.get(conversationKey) ?? 0;
|
||||
if (Date.now() < cooldownUntil) {
|
||||
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
if (
|
||||
getConversationProcessingStartedAt(conversationKey, lane) ===
|
||||
processingStartedAt
|
||||
) {
|
||||
clearConversationProcessing(conversationKey, lane);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
// Split the batch itself — not just route it — so text and media never
|
||||
// share one worker call. A conversation batch commonly mixes plain-text
|
||||
// messages with an image/sticker from someone else; without this split,
|
||||
// ALL of it (including the plain-text messages) would ride along on the
|
||||
// media job and wait for vision analysis to finish. Each sub-batch is now
|
||||
// dispatched to its own pool AND handled independently below, so the text
|
||||
// queue's results land as soon as text analysis completes, full stop.
|
||||
const textMessages = messages.filter((m) => !hasMediaContent(m));
|
||||
const mediaMessages = messages.filter((m) => hasMediaContent(m));
|
||||
|
||||
const jobs: Promise<boolean>[] = [];
|
||||
if (textMessages.length > 0) {
|
||||
jobs.push(runQueueBatch(textWorkerPool, conversationKey, textMessages));
|
||||
}
|
||||
if (mediaMessages.length > 0) {
|
||||
jobs.push(runQueueBatch(mediaWorkerPool, conversationKey, mediaMessages));
|
||||
}
|
||||
|
||||
const outcomes = await Promise.allSettled(jobs);
|
||||
const shouldScheduleNext = outcomes.every(
|
||||
(o) => o.status === "fulfilled" && o.value,
|
||||
const result = await runQueueBatch(
|
||||
lane === "media" ? mediaWorkerPool : textWorkerPool,
|
||||
conversationKey,
|
||||
lane,
|
||||
messages,
|
||||
);
|
||||
|
||||
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
// Release THIS lane's lock now — the other lane (if any) is dispatched
|
||||
// separately by the scheduler and owns its own lock. The old code awaited
|
||||
// BOTH lanes (Promise.allSettled) before releasing the single conversation
|
||||
// lock, so the text sub-batch of a conversation blocked its own lock until
|
||||
// the slow media sub-batch finished. Now each lane is independent: the text
|
||||
// lane frees its lock and re-schedules the moment the text worker returns.
|
||||
if (
|
||||
getConversationProcessingStartedAt(conversationKey, lane) ===
|
||||
processingStartedAt
|
||||
) {
|
||||
clearConversationProcessing(conversationKey, lane);
|
||||
}
|
||||
if (shouldScheduleNext) {
|
||||
|
||||
if (result) {
|
||||
setImmediate(() => {
|
||||
// Dynamic import to avoid circular dependency at module scope
|
||||
// Dynamic import to avoid circular dependency at module scope.
|
||||
// Re-schedule ONLY this lane — the other lane schedules itself.
|
||||
import("./batchScheduler.js").then((m) =>
|
||||
m.scheduleConversationAnalysis(conversationKey),
|
||||
m.scheduleConversationAnalysis(conversationKey, lane),
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
import { type AnalysisLane, splitMessagesByLane } from "./analysisLanes.js";
|
||||
import { computeBudgetOverflowMessages } from "./batchBudget.js";
|
||||
import {
|
||||
pickBatchWithinBudget,
|
||||
processBatch,
|
||||
@@ -9,12 +11,14 @@ import {
|
||||
skipAnalysisUserMessages,
|
||||
} from "./batchProcessor.js";
|
||||
import {
|
||||
clearConversationProcessing,
|
||||
conversationConsecutiveErrors,
|
||||
conversationDebounceTimers,
|
||||
conversationErrorCooldown,
|
||||
conversationProcessing,
|
||||
getConversationProcessingStartedAt,
|
||||
isConversationProcessingLocked,
|
||||
MAX_CONSECUTIVE_ERRORS,
|
||||
setConversationProcessing,
|
||||
} from "./conversationState.js";
|
||||
|
||||
const logger = createChildLogger("batch-scheduler");
|
||||
@@ -23,18 +27,35 @@ const logger = createChildLogger("batch-scheduler");
|
||||
// Scheduling
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Timer key namespaced by lane so one conversation can hold a text timer AND
|
||||
* a media timer independently.
|
||||
*/
|
||||
function timerKey(conversationKey: string, lane: AnalysisLane): string {
|
||||
return `${conversationKey}::${lane}`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Schedules a debounced analysis run for a conversation.
|
||||
*
|
||||
* `lane` optional:
|
||||
* - With a lane: takes that lane's processing lock; if the SAME lane is
|
||||
* already processing, skip. The OTHER lane's lock does not block this one —
|
||||
* text and media of one conversation never block each other.
|
||||
* - Without a lane (recovery worker / whole-conversation): takes BOTH lane
|
||||
* locks (each lock independently) and dispatches both lanes concurrently.
|
||||
* Each lane releases its own lock when its worker job finishes.
|
||||
*
|
||||
* The async work inside setTimeout is wrapped in an explicit .catch() so
|
||||
* DB errors don't produce unhandled promise rejections. Uses a unified
|
||||
* single-timer path: always clear-and-reset one timer per conversation key
|
||||
* single-timer path: always clear-and-reset one timer per conversation+lane
|
||||
* regardless of whether a cooldown is active.
|
||||
*/
|
||||
export function scheduleConversationAnalysis(conversationKey: string): void {
|
||||
if (isConversationProcessingLocked(conversationKey)) {
|
||||
return;
|
||||
}
|
||||
export function scheduleConversationAnalysis(
|
||||
conversationKey: string,
|
||||
lane?: AnalysisLane,
|
||||
): void {
|
||||
const lanesToSchedule: AnalysisLane[] = lane ? [lane] : ["text", "media"];
|
||||
|
||||
const convoCooldown = conversationErrorCooldown.get(conversationKey) ?? 0;
|
||||
const convoErrors = conversationConsecutiveErrors.get(conversationKey) ?? 0;
|
||||
@@ -45,27 +66,49 @@ export function scheduleConversationAnalysis(conversationKey: string): void {
|
||||
}
|
||||
|
||||
// Unified delay: honour the cooldown window if active, otherwise use the
|
||||
// normal debounce interval. Always clear-and-reset so only ONE timer is
|
||||
// ever pending per conversation key regardless of call source.
|
||||
// normal debounce interval. Always clear-and-reset so only ONE timer is
|
||||
// ever pending per conversation+lane regardless of call source.
|
||||
const now = Date.now();
|
||||
const delayMs =
|
||||
convoCooldown > now
|
||||
? convoCooldown - now + 500
|
||||
: config.AI_ANALYSIS_DEBOUNCE_MS;
|
||||
|
||||
const existingTimer = conversationDebounceTimers.get(conversationKey);
|
||||
for (const targetLane of lanesToSchedule) {
|
||||
if (isConversationProcessingLocked(conversationKey, targetLane)) {
|
||||
continue;
|
||||
}
|
||||
scheduleLaneTimer(conversationKey, targetLane, delayMs);
|
||||
}
|
||||
}
|
||||
|
||||
function scheduleLaneTimer(
|
||||
conversationKey: string,
|
||||
lane: AnalysisLane,
|
||||
delayMs: number,
|
||||
): void {
|
||||
const tKey = timerKey(conversationKey, lane);
|
||||
const existingTimer = conversationDebounceTimers.get(tKey);
|
||||
if (existingTimer) {
|
||||
clearTimeout(existingTimer);
|
||||
}
|
||||
|
||||
const timer = setTimeout(() => {
|
||||
conversationDebounceTimers.delete(conversationKey);
|
||||
conversationDebounceTimers.delete(tKey);
|
||||
|
||||
if (isConversationProcessingLocked(conversationKey)) {
|
||||
if (isConversationProcessingLocked(conversationKey, lane)) {
|
||||
logger.warn(
|
||||
{ conversationKey, lane, tKey },
|
||||
"scheduleLaneTimer: lane locked, skipping dispatch",
|
||||
);
|
||||
return;
|
||||
}
|
||||
const processingStartedAt = Date.now();
|
||||
conversationProcessing.set(conversationKey, processingStartedAt);
|
||||
setConversationProcessing(conversationKey, lane, processingStartedAt);
|
||||
logger.debug(
|
||||
{ conversationKey, lane, processingStartedAt },
|
||||
"scheduleLaneTimer: lock acquired, dispatching batch fetch",
|
||||
);
|
||||
|
||||
messageStore
|
||||
.getPendingMessagesByConversation(
|
||||
@@ -73,24 +116,22 @@ export function scheduleConversationAnalysis(conversationKey: string): void {
|
||||
config.AI_ANALYSIS_MAX_BATCH_SIZE,
|
||||
)
|
||||
.then(async (messages: MessageRecord[]) => {
|
||||
if (messages.length === 0) {
|
||||
if (
|
||||
conversationProcessing.get(conversationKey) === processingStartedAt
|
||||
) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
// Filter to THIS lane only. The DB fetch is lane-agnostic (a
|
||||
// conversation key can have both text and media pending); each lane
|
||||
// picks its own subset so text and media never share a worker job.
|
||||
const { [lane]: laneMessages } = splitMessagesByLane(messages);
|
||||
if (laneMessages.length === 0) {
|
||||
// No work for this lane — the other lane (if scheduled) owns the
|
||||
// rest. Clear this lane's lock so the debounce can re-arm.
|
||||
releaseLaneSlot(conversationKey, lane, processingStartedAt);
|
||||
return;
|
||||
}
|
||||
|
||||
const processableMessages = await skipAnalysisUserMessages(
|
||||
await skipAgeRestrictedMessages(messages),
|
||||
await skipAgeRestrictedMessages(laneMessages),
|
||||
);
|
||||
if (processableMessages.length === 0) {
|
||||
if (
|
||||
conversationProcessing.get(conversationKey) === processingStartedAt
|
||||
) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
releaseLaneSlot(conversationKey, lane, processingStartedAt);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -107,6 +148,7 @@ export function scheduleConversationAnalysis(conversationKey: string): void {
|
||||
logger.warn(
|
||||
{
|
||||
conversationKey,
|
||||
lane,
|
||||
messageId: processableMessages[0]?.id,
|
||||
tokenBudget: config.AI_ANALYSIS_MAX_TARGET_TOKENS,
|
||||
},
|
||||
@@ -114,17 +156,76 @@ export function scheduleConversationAnalysis(conversationKey: string): void {
|
||||
);
|
||||
}
|
||||
|
||||
return processBatch(conversationKey, trimmed, processingStartedAt);
|
||||
// Un-claim messages that did NOT make it into the trimmed batch.
|
||||
// getPendingMessagesByConversation() flips EVERY fetched pending row
|
||||
// to `processing`; pickBatchWithinBudget() may then stop early on the
|
||||
// token budget, leaving the tail rows stuck in `processing` forever
|
||||
// (recovery only reverts rows older than 120s, and these keep getting
|
||||
// re-claimed each wave). Return them to `pending` so the next wave
|
||||
// picks them up instead of leaking processing slots.
|
||||
const unclaimed = computeBudgetOverflowMessages(
|
||||
processableMessages,
|
||||
trimmed,
|
||||
);
|
||||
if (unclaimed.length > 0) {
|
||||
const unclaimedRows = await messageStore
|
||||
.updateMessagesAIAnalysisBulk(
|
||||
unclaimed.map((msg) => ({
|
||||
messageId: msg.id,
|
||||
result: {
|
||||
status: "pending",
|
||||
flags: null,
|
||||
score: null,
|
||||
analysis: null,
|
||||
categories: null,
|
||||
severity: null,
|
||||
confidence: null,
|
||||
recommendedAction: null,
|
||||
analyzedAt: null,
|
||||
error: null,
|
||||
},
|
||||
})),
|
||||
)
|
||||
.catch((err: unknown) => {
|
||||
logger.error(
|
||||
{
|
||||
conversationKey,
|
||||
lane,
|
||||
count: unclaimed.length,
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
},
|
||||
"Failed to un-claim budget-overflow messages back to pending",
|
||||
);
|
||||
return null;
|
||||
});
|
||||
if (unclaimedRows) {
|
||||
logger.debug(
|
||||
{
|
||||
conversationKey,
|
||||
lane,
|
||||
unclaimedCount: unclaimed.length,
|
||||
},
|
||||
"Returned budget-overflow messages to pending for next wave",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// processBatch releases THIS lane's lock the moment its worker job
|
||||
// finishes and re-schedules the same lane — independent of the other
|
||||
// lane's (possibly much slower) media batch.
|
||||
return processBatch(
|
||||
conversationKey,
|
||||
lane,
|
||||
trimmed,
|
||||
processingStartedAt,
|
||||
);
|
||||
})
|
||||
.catch((err: unknown) => {
|
||||
if (
|
||||
conversationProcessing.get(conversationKey) === processingStartedAt
|
||||
) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
releaseLaneSlot(conversationKey, lane, processingStartedAt);
|
||||
logger.error(
|
||||
{
|
||||
conversationKey,
|
||||
lane,
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
},
|
||||
"Failed to fetch or dispatch pending messages for scheduled analysis",
|
||||
@@ -132,5 +233,23 @@ export function scheduleConversationAnalysis(conversationKey: string): void {
|
||||
});
|
||||
}, delayMs);
|
||||
|
||||
conversationDebounceTimers.set(conversationKey, timer);
|
||||
conversationDebounceTimers.set(tKey, timer);
|
||||
}
|
||||
|
||||
/**
|
||||
* Clears the processing lock for a lane, but ONLY if this timer still owns it
|
||||
* (processingStartedAt matches). Guards against clearing a newer slot that was
|
||||
* taken after this timer's window expired.
|
||||
*/
|
||||
function releaseLaneSlot(
|
||||
conversationKey: string,
|
||||
lane: AnalysisLane,
|
||||
processingStartedAt: number,
|
||||
): void {
|
||||
if (
|
||||
getConversationProcessingStartedAt(conversationKey, lane) ===
|
||||
processingStartedAt
|
||||
) {
|
||||
clearConversationProcessing(conversationKey, lane);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -2,7 +2,7 @@ import { existsSync } from "node:fs";
|
||||
import { availableParallelism } from "node:os";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { Piscina } from "piscina";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { LRUCache } from "lru-cache";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { LAST_ERROR } from "./moderationState.js";
|
||||
|
||||
/**
|
||||
@@ -24,6 +24,19 @@ import { LAST_ERROR } from "./moderationState.js";
|
||||
* - Alert system: `CircuitBreakerAlert` type, `fireAlert()`, and
|
||||
* `onCircuitBreakerAlert()` for pluggable handler registration.
|
||||
*
|
||||
* ## Processing lanes (2026-09-24)
|
||||
* A conversation batch splits into a **text lane** (messages with no media)
|
||||
* and a **media lane** (messages with attachments/stickers/embeds). The two
|
||||
* lanes are dispatched to separate Piscina pools and MUST NOT block each
|
||||
* other: a fast text sub-batch must be free to finish while the slow
|
||||
* vision/media sub-batch of the SAME conversation is still running.
|
||||
*
|
||||
* The lock is therefore per-lane: `conversationProcessing` maps a
|
||||
* conversation key to its current processing record which carries the lane
|
||||
* name. `isConversationProcessingLocked(key, lane)` reports locked only when
|
||||
* the SAME lane (or all lanes when lane is omitted) is active — a media
|
||||
* sub-batch in flight never blocks scheduling the text sub-batch.
|
||||
*
|
||||
* ## Relationship with moderationState.ts
|
||||
* - `moderationState.ts` owns **infrastructure references** (event broadcaster,
|
||||
* Discord client), the auto-delete guard, the `LAST_ERROR` tracker, and
|
||||
@@ -33,6 +46,14 @@ import { LAST_ERROR } from "./moderationState.js";
|
||||
* - These are **separate concerns** — do not merge them.
|
||||
*/
|
||||
|
||||
/** Processing lanes for conversation analysis. */
|
||||
export type AnalysisLane = "text" | "media";
|
||||
|
||||
export const ANALYSIS_LANES: readonly AnalysisLane[] = [
|
||||
"text",
|
||||
"media",
|
||||
] as const;
|
||||
|
||||
const logger = createChildLogger("conversation-state");
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -60,23 +81,105 @@ export const conversationDebounceTimers = new LRUCache<string, NodeJS.Timeout>({
|
||||
},
|
||||
});
|
||||
|
||||
/** Timestamp of when processing started per conversation key. */
|
||||
export const conversationProcessing = new LRUCache<string, number>({
|
||||
max: 10000,
|
||||
});
|
||||
/**
|
||||
* Per-conversation processing lock, keyed by lane.
|
||||
*
|
||||
* A conversation can hold TWO locks at once — one for its text sub-batch and
|
||||
* one for its media sub-batch — because the two lanes run on separate pools
|
||||
* and finish independently. The value is a partial record of lane →
|
||||
* startedAt; clearing one lane leaves the other lane's lock intact.
|
||||
*/
|
||||
export const conversationProcessing = new LRUCache<
|
||||
string,
|
||||
Partial<Record<AnalysisLane, number>>
|
||||
>({ max: 10000 });
|
||||
|
||||
/**
|
||||
* Locks a conversation for the given lane.
|
||||
* The same conversation can be locked in both lanes simultaneously (text and
|
||||
* media sub-batches run independently); locking an already-locked lane
|
||||
* replaces its startedAt (last writer wins, matching the old single-lock
|
||||
* semantics).
|
||||
*/
|
||||
export function setConversationProcessing(
|
||||
conversationKey: string,
|
||||
lane: AnalysisLane,
|
||||
startedAt: number,
|
||||
): void {
|
||||
const record = conversationProcessing.get(conversationKey) ?? {};
|
||||
conversationProcessing.set(conversationKey, { ...record, [lane]: startedAt });
|
||||
}
|
||||
|
||||
/**
|
||||
* Releases the processing lock for a conversation in a SINGLE lane.
|
||||
* The other lane's lock (if any) is preserved.
|
||||
*/
|
||||
export function clearConversationProcessing(
|
||||
conversationKey: string,
|
||||
lane: AnalysisLane,
|
||||
): void {
|
||||
const record = conversationProcessing.get(conversationKey);
|
||||
if (!record) return;
|
||||
const next = { ...record };
|
||||
delete next[lane];
|
||||
if (Object.keys(next).length === 0) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
} else {
|
||||
conversationProcessing.set(conversationKey, next);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Clears the processing lock for a conversation regardless of lane.
|
||||
* Used by the recovery worker when a lock is stale. If only ONE lane of a
|
||||
* two-lane processing conversation is stale, prefer clearConversationProcessing
|
||||
* with the specific lane to keep the healthy lane's lock intact.
|
||||
*/
|
||||
export function clearConversationProcessingAll(conversationKey: string): void {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the startedAt for a conversation in a lane, or undefined.
|
||||
* Consumers use this to verify a processing slot is still owned by them
|
||||
* before releasing it (guards against clearing a newer slot).
|
||||
*/
|
||||
export function getConversationProcessingStartedAt(
|
||||
conversationKey: string,
|
||||
lane: AnalysisLane,
|
||||
): number | undefined {
|
||||
return conversationProcessing.get(conversationKey)?.[lane];
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Conversation lock helper
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Reports whether the conversation is currently processing.
|
||||
*
|
||||
* When `lane` is provided, only that lane's lock counts — a media sub-batch
|
||||
* in flight does NOT lock the text lane, so the text lane can be scheduled
|
||||
* and vice versa. When `lane` is omitted, any active lane locks it (used by
|
||||
* recovery/individual fallback which must not race ANY batch work).
|
||||
*/
|
||||
export function isConversationProcessingLocked(
|
||||
conversationKey: string,
|
||||
lane?: AnalysisLane,
|
||||
): boolean {
|
||||
const startedAt = conversationProcessing.get(conversationKey);
|
||||
return Boolean(
|
||||
startedAt &&
|
||||
Date.now() - startedAt < config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS,
|
||||
);
|
||||
const now = Date.now();
|
||||
if (lane) {
|
||||
const startedAt = conversationProcessing.get(conversationKey)?.[lane];
|
||||
return Boolean(
|
||||
startedAt && now - startedAt < config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS,
|
||||
);
|
||||
}
|
||||
const record = conversationProcessing.get(conversationKey);
|
||||
if (!record) return false;
|
||||
return ANALYSIS_LANES.some((l) => {
|
||||
const s = record[l];
|
||||
return Boolean(s && now - s < config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS);
|
||||
});
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -86,6 +189,7 @@ export function isConversationProcessingLocked(
|
||||
export type CircuitBreakerAlert = {
|
||||
type: "conversation_cb" | "individual_cb" | "sustained_error";
|
||||
conversationKey?: string;
|
||||
lane?: AnalysisLane;
|
||||
consecutiveErrors: number;
|
||||
message: string;
|
||||
lastError?: string | null;
|
||||
@@ -117,7 +221,10 @@ export function fireAlert(alert: CircuitBreakerAlert): void {
|
||||
// Circuit breaker helpers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export function recordConversationBatchFailure(conversationKey: string): void {
|
||||
export function recordConversationBatchFailure(
|
||||
conversationKey: string,
|
||||
lane?: AnalysisLane,
|
||||
): void {
|
||||
const nextCount =
|
||||
(conversationConsecutiveErrors.get(conversationKey) ?? 0) + 1;
|
||||
conversationConsecutiveErrors.set(conversationKey, nextCount);
|
||||
@@ -130,6 +237,7 @@ export function recordConversationBatchFailure(conversationKey: string): void {
|
||||
fireAlert({
|
||||
type: "conversation_cb",
|
||||
conversationKey,
|
||||
lane,
|
||||
consecutiveErrors: nextCount,
|
||||
message: `Conversation ${conversationKey} circuit breaker triggered after ${nextCount} consecutive errors`,
|
||||
lastError: LAST_ERROR.value,
|
||||
@@ -138,6 +246,9 @@ export function recordConversationBatchFailure(conversationKey: string): void {
|
||||
}
|
||||
}
|
||||
|
||||
export function resetConversationBatchFailures(conversationKey: string): void {
|
||||
export function resetConversationBatchFailures(
|
||||
conversationKey: string,
|
||||
_lane?: AnalysisLane,
|
||||
): void {
|
||||
conversationConsecutiveErrors.delete(conversationKey);
|
||||
}
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { and, desc, eq, sql } from "drizzle-orm";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { getDatabase } from "../../shared/database/drizzle.js";
|
||||
import { messagesTable } from "../../shared/database/schema.js";
|
||||
import { updateChannelCulture } from "./channelCultureStore.js";
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
import OpenAI from "openai";
|
||||
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { cleanContent } from "./textSignals.js";
|
||||
|
||||
const log = createChildLogger("embedding-client");
|
||||
|
||||
@@ -1,11 +1,35 @@
|
||||
// ── Single-pass LLM pipeline exports ─────────────────────────────────────
|
||||
/**
|
||||
* 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 {
|
||||
AnalysisInput,
|
||||
AIRecommendedAction,
|
||||
AISeverity,
|
||||
AIStatus,
|
||||
AnalysisQueueStatus,
|
||||
AnalysisResult,
|
||||
MessageBatch,
|
||||
WorkerConfig,
|
||||
} from "./ai-analysis-worker.js";
|
||||
export { startPendingAIAnalysisWorker } from "./aiAnalyzer.js";
|
||||
export { sanitizeDiscordTokens } from "./discordTokens.js";
|
||||
export { runModerationAnalysis } from "./moderationOrchestrator.js";
|
||||
export { buildSystemPrompt } from "./moderationPrompt.js";
|
||||
} 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";
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { LRUCache } from "lru-cache";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type {
|
||||
AnalysisResult,
|
||||
|
||||
@@ -1,563 +0,0 @@
|
||||
/**
|
||||
* jevAnalyzer.ts
|
||||
*
|
||||
* Jev (TypeSafe System One, `oc/jev-1.13-free` via 9router `/v1/systemone`)
|
||||
* — the PRIMARY analyzer for text-only moderation sub-batches. The existing
|
||||
* LLM (`llmChat`) stays as the fallback for any message Jev cannot decide
|
||||
* confidently (see the acceptance gate) and for media batches (Jev is
|
||||
* decision-only, no image input).
|
||||
*
|
||||
* CRITICAL framing rule (verified 2026-09-23, live probes):
|
||||
* Jev is a System One model — it evaluates typed questions against a STATE.
|
||||
* Feeding it the chat-optimized `SYSTEM_RULES` verbatim INSIDE chat-style
|
||||
* XML (`<messages_to_analyze>`, `<location_context>`, …) makes it
|
||||
* pattern-match the structure and return CONFIDENTLY WRONG verdicts
|
||||
* (flagged clean messages at confidence 0.98 in a probe — would pass any
|
||||
* naive gate and could auto-delete innocent content).
|
||||
*
|
||||
* The state MUST be declarative facts:
|
||||
* OBJEK PENILAIAN / PESAN: `- Pesan "<id>" dari "<user>": "<content>"` /
|
||||
* KEBIJAKAN as statements / KONTEKS as statements
|
||||
* and the questions phrased as "is this true" / "classify this" against
|
||||
* those facts. With that framing the same 4-message probe returned 4/4
|
||||
* correct verdicts at confidence 1.0, including the SARA zero-tolerance
|
||||
* case and the technical-clean case.
|
||||
*
|
||||
* The distilled `JEV_POLICY` below is a compact declarative summary of the
|
||||
* full chat policy (`prompts/rules.ts` SYSTEM_RULES). It is deliberately
|
||||
* kept short (~300 tokens) — the LLM keeps the full 12k-char policy; Jev
|
||||
* triages on the core axes, and anything it can't decide confidently falls
|
||||
* back to the LLM. Keep this block in sync when SYSTEM_RULES changes.
|
||||
*/
|
||||
|
||||
import type { Question, Questions } from "@typesafe-ai/sdk";
|
||||
import { choice, noul, TypeSafeClient } from "@typesafe-ai/sdk";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { incrementCounterBy } from "../gateway-metrics/index.js";
|
||||
import type { AnalysisResult } from "../message-capture/types.js";
|
||||
|
||||
const log = createChildLogger("jev-analyzer");
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Policy + vocab
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Distilled declarative policy for Jev. DERIVED from `SYSTEM_RULES`
|
||||
* (prompts/rules.ts) — update this when the full policy changes. Kept as
|
||||
* factual statements, NOT instructions (System One evaluates truth).
|
||||
*/
|
||||
export const JEV_POLICY = `KEBIJAKAN SERVER (fakta yang berlaku):
|
||||
- Kata vulgar anatomi (kontol, memek, tit, dick, dll) = pelanggaran berat, tanpa kecuali.
|
||||
- SARA / penistaan agama / parodi ayat / mockery tokoh agama / provokasi antar-agama = pelanggaran berat.
|
||||
- Promosi atau diskusi LGBT = pelanggaran berat (zero-tolerance).
|
||||
- Diskusi Israel/Palestina/Yahudi = pelanggaran berat (zero-tolerance).
|
||||
- Hinaan terarah ke orang (harassment), seksisme, ageisme, diskriminasi fisik = pelanggaran.
|
||||
- Konten seksual eksplisit / ajakan seksual / fetish / lolicon-shota = pelanggaran.
|
||||
- Judi, narkoba, scam, doxxing, ancaman kekerasan, self-harm, child safety, konten ilegal = pelanggaran.
|
||||
- Teknik evasi (zalgo, leetspeak, regional indicator, simbol acak) yang menyembunyikan kata terlarang = pelanggaran.
|
||||
- Spam berulang / promosi = pelanggaran ringan.
|
||||
- Memancing konflik (conflict instigation) = pelanggaran ringan.
|
||||
- Username ofensif saja (isi pesan bersih) = peringatan ringan, BUKAN hapus pesan.
|
||||
- Percakapan teknis/normal, slang santai (anjay, wkwk, gaskeun, njir), typo, panggilan akrab (bang, kak, dek), ekspresi religius normal (astaghfirullah, alhamdulillah), istilah anime (waifu, wibu), lirik/kutipan, makian ke benda mati = BUKAN pelanggaran.
|
||||
- Teks acak (kode, log, stack trace, output API, cuplikan UI) = BUKAN pelanggaran.
|
||||
- Setiap pesan dinilai dari isinya sendiri; konteks percakapan dapat memengaruhi interpretasi, bukan menggantikan isi.`;
|
||||
|
||||
/** Choice labels must stay in sync with `AIRecommendedAction` (moderation-types). */
|
||||
export const JEV_ACTIONS = [
|
||||
"none",
|
||||
"monitor",
|
||||
"warn",
|
||||
"review",
|
||||
"delete",
|
||||
"escalate",
|
||||
] as const;
|
||||
|
||||
/** Category choices — the moderation category vocabulary (kept tight). */
|
||||
export const JEV_CATEGORIES = [
|
||||
"none",
|
||||
"harassment",
|
||||
"hate_speech",
|
||||
"sara",
|
||||
"sexual_content",
|
||||
"vulgar_language",
|
||||
"sexual_deviation",
|
||||
"self_harm",
|
||||
"violence",
|
||||
"illegal_content",
|
||||
"gambling",
|
||||
"drugs",
|
||||
"scam",
|
||||
"spam",
|
||||
"conflict_instigation",
|
||||
"offensive_username",
|
||||
"other",
|
||||
] as const;
|
||||
|
||||
export const JEV_STATUSES = ["clean", "warn", "flagged"] as const;
|
||||
export const JEV_SEVERITIES = [
|
||||
"none",
|
||||
"low",
|
||||
"medium",
|
||||
"high",
|
||||
"critical",
|
||||
] as const;
|
||||
|
||||
/** Policy version stamped on every Jev verdict (cache/DB provenance). */
|
||||
export const JEV_POLICY_VERSION = "jev-systemone-2026-09-23";
|
||||
|
||||
/** Lazily-built SDK client (config resolves at first use). */
|
||||
let client: TypeSafeClient | null = null;
|
||||
let clientKey = "";
|
||||
|
||||
function getClient(): TypeSafeClient | null {
|
||||
if (!config.AI_LLM_JEV_API_KEY) return null;
|
||||
if (!client || clientKey !== config.AI_LLM_JEV_API_KEY) {
|
||||
client = new TypeSafeClient({
|
||||
apiKey: config.AI_LLM_JEV_API_KEY,
|
||||
baseURL: config.AI_LLM_JEV_BASE_URL,
|
||||
defaultModel: config.AI_LLM_JEV_MODEL,
|
||||
timeout: config.AI_LLM_JEV_TIMEOUT_MS,
|
||||
retry: { maxRetries: 0 }, // pipeline owns retries/abort
|
||||
logLevel: "warn",
|
||||
});
|
||||
clientKey = config.AI_LLM_JEV_API_KEY;
|
||||
log.info(
|
||||
{
|
||||
baseURL: config.AI_LLM_JEV_BASE_URL,
|
||||
model: config.AI_LLM_JEV_MODEL,
|
||||
},
|
||||
"Jev analyzer client initialized",
|
||||
);
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Question + state builders
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export interface JevTarget {
|
||||
/** Message id — echoed verbatim into analysis/result. */
|
||||
id: string;
|
||||
/** Display name shown to the model (username). */
|
||||
user: string;
|
||||
/** Content to evaluate (truncated by caller). */
|
||||
content: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Build per-message questions, keyed by the message id itself so the
|
||||
* questions read naturally against the declarative state facts
|
||||
* ("pesan <id> melanggar kebijakan server"). Five questions per message.
|
||||
*/
|
||||
export function buildJevQuestions(targets: JevTarget[]): Questions {
|
||||
const questions: Record<string, Question> = {};
|
||||
for (const t of targets) {
|
||||
const k = t.id;
|
||||
questions[`${k}__v`] = noul(
|
||||
`Pesan ${JSON.stringify(t.id)} dari ${JSON.stringify(t.user)} melanggar kebijakan server`,
|
||||
);
|
||||
questions[`${k}__status`] = choice(
|
||||
`Klasifikasi pesan ${JSON.stringify(t.id)} dari ${JSON.stringify(t.user)}`,
|
||||
{
|
||||
clean: "tidak melanggar kebijakan",
|
||||
warn: "pelanggaran ringan",
|
||||
flagged: "melanggar kebijakan",
|
||||
},
|
||||
);
|
||||
questions[`${k}__severity`] = choice(
|
||||
`Tingkat keparahan pelanggaran pesan ${JSON.stringify(t.id)}`,
|
||||
{
|
||||
none: "tidak ada pelanggaran",
|
||||
low: "ringan",
|
||||
medium: "sedang",
|
||||
high: "berat",
|
||||
critical: "kritis/darurat",
|
||||
},
|
||||
);
|
||||
questions[`${k}__category`] = choice(
|
||||
`Kategori utama pelanggaran pesan ${JSON.stringify(t.id)}`,
|
||||
Object.fromEntries(JEV_CATEGORIES.map((c) => [c, null])),
|
||||
);
|
||||
questions[`${k}__action`] = choice(
|
||||
`Tindakan moderasi yang tepat untuk pesan ${JSON.stringify(t.id)}`,
|
||||
{
|
||||
none: "tidak ada tindakan",
|
||||
monitor: "pantau",
|
||||
warn: "beri peringatan",
|
||||
review: "tinjau manual",
|
||||
delete: "hapus pesan",
|
||||
escalate: "eskalasi",
|
||||
},
|
||||
);
|
||||
}
|
||||
return questions as Questions;
|
||||
}
|
||||
|
||||
export interface JevBatchContext {
|
||||
/** Context block (location/conversation) as raw XML or prose — stripped to facts. */
|
||||
contextBlock: string;
|
||||
/** `<web_searches>` XML block (may be ""). */
|
||||
webSearchBlock: string;
|
||||
/** `<term_glossary>` XML block (may be ""). */
|
||||
glossaryBlock: string;
|
||||
/** Raw channel culture summary (may be undefined). */
|
||||
channelCulture?: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Strip XML/HTML tags from a raw block and collapse whitespace so it can be
|
||||
* restated as plain factual prose in the declarative state. Empty after
|
||||
* stripping → omitted from the state.
|
||||
*/
|
||||
function stripToFacts(block: string, maxLength: number): string | null {
|
||||
const cleaned = block
|
||||
.replace(/<[^>]+>/g, " ")
|
||||
.replace(/\s+/g, " ")
|
||||
.trim();
|
||||
if (!cleaned) return null;
|
||||
return JSON.stringify(cleaned.slice(0, maxLength));
|
||||
}
|
||||
|
||||
/**
|
||||
* Build the declarative `state` payload. NO chat/XML scaffolding — plain
|
||||
* factual statements (see the framing rule above; chat-style injection
|
||||
* makes Jev confidently wrong).
|
||||
*/
|
||||
export function buildJevState(
|
||||
targets: JevTarget[],
|
||||
ctx: JevBatchContext,
|
||||
correctedExamples = "",
|
||||
): string {
|
||||
const facts = targets.map(
|
||||
(t) =>
|
||||
`- Pesan ${JSON.stringify(t.id)} dari ${JSON.stringify(t.user)}: ${JSON.stringify(t.content)}`,
|
||||
);
|
||||
const parts = [
|
||||
`OBJEK PENILAIAN: ${targets.length} pesan dari server Discord.`,
|
||||
"PESAN:",
|
||||
...facts,
|
||||
JEV_POLICY,
|
||||
];
|
||||
|
||||
// Conversation/who context as facts (declarative, not instructions).
|
||||
const extraFacts: string[] = [];
|
||||
if (ctx.channelCulture) {
|
||||
extraFacts.push(
|
||||
`KULTUR CHANNEL (fakta): ${JSON.stringify(ctx.channelCulture.slice(0, 800))}`,
|
||||
);
|
||||
}
|
||||
const contextFacts = stripToFacts(ctx.contextBlock, 1200);
|
||||
if (contextFacts) extraFacts.push(`KONTEKS: ${contextFacts}`);
|
||||
const webSearchFacts = stripToFacts(ctx.webSearchBlock, 1500);
|
||||
if (webSearchFacts) extraFacts.push(`HASIL PENCARIAN WEB: ${webSearchFacts}`);
|
||||
const glossaryFacts = stripToFacts(ctx.glossaryBlock, 800);
|
||||
if (glossaryFacts) extraFacts.push(`GLOSARIUM: ${glossaryFacts}`);
|
||||
const correctionFacts = stripToFacts(correctedExamples, 800);
|
||||
if (correctionFacts)
|
||||
extraFacts.push(`KOREKSI SEBELUMNYA: ${correctionFacts}`);
|
||||
if (extraFacts.length > 0) parts.push(...extraFacts);
|
||||
|
||||
return parts.join("\n");
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Acceptance gate + mapper
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/** Shape of the raw `answers` map returned by `systemOne`. */
|
||||
export type JevAnswers = Record<
|
||||
string,
|
||||
| { type: "noul"; noul: number }
|
||||
| {
|
||||
type: "choice";
|
||||
choice: string;
|
||||
confidence: number;
|
||||
probabilities?: Record<string, number>;
|
||||
}
|
||||
>;
|
||||
|
||||
/** Per-message answer subset (nullable until validated — `answersOf`). */
|
||||
interface JevMessageAnswers {
|
||||
v?: { type: "noul"; noul: number };
|
||||
status?: { type: "choice"; choice: string; confidence: number };
|
||||
severity?: { type: "choice"; choice: string };
|
||||
category?: { type: "choice"; choice: string };
|
||||
action?: { type: "choice"; choice: string };
|
||||
}
|
||||
|
||||
/** Reads the per-message answer subset by id, missing → undefined. */
|
||||
function answersOf(answers: JevAnswers, id: string): JevMessageAnswers {
|
||||
return {
|
||||
v: answers[`${id}__v`] as JevMessageAnswers["v"],
|
||||
status: answers[`${id}__status`] as JevMessageAnswers["status"],
|
||||
severity: answers[`${id}__severity`] as JevMessageAnswers["severity"],
|
||||
category: answers[`${id}__category`] as JevMessageAnswers["category"],
|
||||
action: answers[`${id}__action`] as JevMessageAnswers["action"],
|
||||
};
|
||||
}
|
||||
|
||||
/** Optionally-typed accessor for a choice answer's label ("" when missing). */
|
||||
function labelOf(a: { type: "choice"; choice: string } | undefined): string {
|
||||
return a?.type === "choice" ? a.choice : "";
|
||||
}
|
||||
|
||||
/** Set form of the vocab arrays for O(1) membership tests. */
|
||||
const JEV_STATUS_SET = new Set<string>(JEV_STATUSES);
|
||||
const JEV_SEVERITY_SET = new Set<string>(JEV_SEVERITIES);
|
||||
const JEV_CATEGORY_SET = new Set<string>(JEV_CATEGORIES);
|
||||
const JEV_ACTION_SET = new Set<string>(JEV_ACTIONS);
|
||||
|
||||
/**
|
||||
* Decide per-message Jev acceptance. Requires ALL five questions present
|
||||
* with valid labels and CROSS-CONSISTENT semantics:
|
||||
* - status choice confidence >= threshold
|
||||
* - status == clean ⟺ noul < 0.5 (flagged/warn need noul ≥ 0.5)
|
||||
* - severity == none ⟺ status == clean (flagged must have severity)
|
||||
* - action == none ⟺ status == clean; warn must not delete/escalate;
|
||||
* clean must never delete/escalate
|
||||
* - category == none ⟺ status == clean
|
||||
* Anything else → LLM fallback (fail-open).
|
||||
*/
|
||||
export function isJevAccepted(
|
||||
answers: JevAnswers,
|
||||
messageId: string,
|
||||
minConfidence: number,
|
||||
): boolean {
|
||||
const a = answersOf(answers, messageId);
|
||||
if (!a.v || a.v.type !== "noul" || typeof a.v.noul !== "number") return false;
|
||||
if (!a.status || a.status.type !== "choice" || !a.status.choice) return false;
|
||||
if (!a.severity || a.severity.type !== "choice" || !a.severity.choice)
|
||||
return false;
|
||||
if (!a.category || a.category.type !== "choice" || !a.category.choice)
|
||||
return false;
|
||||
if (!a.action || a.action.type !== "choice" || !a.action.choice) return false;
|
||||
|
||||
const { status, severity, category, action } = a;
|
||||
if (
|
||||
typeof status.confidence !== "number" ||
|
||||
status.confidence < minConfidence
|
||||
)
|
||||
return false;
|
||||
|
||||
// Valid label = one of the vocab const arrays. The membership guards
|
||||
// (Set.has) reject anything unknown, then the labels are narrowed via the
|
||||
// const-array includes so the downstream comparisons typecheck.
|
||||
const statusLabel = labelOf(status);
|
||||
const severityLabel = labelOf(severity);
|
||||
const categoryLabel = labelOf(category);
|
||||
const actionLabel = labelOf(action);
|
||||
if (
|
||||
!JEV_STATUS_SET.has(statusLabel) ||
|
||||
!JEV_SEVERITY_SET.has(severityLabel) ||
|
||||
!JEV_CATEGORY_SET.has(categoryLabel) ||
|
||||
!JEV_ACTION_SET.has(actionLabel)
|
||||
)
|
||||
return false; // unknown label — LLM fallback
|
||||
|
||||
const s = statusLabel as (typeof JEV_STATUSES)[number];
|
||||
const sev = severityLabel as (typeof JEV_SEVERITIES)[number];
|
||||
const cat = categoryLabel as (typeof JEV_CATEGORIES)[number];
|
||||
const act = actionLabel as (typeof JEV_ACTIONS)[number];
|
||||
|
||||
const noulVal = a.v.noul;
|
||||
|
||||
// noul ↔ status consistency
|
||||
if (s === "clean" && noulVal >= 0.5) return false;
|
||||
if (s !== "clean" && noulVal < 0.5) return false;
|
||||
// severity ↔ status: clean must be none; flagged/warn must NOT be none
|
||||
if (s === "clean" && sev !== "none") return false;
|
||||
if (s !== "clean" && sev === "none") return false;
|
||||
// action ↔ status: clean must be none; flagged must NOT be none;
|
||||
// warn must not delete/escalate; clean must never delete/escalate
|
||||
if (s === "clean" && act !== "none") return false;
|
||||
if (s === "flagged" && act === "none") return false;
|
||||
if (s === "warn" && (act === "delete" || act === "escalate")) return false;
|
||||
// category ↔ status: clean must be none; flagged must NOT be none
|
||||
if (s === "clean" && cat !== "none") return false;
|
||||
if (s === "flagged" && cat === "none") return false;
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Map accepted Jev answers for one message into the pipeline's `AnalysisResult`.
|
||||
* All values are derived from the model's own typed answers — no fabrication.
|
||||
*/
|
||||
export function mapJevAnswersToResult(
|
||||
answers: JevAnswers,
|
||||
messageId: string,
|
||||
): AnalysisResult {
|
||||
const a = answersOf(answers, messageId);
|
||||
const v = a.v as { type: "noul"; noul: number };
|
||||
const st = a.status as {
|
||||
type: "choice";
|
||||
choice: string;
|
||||
confidence: number;
|
||||
};
|
||||
const sev = labelOf(a.severity);
|
||||
const cat = labelOf(a.category);
|
||||
const act = labelOf(a.action);
|
||||
|
||||
const status = st.choice as (typeof JEV_STATUSES)[number];
|
||||
// Calibrated score: clean → 0; warn → 0.45; flagged → P(violates) clamped.
|
||||
const rawNoul = typeof v.noul === "number" ? v.noul : 0;
|
||||
const score =
|
||||
status === "clean"
|
||||
? 0
|
||||
: status === "warn"
|
||||
? 0.45
|
||||
: Math.min(1, Math.max(0.7, rawNoul));
|
||||
const confidence =
|
||||
typeof st.confidence === "number"
|
||||
? st.confidence
|
||||
: config.AI_LLM_JEV_MIN_CONFIDENCE;
|
||||
|
||||
return {
|
||||
messageId,
|
||||
status,
|
||||
flags: cat === "none" ? [] : [cat],
|
||||
score,
|
||||
analysis:
|
||||
`[Jev] status=${status}, kategori=${cat}, keparahan=${sev}, ` +
|
||||
`keyakinan=${confidence.toFixed(2)}, tindakan=${act}, p_melanggar=${rawNoul.toFixed(2)}`,
|
||||
categories: cat === "none" ? [] : [cat],
|
||||
severity: sev as AnalysisResult["severity"],
|
||||
confidence,
|
||||
recommendedAction: act as AnalysisResult["recommendedAction"],
|
||||
policyVersion: JEV_POLICY_VERSION,
|
||||
evidence: [],
|
||||
};
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Batch entry point (one systemOne call per sub-batch)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export interface JevBatchOutcome {
|
||||
/** Accepted Jev verdicts (keyed by message id). */
|
||||
results: AnalysisResult[];
|
||||
/** Answers the gate rejected for ANY reason (keys = message ids). */
|
||||
rejectedIds: string[];
|
||||
/** Raw `SystemOneResult` (for usage logging / raw passthrough). */
|
||||
raw: unknown;
|
||||
/** Error thrown by the call, if the whole call failed (null = success). */
|
||||
error: string | null;
|
||||
}
|
||||
|
||||
/** True when Jev is configured and enabled (fail-open wrapper). */
|
||||
export function isJevEnabled(): boolean {
|
||||
return (
|
||||
config.AI_LLM_JEV_ENABLED === true && Boolean(config.AI_LLM_JEV_API_KEY)
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Analyze one text sub-batch with Jev. NEVER throws for API-level failures —
|
||||
* returns `{ error }` so the caller falls back to the LLM. Aborts (signal)
|
||||
* propagate as errors too (the caller's timeout must abort the SDK call and
|
||||
* fall back, not hang).
|
||||
*/
|
||||
export async function analyzeBatchWithJev(
|
||||
targets: JevTarget[],
|
||||
ctx: JevBatchContext,
|
||||
signal?: AbortSignal,
|
||||
correctedExamples = "",
|
||||
): Promise<JevBatchOutcome> {
|
||||
const outcome: JevBatchOutcome = {
|
||||
results: [],
|
||||
rejectedIds: [],
|
||||
raw: null,
|
||||
error: null,
|
||||
};
|
||||
if (targets.length === 0) return outcome;
|
||||
|
||||
const jevClient = getClient();
|
||||
if (!jevClient) {
|
||||
outcome.error = "Jev client unavailable (no API key)";
|
||||
return outcome;
|
||||
}
|
||||
|
||||
try {
|
||||
const questions = buildJevQuestions(targets);
|
||||
|
||||
// CONCURRENCY from the skill: one ownership layer owns retries — the SDK
|
||||
// gets retry: { maxRetries: 0 } and the pipeline's timeout/abort layer is
|
||||
// the only retry. The call is wrapped in withLlmConcurrency so Jev down
|
||||
// can't flood the router.
|
||||
const { withLlmConcurrency } = await import("./llmClient.js");
|
||||
const systemOneResult = await withLlmConcurrency(async () => {
|
||||
return await jevClient.systemOne(
|
||||
{
|
||||
model: config.AI_LLM_JEV_MODEL,
|
||||
state: buildJevState(targets, ctx, correctedExamples),
|
||||
questions,
|
||||
},
|
||||
{ signal, timeout: config.AI_LLM_JEV_TIMEOUT_MS },
|
||||
);
|
||||
});
|
||||
|
||||
outcome.raw = systemOneResult;
|
||||
const answers = systemOneResult.answers as unknown as JevAnswers;
|
||||
const minConfidence = config.AI_LLM_JEV_MIN_CONFIDENCE;
|
||||
|
||||
for (const t of targets) {
|
||||
if (isJevAccepted(answers, t.id, minConfidence)) {
|
||||
outcome.results.push(mapJevAnswersToResult(answers, t.id));
|
||||
incrementCounterBy("moderation_jev_decisions", 1, { type: "jev" });
|
||||
} else {
|
||||
outcome.rejectedIds.push(t.id);
|
||||
incrementCounterBy("moderation_jev_decisions", 1, {
|
||||
type: "llm_fallback",
|
||||
});
|
||||
log.debug(
|
||||
{ messageId: t.id, minConfidence },
|
||||
"Jev decision rejected — falling back to LLM",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Usage accounting (same counters as the LLM path).
|
||||
const usage = systemOneResult.usage;
|
||||
if (usage?.input_tokens || usage?.output_tokens) {
|
||||
if (usage.input_tokens) {
|
||||
incrementCounterBy("llm_tokens_total", usage.input_tokens, {
|
||||
model: config.AI_LLM_JEV_MODEL,
|
||||
type: "prompt",
|
||||
label: "jev-batch",
|
||||
});
|
||||
}
|
||||
if (usage.output_tokens) {
|
||||
incrementCounterBy("llm_tokens_total", usage.output_tokens, {
|
||||
model: config.AI_LLM_JEV_MODEL,
|
||||
type: "completion",
|
||||
label: "jev-batch",
|
||||
});
|
||||
}
|
||||
log.info(
|
||||
{
|
||||
targetCount: targets.length,
|
||||
accepted: outcome.results.length,
|
||||
rejected: outcome.rejectedIds.length,
|
||||
model: config.AI_LLM_JEV_MODEL,
|
||||
input_tokens: usage.input_tokens,
|
||||
output_tokens: usage.output_tokens,
|
||||
},
|
||||
"Jev systemone batch usage",
|
||||
);
|
||||
}
|
||||
|
||||
return outcome;
|
||||
} catch (err) {
|
||||
if (err instanceof Error && err.name === "APIUserAbortError") throw err; // real abort — let caller decide
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
outcome.error = msg;
|
||||
log.warn(
|
||||
{ error: msg, targetCount: targets.length },
|
||||
"Jev systemone call failed — falling back to LLM for the whole batch",
|
||||
);
|
||||
return outcome;
|
||||
}
|
||||
}
|
||||
@@ -13,7 +13,7 @@
|
||||
import type { ChatCompletion } from "openai/resources/chat/completions";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { delay, retryWithBackoff } from "@/shared/utils/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { incrementCounterBy } from "../gateway-metrics/index.js";
|
||||
import type { AnalysisResult } from "../message-capture/types.js";
|
||||
import { llmChat } from "./llmClient.js";
|
||||
@@ -58,6 +58,9 @@ export async function callModerationLLM(
|
||||
// callers pass a prompt-derived ceiling so small batches don't reserve a
|
||||
// 16k completion budget (some routers pre-allocate KV cache per max_tokens).
|
||||
maxTokens?: number,
|
||||
// Concurrency lane: "text" (default) uses AI_LLM_MAX_CONCURRENT; "media"
|
||||
// uses AI_LLM_MEDIA_MAX_CONCURRENT.
|
||||
lane: "text" | "media" = "text",
|
||||
): Promise<{
|
||||
results: AnalysisResult[];
|
||||
raw: ChatCompletion | null;
|
||||
@@ -89,13 +92,14 @@ export async function callModerationLLM(
|
||||
jsonResponse: { type: "json_object" },
|
||||
retries: 0,
|
||||
signal,
|
||||
// Router (omniroute) always streams SSE even when the
|
||||
// Router always streams SSE even when the
|
||||
// request omits `stream`. In non-stream mode the OpenAI SDK waits
|
||||
// for the FULL body before parsing, so slow/long upstream streams
|
||||
// hit the 30s/60s timeout and abort mid-generation. Streaming mode
|
||||
// consumes chunks incrementally — timeout only fires on a real
|
||||
// stall. llmClient aggregates the stream into a ChatCompletion.
|
||||
stream: true,
|
||||
lane,
|
||||
});
|
||||
|
||||
if (!completion)
|
||||
|
||||
@@ -10,46 +10,83 @@ import OpenAI from "openai";
|
||||
import pLimit from "p-limit";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { retryWithBackoff } from "@/shared/utils/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
|
||||
const log = createChildLogger("llm-client");
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Concurrency limiter for LLM API calls (inlined from concurrencyLimiter.ts)
|
||||
// Concurrency limiters for LLM API calls (inlined from concurrencyLimiter.ts)
|
||||
// ---------------------------------------------------------------------------
|
||||
//
|
||||
// Split into TWO independent semaphores (2026-09-24): the text lane and the
|
||||
// media lane (vision + media batches) no longer share one global cap. A slow
|
||||
// vision call used to occupy a slot of the SINGLE pLimit(AI_LLM_MAX_CONCURRENT)
|
||||
// semaphore, so a media-heavy burst could starve text inference. Now each lane
|
||||
// has its own cap — media churn can never consume text slots, and vice versa.
|
||||
|
||||
// The limiter is cached per configured concurrency value so it can be tuned
|
||||
// (env / BWS) without a code change and always reflects the current config —
|
||||
// a module-level `pLimit(config.X)` would freeze the cap at import time.
|
||||
let llmSemaphore = pLimit(config.AI_LLM_MAX_CONCURRENT ?? 5);
|
||||
let llmSemaphoreLimit = config.AI_LLM_MAX_CONCURRENT ?? 5;
|
||||
type LlmLane = "text" | "media";
|
||||
|
||||
function getLlmSemaphore() {
|
||||
const wanted = config.AI_LLM_MAX_CONCURRENT ?? 5;
|
||||
if (wanted !== llmSemaphoreLimit) {
|
||||
llmSemaphore = pLimit(wanted);
|
||||
llmSemaphoreLimit = wanted;
|
||||
interface LaneSemaphore {
|
||||
limiter: ReturnType<typeof pLimit>;
|
||||
limit: number;
|
||||
}
|
||||
|
||||
const laneSemaphores: Record<LlmLane, LaneSemaphore> = {
|
||||
text: {
|
||||
limiter: pLimit(config.AI_LLM_MAX_CONCURRENT ?? 5),
|
||||
limit: config.AI_LLM_MAX_CONCURRENT ?? 5,
|
||||
},
|
||||
media: {
|
||||
limiter: pLimit(config.AI_LLM_MEDIA_MAX_CONCURRENT ?? 4),
|
||||
limit: config.AI_LLM_MEDIA_MAX_CONCURRENT ?? 4,
|
||||
},
|
||||
};
|
||||
|
||||
function getLaneSemaphore(lane: LlmLane): ReturnType<typeof pLimit> {
|
||||
const wanted =
|
||||
lane === "media"
|
||||
? (config.AI_LLM_MEDIA_MAX_CONCURRENT ?? 4)
|
||||
: (config.AI_LLM_MAX_CONCURRENT ?? 5);
|
||||
const slot = laneSemaphores[lane];
|
||||
if (wanted !== slot.limit) {
|
||||
slot.limiter = pLimit(wanted);
|
||||
slot.limit = wanted;
|
||||
}
|
||||
return llmSemaphore;
|
||||
return slot.limiter;
|
||||
}
|
||||
|
||||
let activeCount = 0;
|
||||
let pendingCount = 0;
|
||||
|
||||
export async function withLlmConcurrency<T>(fn: () => Promise<T>): Promise<T> {
|
||||
/**
|
||||
* Run `fn` under the per-lane LLM concurrency cap.
|
||||
*
|
||||
* `lane: "text"` uses `AI_LLM_MAX_CONCURRENT`; `lane: "media"` uses
|
||||
* `AI_LLM_MEDIA_MAX_CONCURRENT`. Defaults to "text" so the existing text
|
||||
* moderation path is unchanged.
|
||||
*/
|
||||
export async function withLlmConcurrency<T>(
|
||||
fn: () => Promise<T>,
|
||||
opts: { lane?: LlmLane } = {},
|
||||
): Promise<T> {
|
||||
const lane = opts.lane ?? "text";
|
||||
const maxConcurrent =
|
||||
lane === "media"
|
||||
? (config.AI_LLM_MEDIA_MAX_CONCURRENT ?? 4)
|
||||
: (config.AI_LLM_MAX_CONCURRENT ?? 5);
|
||||
pendingCount++;
|
||||
log.debug(
|
||||
{ activeCount, pendingCount, maxConcurrent: config.AI_LLM_MAX_CONCURRENT },
|
||||
{ activeCount, pendingCount, maxConcurrent, lane },
|
||||
"Queuing LLM request",
|
||||
);
|
||||
|
||||
return getLlmSemaphore()(async () => {
|
||||
return getLaneSemaphore(lane)(async () => {
|
||||
pendingCount--;
|
||||
activeCount++;
|
||||
|
||||
if (activeCount >= (config.AI_LLM_MAX_CONCURRENT ?? 5)) {
|
||||
if (activeCount >= maxConcurrent) {
|
||||
log.warn(
|
||||
{ activeCount, maxConcurrent: config.AI_LLM_MAX_CONCURRENT },
|
||||
{ activeCount, maxConcurrent, lane },
|
||||
"LLM concurrency limit reached",
|
||||
);
|
||||
}
|
||||
@@ -94,7 +131,7 @@ type LLMResponseChunk = {
|
||||
* `delta.content`; falls back to reasoning fields so reasoning-only models
|
||||
* still produce usable aggregated text. Providers differ in the field name:
|
||||
* - DeepSeek-style / Cloudflare gemma → `delta.reasoning_content`
|
||||
* - mimo (via omniroute) streams reasoning in `delta.reasoning` +
|
||||
* - mimo (via the router) streams reasoning in `delta.reasoning` +
|
||||
* `delta.reasoning_details[].text` (content:"") — without these fallbacks
|
||||
* vision aggregation came back empty ("Vision API null response").
|
||||
* Exported for unit tests.
|
||||
@@ -177,6 +214,11 @@ export interface LlmCallOpts {
|
||||
* so a single large-image call isn't killed early by the shared default.
|
||||
*/
|
||||
timeout?: number;
|
||||
/**
|
||||
* Concurrency lane. "text" uses AI_LLM_MAX_CONCURRENT; "media" (vision,
|
||||
* media batches) uses AI_LLM_MEDIA_MAX_CONCURRENT. Defaults to "text".
|
||||
*/
|
||||
lane?: "text" | "media";
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -225,7 +267,7 @@ export function buildLlmParams(
|
||||
reasoning: { enabled: false },
|
||||
// vLLM / Qwen / litellm
|
||||
chat_template_kwargs: { enable_thinking: false },
|
||||
// Anthropic / Claude-format (omniroute exposes thinkingFormat
|
||||
// Anthropic / Claude-format (router exposes thinkingFormat
|
||||
// "claude-adaptive" / "claude-budget" on its reasoning models)
|
||||
thinking: { type: "disabled" },
|
||||
} as Record<string, unknown>);
|
||||
@@ -255,78 +297,84 @@ export async function llmChat(
|
||||
|
||||
return retryWithBackoff(
|
||||
async () => {
|
||||
return withLlmConcurrency(async () => {
|
||||
const execute = async (
|
||||
currentParams: OpenAI.Chat.Completions.ChatCompletionCreateParams,
|
||||
) => {
|
||||
const response = await client.chat.completions.create(currentParams, {
|
||||
signal,
|
||||
...(opts.timeout ? { timeout: opts.timeout } : {}),
|
||||
});
|
||||
if (currentParams.stream) {
|
||||
let content = "";
|
||||
let finishReason = "stop";
|
||||
for await (const chunk of response as unknown as AsyncIterable<LLMResponseChunk>) {
|
||||
const choice = chunk?.choices?.[0];
|
||||
content += extractChunkText(chunk);
|
||||
const fr = choice?.finish_reason || chunk?.finish_reason;
|
||||
if (fr) finishReason = fr;
|
||||
}
|
||||
return {
|
||||
id: "stream-aggregated",
|
||||
choices: [
|
||||
{
|
||||
message: { role: "assistant", content, refusal: null },
|
||||
finish_reason: finishReason,
|
||||
index: 0,
|
||||
logprobs: null,
|
||||
},
|
||||
],
|
||||
created: Math.floor(Date.now() / 1000),
|
||||
model: currentParams.model,
|
||||
object: "chat.completion",
|
||||
} as OpenAI.Chat.Completions.ChatCompletion;
|
||||
}
|
||||
return response as OpenAI.Chat.Completions.ChatCompletion;
|
||||
};
|
||||
|
||||
try {
|
||||
return await execute(params);
|
||||
} catch (error: any) {
|
||||
const rawResponse =
|
||||
error.error || error.body || error.response?.data || "N/A";
|
||||
const errorStr = (
|
||||
JSON.stringify(rawResponse) + String(error.message)
|
||||
).toLowerCase();
|
||||
|
||||
// Auto-fallback: If provider strictly demands streaming (400 Bad Request on stream params)
|
||||
if (
|
||||
error.status === 400 &&
|
||||
errorStr.includes("stream") &&
|
||||
!params.stream
|
||||
) {
|
||||
log.warn(
|
||||
{ model },
|
||||
"Provider rejected non-streaming request. Fallback to stream: true initiated.",
|
||||
return withLlmConcurrency(
|
||||
async () => {
|
||||
const execute = async (
|
||||
currentParams: OpenAI.Chat.Completions.ChatCompletionCreateParams,
|
||||
) => {
|
||||
const response = await client.chat.completions.create(
|
||||
currentParams,
|
||||
{
|
||||
signal,
|
||||
...(opts.timeout ? { timeout: opts.timeout } : {}),
|
||||
},
|
||||
);
|
||||
(
|
||||
params as unknown as OpenAI.Chat.Completions.ChatCompletionCreateParamsStreaming
|
||||
).stream = true;
|
||||
return await execute(params);
|
||||
}
|
||||
if (currentParams.stream) {
|
||||
let content = "";
|
||||
let finishReason = "stop";
|
||||
for await (const chunk of response as unknown as AsyncIterable<LLMResponseChunk>) {
|
||||
const choice = chunk?.choices?.[0];
|
||||
content += extractChunkText(chunk);
|
||||
const fr = choice?.finish_reason || chunk?.finish_reason;
|
||||
if (fr) finishReason = fr;
|
||||
}
|
||||
return {
|
||||
id: "stream-aggregated",
|
||||
choices: [
|
||||
{
|
||||
message: { role: "assistant", content, refusal: null },
|
||||
finish_reason: finishReason,
|
||||
index: 0,
|
||||
logprobs: null,
|
||||
},
|
||||
],
|
||||
created: Math.floor(Date.now() / 1000),
|
||||
model: currentParams.model,
|
||||
object: "chat.completion",
|
||||
} as OpenAI.Chat.Completions.ChatCompletion;
|
||||
}
|
||||
return response as OpenAI.Chat.Completions.ChatCompletion;
|
||||
};
|
||||
|
||||
log.error(
|
||||
{
|
||||
error: error.message,
|
||||
status: error.status,
|
||||
rawResponse,
|
||||
model,
|
||||
},
|
||||
"LLM API request failed",
|
||||
);
|
||||
throw error;
|
||||
}
|
||||
});
|
||||
try {
|
||||
return await execute(params);
|
||||
} catch (error: any) {
|
||||
const rawResponse =
|
||||
error.error || error.body || error.response?.data || "N/A";
|
||||
const errorStr = (
|
||||
JSON.stringify(rawResponse) + String(error.message)
|
||||
).toLowerCase();
|
||||
|
||||
// Auto-fallback: If provider strictly demands streaming (400 Bad Request on stream params)
|
||||
if (
|
||||
error.status === 400 &&
|
||||
errorStr.includes("stream") &&
|
||||
!params.stream
|
||||
) {
|
||||
log.warn(
|
||||
{ model },
|
||||
"Provider rejected non-streaming request. Fallback to stream: true initiated.",
|
||||
);
|
||||
(
|
||||
params as unknown as OpenAI.Chat.Completions.ChatCompletionCreateParamsStreaming
|
||||
).stream = true;
|
||||
return await execute(params);
|
||||
}
|
||||
|
||||
log.error(
|
||||
{
|
||||
error: error.message,
|
||||
status: error.status,
|
||||
rawResponse,
|
||||
model,
|
||||
},
|
||||
"LLM API request failed",
|
||||
);
|
||||
throw error;
|
||||
}
|
||||
},
|
||||
{ lane: opts.lane ?? "text" },
|
||||
);
|
||||
},
|
||||
{
|
||||
retries,
|
||||
@@ -356,7 +404,7 @@ export async function llmVision(
|
||||
promptText: string,
|
||||
imageUrl: { url: string },
|
||||
): Promise<string | null> {
|
||||
const params = {
|
||||
const params: LlmCallOpts = {
|
||||
messages: [
|
||||
{
|
||||
role: "user" as const,
|
||||
@@ -372,6 +420,7 @@ export async function llmVision(
|
||||
top_p: 0.9,
|
||||
retries: 0,
|
||||
timeout: config.AI_LLM_VISION_ANALYSIS_TIMEOUT_MS ?? 60_000,
|
||||
lane: "media",
|
||||
};
|
||||
|
||||
// Streaming first (the router always streams SSE; a non-stream request
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
* moderationOrchestrator.ts.
|
||||
*/
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type {
|
||||
AnalysisResult,
|
||||
AttachmentRecord,
|
||||
@@ -108,6 +108,7 @@ export async function runMediaBatch(
|
||||
`media-batch:${targetIds.length}msgs`,
|
||||
abortController.signal,
|
||||
dynamicMaxTokens,
|
||||
"media",
|
||||
);
|
||||
log.info(
|
||||
{ mediaCount: targets.length, resultCount: result.results.length },
|
||||
|
||||
@@ -7,7 +7,7 @@
|
||||
|
||||
import { LRUCache } from "lru-cache";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { incrementCounterBy } from "../gateway-metrics/index.js";
|
||||
import { extractMessageMediaEvidence } from "../message-capture/messageMetadata.js";
|
||||
import type {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import type { Client } from "discord.js-selfbot-v13";
|
||||
import { LRUCache } from "lru-cache";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import type { EventBroadcaster } from "../event-broadcaster/index.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
import { attemptAutoDeleteFlaggedMessage } from "./autoDeleteManager.js";
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -13,7 +13,7 @@
|
||||
import { createHash } from "node:crypto";
|
||||
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
|
||||
const log = createChildLogger("qdrant");
|
||||
|
||||
|
||||
@@ -0,0 +1,195 @@
|
||||
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.
|
||||
* Kept in lockstep with the batch processing timeout
|
||||
* (AI_ANALYSIS_PROCESSING_TIMEOUT_MS, default 120s): a row sitting past the
|
||||
* batch budget is a leak, not a legitimate slow batch. messagesCleanup's
|
||||
* default was lowered 300s→120s in 2026-08-24; this constant was missed and
|
||||
* stayed at 300s — messages looked stuck for up to 5 minutes before recovery
|
||||
* touched them.
|
||||
*/
|
||||
const STUCK_PROCESSING_AGE_MS = config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS;
|
||||
|
||||
/**
|
||||
* 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();
|
||||
|
||||
// Revert stuck `processing` messages back to `pending` unconditionally.
|
||||
// The old `if (conversationProcessing.size > 0)` guard skipped the revert
|
||||
// when the in-memory lock map was empty (fresh boot, or every lock was
|
||||
// already pruned) — precisely the moment stranded `processing` rows from
|
||||
// a previous process still need rescuing. `revertStuckProcessingMessages`
|
||||
// is a cheap UPDATE..RETURNING keyed on ai_status + age, safe to run
|
||||
// every interval; it matches 0 rows when there is nothing to do.
|
||||
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,5 +1,5 @@
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { executeAll, executeGet } from "../../shared/database/drizzle.js";
|
||||
import { uploadToTele } from "../../shared/uploader.js";
|
||||
|
||||
|
||||
@@ -34,7 +34,7 @@ import { LRUCache } from "lru-cache";
|
||||
import pLimit from "p-limit";
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { delay } from "@/shared/utils/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { cacheGet, cacheSet, makeCacheKey } from "./cacheStore.js";
|
||||
import { escapeXml } from "./moderationBuilders.js";
|
||||
import {
|
||||
|
||||
@@ -7,7 +7,7 @@
|
||||
*/
|
||||
import { createChildLogger } from "@/shared/logger/index";
|
||||
import { delay } from "@/shared/utils/index";
|
||||
import { config } from "../../shared/config/config.js";
|
||||
import { config } from "../../shared/config/index.js";
|
||||
import { resizeImageForVision } from "../attachment-upload/imageResizer.js";
|
||||
import type {
|
||||
AnalysisResult,
|
||||
@@ -15,12 +15,6 @@ import type {
|
||||
} from "../message-capture/types.js";
|
||||
import { getChannelCulture } from "./channelCultureStore.js";
|
||||
import { estimateTokens } from "./conversationContext.js";
|
||||
import {
|
||||
analyzeBatchWithJev,
|
||||
isJevEnabled,
|
||||
JEV_POLICY_VERSION,
|
||||
type JevTarget,
|
||||
} from "./jevAnalyzer.js";
|
||||
import type { ModerationPromptContent, RetryState } from "./llmCaller.js";
|
||||
import { callModerationLLM } from "./llmCaller.js";
|
||||
import { analyzeSingleMediaImage } from "./mediaAnalysisClient.js";
|
||||
@@ -56,16 +50,18 @@ interface UrlFetchResult {
|
||||
text: Map<string, string>;
|
||||
image: Map<string, { data: Buffer; mimeType: string }>;
|
||||
title: Map<string, string>;
|
||||
/** URL → reason the fetch failed (HTTP status / throw / unsupported). */
|
||||
error: Map<string, string>;
|
||||
}
|
||||
|
||||
/** Provider-reported token usage from a raw LLM/Jev payload (may be absent). */
|
||||
/** Provider-reported token usage from a raw LLM payload (may be absent). */
|
||||
interface TokenUsage {
|
||||
prompt_tokens: number;
|
||||
completion_tokens: number;
|
||||
total_tokens: number;
|
||||
}
|
||||
|
||||
/** Read provider-reported token usage from either the LLM or Jev raw payload. */
|
||||
/** Read provider-reported token usage from the raw LLM payload. */
|
||||
function extractUsage(raw: unknown): TokenUsage | undefined {
|
||||
return (raw as { usage?: TokenUsage } | null)?.usage ?? undefined;
|
||||
}
|
||||
@@ -175,6 +171,7 @@ export async function runTextOnlyBatch(
|
||||
text: new Map(),
|
||||
image: new Map(),
|
||||
title: new Map(),
|
||||
error: new Map(),
|
||||
} satisfies UrlFetchResult;
|
||||
}
|
||||
const results = await Promise.allSettled(
|
||||
@@ -183,9 +180,13 @@ export async function runTextOnlyBatch(
|
||||
const textMap = new Map<string, string>();
|
||||
const imageMap = new Map<string, { data: Buffer; mimeType: string }>();
|
||||
const titleMap = new Map<string, string>();
|
||||
const errorMap = new Map<string, string>();
|
||||
for (let i = 0; i < urlArr.length; i++) {
|
||||
const r = results[i];
|
||||
if (r.status !== "fulfilled") continue;
|
||||
if (r.status !== "fulfilled") {
|
||||
errorMap.set(urlArr[i], "fetch threw");
|
||||
continue;
|
||||
}
|
||||
const v = r.value;
|
||||
if (v.type === "text" && v.textContent) {
|
||||
textMap.set(urlArr[i], v.textContent);
|
||||
@@ -194,9 +195,11 @@ export async function runTextOnlyBatch(
|
||||
// Direct image link (or og:image followed from an HTML page) —
|
||||
// kept for vision analysis below.
|
||||
imageMap.set(urlArr[i], { data: v.data, mimeType: v.mimeType });
|
||||
} else {
|
||||
errorMap.set(urlArr[i], v.error ?? "unsupported content");
|
||||
}
|
||||
}
|
||||
return { text: textMap, image: imageMap, title: titleMap };
|
||||
return { text: textMap, image: imageMap, title: titleMap, error: errorMap };
|
||||
})();
|
||||
|
||||
const webSearchPromise = (async () => {
|
||||
@@ -230,6 +233,7 @@ export async function runTextOnlyBatch(
|
||||
glossaryPromise,
|
||||
]);
|
||||
const urlFetchMap = urlFetchMaps.text;
|
||||
const urlFetchErrors = urlFetchMaps.error;
|
||||
|
||||
// Deduplicate identical short messages
|
||||
const shortContentGroups = new Map<string, MessageRecord[]>();
|
||||
@@ -374,10 +378,16 @@ export async function runTextOnlyBatch(
|
||||
const urlContexts = msgUrls
|
||||
.map((url) => {
|
||||
const ft = urlFetchMap.get(url);
|
||||
if (!ft) return null;
|
||||
const title = urlTitles.get(url);
|
||||
const titleAttr = title ? ` title="${escapeXml(title)}"` : "";
|
||||
return `<web_content url="${escapeXml(url)}"${titleAttr}>${escapeXml(ft)}</web_content>`;
|
||||
if (ft) {
|
||||
const title = urlTitles.get(url);
|
||||
const titleAttr = title ? ` title="${escapeXml(title)}"` : "";
|
||||
return `<web_content url="${escapeXml(url)}"${titleAttr}>${escapeXml(ft)}</web_content>`;
|
||||
}
|
||||
const fetchError = urlFetchErrors.get(url);
|
||||
if (fetchError) {
|
||||
return `<web_content url="${escapeXml(url)}" fetch_error="true">Konten tidak dapat diambil otomatis (${escapeXml(fetchError)}). Analisis hanya dari teks pesan; JANGAN mengarang isi halaman.</web_content>`;
|
||||
}
|
||||
return null;
|
||||
})
|
||||
.filter(Boolean)
|
||||
.join("\n");
|
||||
@@ -419,7 +429,7 @@ export async function runTextOnlyBatch(
|
||||
results: [],
|
||||
raw: null,
|
||||
};
|
||||
// Per-sub-batch verdicts before fan-out (Jev + LLM fallback merged).
|
||||
// Per-sub-batch verdicts before fan-out.
|
||||
let subBatchResults: AnalysisResult[] = [];
|
||||
try {
|
||||
// Output budget scales with the prompt: the JSON verdict block is
|
||||
@@ -440,78 +450,14 @@ export async function runTextOnlyBatch(
|
||||
Math.max(2048, Math.ceil(subBatchPromptEstimate * 1.5)),
|
||||
);
|
||||
|
||||
// ── Jev-first (TypeSafe System One) with LLM fallback ────────────────
|
||||
// Jev is the PRIMARY text analyzer: ONE systemOne call per sub-batch
|
||||
// (5 typed questions × N messages, evaluated in parallel by Jev).
|
||||
// Verdicts that pass the acceptance gate are used directly; anything
|
||||
// Jev rejects (low confidence / inconsistent) and any Jev API failure
|
||||
// falls back to the existing LLM call — fail-open, never dead.
|
||||
if (isJevEnabled()) {
|
||||
const jevTargets: JevTarget[] = batch.map((msg) => ({
|
||||
id: msg.id,
|
||||
user: resolveDisplayName(msg),
|
||||
content: analysisContentOf(msg),
|
||||
}));
|
||||
const jevOutcome = await analyzeBatchWithJev(
|
||||
jevTargets,
|
||||
{
|
||||
contextBlock,
|
||||
webSearchBlock: buildWebSearchBlock(webSearchResults),
|
||||
glossaryBlock,
|
||||
channelCulture: channelCultureObj?.culture_summary,
|
||||
},
|
||||
abortController.signal,
|
||||
correctedExamples,
|
||||
);
|
||||
|
||||
subBatchResults.push(...jevOutcome.results);
|
||||
if (jevOutcome.results.length > 0) {
|
||||
log.info(
|
||||
{
|
||||
subBatch: i + 1,
|
||||
accepted: jevOutcome.results.length,
|
||||
rejected: jevOutcome.rejectedIds.length,
|
||||
},
|
||||
"Jev analyzed sub-batch — accepted verdicts kept, rejected go to LLM",
|
||||
);
|
||||
}
|
||||
|
||||
// Which targets still need the LLM?
|
||||
const coveredIds = new Set(subBatchResults.map((r) => r.messageId));
|
||||
const llmTargets = batch.filter((m) => !coveredIds.has(m.id));
|
||||
|
||||
if (llmTargets.length > 0) {
|
||||
const llmResult = await callModerationLLM(
|
||||
(state) => buildContent(state, llmTargets),
|
||||
llmTargets.map((m) => m.id),
|
||||
`text-batch-${i + 1}-jev-fallback`,
|
||||
abortController.signal,
|
||||
dynamicMaxTokens,
|
||||
);
|
||||
subBatchResults.push(...llmResult.results);
|
||||
batchResult = llmResult;
|
||||
logModerationAnalysis(
|
||||
llmTargets.map((m) => m.id),
|
||||
config.AI_LLM_MODEL,
|
||||
llmResult.results,
|
||||
0,
|
||||
extractUsage(llmResult.raw),
|
||||
);
|
||||
} else {
|
||||
// Jev accepted everything — no LLM usage to attribute.
|
||||
batchResult = { results: subBatchResults, raw: null };
|
||||
}
|
||||
} else {
|
||||
// Jev disabled / unconfigured — pure LLM path (unchanged).
|
||||
batchResult = await callModerationLLM(
|
||||
buildContent,
|
||||
targetIds,
|
||||
`text-batch-${i + 1}`,
|
||||
abortController.signal,
|
||||
dynamicMaxTokens,
|
||||
);
|
||||
subBatchResults = batchResult.results;
|
||||
}
|
||||
batchResult = await callModerationLLM(
|
||||
buildContent,
|
||||
targetIds,
|
||||
`text-batch-${i + 1}`,
|
||||
abortController.signal,
|
||||
dynamicMaxTokens,
|
||||
);
|
||||
subBatchResults = batchResult.results;
|
||||
} catch (err: unknown) {
|
||||
const isAbort =
|
||||
(err instanceof Error && err.name === "AbortError") ||
|
||||
@@ -530,8 +476,7 @@ export async function runTextOnlyBatch(
|
||||
clearTimeout(timeoutId);
|
||||
}
|
||||
|
||||
// Fan-out results for deduplicated messages (applies to Jev + LLM
|
||||
// verdicts alike — they only consume AnalysisResult[]).
|
||||
// Fan-out results for deduplicated messages.
|
||||
const fannedOutResults =
|
||||
groupMapping.size > 0
|
||||
? subBatchResults.flatMap((result) => {
|
||||
@@ -548,9 +493,7 @@ export async function runTextOnlyBatch(
|
||||
if (subBatchResults.length > 0) {
|
||||
logModerationAnalysis(
|
||||
targetIds,
|
||||
subBatchResults.every((r) => r.policyVersion === JEV_POLICY_VERSION)
|
||||
? config.AI_LLM_JEV_MODEL
|
||||
: config.AI_LLM_MODEL,
|
||||
config.AI_LLM_MODEL,
|
||||
subBatchResults,
|
||||
0,
|
||||
extractUsage(batchResult.raw),
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user