Author SHA1 Message Date
asepharyana 045cdf1f75 fix(gateway): align stuck-recovery threshold with batch timeout (300s→120s)
recovery-worker used a hardcoded STUCK_PROCESSING_AGE_MS=300_000 while
messagesCleanup's default and AI_ANALYSIS_PROCESSING_TIMEOUT_MS are both
120s. Rows stuck between 2 and 5 minutes were never reverted by the
recovery worker — they looked permanently stuck (and accumulated under
load) even though the batch budget had long passed. Now derives the
threshold from config so the two knobs can never drift again.
2026-09-24 17:52:17 +07:00
asepharyana deed7bdfb0 fix(gateway): destroy AI worker pools on graceful shutdown
Piscina worker threads outlive process.exit() and linger as orphaned
processes holding DB connections/locks after a deploy restart. Two live
gateways then fight over the same messages table rows (one claims
processing, the other reverts), which left messages stuck in
ai_status='processing' forever.

Destroy both worker pools before closing the DB, with a 5s fallback so
a hung vision job cannot block shutdown indefinitely.
2026-09-24 17:26:54 +07:00
asepharyana c4d9ade85e debug(gateway): log lane-lock skip + dispatch in scheduleLaneTimer 2026-09-24 17:11:04 +07:00
asepharyana 80daa9f045 fix(gateway): un-claim budget-overflow messages stuck in processing
getPendingMessagesByConversation() flips every fetched pending row to
'processing', then pickBatchWithinBudget() may stop early on the token
budget. The tail rows that did NOT make the batch were never un-claimed,
so they stayed 'processing' forever — the recovery worker reverted them
(120s) only for the next wave to re-claim them, an infinite loop of
stuck messages that never get analyzed (saw 22 rows, some recycled for
40+ minutes).

- add computeBudgetOverflowMessages() pure helper (batchBudget.ts)
- batchScheduler un-claims overflow rows back to 'pending' before
  dispatching the trimmed batch
- recovery-worker now reverts stuck processing unconditionally (the old
  'conversationProcessing.size > 0' guard skipped the revert when the
  in-memory lock map was empty, e.g. fresh boot — exactly when stranded
  rows from a previous process need rescuing)
- 3 regression tests for the overflow helper
2026-09-24 16:49:27 +07:00
asepharyana ef7708bf7d feat(gmw): route all LLM traffic through 9router
GMW moves off omniroute (100.121.180.82:20128) and off the direct NVIDIA
vision endpoint onto 9router, which runs on the same host as both services
(127.0.0.1:4014) — loopback avoids the TLS/proxy hop and localhost calls
bypass 9router's remote-key guard.

- gateway + backend: AI_LLM_BASE_URL default -> http://127.0.0.1:4014/v1
- drop stale 'omniroute' router references from comments/docs now that the
  active router is 9router (llmClient, llmCaller, ARCHITECTURE, AGENTS)

Verified against 9router before wiring: model 'text' -> gemini-3.5-flash-lite
(SSE, as the pipeline expects), 'multimodal' -> nemotron-3-nano-omni answers
image input, and gemini/gemini-embedding-001 returns 3072 dims — matching the
existing Qdrant collections (no reindex needed). The GMW key is already
registered in 9router's apiKeys table.

typecheck + lint + tests green (gateway 138, backend 37 excluding e2e).
2026-09-24 16:05:36 +07:00
asepharyana 750f3aa598 docs(gateway): correct metrics port — 4016 was wrong
The metrics server binds METRICS_PORT (code default 9090; this host runs it
on 4018 — 4016 is occupied by another process). Docs said 4016 in three
places, which is not what the service does.
2026-09-24 15:50:15 +07:00
asepharyana c57ee12da1 docs(gateway): consolidate README/ARCHITECTURE, drop stale MODULE_STRUCTURE
README.md was the extraction-era document (referenced winston, mock-crc.ts,
llmModerationClient.ts, indonesianTextNormalizer.ts — all long gone) and
duplicated ARCHITECTURE.md. Rewritten as a short run-the-service guide;
layout/design lives only in ARCHITECTURE.md.

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

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

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

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

