fix(gateway): separate text and media lanes in AI analysis queue

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:
asepharyana
2026-09-24 14:06:58 +07:00
parent f9fecfc144
commit 6d0d7b34a3
15 changed files with 692 additions and 205 deletions
+7 -1
View File
@@ -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
+32 -18
View File
@@ -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;
@@ -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);
}
}
@@ -24,6 +24,19 @@ import { LAST_ERROR } from "./moderationState.js";
* - Alert system: `CircuitBreakerAlert` type, `fireAlert()`, and
* `onCircuitBreakerAlert()` for pluggable handler registration.
*
* ## Processing lanes (2026-09-24)
* A conversation batch splits into a **text lane** (messages with no media)
* and a **media lane** (messages with attachments/stickers/embeds). The two
* lanes are dispatched to separate Piscina pools and MUST NOT block each
* other: a fast text sub-batch must be free to finish while the slow
* vision/media sub-batch of the SAME conversation is still running.
*
* The lock is therefore per-lane: `conversationProcessing` maps a
* conversation key to its current processing record which carries the lane
* name. `isConversationProcessingLocked(key, lane)` reports locked only when
* the SAME lane (or all lanes when lane is omitted) is active — a media
* sub-batch in flight never blocks scheduling the text sub-batch.
*
* ## Relationship with moderationState.ts
* - `moderationState.ts` owns **infrastructure references** (event broadcaster,
* Discord client), the auto-delete guard, the `LAST_ERROR` tracker, and
@@ -33,6 +46,14 @@ import { LAST_ERROR } from "./moderationState.js";
* - These are **separate concerns** — do not merge them.
*/
/** Processing lanes for conversation analysis. */
export type AnalysisLane = "text" | "media";
export const ANALYSIS_LANES: readonly AnalysisLane[] = [
"text",
"media",
] as const;
const logger = createChildLogger("conversation-state");
// ---------------------------------------------------------------------------
@@ -60,23 +81,105 @@ export const conversationDebounceTimers = new LRUCache<string, NodeJS.Timeout>({
},
});
/** Timestamp of when processing started per conversation key. */
export const conversationProcessing = new LRUCache<string, number>({
max: 10000,
});
/**
* Per-conversation processing lock, keyed by lane.
*
* A conversation can hold TWO locks at once — one for its text sub-batch and
* one for its media sub-batch — because the two lanes run on separate pools
* and finish independently. The value is a partial record of lane →
* startedAt; clearing one lane leaves the other lane's lock intact.
*/
export const conversationProcessing = new LRUCache<
string,
Partial<Record<AnalysisLane, number>>
>({ max: 10000 });
/**
* Locks a conversation for the given lane.
* The same conversation can be locked in both lanes simultaneously (text and
* media sub-batches run independently); locking an already-locked lane
* replaces its startedAt (last writer wins, matching the old single-lock
* semantics).
*/
export function setConversationProcessing(
conversationKey: string,
lane: AnalysisLane,
startedAt: number,
): void {
const record = conversationProcessing.get(conversationKey) ?? {};
conversationProcessing.set(conversationKey, { ...record, [lane]: startedAt });
}
/**
* Releases the processing lock for a conversation in a SINGLE lane.
* The other lane's lock (if any) is preserved.
*/
export function clearConversationProcessing(
conversationKey: string,
lane: AnalysisLane,
): void {
const record = conversationProcessing.get(conversationKey);
if (!record) return;
const next = { ...record };
delete next[lane];
if (Object.keys(next).length === 0) {
conversationProcessing.delete(conversationKey);
} else {
conversationProcessing.set(conversationKey, next);
}
}
/**
* Clears the processing lock for a conversation regardless of lane.
* Used by the recovery worker when a lock is stale. If only ONE lane of a
* two-lane processing conversation is stale, prefer clearConversationProcessing
* with the specific lane to keep the healthy lane's lock intact.
*/
export function clearConversationProcessingAll(conversationKey: string): void {
conversationProcessing.delete(conversationKey);
}
/**
* Returns the startedAt for a conversation in a lane, or undefined.
* Consumers use this to verify a processing slot is still owned by them
* before releasing it (guards against clearing a newer slot).
*/
export function getConversationProcessingStartedAt(
conversationKey: string,
lane: AnalysisLane,
): number | undefined {
return conversationProcessing.get(conversationKey)?.[lane];
}
// ---------------------------------------------------------------------------
// Conversation lock helper
// ---------------------------------------------------------------------------
/**
* Reports whether the conversation is currently processing.
*
* When `lane` is provided, only that lane's lock counts — a media sub-batch
* in flight does NOT lock the text lane, so the text lane can be scheduled
* and vice versa. When `lane` is omitted, any active lane locks it (used by
* recovery/individual fallback which must not race ANY batch work).
*/
export function isConversationProcessingLocked(
conversationKey: string,
lane?: AnalysisLane,
): boolean {
const startedAt = conversationProcessing.get(conversationKey);
return Boolean(
startedAt &&
Date.now() - startedAt < config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS,
);
const now = Date.now();
if (lane) {
const startedAt = conversationProcessing.get(conversationKey)?.[lane];
return Boolean(
startedAt && now - startedAt < config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS,
);
}
const record = conversationProcessing.get(conversationKey);
if (!record) return false;
return ANALYSIS_LANES.some((l) => {
const s = record[l];
return Boolean(s && now - s < config.AI_ANALYSIS_PROCESSING_TIMEOUT_MS);
});
}
// ---------------------------------------------------------------------------
@@ -86,6 +189,7 @@ export function isConversationProcessingLocked(
export type CircuitBreakerAlert = {
type: "conversation_cb" | "individual_cb" | "sustained_error";
conversationKey?: string;
lane?: AnalysisLane;
consecutiveErrors: number;
message: string;
lastError?: string | null;
@@ -117,7 +221,10 @@ export function fireAlert(alert: CircuitBreakerAlert): void {
// Circuit breaker helpers
// ---------------------------------------------------------------------------
export function recordConversationBatchFailure(conversationKey: string): void {
export function recordConversationBatchFailure(
conversationKey: string,
lane?: AnalysisLane,
): void {
const nextCount =
(conversationConsecutiveErrors.get(conversationKey) ?? 0) + 1;
conversationConsecutiveErrors.set(conversationKey, nextCount);
@@ -130,6 +237,7 @@ export function recordConversationBatchFailure(conversationKey: string): void {
fireAlert({
type: "conversation_cb",
conversationKey,
lane,
consecutiveErrors: nextCount,
message: `Conversation ${conversationKey} circuit breaker triggered after ${nextCount} consecutive errors`,
lastError: LAST_ERROR.value,
@@ -138,6 +246,9 @@ export function recordConversationBatchFailure(conversationKey: string): void {
}
}
export function resetConversationBatchFailures(conversationKey: string): void {
export function resetConversationBatchFailures(
conversationKey: string,
_lane?: AnalysisLane,
): void {
conversationConsecutiveErrors.delete(conversationKey);
}
@@ -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;
}
const laneSemaphores: Record<LlmLane, LaneSemaphore> = {
text: {
limiter: pLimit(config.AI_LLM_MAX_CONCURRENT ?? 5),
limit: config.AI_LLM_MAX_CONCURRENT ?? 5,
},
media: {
limiter: pLimit(config.AI_LLM_MEDIA_MAX_CONCURRENT ?? 4),
limit: config.AI_LLM_MEDIA_MAX_CONCURRENT ?? 4,
},
};
function getLaneSemaphore(lane: LlmLane): ReturnType<typeof pLimit> {
const wanted =
lane === "media"
? (config.AI_LLM_MEDIA_MAX_CONCURRENT ?? 4)
: (config.AI_LLM_MAX_CONCURRENT ?? 5);
const slot = laneSemaphores[lane];
if (wanted !== slot.limit) {
slot.limiter = pLimit(wanted);
slot.limit = wanted;
}
return llmSemaphore;
return slot.limiter;
}
let activeCount = 0;
let pendingCount = 0;
export async function withLlmConcurrency<T>(fn: () => Promise<T>): Promise<T> {
/**
* Run `fn` under the per-lane LLM concurrency cap.
*
* `lane: "text"` uses `AI_LLM_MAX_CONCURRENT`; `lane: "media"` uses
* `AI_LLM_MEDIA_MAX_CONCURRENT`. Defaults to "text" so the existing text
* moderation path is unchanged.
*/
export async function withLlmConcurrency<T>(
fn: () => Promise<T>,
opts: { lane?: LlmLane } = {},
): Promise<T> {
const lane = opts.lane ?? "text";
const maxConcurrent =
lane === "media"
? (config.AI_LLM_MEDIA_MAX_CONCURRENT ?? 4)
: (config.AI_LLM_MAX_CONCURRENT ?? 5);
pendingCount++;
log.debug(
{ activeCount, pendingCount, maxConcurrent: config.AI_LLM_MAX_CONCURRENT },
{ activeCount, pendingCount, maxConcurrent, lane },
"Queuing LLM request",
);
return getLlmSemaphore()(async () => {
return getLaneSemaphore(lane)(async () => {
pendingCount--;
activeCount++;
if (activeCount >= (config.AI_LLM_MAX_CONCURRENT ?? 5)) {
if (activeCount >= maxConcurrent) {
log.warn(
{ activeCount, maxConcurrent: config.AI_LLM_MAX_CONCURRENT },
{ activeCount, maxConcurrent, lane },
"LLM concurrency limit reached",
);
}
@@ -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<LLMResponseChunk>) {
const choice = chunk?.choices?.[0];
content += extractChunkText(chunk);
const fr = choice?.finish_reason || chunk?.finish_reason;
if (fr) finishReason = fr;
}
return {
id: "stream-aggregated",
choices: [
{
message: { role: "assistant", content, refusal: null },
finish_reason: finishReason,
index: 0,
logprobs: null,
},
],
created: Math.floor(Date.now() / 1000),
model: currentParams.model,
object: "chat.completion",
} as OpenAI.Chat.Completions.ChatCompletion;
}
return response as OpenAI.Chat.Completions.ChatCompletion;
};
try {
return await execute(params);
} catch (error: any) {
const rawResponse =
error.error || error.body || error.response?.data || "N/A";
const errorStr = (
JSON.stringify(rawResponse) + String(error.message)
).toLowerCase();
// Auto-fallback: If provider strictly demands streaming (400 Bad Request on stream params)
if (
error.status === 400 &&
errorStr.includes("stream") &&
!params.stream
) {
log.warn(
{ model },
"Provider rejected non-streaming request. Fallback to stream: true initiated.",
return withLlmConcurrency(
async () => {
const execute = async (
currentParams: OpenAI.Chat.Completions.ChatCompletionCreateParams,
) => {
const response = await client.chat.completions.create(
currentParams,
{
signal,
...(opts.timeout ? { timeout: opts.timeout } : {}),
},
);
(
params as unknown as OpenAI.Chat.Completions.ChatCompletionCreateParamsStreaming
).stream = true;
return await execute(params);
}
if (currentParams.stream) {
let content = "";
let finishReason = "stop";
for await (const chunk of response as unknown as AsyncIterable<LLMResponseChunk>) {
const choice = chunk?.choices?.[0];
content += extractChunkText(chunk);
const fr = choice?.finish_reason || chunk?.finish_reason;
if (fr) finishReason = fr;
}
return {
id: "stream-aggregated",
choices: [
{
message: { role: "assistant", content, refusal: null },
finish_reason: finishReason,
index: 0,
logprobs: null,
},
],
created: Math.floor(Date.now() / 1000),
model: currentParams.model,
object: "chat.completion",
} as OpenAI.Chat.Completions.ChatCompletion;
}
return response as OpenAI.Chat.Completions.ChatCompletion;
};
log.error(
{
error: error.message,
status: error.status,
rawResponse,
model,
},
"LLM API request failed",
);
throw error;
}
});
try {
return await execute(params);
} catch (error: any) {
const rawResponse =
error.error || error.body || error.response?.data || "N/A";
const errorStr = (
JSON.stringify(rawResponse) + String(error.message)
).toLowerCase();
// Auto-fallback: If provider strictly demands streaming (400 Bad Request on stream params)
if (
error.status === 400 &&
errorStr.includes("stream") &&
!params.stream
) {
log.warn(
{ model },
"Provider rejected non-streaming request. Fallback to stream: true initiated.",
);
(
params as unknown as OpenAI.Chat.Completions.ChatCompletionCreateParamsStreaming
).stream = true;
return await execute(params);
}
log.error(
{
error: error.message,
status: error.status,
rawResponse,
model,
},
"LLM API request failed",
);
throw error;
}
},
{ lane: opts.lane ?? "text" },
);
},
{
retries,
@@ -356,7 +404,7 @@ export async function llmVision(
promptText: string,
imageUrl: { url: string },
): Promise<string | null> {
const params = {
const params: LlmCallOpts = {
messages: [
{
role: "user" as const,
@@ -372,6 +420,7 @@ export async function llmVision(
top_p: 0.9,
retries: 0,
timeout: config.AI_LLM_VISION_ANALYSIS_TIMEOUT_MS ?? 60_000,
lane: "media",
};
// Streaming first (the router always streams SSE; a non-stream request
@@ -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");
});
});