diff --git a/src/lib/state-store-registrations.ts b/src/lib/state-store-registrations.ts index 13a22bfce0..849f145848 100644 --- a/src/lib/state-store-registrations.ts +++ b/src/lib/state-store-registrations.ts @@ -21,7 +21,7 @@ import { } from "../combos/failover"; import { reconcileComboWarningMemos } from "../combos/request"; import { reconcileComboRotationState } from "../combos/resolve"; -import { reconcileComboRecall } from "../server/responses/combo-session-recall"; +import { reconcileComboRecall, sweepExpiredComboRecall } from "../server/responses/combo-session-recall"; import { listLiveComboTargetKeys } from "../combos/types"; import { listLiveConfigOwnershipRoots, @@ -112,7 +112,11 @@ export const STATE_STORE_REGISTRATIONS = [ { name: "model-cache-history", reconcileGeneration: reconcileModelCacheGeneration }, { name: "pool-rotation", reconcileGeneration: reconcilePoolRotationState }, { name: "combo-rotation", reconcileGeneration: reconcileComboRotationState }, - { name: "combo-session-recall", reconcileGeneration: reconcileComboRecall }, + { + name: "combo-session-recall", + sweepExpired: sweepExpiredComboRecall, + reconcileGeneration: reconcileComboRecall, + }, { name: "guardian-backoff", reconcileGeneration: reconcileGuardianBackoff }, { name: "codex-reauth", reconcileGeneration: reconcileCodexReauthState }, { name: "oauth-reauth", reconcileGeneration: reconcileOAuthReauthState }, diff --git a/src/server/responses/combo-session-recall.ts b/src/server/responses/combo-session-recall.ts index 84dfd8d466..aafece9e04 100644 --- a/src/server/responses/combo-session-recall.ts +++ b/src/server/responses/combo-session-recall.ts @@ -7,12 +7,16 @@ interface ComboRecallEntry { comboId: string; target: Pick; responseModel: string; + responseModelBytes: number; at: number; } const RECALL_CAPACITY = 256; const RECALL_TTL_MS = 30 * 60 * 1000; +const RECALL_MODEL_MAX_BYTES = 1024; +const RECALL_MODEL_TOTAL_BYTES = 64 * 1024; const recall = new Map(); +let retainedModelBytes = 0; let lastReconciledGeneration = 0; let liveOwners: Pick | undefined; @@ -22,6 +26,22 @@ function ownsEntry(context: Pick RECALL_MODEL_MAX_BYTES) return undefined; + const bytes = new TextEncoder().encode(model).byteLength; + return bytes <= RECALL_MODEL_MAX_BYTES ? bytes : undefined; +} + export function rememberComboForLane( lane: string | undefined, comboId: string, @@ -30,16 +50,25 @@ export function rememberComboForLane( writerGeneration: number, ): void { if (!lane || !comboId || !responseModel.trim()) return; + const modelBytes = boundedModelBytes(responseModel); + if (modelBytes === undefined) return; // Reject even a same-named recreated owner: its previous in-flight turn is obsolete. if (writerGeneration < Math.max(lastReconciledGeneration, captureConfigGeneration())) return; - const entry = { comboId, target: { provider: target.provider, model: target.model }, responseModel, at: Date.now() }; + const entry = { + comboId, + target: { provider: target.provider, model: target.model }, + responseModel, + responseModelBytes: modelBytes, + at: Date.now(), + }; if (liveOwners && !ownsEntry(liveOwners, entry)) return; - recall.delete(lane); + deleteEntry(lane); recall.set(lane, entry); - while (recall.size > RECALL_CAPACITY) { + retainedModelBytes += modelBytes; + while (recall.size > RECALL_CAPACITY || retainedModelBytes > RECALL_MODEL_TOTAL_BYTES) { const oldest = recall.keys().next().value; if (oldest === undefined) break; - recall.delete(oldest); + deleteEntry(oldest); } } @@ -57,7 +86,7 @@ export function recallComboForLane( || !Object.hasOwn(config.providers, entry.target.provider) || !provider || provider.disabled === true || !combo?.targets.some(target => targetKey(target) === targetKey(entry.target))) { - recall.delete(lane); + deleteEntry(lane); return undefined; } return entry.responseModel === model ? entry.comboId : undefined; @@ -74,16 +103,27 @@ export function reconcileComboRecall(context: GenerationContext): number { let removed = 0; for (const [lane, entry] of recall) { if (!ownsEntry(context, entry) || Date.now() - entry.at >= RECALL_TTL_MS) { - recall.delete(lane); + deleteEntry(lane); removed += 1; } } return removed; } +export function sweepExpiredComboRecall(now: number): number { + let removed = 0; + for (const [lane, entry] of recall) { + if (now - entry.at < RECALL_TTL_MS) continue; + deleteEntry(lane); + removed += 1; + } + return removed; +} + /** Test-only reset, alongside the combo rotation/cooldown resets. */ export function clearComboRecallForTests(): void { recall.clear(); + retainedModelBytes = 0; lastReconciledGeneration = 0; liveOwners = undefined; } diff --git a/tests/oauth/state-store-sweeper.test.ts b/tests/oauth/state-store-sweeper.test.ts index b162e2214b..0420cc8951 100644 --- a/tests/oauth/state-store-sweeper.test.ts +++ b/tests/oauth/state-store-sweeper.test.ts @@ -188,6 +188,39 @@ describe("state-store sweeper", () => { expect(reconcileLiveStateStores()).toEqual({ storesVisited: 1, rowsRemoved: 2 }); }); + test("registered combo recall sweep removes dormant expired lanes", () => { + registerStateStore(STATE_STORE_REGISTRATIONS.find(row => row.name === "combo-session-recall")!); + const config: OcxConfig = { + port: 0, defaultProvider: "a", + providers: { a: { adapter: "openai-chat", baseUrl: "https://a.example/v1" } }, + combos: { first: { targets: [{ provider: "a", model: "m1" }] } }, + }; + rememberComboForLane("lane", "first", { provider: "a", model: "m1" }, "visible", captureConfigGeneration()); + + expect(sweepExpired(Date.now() + 30 * 60 * 1_000)).toEqual({ storesVisited: 1, rowsRemoved: 1 }); + expect(recallComboForLane(config, "lane", "visible")).toBeUndefined(); + }); + + test("combo recall bounds individual and aggregate upstream model bytes", () => { + const config: OcxConfig = { + port: 0, defaultProvider: "a", + providers: { a: { adapter: "openai-chat", baseUrl: "https://a.example/v1" } }, + combos: { first: { targets: [{ provider: "a", model: "m1" }] } }, + }; + const target = { provider: "a", model: "m1" }; + rememberComboForLane("oversized-ascii", "first", target, "x".repeat(1025), captureConfigGeneration()); + rememberComboForLane("oversized-utf8", "first", target, "é".repeat(513), captureConfigGeneration()); + expect(recallComboForLane(config, "oversized-ascii", "x".repeat(1025))).toBeUndefined(); + expect(recallComboForLane(config, "oversized-utf8", "é".repeat(513))).toBeUndefined(); + + for (let index = 0; index < 65; index += 1) { + rememberComboForLane(`lane-${index}`, "first", target, `${index}`.padEnd(1024, "x"), captureConfigGeneration()); + } + expect(recallComboForLane(config, "lane-0", "0".padEnd(1024, "x"))).toBeUndefined(); + expect(recallComboForLane(config, "lane-1", "1".padEnd(1024, "x"))).toBe("first"); + expect(recallComboForLane(config, "lane-64", "64".padEnd(1024, "x"))).toBe("first"); + }); + test("combo recall watermark rejects writers after a partially failed generation", () => { registerStateStore(STATE_STORE_REGISTRATIONS.find(row => row.name === "combo-session-recall")!); const warning = spyOn(console, "warn").mockImplementation(() => {});