Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions src/lib/state-store-registrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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 },
Expand Down
52 changes: 46 additions & 6 deletions src/server/responses/combo-session-recall.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,16 @@ interface ComboRecallEntry {
comboId: string;
target: Pick<OcxComboTarget, "provider" | "model">;
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<string, ComboRecallEntry>();
let retainedModelBytes = 0;
let lastReconciledGeneration = 0;
let liveOwners: Pick<GenerationContext, "comboIds" | "comboTargets" | "providerNames"> | undefined;

Expand All @@ -22,6 +26,22 @@ function ownsEntry(context: Pick<GenerationContext, "comboIds" | "comboTargets"
&& context.comboTargets.has(`${entry.comboId}::${targetKey(entry.target)}`);
}

function deleteEntry(lane: string): boolean {
const entry = recall.get(lane);
if (!entry) return false;
retainedModelBytes -= entry.responseModelBytes;
recall.delete(lane);
return true;
}

function boundedModelBytes(model: string): number | undefined {
// UTF-8 is never shorter than the JS string's code-unit length. Reject oversized
// upstream values before encoding so validation itself cannot make another large copy.
if (model.length > 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,
Expand All @@ -30,16 +50,25 @@ export function rememberComboForLane(
writerGeneration: number,
): void {
if (!lane || !comboId || !responseModel.trim()) return;
const modelBytes = boundedModelBytes(responseModel);
if (modelBytes === undefined) return;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Clear stale recall when rejecting an oversized model

When a lane already has a cached entry and its next completed combo response reports a model over 1 KiB, this early return leaves the previous response registered as the lane's latest result. If the client later submits that previous bare model for compaction, recallComboForLane can incorrectly route it back through the old combo even though a newer response completed on the lane. After validating the writer generation and owner, invalidate the lane when the new model cannot be retained, and add a valid-to-oversized transition test.

Useful? React with 👍 / 👎.

// 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);
}
}

Expand All @@ -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;
Expand All @@ -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;
}
33 changes: 33 additions & 0 deletions tests/oauth/state-store-sweeper.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(() => {});
Expand Down
Loading