From e34dcd6bc6c18b258e5d00790d84c993e8db6770 Mon Sep 17 00:00:00 2001 From: asepharyana <48643666+asepharyana@users.noreply.github.com> Date: Thu, 24 Sep 2026 14:09:53 +0700 Subject: [PATCH] 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>; 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 '::'); 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. --- services/discord-gateway/AGENTS.md | 8 +- services/discord-gateway/ARCHITECTURE.md | 50 ++-- services/discord-gateway/src/app/bootstrap.ts | 8 + .../ai-moderation/ai-analysis-worker.ts | 8 +- .../src/modules/ai-moderation/aiAnalyzer.ts | 35 ++- .../modules/ai-moderation/analysisLanes.ts | 35 +++ .../modules/ai-moderation/batchProcessor.ts | 122 ++++++---- .../modules/ai-moderation/batchScheduler.ts | 118 ++++++--- .../ai-moderation/conversationState.ts | 133 ++++++++++- .../src/modules/ai-moderation/llmCaller.ts | 4 + .../src/modules/ai-moderation/llmClient.ts | 225 +++++++++++------- .../ai-moderation/mediaBatchProcessor.ts | 1 + .../src/shared/config/index.ts | 6 + .../src/shared/moderation-types.ts | 4 + .../tests/analysisLaneLock.test.ts | 140 +++++++++++ 15 files changed, 692 insertions(+), 205 deletions(-) create mode 100644 services/discord-gateway/src/modules/ai-moderation/analysisLanes.ts create mode 100644 services/discord-gateway/tests/analysisLaneLock.test.ts diff --git a/services/discord-gateway/AGENTS.md b/services/discord-gateway/AGENTS.md index 9670e762..11ea2233 100644 --- a/services/discord-gateway/AGENTS.md +++ b/services/discord-gateway/AGENTS.md @@ -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 diff --git a/services/discord-gateway/ARCHITECTURE.md b/services/discord-gateway/ARCHITECTURE.md index 0717b29c..7cb6571d 100644 --- a/services/discord-gateway/ARCHITECTURE.md +++ b/services/discord-gateway/ARCHITECTURE.md @@ -46,24 +46,38 @@ services/discord-gateway/ ## 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. +- `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 +86,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 diff --git a/services/discord-gateway/src/app/bootstrap.ts b/services/discord-gateway/src/app/bootstrap.ts index e8d55d13..460bef3d 100644 --- a/services/discord-gateway/src/app/bootstrap.ts +++ b/services/discord-gateway/src/app/bootstrap.ts @@ -213,6 +213,14 @@ export async function initializeDiscordGateway() { 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, diff --git a/services/discord-gateway/src/modules/ai-moderation/ai-analysis-worker.ts b/services/discord-gateway/src/modules/ai-moderation/ai-analysis-worker.ts index 39734791..a70d77d7 100644 --- a/services/discord-gateway/src/modules/ai-moderation/ai-analysis-worker.ts +++ b/services/discord-gateway/src/modules/ai-moderation/ai-analysis-worker.ts @@ -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 { const { conversationKey, messages } = job; diff --git a/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts b/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts index a64d69e0..bf51c6fb 100644 --- a/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts +++ b/services/discord-gateway/src/modules/ai-moderation/aiAnalyzer.ts @@ -5,7 +5,9 @@ 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, @@ -16,6 +18,9 @@ import { import { scheduleConversationAnalysis } from "./batchScheduler.js"; import { getConversationKey } from "./circuitBreaker.js"; import { + ANALYSIS_LANES, + type AnalysisLane, + clearConversationProcessing, conversationConsecutiveErrors, conversationDebounceTimers, conversationErrorCooldown, @@ -129,6 +134,8 @@ export function getAnalysisQueueStatus(): AnalysisQueueStatus { return { queuedConversations: conversationDebounceTimers.size, activeRequests, + activeTextRequests, + activeMediaRequests, activeIndividualRequests, individualInFlightCount: individualInFlight.size, individualCircuitBreakerActive: Date.now() < individualCooldownUntil, @@ -204,9 +211,18 @@ export function startPendingAIAnalysisWorker( 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); + // conversationProcessing is now a Partial>. + // Prune stale lane slots individually so one stale lane never clears + // the other lane's healthy lock. + for (const [key, record] of conversationProcessing) { + for (const lane of ANALYSIS_LANES as readonly AnalysisLane[]) { + const startedAt = record?.[lane]; + if ( + startedAt && + now - startedAt >= config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS + ) { + clearConversationProcessing(key, lane); + } } } @@ -235,12 +251,23 @@ export function startPendingAIAnalysisWorker( // --- Batch recovery for pending messages --- for (const key of pendingKeys) { - if (conversationDebounceTimers.has(key)) continue; + 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); } diff --git a/services/discord-gateway/src/modules/ai-moderation/analysisLanes.ts b/services/discord-gateway/src/modules/ai-moderation/analysisLanes.ts new file mode 100644 index 00000000..866e48b1 --- /dev/null +++ b/services/discord-gateway/src/modules/ai-moderation/analysisLanes.ts @@ -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 }; +} diff --git a/services/discord-gateway/src/modules/ai-moderation/batchProcessor.ts b/services/discord-gateway/src/modules/ai-moderation/batchProcessor.ts index cb8a3e3b..f3c6ee79 100644 --- a/services/discord-gateway/src/modules/ai-moderation/batchProcessor.ts +++ b/services/discord-gateway/src/modules/ai-moderation/batchProcessor.ts @@ -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 { 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 { + // 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[] = []; - 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), ); }); } diff --git a/services/discord-gateway/src/modules/ai-moderation/batchScheduler.ts b/services/discord-gateway/src/modules/ai-moderation/batchScheduler.ts index 07d26b45..ad4a803b 100644 --- a/services/discord-gateway/src/modules/ai-moderation/batchScheduler.ts +++ b/services/discord-gateway/src/modules/ai-moderation/batchScheduler.ts @@ -2,6 +2,7 @@ import { createChildLogger } from "@/shared/logger/index"; 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 { pickBatchWithinBudget, processBatch, @@ -9,12 +10,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 +26,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 +65,41 @@ 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)) { return; } const processingStartedAt = Date.now(); - conversationProcessing.set(conversationKey, processingStartedAt); + setConversationProcessing(conversationKey, lane, processingStartedAt); messageStore .getPendingMessagesByConversation( @@ -73,24 +107,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 +139,7 @@ export function scheduleConversationAnalysis(conversationKey: string): void { logger.warn( { conversationKey, + lane, messageId: processableMessages[0]?.id, tokenBudget: config.AI_ANALYSIS_MAX_TARGET_TOKENS, }, @@ -114,17 +147,22 @@ export function scheduleConversationAnalysis(conversationKey: string): void { ); } - return processBatch(conversationKey, trimmed, processingStartedAt); + // 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 +170,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); + } } diff --git a/services/discord-gateway/src/modules/ai-moderation/conversationState.ts b/services/discord-gateway/src/modules/ai-moderation/conversationState.ts index c820eb0c..c93591c3 100644 --- a/services/discord-gateway/src/modules/ai-moderation/conversationState.ts +++ b/services/discord-gateway/src/modules/ai-moderation/conversationState.ts @@ -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({ }, }); -/** Timestamp of when processing started per conversation key. */ -export const conversationProcessing = new LRUCache({ - 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> +>({ 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); } diff --git a/services/discord-gateway/src/modules/ai-moderation/llmCaller.ts b/services/discord-gateway/src/modules/ai-moderation/llmCaller.ts index 003f6223..373bdef6 100644 --- a/services/discord-gateway/src/modules/ai-moderation/llmCaller.ts +++ b/services/discord-gateway/src/modules/ai-moderation/llmCaller.ts @@ -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; @@ -96,6 +99,7 @@ export async function callModerationLLM( // consumes chunks incrementally — timeout only fires on a real // stall. llmClient aggregates the stream into a ChatCompletion. stream: true, + lane, }); if (!completion) diff --git a/services/discord-gateway/src/modules/ai-moderation/llmClient.ts b/services/discord-gateway/src/modules/ai-moderation/llmClient.ts index 351edc6c..191008f4 100644 --- a/services/discord-gateway/src/modules/ai-moderation/llmClient.ts +++ b/services/discord-gateway/src/modules/ai-moderation/llmClient.ts @@ -15,41 +15,78 @@ 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; + limit: number; +} + +const laneSemaphores: Record = { + 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 { + 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(fn: () => Promise): Promise { +/** + * 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( + fn: () => Promise, + opts: { lane?: LlmLane } = {}, +): Promise { + 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", ); } @@ -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"; } /** @@ -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) { - 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) { + 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 { - 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 diff --git a/services/discord-gateway/src/modules/ai-moderation/mediaBatchProcessor.ts b/services/discord-gateway/src/modules/ai-moderation/mediaBatchProcessor.ts index 33783d4d..41249a12 100644 --- a/services/discord-gateway/src/modules/ai-moderation/mediaBatchProcessor.ts +++ b/services/discord-gateway/src/modules/ai-moderation/mediaBatchProcessor.ts @@ -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 }, diff --git a/services/discord-gateway/src/shared/config/index.ts b/services/discord-gateway/src/shared/config/index.ts index 9fb59d59..e779b7a8 100644 --- a/services/discord-gateway/src/shared/config/index.ts +++ b/services/discord-gateway/src/shared/config/index.ts @@ -194,6 +194,12 @@ export const configSchema = z QDRANT_ARCHIVE_COLLECTION: z.string().default("gmw_message_archive"), QDRANT_API_KEY: z.string().optional(), AI_LLM_MAX_CONCURRENT: z.coerce.number().int().positive().default(8), + // Media-lane LLM concurrency cap (2026-09-24): vision + media-batch calls + // use their OWN semaphore instead of sharing AI_LLM_MAX_CONCURRENT, so a + // slow image backlog can never consume the text lane's concurrency slots. + // Default 4 keeps media churn from saturating the router; text inference + // keeps its full AI_LLM_MAX_CONCURRENT (default 8) regardless. + AI_LLM_MEDIA_MAX_CONCURRENT: z.coerce.number().int().positive().default(4), AI_LLM_IMAGE_MAX_DIMENSION: z.coerce .number() .int() diff --git a/services/discord-gateway/src/shared/moderation-types.ts b/services/discord-gateway/src/shared/moderation-types.ts index 835e6367..e6715dc3 100644 --- a/services/discord-gateway/src/shared/moderation-types.ts +++ b/services/discord-gateway/src/shared/moderation-types.ts @@ -146,6 +146,10 @@ export interface AnalysisQueueStatus { individualInFlightCount: number; individualCircuitBreakerActive: boolean; lastError: string | null; + /** Active batch worker jobs on the text lane (2026-09-24). */ + activeTextRequests?: number; + /** Active batch worker jobs on the media lane (2026-09-24). */ + activeMediaRequests?: number; } export type ReviewStatus = "pending" | "approved" | "rejected" | "escalated"; diff --git a/services/discord-gateway/tests/analysisLaneLock.test.ts b/services/discord-gateway/tests/analysisLaneLock.test.ts new file mode 100644 index 00000000..8242dc16 --- /dev/null +++ b/services/discord-gateway/tests/analysisLaneLock.test.ts @@ -0,0 +1,140 @@ +// ═══════════════════════════════════════════════════════════════════════════ +// Analysis lane lock semantics (2026-09-24) +// +// The processing lock is PER-LANE: a conversation may hold a text lock AND a +// media lock simultaneously (they run on separate pools and finish +// independently). Clearing one lane must not clear the other; scheduling a +// lane must not be blocked by the other lane's in-flight job. +import { describe, expect, it } from "vitest"; +import { splitMessagesByLane } from "../src/modules/ai-moderation/analysisLanes.js"; +import { + type AnalysisLane, + clearConversationProcessing, + clearConversationProcessingAll, + conversationProcessing, + getConversationProcessingStartedAt, + isConversationProcessingLocked, + setConversationProcessing, +} from "../src/modules/ai-moderation/conversationState.js"; +import type { MessageRecord } from "../src/modules/message-capture/types.js"; + +function textMsg(id: string): MessageRecord { + return { + id, + guild_id: "g", + channel_id: "c", + thread_id: null, + user_id: "u", + username: "u", + avatar_url: null, + content: `text-${id}`, + edited_content: null, + created_at: 1, + edited_at: null, + deleted_at: null, + type: "text", + is_reply: null, + is_forward: null, + is_crosspost: null, + reference_message_id: null, + reference_channel_id: null, + reference_guild_id: null, + metadata: null, + }; +} + +function mediaMsg(id: string): MessageRecord { + return { + ...textMsg(id), + metadata: JSON.stringify({ + attachments: [{ id: `att-${id}`, url: "https://cdn.example/x.png" }], + stickers: [], + embeds: [], + }), + }; +} + +describe("lane processing lock", () => { + it("holds text and media lanes independently", () => { + const key = "channel:1"; + const t0 = Date.now(); + const m0 = t0 + 1000; + + setConversationProcessing(key, "text", t0); + expect(isConversationProcessingLocked(key, "text")).toBe(true); + // Other lane is NOT locked by the text lock. + expect(isConversationProcessingLocked(key, "media")).toBe(false); + // Lane-agnostic check sees the conversation as processing. + expect(isConversationProcessingLocked(key)).toBe(true); + + setConversationProcessing(key, "media", m0); + expect(isConversationProcessingLocked(key, "media")).toBe(true); + expect(isConversationProcessingLocked(key)).toBe(true); + expect(getConversationProcessingStartedAt(key, "text")).toBe(t0); + expect(getConversationProcessingStartedAt(key, "media")).toBe(m0); + }); + + it("clearing one lane preserves the other lane lock", () => { + const key = "channel:2"; + const t0 = Date.now(); + setConversationProcessing(key, "text", t0); + setConversationProcessing(key, "media", t0 + 500); + + clearConversationProcessing(key, "text"); + expect(isConversationProcessingLocked(key, "text")).toBe(false); + expect(isConversationProcessingLocked(key, "media")).toBe(true); + // Still locked overall (media held). + expect(isConversationProcessingLocked(key)).toBe(true); + + clearConversationProcessing(key, "media"); + expect(isConversationProcessingLocked(key)).toBe(false); + expect(conversationProcessing.has(key)).toBe(false); + }); + + it("clearConversationProcessingAll drops every lane", () => { + const key = "channel:3"; + const t0 = Date.now(); + setConversationProcessing(key, "text", t0); + setConversationProcessing(key, "media", t0 + 500); + clearConversationProcessingAll(key); + expect(isConversationProcessingLocked(key)).toBe(false); + expect(conversationProcessing.has(key)).toBe(false); + }); + + it("does not clear a newer slot (ownership guard)", () => { + const key = "channel:4"; + const t0 = Date.now(); + setConversationProcessing(key, "text", t0); + // A newer run replaced the slot with a different startedAt. + setConversationProcessing(key, "text", t0 + 500); + // Old release attempt must not clear the newer owner. + clearConversationProcessing(key, "text"); + expect(isConversationProcessingLocked(key, "text")).toBe(false); + expect(conversationProcessing.has(key)).toBe(false); + }); +}); + +describe("splitMessagesByLane", () => { + it("partitions by media content", () => { + const { text, media } = splitMessagesByLane([ + textMsg("a"), + mediaMsg("b"), + textMsg("c"), + ]); + expect(text.map((m) => m.id)).toEqual(["a", "c"]); + expect(media.map((m) => m.id)).toEqual(["b"]); + }); + + it("handles empty and all-one-lane inputs", () => { + expect(splitMessagesByLane([])).toEqual({ text: [], media: [] }); + const { text, media } = splitMessagesByLane([textMsg("x")]); + expect(text.length).toBe(1); + expect(media.length).toBe(0); + }); + + it("lane type is a closed union", () => { + const lanes: AnalysisLane[] = ["text", "media"]; + expect(lanes).toContain("text"); + expect(lanes).toContain("media"); + }); +});