feat(gateway): implement separate Piscina pools for text and media analysis to optimize processing
This commit is contained in:
@@ -193,6 +193,7 @@ export const configSchema = z
|
||||
.positive()
|
||||
.default(50),
|
||||
PISCINA_MAX_THREADS: z.coerce.number().int().positive().optional(),
|
||||
PISCINA_MEDIA_MAX_THREADS: z.coerce.number().int().positive().optional(),
|
||||
|
||||
// ── Voice Transcription ────────────────────────────────────────────────
|
||||
AI_VOICE_TRANSCRIPTION_ENABLED: z
|
||||
|
||||
@@ -77,20 +77,26 @@ handles a whole batch (text + media split internally, parallel paths).
|
||||
|
||||
- Main thread owns the LLM semaphore (`AI_LLM_MAX_CONCURRENT`, default 5) via
|
||||
`llmClient.withLlmConcurrency`.
|
||||
- Piscina pool (`PISCINA_MAX_THREADS`, default 4) runs the heavy LLM work off
|
||||
the event loop; **each worker thread initializes its own pg Pool** (min 0,
|
||||
grows to `POSTGRES_POOL_MAX`). See "Memory & connections" below.
|
||||
- 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.
|
||||
|
||||
## Memory & DB connections
|
||||
|
||||
`MemoryMax=1G` (raised from 512M — live RSS sits at ~500 MiB, peak 508 MiB,
|
||||
so 512M left ~2% headroom and risked an OOM-kill restart). Host has 8 GB free.
|
||||
|
||||
`POSTGRES_POOL_MIN=0` (default). The gateway = main process + up to 4 Piscina
|
||||
worker threads, each with its own pg Pool. With min:0 the pools stay empty
|
||||
until a query runs and drop idle clients afterward, instead of holding
|
||||
`(1 main + 4 workers) × 2 = 10` permanently-open idle connections against
|
||||
PgBouncer. The pool still grows on demand up to `POSTGRES_POOL_MAX`.
|
||||
`POSTGRES_POOL_MIN=0` (default). The gateway = main process + up to 4 text
|
||||
Piscina worker threads + up to 2 media Piscina worker threads, each with its
|
||||
own pg Pool. With min:0 the pools stay empty until a query runs and drop
|
||||
idle clients afterward, instead of holding `(1 main + 4 text + 2 media) × 2
|
||||
= 14` permanently-open idle connections against PgBouncer. The pool still
|
||||
grows on demand up to `POSTGRES_POOL_MAX`.
|
||||
|
||||
## Event channels (Redis pub/sub)
|
||||
|
||||
|
||||
@@ -7,7 +7,10 @@ import {
|
||||
getAnalysisQueueStatus,
|
||||
startPendingAIAnalysisWorker,
|
||||
} from "../modules/ai-moderation/aiAnalyzer.js";
|
||||
import { workerPool } from "../modules/ai-moderation/circuitBreaker.js";
|
||||
import {
|
||||
mediaWorkerPool,
|
||||
textWorkerPool,
|
||||
} from "../modules/ai-moderation/circuitBreaker.js";
|
||||
import { registerChannelTopicCapture } from "../modules/channel-topic/index.js";
|
||||
import { CommandHandler } from "../modules/command-handler/commandHandler.js";
|
||||
import {
|
||||
@@ -396,12 +399,25 @@ export async function initializeDiscordGateway() {
|
||||
if (typeof status.lastError === "string") {
|
||||
setGauge("ai_analysis_last_error_present", status.lastError ? 1 : 0);
|
||||
}
|
||||
const pool = workerPool as unknown as {
|
||||
_poolState?: { size: number; active: number };
|
||||
};
|
||||
if (pool._poolState) {
|
||||
setGauge("ai_analysis_worker_threads", pool._poolState.size);
|
||||
setGauge("ai_analysis_worker_threads_active", pool._poolState.active);
|
||||
type PoolState = { _poolState?: { size: number; active: number } };
|
||||
const textPool = textWorkerPool as unknown as PoolState;
|
||||
const mediaPool = mediaWorkerPool as unknown as PoolState;
|
||||
// Reported per queue (2026-08-31 text/media pool split) so the text
|
||||
// and media backlogs are distinguishable in dashboards/alerts instead
|
||||
// of one combined "worker threads" number.
|
||||
if (textPool._poolState) {
|
||||
setGauge("ai_analysis_worker_threads_text", textPool._poolState.size);
|
||||
setGauge(
|
||||
"ai_analysis_worker_threads_active_text",
|
||||
textPool._poolState.active,
|
||||
);
|
||||
}
|
||||
if (mediaPool._poolState) {
|
||||
setGauge("ai_analysis_worker_threads_media", mediaPool._poolState.size);
|
||||
setGauge(
|
||||
"ai_analysis_worker_threads_active_media",
|
||||
mediaPool._poolState.active,
|
||||
);
|
||||
}
|
||||
} catch (err) {
|
||||
logger.warn({ error: String(err) }, "AI metrics collector failed");
|
||||
|
||||
@@ -5,7 +5,7 @@ import { messageStore } from "../message-capture/messageStore.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
import { pickBatchWithinBudget as pickBatchWithinBudgetPure } from "./batchBudget.js";
|
||||
import { partitionBatchOutcome } from "./batchOutcomeClassifier.js";
|
||||
import { workerPool } from "./circuitBreaker.js";
|
||||
import { mediaWorkerPool, textWorkerPool } from "./circuitBreaker.js";
|
||||
import { estimateTokens } from "./conversationContext.js";
|
||||
import {
|
||||
conversationErrorCooldown,
|
||||
@@ -14,6 +14,7 @@ import {
|
||||
resetConversationBatchFailures,
|
||||
} from "./conversationState.js";
|
||||
import { enqueueIndividualFallbacks } from "./individualFallbackProcessor.js";
|
||||
import { hasMediaContent } from "./mediaAnalysisClient.js";
|
||||
import {
|
||||
broadcastAnalysisCompleted,
|
||||
LAST_ERROR,
|
||||
@@ -121,29 +122,31 @@ export async function skipAgeRestrictedMessages(
|
||||
// Batch pipeline
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export async function processBatch(
|
||||
/**
|
||||
* Runs one worker job (either the text-only or the 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.
|
||||
*/
|
||||
async function runQueueBatch(
|
||||
pool: typeof textWorkerPool,
|
||||
conversationKey: string,
|
||||
messages: MessageRecord[],
|
||||
processingStartedAt: number,
|
||||
): Promise<void> {
|
||||
if (messages.length === 0) {
|
||||
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
return;
|
||||
}
|
||||
const cooldownUntil = conversationErrorCooldown.get(conversationKey) ?? 0;
|
||||
if (Date.now() < cooldownUntil) {
|
||||
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
): Promise<boolean> {
|
||||
activeRequests++;
|
||||
let shouldScheduleNext = false;
|
||||
try {
|
||||
const result = (await workerPool.run({
|
||||
const result = (await pool.run({
|
||||
type: "batch",
|
||||
conversationKey,
|
||||
messages,
|
||||
@@ -191,7 +194,7 @@ export async function processBatch(
|
||||
},
|
||||
"Batch analysis failed, will retry after cooldown",
|
||||
);
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
|
||||
// Batch succeeded -- partition per-message outcome (2026-08-25).
|
||||
@@ -281,20 +284,13 @@ export async function processBatch(
|
||||
conversationErrorCooldown.set(conversationKey, newCooldown);
|
||||
}
|
||||
|
||||
// Release the processing lock immediately so the cooldown timer controls retry
|
||||
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
|
||||
// Do NOT schedule next -- let the cooldown gate it
|
||||
shouldScheduleNext = false;
|
||||
return false;
|
||||
}
|
||||
|
||||
if (apiFailedMessages.length === 0) {
|
||||
resetConversationBatchFailures(conversationKey);
|
||||
conversationErrorCooldown.delete(conversationKey);
|
||||
}
|
||||
shouldScheduleNext = true;
|
||||
resetConversationBatchFailures(conversationKey);
|
||||
conversationErrorCooldown.delete(conversationKey);
|
||||
return true;
|
||||
} catch (error) {
|
||||
recordConversationBatchFailure(conversationKey);
|
||||
|
||||
@@ -326,18 +322,63 @@ export async function processBatch(
|
||||
},
|
||||
"Analysis worker failed, will retry after cooldown",
|
||||
);
|
||||
return false;
|
||||
} finally {
|
||||
activeRequests--;
|
||||
}
|
||||
}
|
||||
|
||||
export async function processBatch(
|
||||
conversationKey: string,
|
||||
messages: MessageRecord[],
|
||||
processingStartedAt: number,
|
||||
): Promise<void> {
|
||||
if (messages.length === 0) {
|
||||
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
if (shouldScheduleNext) {
|
||||
setImmediate(() => {
|
||||
// Dynamic import to avoid circular dependency at module scope
|
||||
import("./batchScheduler.js").then((m) =>
|
||||
m.scheduleConversationAnalysis(conversationKey),
|
||||
);
|
||||
});
|
||||
return;
|
||||
}
|
||||
const cooldownUntil = conversationErrorCooldown.get(conversationKey) ?? 0;
|
||||
if (Date.now() < cooldownUntil) {
|
||||
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
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,
|
||||
);
|
||||
|
||||
if (conversationProcessing.get(conversationKey) === processingStartedAt) {
|
||||
conversationProcessing.delete(conversationKey);
|
||||
}
|
||||
if (shouldScheduleNext) {
|
||||
setImmediate(() => {
|
||||
// Dynamic import to avoid circular dependency at module scope
|
||||
import("./batchScheduler.js").then((m) =>
|
||||
m.scheduleConversationAnalysis(conversationKey),
|
||||
);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,7 +6,16 @@ import { config } from "../../shared/config/config.js";
|
||||
import type { MessageRecord } from "../message-capture/types.js";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Piscina worker pool (shared by batch + individual pipelines)
|
||||
// Piscina worker pools (2026-08-31: split text vs media)
|
||||
//
|
||||
// Both the batch and individual-fallback pipelines used to share ONE pool.
|
||||
// A conversation batch containing images/stickers/embeds routes through
|
||||
// runMediaBatch (download + vision + LLM — tens of seconds per batch), and
|
||||
// Piscina hands each worker thread exactly one task at a time. With a small
|
||||
// fixed thread count, a handful of slow media batches could occupy every
|
||||
// thread and leave fast text-only batches for OTHER conversations queued
|
||||
// behind them for the whole media duration. Text and media now get their
|
||||
// own dedicated pools so a media backlog can never starve text analysis.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
function getAnalysisWorkerUrl(): URL {
|
||||
@@ -25,12 +34,32 @@ function getAnalysisWorkerUrl(): URL {
|
||||
return candidates[2];
|
||||
}
|
||||
|
||||
export const workerPool = new Piscina({
|
||||
filename: fileURLToPath(getAnalysisWorkerUrl()),
|
||||
const analysisWorkerFilename = fileURLToPath(getAnalysisWorkerUrl());
|
||||
|
||||
/** Dedicated pool for text-only batch/individual analysis jobs. */
|
||||
export const textWorkerPool = new Piscina({
|
||||
filename: analysisWorkerFilename,
|
||||
execArgv: process.execArgv,
|
||||
maxThreads: config.PISCINA_MAX_THREADS ?? availableParallelism(),
|
||||
});
|
||||
|
||||
/**
|
||||
* Dedicated pool for jobs whose batch contains at least one message with
|
||||
* media (attachments/stickers/embeds). Kept separate so slow image/vision
|
||||
* analysis never blocks the text pool above.
|
||||
*/
|
||||
export const mediaWorkerPool = new Piscina({
|
||||
filename: analysisWorkerFilename,
|
||||
execArgv: process.execArgv,
|
||||
maxThreads: config.PISCINA_MEDIA_MAX_THREADS ?? 2,
|
||||
});
|
||||
|
||||
/**
|
||||
* @deprecated Use `textWorkerPool` or `mediaWorkerPool` directly. Kept as an
|
||||
* alias to the text pool only for anything not yet migrated.
|
||||
*/
|
||||
export const workerPool = textWorkerPool;
|
||||
|
||||
/**
|
||||
* Gets the conversation key for a message (thread_id or channel_id).
|
||||
*/
|
||||
|
||||
@@ -6,9 +6,14 @@ import type {
|
||||
AnalysisResult,
|
||||
MessageRecord,
|
||||
} from "../message-capture/types.js";
|
||||
import { getConversationKey, workerPool } from "./circuitBreaker.js";
|
||||
import {
|
||||
getConversationKey,
|
||||
mediaWorkerPool,
|
||||
textWorkerPool,
|
||||
} from "./circuitBreaker.js";
|
||||
import { fireAlert } from "./conversationState.js";
|
||||
import { classifyIndividualWorkerResult } from "./fallbackResultClassifier.js";
|
||||
import { hasMediaContent } from "./mediaAnalysisClient.js";
|
||||
import {
|
||||
broadcastAnalysisCompleted,
|
||||
LAST_ERROR,
|
||||
@@ -73,8 +78,12 @@ async function processIndividualFallback(
|
||||
let exhaustedOnIncomplete = false;
|
||||
|
||||
try {
|
||||
// Run the LLM-heavy work in the worker thread
|
||||
const workerResult = (await workerPool.run({
|
||||
// Run the LLM-heavy work in the worker thread — same text/media pool
|
||||
// split as the batch pipeline (see circuitBreaker.ts), so a single
|
||||
// media message falling back to individual analysis can't queue behind
|
||||
// (or block) text-only fallbacks, and vice versa.
|
||||
const pool = hasMediaContent(message) ? mediaWorkerPool : textWorkerPool;
|
||||
const workerResult = (await pool.run({
|
||||
type: "individual",
|
||||
message,
|
||||
skipNormalAnalysis: false,
|
||||
|
||||
@@ -319,11 +319,20 @@ export const configSchema = z
|
||||
.int()
|
||||
.positive()
|
||||
.default(50),
|
||||
// Worker pool size. Default 4 (not availableParallelism) because each
|
||||
// Piscina thread owns its own pLimit(5) semaphore — on big VPSes
|
||||
// availableParallelism × 5 concurrent LLM calls would overwhelm the
|
||||
// router. Keep threads modest; concurrency is capped per-thread anyway.
|
||||
// Text-analysis worker pool size (2026-08-31: split from the media pool
|
||||
// below so a slow image/vision batch can never occupy every thread and
|
||||
// starve the far more common text-only batches). Default 4 (not
|
||||
// availableParallelism) because each Piscina thread owns its own
|
||||
// pLimit(5) semaphore — on big VPSes availableParallelism × 5 concurrent
|
||||
// LLM calls would overwhelm the router. Keep threads modest; concurrency
|
||||
// is capped per-thread anyway.
|
||||
PISCINA_MAX_THREADS: z.coerce.number().int().positive().default(4),
|
||||
// Media-analysis worker pool size — dedicated threads for batches that
|
||||
// contain images/stickers/embeds (download + vision + LLM, much slower
|
||||
// than text). Kept small since media batches are less frequent and each
|
||||
// one is long-running; sized independently from PISCINA_MAX_THREADS so
|
||||
// tuning one never starves the other.
|
||||
PISCINA_MEDIA_MAX_THREADS: z.coerce.number().int().positive().default(2),
|
||||
|
||||
// ── Voice Transcription ────────────────────────────────────────────────
|
||||
AI_VOICE_TRANSCRIPTION_ENABLED: z
|
||||
|
||||
Reference in New Issue
Block a user