No behavior change. typecheck + lint + 138 tests green.
2026-09-24 14:58:19 +07:00
asepharyana 8743fcc0b5 fix(ci): recover attic client binary in VPS-hop fallback path 2026-09-24 14:31:26 +07:00
asepharyana e34dcd6bc6 fix(gateway): separate text and media lanes in AI analysis queue (#86)
Image messages previously blocked the whole analysis pipeline:
- conversationProcessing was a single lock per conversation; processBatch
  awaited BOTH text and media worker jobs before releasing it, so a fast
  text verdict sat unused until the slow vision/media batch finished
- one global LLM semaphore (AI_LLM_MAX_CONCURRENT) was shared by text and
  media, so a vision backlog could starve text inference
- recovery worker gated on conversationProcessing.size

Now the queue is split into independent text/media lanes:
- conversationProcessing maps key -> Partial<Record<lane, startedAt>>;
  each lane holds its own lock and frees it the moment ITS worker job
  resolves (ownership-guarded clear prevents stale timers clearing newer
  slots)
- two LLM semaphores: AI_LLM_MAX_CONCURRENT (text, default 8) and
  AI_LLM_MEDIA_MAX_CONCURRENT (media, default 4) via
  withLlmConcurrency(fn, { lane })
- batchScheduler schedules per conversation+lane (timer keys
  '<key>::<lane>'); splitMessagesByLane/laneOfMessage moved to pure
  analysisLanes.ts (unit-testable without Piscina)
- ai-analysis-worker batch jobs carry a lane field; per-lane active
  request gauges (active_text_requests / active_media_requests)
- added tests/analysisLaneLock.test.ts (7 tests: independent lane locks,
  preserving other-lane lock, clear-all, ownership guard, lane split)

Docs: ARCHITECTURE.md + AGENTS.md concurrency model updated.
typecheck/lint/test(138)/build all green.
2026-09-24 14:09:53 +07:00
asepharyana f9fecfc144 chore: clean up dead barrels, duplicate config, and orphaned frontend components
Gateway:
- Remove dead barrels (ai-moderation/index, attachment-upload/index, message-capture/index) — all consumers import files directly
- Remove orphaned schema/ split dir (analytics/cache/messages/meta) — schema.ts is monolithic
- Merge duplicate config singleton: delete shared/config/config.ts, point all 44 imports at shared/config/index

Backend:
- Remove dead commandHelper.ts (voice-era fallback), ws/index.ts barrel, health.schema.ts, moderationMetrics.ts, analysis.schema.ts (0 importers; metrics/handlers route directly)

Frontend:
- Remove orphaned CategoryDrilldown/CoverageTiles/TopicTrends, primitives/slot, use-mobile, use-mounted
- Remove unused charts donut/sparkline (TopicTrends was only consumer)

Kept (verified active): shared/database/index.ts facade (11 importers), hooks/index + lib/api/index barrels (10 importers), orpc/ws.ts, charts/index.ts barrel.
Verified: tsc + biome + vitest per service (backend e2e 3 failures pre-existing on main); frontend next build 8 routes.
2026-09-24 13:23:38 +07:00
asepharyana a59f3132ee fix: add findings and progress documentation to .gitignore 2026-09-24 12:11:21 +07:00
asepharyana c7f53e4f7e feat(gateway): remove Jev (System One) analyzer, restore LLM-only text moderation
Jev (oc/jev-1.13-free via 9router /v1/systemone) added as primary text
analyzer was underperforming. Delete the whole feature:
- jevAnalyzer.ts + its unit & live-smoke tests
- Jev-first branch in textBatchProcessor, restore pure callModerationLLM path
- AI_LLM_JEV_* config vars (zod) and .env.example entries
- @typesafe-ai/sdk dependency (+ lockfile)

Behavior: text moderation is LLM-only again, exactly as before the
Jev feature; AGENTS.md invariant 'LLM is the only judge' holds.
2026-09-24 12:06:36 +07:00
asepharyana 8583bcdf17 fe(fe): enrich dashboard Top Reacted Messages
Show relative time, always-on username with channel context, and
collapse top emojis to two; add VIEW ALL toggle to reveal all 20
fetched reactions instead of the top 4.
2026-09-23 23:56:02 +07:00
asepharyana b12eb0a038 fe(fe): surface moderation explainability in live feed
Show confidence bar, flag chips, status dot, and error text per action
in the Live Stream Audit Log; drop stale media comment on SkeletonHero.
2026-09-23 23:46:22 +07:00
asepharyana b72423c64d fix(fe): remove remaining voice/recording/media leftovers from navigation, dashboard, and ws layer 2026-09-23 23:21:05 +07:00
asepharyana 6b34fc97ec Merge branch 'fix/remove-voice-recording' 2026-09-23 22:13:00 +07:00
dependabot[bot] 8ce8978755 build(deps-dev): bump tsx (#84)
Bumps the development group in /services/discord-gateway with 1 update: [tsx](https://github.com/privatenumber/tsx).


Updates `tsx` from 4.23.13 to 4.23.15
- [Release notes](https://github.com/privatenumber/tsx/releases)
- [Changelog](https://github.com/privatenumber/tsx/blob/master/release.config.cjs)
- [Commits](https://github.com/privatenumber/tsx/compare/v4.23.13...v4.23.15)

---
updated-dependencies:
- dependency-name: tsx
  dependency-version: 4.23.15
  dependency-type: direct:development
  update-type: version-update:semver-patch
  dependency-group: development
...

Signed-off-by: dependabot[bot] <support@github.com>
Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
2026-09-23 22:10:04 +07:00
mytheclipsebotreview[bot] 869cad88e4 Auto-merge PR #83
build(deps-dev): bump tsx from 4.23.13 to 4.23.15 in /services/backend in the development group
2026-09-23 14:41:14 +00:00
mytheclipsebotreview[bot] f136cf3f6b Auto-merge PR #82
build(deps): bump openai from 7.19.0 to 7.20.0 in /services/discord-gateway in the production group
2026-09-23 14:31:27 +00:00
dependabot[bot] cd3ee5b5d8 build(deps-dev): bump tsx in /services/backend in the development group
Bumps the development group in /services/backend with 1 update: [tsx](https://github.com/privatenumber/tsx).


Updates `tsx` from 4.23.13 to 4.23.15
- [Release notes](https://github.com/privatenumber/tsx/releases)
- [Changelog](https://github.com/privatenumber/tsx/blob/master/release.config.cjs)
- [Commits](https://github.com/privatenumber/tsx/compare/v4.23.13...v4.23.15)

---
updated-dependencies:
- dependency-name: tsx
  dependency-version: 4.23.15
  dependency-type: direct:development
  update-type: version-update:semver-patch
  dependency-group: development
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-09-23 14:24:58 +00:00
dependabot[bot] 93288afa0b build(deps): bump openai
Bumps the production group in /services/discord-gateway with 1 update: [openai](https://github.com/openai/openai-node).


Updates `openai` from 7.19.0 to 7.20.0
- [Release notes](https://github.com/openai/openai-node/releases)
- [Changelog](https://github.com/openai/openai-node/blob/main/CHANGELOG.md)
- [Commits](https://github.com/openai/openai-node/compare/v7.19.0...v7.20.0)

---
updated-dependencies:
- dependency-name: openai
  dependency-version: 7.20.0
  dependency-type: direct:production
  update-type: version-update:semver-minor
  dependency-group: production
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-09-23 14:24:44 +00:00
110 changed files with 1844 additions and 3307 deletions
-6
View File
@@ -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)
+10 -2
View File
@@ -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)"
}
+4
View File
@@ -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();
}
+4 -4
View File
@@ -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(),
-8
View File
@@ -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";
+8 -2
View File
@@ -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
+102 -49
View File
@@ -5,10 +5,9 @@ messages/attachments/reactions/threads/presence, runs LLM-based AI
moderation, and publishes everything to Redis pub/sub for the backend to
consume. The backend serves the HTTP/WS API to the frontend.
> NOTE: this doc is the source of truth for the module layout. The older
> `MODULE_STRUCTURE.md` was stale (referenced `winston`, `mock-crc.ts`,
> `indonesianTextNormalizer.ts`, and `aiAnalysisWorker.ts`/`llmModerationClient.ts`
> which were renamed/merged). If they disagree, this file wins.
> NOTE: this doc is the source of truth for the module layout. The old
> `MODULE_STRUCTURE.md` was a stale duplicate and has been removed. `README.md`
> only covers how to run the service.
## Top-level layout
@@ -16,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.
+50 -306
View File
@@ -1,319 +1,63 @@
# Discord Gateway Service - Extraction Complete
# Discord Gateway
## Overview
Event-driven selfbot service: captures Discord events, runs LLM moderation,
publishes everything to Redis for the backend to consume.
Successfully extracted Discord Gateway service with **Modular MVC + Event-Driven Architecture** using Redis pub/sub for inter-service communication.
> Architecture, invariants and the AI pipeline are documented in
> **`ARCHITECTURE.md`** — that file is the source of truth. This README only
> covers how to run it.
## Directory Structure
## Commands
```
services/discord-gateway/
├── src/
│ ├── app/
│ │ ├── bootstrap.ts # Service initialization (Discord client, DB, Redis)
│ │ └── shutdown.ts # Graceful shutdown handler
│ ├── shared/ # Shared infrastructure layer
│ │ ├── config/
│ │ │ └── config.ts # Zod-validated environment config
│ │ ├── database/
│ │ │ ├── schema.ts # Drizzle ORM schema
│ │ │ ├── drizzle.ts # PostgreSQL connection
│ │ │ ├── migrate.ts # Migration runner
│ │ ├── errors/
│ │ │ └── errors.ts # Custom error classes
│ │ ├── logger/
│ │ │ ├── logger.ts # Winston logger wrapper
│ │ │ └── serialization.ts # Log serialization
│ │ ├── utils/
│ │ │ └── retry.ts # Retry with exponential backoff
│ │ └── discord/
│ │ └── clientOptions.ts # Discord.js client config
│ ├── modules/ # Feature modules (Modular MVC)
│ │ ├── message-capture/ # Controller-Service-Repository
│ │ │ ├── messageCapture.ts # Controller: Discord event listeners
│ │ │ ├── messageStore.ts # Repository: DB operations
│ │ │ ├── messageMetadata.ts # Service: Metadata extraction
│ │ │ ├── types.ts # Domain types
│ │ │ └── index.ts # Module exports
│ │ ├── ai-moderation/ # Controller-Service-Repository
│ │ │ ├── aiAnalyzer.ts # Controller: Analysis orchestration
│ │ │ ├── llmModerationClient.ts # Service: LLM API client
│ │ │ ├── aiAnalysisWorker.ts # Service: Worker pool
│ │ │ ├── indonesianTextNormalizer.ts # Service: Text normalization
│ │ │ ├── moderationPrompt.ts # Service: Prompt generation
│ │ │ └── index.ts # Module exports
│ │ ├── attachment-upload/ # Controller-Service-Repository
│ │ │ ├── attachmentUploader.ts # Service: Upload orchestration
│ │ │ ├── imageResizer.ts # Service: Image resizing
│ │ │ └── index.ts # Module exports
│ │ └── event-broadcaster/ # Event-driven layer
│ │ ├── eventBroadcaster.ts # Service: Redis pub/sub publisher
│ │ ├── eventTypes.ts # Domain: Event type definitions
│ │ └── index.ts # Module exports
│ ├── mock-crc.ts # CRC polyfill for discord.js
│ └── index.ts # Service entry point
├── ARCHITECTURE.md # Detailed architecture documentation
├── package.json # Service dependencies
└── tsconfig.json # TypeScript configuration (inherited)
```bash
pnpm install
pnpm typecheck # tsc --noEmit
pnpm lint # biome check --diagnostic-level=error .
pnpm test # vitest run (138 tests)
pnpm build # tsc — CI/prod builds run this inside nix, which also
# runs scripts/fix-imports.mjs to rewrite @/ aliases and
# extensionless imports for Node ESM
pnpm dev # tsx watch src/index.ts
pnpm start # node dist/index.js
```
## Architecture Patterns
Deployment is CI-only: `nix build .#discord-gateway` → Attic cache → systemd
restart on the VPS. Do not build/hand-copy the artifact.
### 1. Modular MVC Structure
Each feature module follows **Controller-Service-Repository** pattern:
**Message Capture Module**:
- **Controller** (`messageCapture.ts`): Listens to Discord events (messageCreate, messageUpdate, messageDelete)
- **Service** (`messageMetadata.ts`): Extracts and normalizes message metadata
- **Repository** (`messageStore.ts`): Database CRUD operations
**AI Moderation Module**:
- **Controller** (`aiAnalyzer.ts`): Orchestrates analysis workflow
- **Service** (`llmModerationClient.ts`): LLM API integration
- **Service** (`aiAnalysisWorker.ts`): Worker pool management
- **Service** (`indonesianTextNormalizer.ts`): Text preprocessing
**Attachment Upload Module**:
- **Service** (`attachmentUploader.ts`): Upload orchestration
- **Service** (`imageResizer.ts`): Image processing
### 2. Event-Driven Architecture
**Redis Pub/Sub** replaces WebSocket broadcaster:
## Layout
```
Discord Events → Discord Gateway Service → Redis Pub/Sub → Backend Service
↓
Event Channels:
- discord:message:created
- discord:message:updated
- discord:message:deleted
- discord:message:analyzed
- discord:attachment:created
- discord:attachment:uploaded
- discord:analysis:queue_status
src/
├── index.ts # entry → initializeDiscordGateway()
├── app/ # process lifecycle
│ ├── bootstrap.ts # startup order: config → DB → services → metrics → login
│ ├── lifecycle.ts # everything wired on the Discord 'ready' hook
│ ├── process-guards.ts # SIGINT/SIGTERM + uncaught error policy
│ ├── metrics-collector.ts # AI pipeline Prometheus gauges
│ ├── shutdown.ts # graceful shutdown sequence
│ └── retention.ts # expired-record cleanup scheduler
├── shared/ # infrastructure — never imports from modules/
│ ├── config/ database/ logger/ errors/ utils/
│ ├── discord/clientOptions.ts
│ ├── redis-channels.ts # canonical Redis channel + command constants
│ └── moderation-types.ts # domain types shared across services
└── modules/ # feature modules (each exposes an index.ts facade)
├── ai-moderation/ # LLM moderation pipeline (largest module)
├── message-capture/ # Discord listeners + message/attachment DB
├── attachment-upload/ # download → resize → upload
├── event-broadcaster/ # Redis pub/sub publisher
├── command-handler/ # backend → gateway commands over Redis
├── gateway-metrics/ # Prometheus /metrics (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`).
+1 -2
View File
@@ -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"
}
+14 -23
View File
@@ -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
+72 -196
View File
@@ -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";
@@ -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";
@@ -58,14 +52,14 @@ interface UrlFetchResult {
title: 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;
}
@@ -419,7 +413,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 +434,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 +460,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 +477,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),
@@ -1,6 +1,6 @@
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";
import { executeAll, executeGet } from "../../shared/database/drizzle.js";
import { findBestEmbeddingMatch } from "./embeddingClient.js";
import {

Some files were not shown because too many files have changed in this diff Show More