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.
36 lines
1.3 KiB
TypeScript
36 lines
1.3 KiB
TypeScript
/**
|
|
* 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 };
|
|
}
|