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.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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<Record<lane, startedAt>>.
|
||||
// Prune stale lane slots individually so one stale lane never clears
|
||||
// the other lane's healthy lock.
|
||||
for (const [key, record] of conversationProcessing) {
|
||||
for (const lane of ANALYSIS_LANES as readonly AnalysisLane[]) {
|
||||
const startedAt = record?.[lane];
|
||||
if (
|
||||
startedAt &&
|
||||
now - startedAt >= config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS
|
||||
) {
|
||||
clearConversationProcessing(key, lane);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
@@ -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 };
|
||||
}
|
||||
@@ -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),
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
@@ -46,26 +66,40 @@ 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.
|
||||
// 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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,24 +81,106 @@ 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);
|
||||
const now = Date.now();
|
||||
if (lane) {
|
||||
const startedAt = conversationProcessing.get(conversationKey)?.[lane];
|
||||
return Boolean(
|
||||
startedAt &&
|
||||
Date.now() - startedAt < config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS,
|
||||
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);
|
||||
});
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Alert system
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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<typeof pLimit>;
|
||||
limit: number;
|
||||
}
|
||||
return llmSemaphore;
|
||||
|
||||
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 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",
|
||||
);
|
||||
}
|
||||
@@ -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,14 +297,18 @@ export async function llmChat(
|
||||
|
||||
return retryWithBackoff(
|
||||
async () => {
|
||||
return withLlmConcurrency(async () => {
|
||||
return withLlmConcurrency(
|
||||
async () => {
|
||||
const execute = async (
|
||||
currentParams: OpenAI.Chat.Completions.ChatCompletionCreateParams,
|
||||
) => {
|
||||
const response = await client.chat.completions.create(currentParams, {
|
||||
const response = await client.chat.completions.create(
|
||||
currentParams,
|
||||
{
|
||||
signal,
|
||||
...(opts.timeout ? { timeout: opts.timeout } : {}),
|
||||
});
|
||||
},
|
||||
);
|
||||
if (currentParams.stream) {
|
||||
let content = "";
|
||||
let finishReason = "stop";
|
||||
@@ -326,7 +372,9 @@ export async function llmChat(
|
||||
);
|
||||
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
|
||||
|
||||
@@ -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 },
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user