fix(moderation): audit fixes — media/vision timeout 120s, Qdrant retry w/ backoff, JSON repair in LLM caller
- AI_LLM_MEDIA_ANALYSIS_TIMEOUT_MS 60s→120s + vision 60s→120s: vision model via router regularly exceeded 60s, dropping media batches into the individual-fallback chain then exhausting into ai_status=error. - Qdrant upserts: retryWithRetry() wraps PUT /points with exponential backoff (3 attempts, jitter) for transient 408/abort/ECONNRESET — the 41 six-hour 'Qdrant upsert failed — semantic entry skipped' warnings were single-hop timeouts on a healthy-but-loaded Qdrant. - LLM caller: on parse failure, attempt extractJson() structural repair of the raw content (models with thinking disabled sometimes emit JSON as plain text) before giving up and re-requesting.
This commit is contained in:
@@ -17,7 +17,10 @@ import { config } from "../../shared/config/config.js";
|
||||
import { incrementCounterBy } from "../gateway-metrics/index.js";
|
||||
import type { AnalysisResult } from "../message-capture/types.js";
|
||||
import { llmChat } from "./llmClient.js";
|
||||
import { parseModerationResponse } from "./moderationResponseParser.js";
|
||||
import {
|
||||
extractJson,
|
||||
parseModerationResponse,
|
||||
} from "./moderationResponseParser.js";
|
||||
import { logModerationError } from "./responseLogger.js";
|
||||
|
||||
const log = createChildLogger("llm-caller");
|
||||
@@ -129,6 +132,35 @@ export async function callModerationLLM(
|
||||
},
|
||||
`Failed to parse moderation response (${label})`,
|
||||
);
|
||||
// Attempt cheap structural repair before re-throwing: models with
|
||||
// thinking disabled sometimes emit the JSON object as plain text
|
||||
// (no code fence, outer prose). extractJson walks the raw content
|
||||
// for a balanced {…} / […]. If that yields a parseable object,
|
||||
// use it — the parse is data-driven, never injected into the
|
||||
// prompt, so it cannot leak into future calls.
|
||||
try {
|
||||
const repaired = extractJson(rawContent);
|
||||
if (repaired) {
|
||||
return {
|
||||
parsed: parseModerationResponse(
|
||||
JSON.stringify(repaired),
|
||||
targetIds,
|
||||
),
|
||||
result: completion,
|
||||
};
|
||||
}
|
||||
} catch (repairError: any) {
|
||||
log.warn(
|
||||
{
|
||||
error:
|
||||
repairError instanceof Error
|
||||
? repairError.message
|
||||
: String(repairError),
|
||||
label,
|
||||
},
|
||||
"JSON repair attempt failed — falling through to retry",
|
||||
);
|
||||
}
|
||||
throw parseError;
|
||||
}
|
||||
} catch (apiError: any) {
|
||||
|
||||
@@ -78,6 +78,10 @@ export async function runMediaBatch(
|
||||
const userContent = userBlocks.join("\n\n");
|
||||
|
||||
const perMsgTimeout = config.AI_LLM_MEDIA_ANALYSIS_TIMEOUT_MS ?? 60000;
|
||||
// Batch budget = per-message budget × message count, capped at 5 minutes
|
||||
// absolute so a single slow vision call cannot stall the whole pipeline
|
||||
// for a huge batch (the abandonment budget is a hard safety net; messages
|
||||
// that time out are routed to the individual fallback queue anyway).
|
||||
const batchTimeout = Math.min(
|
||||
Math.max(perMsgTimeout, perMsgTimeout * targets.length),
|
||||
300_000,
|
||||
|
||||
@@ -107,6 +107,54 @@ async function request(
|
||||
}
|
||||
}
|
||||
|
||||
/** True for transient errors worth retrying (408 timeout, ECONNRESET, aborts). */
|
||||
function isTransientQdrantError(error: unknown): boolean {
|
||||
if (error instanceof Error && error.name === "AbortError") return true;
|
||||
const msg = error instanceof Error ? error.message : String(error);
|
||||
return (
|
||||
msg.includes("-> 408") ||
|
||||
msg.includes("aborted") ||
|
||||
msg.includes("ECONNRESET") ||
|
||||
msg.includes("ETIMEDOUT") ||
|
||||
msg.includes("fetch failed")
|
||||
);
|
||||
}
|
||||
|
||||
/** Retry with exponential backoff around `request`, abort-aware. */
|
||||
async function requestWithRetry(
|
||||
method: string,
|
||||
path: string,
|
||||
body?: unknown,
|
||||
timeoutMs = 10_000,
|
||||
attempts = 3,
|
||||
): Promise<unknown> {
|
||||
let lastError: unknown;
|
||||
for (let attempt = 0; attempt < attempts; attempt++) {
|
||||
if (attempt > 0) {
|
||||
// Exponential backoff: 500ms → 1s → 2s (jittered ±20%).
|
||||
const base = 500 * 2 ** (attempt - 1);
|
||||
const delayMs = base + Math.floor(Math.random() * base * 0.2);
|
||||
await new Promise((resolve) => setTimeout(resolve, delayMs));
|
||||
}
|
||||
try {
|
||||
return await request(method, path, body, timeoutMs);
|
||||
} catch (error) {
|
||||
lastError = error;
|
||||
if (!isTransientQdrantError(error)) throw error;
|
||||
log.debug(
|
||||
{
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
attempt: attempt + 1,
|
||||
method,
|
||||
path,
|
||||
},
|
||||
"Qdrant transient error — retrying with backoff",
|
||||
);
|
||||
}
|
||||
}
|
||||
throw lastError;
|
||||
}
|
||||
|
||||
/** Deterministic uint64 point id from the exact-hash cache key. */
|
||||
export function qdrantPointId(cacheKey: string): number {
|
||||
const digest = createHash("sha256").update(cacheKey).digest();
|
||||
@@ -185,7 +233,7 @@ export async function upsertQdrantPoint(
|
||||
): Promise<boolean> {
|
||||
try {
|
||||
if (!(await ensureQdrantCollection(vector.length))) return false;
|
||||
await request(
|
||||
await requestWithRetry(
|
||||
"PUT",
|
||||
`/collections/${collectionName()}/points`,
|
||||
{
|
||||
@@ -193,6 +241,7 @@ export async function upsertQdrantPoint(
|
||||
wait: true,
|
||||
},
|
||||
30_000,
|
||||
3,
|
||||
);
|
||||
return true;
|
||||
} catch (error) {
|
||||
@@ -470,7 +519,7 @@ export async function upsertQdrantPointV2(
|
||||
): Promise<boolean> {
|
||||
try {
|
||||
if (!(await ensureQdrantCollectionV2(name, vector.length))) return false;
|
||||
await request(
|
||||
await requestWithRetry(
|
||||
"PUT",
|
||||
`/collections/${name}/points`,
|
||||
{
|
||||
@@ -478,6 +527,7 @@ export async function upsertQdrantPointV2(
|
||||
wait: true,
|
||||
},
|
||||
30_000,
|
||||
3,
|
||||
);
|
||||
return true;
|
||||
} catch (error) {
|
||||
|
||||
@@ -105,6 +105,14 @@ export const configSchema = z
|
||||
// to POSTGRES_POOL_MAX; min:0 only drops idle clients after
|
||||
// idleTimeoutMillis. This both trims RSS and frees PgBouncer slots.
|
||||
POSTGRES_POOL_MIN: z.coerce.number().int().min(0).default(0),
|
||||
// Ceiling for the gateway's per-process pg Pool. Each Piscina worker
|
||||
// thread owns its own pool (main + 4 workers = 5 pools), so this value
|
||||
// is the per-thread cap. Kept at 10 (2026-09-09 audit): the real
|
||||
// bottleneck is PgBouncer's per-(user,db) default_pool_size on imrnes —
|
||||
// raising this ceiling without raising the Bouncer pool just makes more
|
||||
// clients queue at the same 20 slots. pool_mode=session means each pg
|
||||
// Pool client occupies a Bouncer slot for the whole transaction; min:0
|
||||
// + idleTimeoutMillis frees idle slots automatically.
|
||||
POSTGRES_POOL_MAX: z.coerce.number().int().positive().default(10),
|
||||
|
||||
// ── Redis ────────────────────────────────────────────────────────────
|
||||
@@ -211,16 +219,18 @@ export const configSchema = z
|
||||
.number()
|
||||
.int()
|
||||
.positive()
|
||||
.default(60000),
|
||||
.default(120_000),
|
||||
// Standalone image/sticker/emoji vision analysis (analyzeSingleMediaImage
|
||||
// → llmVision → llmChat). Decoupled from the media *batch* timeout above so
|
||||
// a single vision call can be tuned independently. 1 minute by default —
|
||||
// vision models (especially behind a router) need headroom for large images.
|
||||
// a single vision call can be tuned independently. 2 minutes by default —
|
||||
// vision models (especially behind a router) need headroom for large images
|
||||
// and the media batch budget grew to 120s (2026-09-09) so single-image calls
|
||||
// must not be the bottleneck in the fallback chain.
|
||||
AI_LLM_VISION_ANALYSIS_TIMEOUT_MS: z.coerce
|
||||
.number()
|
||||
.int()
|
||||
.positive()
|
||||
.default(60000),
|
||||
.default(120_000),
|
||||
// Text-only moderation batches are cheaper than media (no downloads /
|
||||
// vision pre-pass), so they get their own (shorter) timeout instead of
|
||||
// being tied to the media budget.
|
||||
|
||||
Reference in New Issue
Block a user