fix(moderation): patch 3 additional aiAnalyzer vulnerabilities
Infinite recovery loop (#new):
- Inside retryWithBackoff callback in processIndividualFallback, detect
'analysis_incomplete' in the LLM result and throw to trigger backoff.
- Track exhaustedOnIncomplete flag across retries.
- On final exhaustion: write terminal flag 'individual_analysis_exhausted'
to DB so the recovery query (which only looks for 'analysis_incomplete')
never picks this message up again.
- Transient failures (network/parse) are NOT written as exhausted; they
remain as 'analysis_incomplete' and are retried via the CB-throttled
recovery cycle.
Token budget zero-result deadlock (#10):
- If pickBatchWithinBudget returns [] because every candidate message
individually exceeds AI_ANALYSIS_MAX_TARGET_TOKENS, fall back to
messages.slice(0,1) so at least the first message is processed.
- Without this, messages would be permanently stuck as 'pending' because
every recovery tick would fetch them, trim to 0, and exit silently.
- Uses messages.slice(0,1) instead of messages[0]! to avoid the
forbidden noNonNullAssertion lint rule.
Stale state map memory leak (#9):
- startPendingAIAnalysisWorker now prunes conversationErrorCooldown and
conversationProcessing on every recovery interval tick.
- Cooldown entries past their expiry timestamp are deleted.
- Processing entries older than AI_ANALYSIS_PROCESSING_TIMEOUT_MS are
deleted (these represent stale locks from crashed processing runs).
- Prevents unbounded Map growth for long-running bots with many channels.
Batch/individual scheduling collision (#8):
- Build incompleteKeySet (Set<string>) from incompleteKeys before the
batch recovery loop.
- Batch recovery loop skips any key present in incompleteKeySet so a
conversation that has both 'pending' and 'analysis_incomplete' messages
is only targeted by the individual pipeline, not both simultaneously.
- Avoids the DB last-write-wins race where batch and individual pipelines
both update the same message rows concurrently.
This commit is contained in:
+118
-15
@@ -178,6 +178,15 @@ function isConversationProcessingLocked(conversationKey: string): boolean {
|
|||||||
*
|
*
|
||||||
* FIX #1+#5: Increments the individual circuit breaker on failure so a
|
* FIX #1+#5: Increments the individual circuit breaker on failure so a
|
||||||
* sustained outage stops hammering the LLM endpoint.
|
* sustained outage stops hammering the LLM endpoint.
|
||||||
|
*
|
||||||
|
* Infinite-loop prevention: if the LLM consistently drops the single target
|
||||||
|
* message across all retries (analysis_incomplete), we write a terminal flag
|
||||||
|
* 'individual_analysis_exhausted' to DB instead of 'analysis_incomplete'.
|
||||||
|
* The recovery worker only queries for 'analysis_incomplete', so exhausted
|
||||||
|
* messages are permanently excluded from the reprocessing loop.
|
||||||
|
* Transient failures (network/parse/DB) are NOT written as exhausted — they
|
||||||
|
* stay as 'analysis_incomplete' so the circuit-breaker-throttled recovery
|
||||||
|
* cycle can retry them later.
|
||||||
*/
|
*/
|
||||||
async function processIndividualFallback(
|
async function processIndividualFallback(
|
||||||
message: MessageRecord,
|
message: MessageRecord,
|
||||||
@@ -192,6 +201,10 @@ async function processIndividualFallback(
|
|||||||
(individualInFlightByConversation.get(conversationKey) ?? 0) + 1,
|
(individualInFlightByConversation.get(conversationKey) ?? 0) + 1,
|
||||||
);
|
);
|
||||||
|
|
||||||
|
// Track whether all retries were exhausted specifically because the LLM
|
||||||
|
// consistently returned no result for this message (vs. a transient error).
|
||||||
|
let exhaustedOnIncomplete = false;
|
||||||
|
|
||||||
try {
|
try {
|
||||||
const contextBefore = await getConversationContextBefore({
|
const contextBefore = await getConversationContextBefore({
|
||||||
channelId: message.channel_id,
|
channelId: message.channel_id,
|
||||||
@@ -213,12 +226,30 @@ async function processIndividualFallback(
|
|||||||
]);
|
]);
|
||||||
|
|
||||||
const analysisResult = await retryWithBackoff(
|
const analysisResult = await retryWithBackoff(
|
||||||
() =>
|
async () => {
|
||||||
runModerationAnalysis({
|
const result = await runModerationAnalysis({
|
||||||
targets: [message],
|
targets: [message],
|
||||||
contextText: contextLines.join("\n"),
|
contextText: contextLines.join("\n"),
|
||||||
attachments,
|
attachments,
|
||||||
}),
|
});
|
||||||
|
|
||||||
|
// If the LLM still dropped our only target, convert to a retryable
|
||||||
|
// throw so backoff kicks in. Track this so the catch block can
|
||||||
|
// distinguish it from a transient network/parse failure.
|
||||||
|
const stillIncomplete = result.results.some((r) =>
|
||||||
|
r.flags.includes("analysis_incomplete"),
|
||||||
|
);
|
||||||
|
if (stillIncomplete) {
|
||||||
|
exhaustedOnIncomplete = true;
|
||||||
|
throw new Error(
|
||||||
|
`LLM returned no result for single-target message ${messageId} — will retry with backoff`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
// Got a real result — clear the incomplete flag.
|
||||||
|
exhaustedOnIncomplete = false;
|
||||||
|
return result;
|
||||||
|
},
|
||||||
{
|
{
|
||||||
retries: 2,
|
retries: 2,
|
||||||
minTimeout: 2000,
|
minTimeout: 2000,
|
||||||
@@ -269,14 +300,50 @@ async function processIndividualFallback(
|
|||||||
}
|
}
|
||||||
|
|
||||||
lastError = error instanceof Error ? error.message : String(error);
|
lastError = error instanceof Error ? error.message : String(error);
|
||||||
logger.error(
|
|
||||||
{
|
// Infinite-loop prevention: if all retries were exhausted because the LLM
|
||||||
messageId,
|
// consistently dropped this specific message (not a transient error),
|
||||||
error: lastError,
|
// overwrite the DB entry with a terminal flag that the recovery query
|
||||||
stack: error instanceof Error ? error.stack : undefined,
|
// does NOT match. This permanently removes it from the recovery loop
|
||||||
},
|
// while keeping it visible as an error in the dashboard.
|
||||||
"Individual fallback analysis failed",
|
if (exhaustedOnIncomplete) {
|
||||||
);
|
await updateMessagesAIAnalysisBulk([
|
||||||
|
{
|
||||||
|
messageId,
|
||||||
|
result: {
|
||||||
|
status: "error",
|
||||||
|
flags: JSON.stringify(["individual_analysis_exhausted"]),
|
||||||
|
score: 0,
|
||||||
|
raw: null,
|
||||||
|
analysis:
|
||||||
|
"Individual fallback exhausted all retries: LLM consistently dropped this message even in single-target mode",
|
||||||
|
analyzedAt: Date.now(),
|
||||||
|
error: lastError,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
]).catch((dbErr) => {
|
||||||
|
logger.error(
|
||||||
|
{ messageId, error: String(dbErr) },
|
||||||
|
"Failed to write terminal exhausted status — message may re-enter recovery loop",
|
||||||
|
);
|
||||||
|
});
|
||||||
|
logger.warn(
|
||||||
|
{ messageId },
|
||||||
|
"Individual fallback exhausted — marked as individual_analysis_exhausted to stop recovery loop",
|
||||||
|
);
|
||||||
|
} else {
|
||||||
|
// Transient failure (network/parse/DB): do NOT write terminal status.
|
||||||
|
// Message stays as error/analysis_incomplete in DB and will be retried
|
||||||
|
// by the recovery worker, subject to the individual circuit breaker.
|
||||||
|
logger.error(
|
||||||
|
{
|
||||||
|
messageId,
|
||||||
|
error: lastError,
|
||||||
|
stack: error instanceof Error ? error.stack : undefined,
|
||||||
|
},
|
||||||
|
"Individual fallback analysis failed (transient) — will be retried by recovery worker",
|
||||||
|
);
|
||||||
|
}
|
||||||
} finally {
|
} finally {
|
||||||
activeIndividualRequests--;
|
activeIndividualRequests--;
|
||||||
individualInFlight.delete(messageId);
|
individualInFlight.delete(messageId);
|
||||||
@@ -562,12 +629,27 @@ function scheduleConversationAnalysis(conversationKey: string): void {
|
|||||||
|
|
||||||
// FIX #6: trim to token budget before sending to LLM.
|
// FIX #6: trim to token budget before sending to LLM.
|
||||||
// 50 tokens overhead accounts for JSON structure + id/username fields.
|
// 50 tokens overhead accounts for JSON structure + id/username fields.
|
||||||
const trimmed = pickBatchWithinBudget(
|
let trimmed = pickBatchWithinBudget(
|
||||||
messages,
|
messages,
|
||||||
config.AI_ANALYSIS_MAX_TARGET_TOKENS,
|
config.AI_ANALYSIS_MAX_TARGET_TOKENS,
|
||||||
50,
|
50,
|
||||||
);
|
);
|
||||||
if (trimmed.length === 0) return;
|
|
||||||
|
// FIX #10: if every message individually exceeds the token budget,
|
||||||
|
// pickBatchWithinBudget returns [] — which would leave them permanently
|
||||||
|
// stuck as `pending`. Fall back to the first message alone so at
|
||||||
|
// least one makes progress; the rest will be processed in later ticks.
|
||||||
|
if (trimmed.length === 0 && messages.length > 0) {
|
||||||
|
trimmed = messages.slice(0, 1);
|
||||||
|
logger.warn(
|
||||||
|
{
|
||||||
|
conversationKey,
|
||||||
|
messageId: messages[0]?.id,
|
||||||
|
tokenBudget: config.AI_ANALYSIS_MAX_TARGET_TOKENS,
|
||||||
|
},
|
||||||
|
"All messages exceed token budget — processing first message alone to avoid stuck-pending deadlock",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
return processBatch(conversationKey, trimmed);
|
return processBatch(conversationKey, trimmed);
|
||||||
})
|
})
|
||||||
@@ -652,20 +734,41 @@ export function startPendingAIAnalysisWorker(): void {
|
|||||||
getConversationKeysWithIncompleteAnalysis(50),
|
getConversationKeysWithIncompleteAnalysis(50),
|
||||||
])
|
])
|
||||||
.then(([pendingKeys, incompleteKeys]) => {
|
.then(([pendingKeys, incompleteKeys]) => {
|
||||||
|
const now = Date.now();
|
||||||
|
|
||||||
|
// FIX #9: Prune stale entries from state maps to prevent unbounded
|
||||||
|
// memory growth from channels/threads that are no longer active.
|
||||||
|
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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// FIX #8: Build a set of keys already targeted for individual recovery
|
||||||
|
// so the batch loop below skips them, preventing a race where batch
|
||||||
|
// scheduling and individual scheduling collide on the same conversation.
|
||||||
|
const incompleteKeySet = new Set(incompleteKeys);
|
||||||
|
|
||||||
// --- Batch recovery for `pending` messages ---
|
// --- Batch recovery for `pending` messages ---
|
||||||
for (const key of pendingKeys) {
|
for (const key of pendingKeys) {
|
||||||
if (conversationDebounceTimers.has(key)) continue;
|
if (conversationDebounceTimers.has(key)) continue;
|
||||||
if (isConversationProcessingLocked(key)) continue;
|
if (isConversationProcessingLocked(key)) continue;
|
||||||
// FIX #4: skip if individual fallback already running for this conversation.
|
// FIX #4: skip if individual fallback already running for this conversation.
|
||||||
if (individualInFlightByConversation.has(key)) continue;
|
if (individualInFlightByConversation.has(key)) continue;
|
||||||
|
// FIX #8: skip if this conversation also needs individual recovery
|
||||||
|
// (batch processing would conflict with in-flight individual work).
|
||||||
|
if (incompleteKeySet.has(key)) continue;
|
||||||
const cooldownUntil = conversationErrorCooldown.get(key);
|
const cooldownUntil = conversationErrorCooldown.get(key);
|
||||||
if (cooldownUntil && Date.now() < cooldownUntil) continue;
|
if (cooldownUntil && now < cooldownUntil) continue;
|
||||||
scheduleConversationAnalysis(key);
|
scheduleConversationAnalysis(key);
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- Individual recovery for `error/analysis_incomplete` messages ---
|
// --- Individual recovery for `error/analysis_incomplete` messages ---
|
||||||
// Circuit breaker check: no point iterating if individual CB is active.
|
// Circuit breaker check: no point iterating if individual CB is active.
|
||||||
if (Date.now() >= individualCooldownUntil) {
|
if (now >= individualCooldownUntil) {
|
||||||
const promises: Promise<void>[] = [];
|
const promises: Promise<void>[] = [];
|
||||||
for (const key of incompleteKeys) {
|
for (const key of incompleteKeys) {
|
||||||
// Skip if individual work is already running for this conversation.
|
// Skip if individual work is already running for this conversation.
|
||||||
|
|||||||
Reference in New Issue
Block a user