diff --git a/services/discord-gateway/src/modules/ai-moderation/batchBudget.ts b/services/discord-gateway/src/modules/ai-moderation/batchBudget.ts index 6f339e87..b8ea9200 100644 --- a/services/discord-gateway/src/modules/ai-moderation/batchBudget.ts +++ b/services/discord-gateway/src/modules/ai-moderation/batchBudget.ts @@ -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)); +} diff --git a/services/discord-gateway/src/modules/ai-moderation/batchScheduler.ts b/services/discord-gateway/src/modules/ai-moderation/batchScheduler.ts index ad4a803b..8cb7f90d 100644 --- a/services/discord-gateway/src/modules/ai-moderation/batchScheduler.ts +++ b/services/discord-gateway/src/modules/ai-moderation/batchScheduler.ts @@ -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. diff --git a/services/discord-gateway/src/modules/ai-moderation/recovery-worker.ts b/services/discord-gateway/src/modules/ai-moderation/recovery-worker.ts index f0d77130..52156bbb 100644 --- a/services/discord-gateway/src/modules/ai-moderation/recovery-worker.ts +++ b/services/discord-gateway/src/modules/ai-moderation/recovery-worker.ts @@ -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), diff --git a/services/discord-gateway/tests/batchBudget.test.ts b/services/discord-gateway/tests/batchBudget.test.ts index 8b4b5f6e..88c7c0c7 100644 --- a/services/discord-gateway/tests/batchBudget.test.ts +++ b/services/discord-gateway/tests/batchBudget.test.ts @@ -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([]); + }); +});