fix(gateway): un-claim budget-overflow messages stuck in processing

getPendingMessagesByConversation() flips every fetched pending row to
'processing', then pickBatchWithinBudget() may stop early on the token
budget. The tail rows that did NOT make the batch were never un-claimed,
so they stayed 'processing' forever — the recovery worker reverted them
(120s) only for the next wave to re-claim them, an infinite loop of
stuck messages that never get analyzed (saw 22 rows, some recycled for
40+ minutes).

- add computeBudgetOverflowMessages() pure helper (batchBudget.ts)
- batchScheduler un-claims overflow rows back to 'pending' before
  dispatching the trimmed batch
- recovery-worker now reverts stuck processing unconditionally (the old
  'conversationProcessing.size > 0' guard skipped the revert when the
  in-memory lock map was empty, e.g. fresh boot — exactly when stranded
  rows from a previous process need rescuing)
- 3 regression tests for the overflow helper
This commit is contained in:
asepharyana
2026-09-24 16:49:27 +07:00
parent ef7708bf7d
commit 80daa9f045
4 changed files with 148 additions and 14 deletions
@@ -43,3 +43,22 @@ export function pickBatchWithinBudget(
return batch;
}
/**
* Returns the messages that were fetched/claimed but did NOT make it into the
* trimmed batch (i.e. the tail past the token budget).
*
* The DB claim step flips every fetched pending row to `processing`; the batch
* trim may then stop early on the token budget. Those tail rows would stay
* stuck in `processing` forever unless the caller explicitly un-claims them —
* this helper identifies exactly which rows that is, so the caller can write
* them back to `pending` for the next wave.
*/
export function computeBudgetOverflowMessages(
claimed: MessageRecord[],
trimmed: MessageRecord[],
): MessageRecord[] {
if (trimmed.length === 0) return claimed;
const trimmedIds = new Set(trimmed.map((m) => m.id));
return claimed.filter((m) => !trimmedIds.has(m.id));
}
@@ -3,6 +3,7 @@ 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 { computeBudgetOverflowMessages } from "./batchBudget.js";
import {
pickBatchWithinBudget,
processBatch,
@@ -147,6 +148,60 @@ function scheduleLaneTimer(
);
}
// Un-claim messages that did NOT make it into the trimmed batch.
// getPendingMessagesByConversation() flips EVERY fetched pending row
// to `processing`; pickBatchWithinBudget() may then stop early on the
// token budget, leaving the tail rows stuck in `processing` forever
// (recovery only reverts rows older than 120s, and these keep getting
// re-claimed each wave). Return them to `pending` so the next wave
// picks them up instead of leaking processing slots.
const unclaimed = computeBudgetOverflowMessages(
processableMessages,
trimmed,
);
if (unclaimed.length > 0) {
const unclaimedRows = await messageStore
.updateMessagesAIAnalysisBulk(
unclaimed.map((msg) => ({
messageId: msg.id,
result: {
status: "pending",
flags: null,
score: null,
analysis: null,
categories: null,
severity: null,
confidence: null,
recommendedAction: null,
analyzedAt: null,
error: null,
},
})),
)
.catch((err: unknown) => {
logger.error(
{
conversationKey,
lane,
count: unclaimed.length,
error: err instanceof Error ? err.message : String(err),
},
"Failed to un-claim budget-overflow messages back to pending",
);
return null;
});
if (unclaimedRows) {
logger.debug(
{
conversationKey,
lane,
unclaimedCount: unclaimed.length,
},
"Returned budget-overflow messages to pending for next wave",
);
}
}
// 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.
@@ -48,18 +48,21 @@ export function startRecoveryWorker(): void {
setInterval(() => {
runCachePruneIfDue();
// Only revert stuck processing messages if there's active processing.
// Avoids a DB query every recovery interval when the pipeline is idle.
if (conversationProcessing.size > 0) {
messageStore
.revertStuckProcessingMessages(STUCK_PROCESSING_AGE_MS)
.catch((err: unknown) => {
logger.error(
{ error: String(err) },
"Failed to run stuck processing recovery",
);
});
}
// Revert stuck `processing` messages back to `pending` unconditionally.
// The old `if (conversationProcessing.size > 0)` guard skipped the revert
// when the in-memory lock map was empty (fresh boot, or every lock was
// already pruned) — precisely the moment stranded `processing` rows from
// a previous process still need rescuing. `revertStuckProcessingMessages`
// is a cheap UPDATE..RETURNING keyed on ai_status + age, safe to run
// every interval; it matches 0 rows when there is nothing to do.
messageStore
.revertStuckProcessingMessages(STUCK_PROCESSING_AGE_MS)
.catch((err: unknown) => {
logger.error(
{ error: String(err) },
"Failed to run stuck processing recovery",
);
});
Promise.all([
messageStore.getPendingConversationKeys(500),
@@ -1,5 +1,6 @@
import { describe, expect, it } from "vitest";
import {
computeBudgetOverflowMessages,
pickBatchWithinBudget,
type TokenEstimator,
} from "../src/modules/ai-moderation/batchBudget.js";
@@ -10,6 +11,8 @@ import type { MessageRecord } from "../src/modules/message-capture/types.js";
// the call site — see batchProcessor.ts).
const estimate: TokenEstimator = (text: string) => text.length;
const TOKENS_PER_MESSAGE = 50;
function msg(id: string, content: string, createdAt: number): MessageRecord {
return {
id,
@@ -36,8 +39,6 @@ function msg(id: string, content: string, createdAt: number): MessageRecord {
}
describe("pickBatchWithinBudget", () => {
const TOKENS_PER_MESSAGE = 50;
it("returns a contiguous chronological prefix — no gaps mid-timeline", () => {
// sizes: 100, 100, 400 (overflow), 10
const messages = [
@@ -89,3 +90,59 @@ describe("pickBatchWithinBudget", () => {
).toEqual([]);
});
});
describe("computeBudgetOverflowMessages", () => {
it("returns the tail that did NOT fit the token budget", () => {
const messages = [
msg("m1", "a".repeat(100), 1),
msg("m2", "b".repeat(100), 2),
msg("m3", "c".repeat(900), 3),
msg("m4", "d".repeat(10), 4),
];
// 500 budget: m1(150)+m2(150)=300 fits, m3(950) overflows → batch = [m1,m2]
const batch = pickBatchWithinBudget(
messages,
500,
TOKENS_PER_MESSAGE,
estimate,
);
const overflow = computeBudgetOverflowMessages(messages, batch);
// m3 + m4 were claimed `processing` by the DB fetch but never processed —
// they must be identified so the scheduler can return them to `pending`.
expect(overflow.map((m) => m.id)).toEqual(["m3", "m4"]);
});
it("returns everything when the batch is empty (all overflow)", () => {
const messages = [
msg("m1", "x".repeat(1000), 1),
msg("m2", "y".repeat(1000), 2),
];
const batch = pickBatchWithinBudget(
messages,
500,
TOKENS_PER_MESSAGE,
estimate,
);
expect(batch).toHaveLength(0);
expect(
computeBudgetOverflowMessages(messages, batch).map((m) => m.id),
).toEqual(["m1", "m2"]);
});
it("returns nothing when every claimed message was processed", () => {
const messages = [
msg("m1", "a".repeat(100), 1),
msg("m2", "b".repeat(100), 2),
];
const batch = pickBatchWithinBudget(
messages,
500,
TOKENS_PER_MESSAGE,
estimate,
);
expect(batch).toHaveLength(2);
expect(
computeBudgetOverflowMessages(messages, batch).map((m) => m.id),
).toEqual([]);
});
});