From e03ba9cbc4ccf5d2389dc1c6cd7ea14824c618c6 Mon Sep 17 00:00:00 2001 From: testikun Date: Thu, 3 Sep 2026 14:42:19 +0800 Subject: [PATCH 01/11] fix(runtime): deduplicate in-flight context compaction Generated-by: OpenAI Codex --- .../src/__tests__/ai-sdk-backend.test.ts | 67 +++++++++++++++++++ packages/runtime/src/ai-sdk-compaction.ts | 43 +++++++++++- 2 files changed, 109 insertions(+), 1 deletion(-) diff --git a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts index f9693fd817..9b7f91ca3b 100644 --- a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts +++ b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts @@ -4527,6 +4527,73 @@ describe('AiSdkBackend model history', () => { assert.equal(result.contextBudget?.compactionDecisions?.[0]?.decision, 'replaced'); }); + test('coalesces identical in-flight compactHistory requests', async () => { + const recorded: HistoryCompactCheckpoint[] = []; + let summarizeCalls = 0; + let releaseSummary!: () => void; + const summaryReady = new Promise((resolve) => { + releaseSummary = resolve; + }); + const backend = createTestAiSdkBackend({ + sessionId: 'session-1', + header: header(), + appendMessage: async () => {}, + connection: connection(), + apiKey: 'sk-test', + modelId: 'mock-model-id', + modelFactory: () => completionModel(), + tools: [], + newId: idGenerator(), + now: monotonicClock(), + contextBudget: { + name: 'in-flight-dedup-test', + maxHistoryEstimatedTokens: 10_000, + charsPerToken: 1, + }, + summarizeHistoryCompact: async () => { + summarizeCalls += 1; + await summaryReady; + return structuredSummary('IN_FLIGHT_DEDUP_SUMMARY'); + }, + recordHistoryCompactCheckpoint: (checkpoint) => { + recorded.push(checkpoint); + }, + }); + const runtimeContext = [ + runtimeTextEvent({ + id: 'dedup-old-user', + turnId: 'dedup-old-turn', + role: 'user', + author: 'user', + text: 'old context '.repeat(100), + }), + runtimeTextEvent({ + id: 'dedup-old-agent', + turnId: 'dedup-old-turn', + role: 'model', + author: 'agent', + text: 'old response '.repeat(100), + }), + ]; + const first = backend.compactHistory({ + turnId: 'dedup-compact-1', + runId: 'run-dedup-1', + runtimeContext, + }); + await new Promise((resolve) => setImmediate(resolve)); + const second = backend.compactHistory({ + turnId: 'dedup-compact-2', + runId: 'run-dedup-2', + runtimeContext, + }); + + assert.equal(summarizeCalls, 1); + releaseSummary(); + const [firstResult, secondResult] = await Promise.all([first, second]); + assert.deepEqual(secondResult, firstResult); + assert.equal(recorded.length, 1); + }); + test('manual compactHistory compacts one completed turn with multiple agent steps', async () => { const recorded: HistoryCompactCheckpoint[] = []; const backend = createTestAiSdkBackend({ diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 4199e4541a..311a57787c 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -210,6 +210,18 @@ export class AiSdkCompaction { ) => Promise; private readonly canReplayProviderNative: (plan: RuntimeEventModelReplayPlan) => boolean; private historyCompactAbortController: AbortController | null = null; + /** + * Exact duplicate compaction requests share one physical summarizer call. + * Automatic capacity checks and explicit/manual entry can converge while a + * checkpoint is still being written; dispatching both would race the same + * source prefix and charge the provider twice. The key is derived from the + * source/configuration rather than the issuing turn id so callers that are + * otherwise asking for the same fold share the result. + */ + private readonly inFlightHistoryCompactions = new Map< + string, + Promise + >(); /** * Session-scoped circuit for exact malformed compaction inputs. A retry or * regeneration on the same backend must not dispatch the same doomed call; @@ -284,7 +296,25 @@ export class AiSdkCompaction { this.historyCompactAbortController?.abort(); } - public async compactHistory( + public compactHistory( + input: Omit & { runId: string | undefined }, + automaticMemoryBoundary?: HistoryCompactMemoryExtractionBoundary, + ): Promise { + const key = historyCompactRequestKey(input); + const existing = this.inFlightHistoryCompactions.get(key); + if (existing) return existing; + const pending = this.compactHistoryOnce(input, automaticMemoryBoundary); + this.inFlightHistoryCompactions.set(key, pending); + const clear = () => { + if (this.inFlightHistoryCompactions.get(key) === pending) { + this.inFlightHistoryCompactions.delete(key); + } + }; + void pending.then(clear, clear); + return pending; + } + + private async compactHistoryOnce( input: Omit & { runId: string | undefined }, automaticMemoryBoundary?: HistoryCompactMemoryExtractionBoundary, ): Promise { @@ -1382,6 +1412,17 @@ function sha256(text: string): string { return createHash('sha256').update(text).digest('hex'); } +function historyCompactRequestKey( + input: Omit & { runId: string | undefined }, +): string { + return sha256( + stableStringifyForSignature({ + runtimeContext: input.runtimeContext, + runtimeContextRunHeaders: input.runtimeContextRunHeaders ?? [], + }), + ); +} + function modelMessageSignature(message: ModelMessage): string { return sha256(stableStringifyForSignature(message)); } From 69531f2b3fe09fe4c90136d3679f9b53da4beb31 Mon Sep 17 00:00:00 2001 From: testikun Date: Thu, 3 Sep 2026 15:04:39 +0800 Subject: [PATCH 02/11] fix(runtime): key compaction deduplication by effective history Generated-by: OpenAI Codex --- .../runtime/src/__tests__/ai-sdk-backend.test.ts | 11 ++++++++++- packages/runtime/src/ai-sdk-compaction.ts | 14 +++++++++++--- 2 files changed, 21 insertions(+), 4 deletions(-) diff --git a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts index 9b7f91ca3b..3d7d6b312f 100644 --- a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts +++ b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts @@ -4584,7 +4584,16 @@ describe('AiSdkBackend model history', () => { const second = backend.compactHistory({ turnId: 'dedup-compact-2', runId: 'run-dedup-2', - runtimeContext, + runtimeContext: [ + ...runtimeContext, + runtimeTextEvent({ + id: 'dedup-current-turn', + turnId: 'dedup-compact-2', + role: 'user', + author: 'user', + text: 'current turn content must not affect the fold key', + }), + ], }); assert.equal(summarizeCalls, 1); diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 311a57787c..679c4ba63b 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -300,7 +300,7 @@ export class AiSdkCompaction { input: Omit & { runId: string | undefined }, automaticMemoryBoundary?: HistoryCompactMemoryExtractionBoundary, ): Promise { - const key = historyCompactRequestKey(input); + const key = historyCompactRequestKey(input, automaticMemoryBoundary); const existing = this.inFlightHistoryCompactions.get(key); if (existing) return existing; const pending = this.compactHistoryOnce(input, automaticMemoryBoundary); @@ -1414,11 +1414,19 @@ function sha256(text: string): string { function historyCompactRequestKey( input: Omit & { runId: string | undefined }, + automaticMemoryBoundary: HistoryCompactMemoryExtractionBoundary | undefined, ): string { + const runtimeContext = input.runtimeContext + .filter((event) => event.turnId !== input.turnId) + .filter(isHistoryCompactContentEvent); + const sourceRunIds = new Set(runtimeContext.map((event) => event.runId)); return sha256( stableStringifyForSignature({ - runtimeContext: input.runtimeContext, - runtimeContextRunHeaders: input.runtimeContextRunHeaders ?? [], + runtimeContext, + runtimeContextRunHeaders: (input.runtimeContextRunHeaders ?? []).filter((header) => + sourceRunIds.has(header.runId), + ), + automaticMemoryBoundary, }), ); } From b0b3d256ea611b3aaa3a248609c09ae2f6cc51c6 Mon Sep 17 00:00:00 2001 From: testikun Date: Thu, 3 Sep 2026 22:33:24 +0800 Subject: [PATCH 03/11] fix(runtime): deduplicate effective history compaction calls Generated-by: OpenAI Codex --- .../src/__tests__/ai-sdk-backend.test.ts | 17 +++- packages/runtime/src/ai-sdk-compaction.ts | 77 +++++++------------ .../history-compact-checkpoint-coordinator.ts | 42 +++++++++- 3 files changed, 84 insertions(+), 52 deletions(-) diff --git a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts index 3d7d6b312f..436eab87fe 100644 --- a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts +++ b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts @@ -4596,11 +4596,22 @@ describe('AiSdkBackend model history', () => { ], }); - assert.equal(summarizeCalls, 1); releaseSummary(); const [firstResult, secondResult] = await Promise.all([first, second]); - assert.deepEqual(secondResult, firstResult); - assert.equal(recorded.length, 1); + assert.equal(summarizeCalls, 1); + assert.equal(secondResult.outcome.kind, 'compacted'); + assert.equal(firstResult.outcome.kind, 'compacted'); + assert.equal(recorded.length, 2); + const firstCheckpoint = recorded[0]; + const secondCheckpoint = recorded[1]; + assert.ok(firstCheckpoint); + assert.ok(secondCheckpoint); + assert.equal(firstCheckpoint.version, 2); + assert.equal(secondCheckpoint.version, 2); + if (firstCheckpoint.version === 2 && secondCheckpoint.version === 2) { + assert.equal(secondCheckpoint.summary, firstCheckpoint.summary); + assert.deepEqual(secondCheckpoint.coverage, firstCheckpoint.coverage); + } }); test('manual compactHistory compacts one completed turn with multiple agent steps', async () => { diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 679c4ba63b..21712368de 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -209,18 +209,16 @@ export class AiSdkCompaction { providerReasoningReplayEventIds: ReadonlySet, ) => Promise; private readonly canReplayProviderNative: (plan: RuntimeEventModelReplayPlan) => boolean; - private historyCompactAbortController: AbortController | null = null; + private readonly historyCompactAbortControllers = new Set(); /** - * Exact duplicate compaction requests share one physical summarizer call. - * Automatic capacity checks and explicit/manual entry can converge while a - * checkpoint is still being written; dispatching both would race the same - * source prefix and charge the provider twice. The key is derived from the - * source/configuration rather than the issuing turn id so callers that are - * otherwise asking for the same fold share the result. + * Exact duplicate summary inputs share one physical provider call. The map + * lives at the summarizer boundary, where the existing effective-history + * fingerprint already describes the request that spends provider budget. + * Each caller still completes its own checkpoint/turn bookkeeping. */ - private readonly inFlightHistoryCompactions = new Map< + private readonly inFlightHistorySummaries = new Map< string, - Promise + Promise >(); /** * Session-scoped circuit for exact malformed compaction inputs. A retry or @@ -293,25 +291,14 @@ export class AiSdkCompaction { /** Abort an in-flight manual history compaction (called by AiSdkBackend.stop). */ public abortHistoryCompact(): void { - this.historyCompactAbortController?.abort(); + for (const controller of this.historyCompactAbortControllers) controller.abort(); } - public compactHistory( + public async compactHistory( input: Omit & { runId: string | undefined }, automaticMemoryBoundary?: HistoryCompactMemoryExtractionBoundary, ): Promise { - const key = historyCompactRequestKey(input, automaticMemoryBoundary); - const existing = this.inFlightHistoryCompactions.get(key); - if (existing) return existing; - const pending = this.compactHistoryOnce(input, automaticMemoryBoundary); - this.inFlightHistoryCompactions.set(key, pending); - const clear = () => { - if (this.inFlightHistoryCompactions.get(key) === pending) { - this.inFlightHistoryCompactions.delete(key); - } - }; - void pending.then(clear, clear); - return pending; + return this.compactHistoryOnce(input, automaticMemoryBoundary); } private async compactHistoryOnce( @@ -319,7 +306,7 @@ export class AiSdkCompaction { automaticMemoryBoundary?: HistoryCompactMemoryExtractionBoundary, ): Promise { const historyCompactAbortController = new AbortController(); - this.historyCompactAbortController = historyCompactAbortController; + this.historyCompactAbortControllers.add(historyCompactAbortController); try { const policy = this.input.contextBudget; const summarizer = this.input.summarizeHistoryCompact; @@ -421,7 +408,13 @@ export class AiSdkCompaction { source: { foldedRuntimeEvents: [...coveredRuntimeEvents], ...(input.runtimeContextInvocations - ? { invocations: input.runtimeContextInvocations } + ? { + invocations: input.runtimeContextInvocations.filter((invocation) => + coveredRuntimeEvents.some( + (event) => event.invocationId === invocation.invocationId, + ), + ), + } : {}), }, newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], @@ -496,9 +489,7 @@ export class AiSdkCompaction { }), }; } finally { - if (this.historyCompactAbortController === historyCompactAbortController) { - this.historyCompactAbortController = null; - } + this.historyCompactAbortControllers.delete(historyCompactAbortController); } } @@ -544,8 +535,13 @@ export class AiSdkCompaction { const priorFailure = this.malformedSummaryFailures.get(fingerprint); if (priorFailure) throw new HistoryCompactSummarizerError(priorFailure); + const existing = this.inFlightHistorySummaries.get(fingerprint); + if (existing) return existing; + + const pending = Promise.resolve().then(() => summarizer(input)); + this.inFlightHistorySummaries.set(fingerprint, pending); try { - return await Promise.resolve(summarizer(input)); + return await pending; } catch (error) { if ( error instanceof HistoryCompactSummarizerError && @@ -560,6 +556,10 @@ export class AiSdkCompaction { } } throw error; + } finally { + if (this.inFlightHistorySummaries.get(fingerprint) === pending) { + this.inFlightHistorySummaries.delete(fingerprint); + } } } @@ -1412,25 +1412,6 @@ function sha256(text: string): string { return createHash('sha256').update(text).digest('hex'); } -function historyCompactRequestKey( - input: Omit & { runId: string | undefined }, - automaticMemoryBoundary: HistoryCompactMemoryExtractionBoundary | undefined, -): string { - const runtimeContext = input.runtimeContext - .filter((event) => event.turnId !== input.turnId) - .filter(isHistoryCompactContentEvent); - const sourceRunIds = new Set(runtimeContext.map((event) => event.runId)); - return sha256( - stableStringifyForSignature({ - runtimeContext, - runtimeContextRunHeaders: (input.runtimeContextRunHeaders ?? []).filter((header) => - sourceRunIds.has(header.runId), - ), - automaticMemoryBoundary, - }), - ); -} - function modelMessageSignature(message: ModelMessage): string { return sha256(stableStringifyForSignature(message)); } diff --git a/packages/runtime/src/history-compact-checkpoint-coordinator.ts b/packages/runtime/src/history-compact-checkpoint-coordinator.ts index 05533e500f..3478725961 100644 --- a/packages/runtime/src/history-compact-checkpoint-coordinator.ts +++ b/packages/runtime/src/history-compact-checkpoint-coordinator.ts @@ -90,9 +90,17 @@ export class HistoryCompactCheckpointCoordinator { .catch(() => {}) .then(async () => { const durableCheckpoint = await this.load(sessionId); - if (!canReplaceHistoryCompactCheckpoint(durableCheckpoint, checkpoint)) { + const sameEffectiveCheckpoint = hasSameEffectiveCoverage(durableCheckpoint, checkpoint); + if ( + !sameEffectiveCheckpoint && + !canReplaceHistoryCompactCheckpoint(durableCheckpoint, checkpoint) + ) { throw new Error('History compact checkpoint was superseded before persistence'); } + // A concurrent caller may have received the same effective checkpoint + // from the summarizer coalescer. Keep its run-local ledger event too so + // the rider Turn retains provenance even though the session checkpoint + // itself is already current. await run.recordHistoryCompactCheckpoint(checkpoint); this.checkpoints.set(sessionId, checkpoint); this.scheduleCleanup(sessionId, checkpoint); @@ -155,3 +163,35 @@ export class HistoryCompactCheckpointCoordinator { this.cleanups.set(sessionId, tracked); } } + +function hasSameEffectiveCoverage( + current: HistoryCompactCheckpoint | undefined, + candidate: HistoryCompactCheckpoint, +): boolean { + if (!current) return false; + const currentContent = checkpointContent(current); + const candidateContent = checkpointContent(candidate); + return ( + current.sessionId === candidate.sessionId && + current.version === candidate.version && + current.highWaterName === candidate.highWaterName && + current.phase === candidate.phase && + current.coverage.eventCount === candidate.coverage.eventCount && + current.coverage.turnCount === candidate.coverage.turnCount && + current.coverage.sourceDigest === candidate.coverage.sourceDigest && + current.coverage.through.runId === candidate.coverage.through.runId && + current.coverage.through.turnId === candidate.coverage.through.turnId && + current.coverage.through.runtimeEventId === candidate.coverage.through.runtimeEventId && + JSON.stringify(current.source) === JSON.stringify(candidate.source) && + JSON.stringify(current.headAnchor) === JSON.stringify(candidate.headAnchor) && + JSON.stringify(current.memoryExtractionBoundary) === + JSON.stringify(candidate.memoryExtractionBoundary) && + JSON.stringify(currentContent) === JSON.stringify(candidateContent) + ); +} + +function checkpointContent(checkpoint: HistoryCompactCheckpoint): unknown { + return checkpoint.version === 2 + ? { summary: checkpoint.summary, summaryFormat: checkpoint.summaryFormat } + : { providerState: checkpoint.providerState }; +} From da997c4c314b761cf352d2f75be53f1ed2858ced Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 09:33:42 +0800 Subject: [PATCH 04/11] fix(runtime): scope compaction coalescing by checkpoint intent Generated-by: OpenAI Codex --- packages/runtime/src/ai-sdk-compaction.ts | 104 +++++++++++++--------- 1 file changed, 64 insertions(+), 40 deletions(-) diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 21712368de..3599164b01 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -57,6 +57,7 @@ import { matchHistoryCompactCheckpointPrefix, projectHistoryCompactCheckpointReplay, type HistoryCompactCheckpoint, + type HistoryCompactCheckpointHeadAnchor, type HistoryCompactMemoryExtractionBoundary, type HistoryCompactProviderState, } from './history-compact-checkpoint.js'; @@ -401,27 +402,36 @@ export class AiSdkCompaction { ...(automaticMemoryBoundary ? { memoryExtractionBoundary: automaticMemoryBoundary } : {}), ...(previousCheckpoint ? { previousCheckpoint } : {}), summarize: async ({ coveredRuntimeEvents, newlyFoldedRuntimeEvents, previousCheckpoint }) => - await this.summarizeWithFailureCircuit(summarizer, { - sessionId: this.sessionId, - turnId: input.turnId, - runId: input.runId, - source: { - foldedRuntimeEvents: [...coveredRuntimeEvents], - ...(input.runtimeContextInvocations - ? { - invocations: input.runtimeContextInvocations.filter((invocation) => - coveredRuntimeEvents.some( - (event) => event.invocationId === invocation.invocationId, + await this.summarizeWithFailureCircuit( + summarizer, + { + sessionId: this.sessionId, + turnId: input.turnId, + runId: input.runId, + source: { + foldedRuntimeEvents: [...coveredRuntimeEvents], + ...(input.runtimeContextInvocations + ? { + invocations: input.runtimeContextInvocations.filter((invocation) => + coveredRuntimeEvents.some( + (event) => event.invocationId === invocation.invocationId, + ), ), - ), - } + } + : {}), + }, + newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], + ...(previousCheckpoint ? { previousCheckpoint } : {}), + abortSignal: historyCompactAbortController.signal, + ...(tracker ? { providerRequestTracker: tracker } : {}), + }, + { + phase: 'pre_turn', + ...(automaticMemoryBoundary + ? { memoryExtractionBoundary: automaticMemoryBoundary } : {}), }, - newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], - ...(previousCheckpoint ? { previousCheckpoint } : {}), - abortSignal: historyCompactAbortController.signal, - ...(tracker ? { providerRequestTracker: tracker } : {}), - }), + ), }); if (historyCompactAbortController.signal.aborted) { return { outcome: { kind: 'failed', reason: 'aborted' } }; @@ -500,6 +510,11 @@ export class AiSdkCompaction { private async summarizeWithFailureCircuit( summarizer: HistoryCompactSummarizer, input: HistoryCompactSummaryInput, + checkpointIntent?: { + phase?: 'pre_turn' | 'mid_turn'; + headAnchor?: HistoryCompactCheckpointHeadAnchor; + memoryExtractionBoundary?: HistoryCompactMemoryExtractionBoundary; + }, ): Promise { const foldedRunIds = new Set(input.source.foldedRuntimeEvents.map((event) => event.runId)); const sourceRunRoutes = input.source.invocations @@ -530,6 +545,7 @@ export class AiSdkCompaction { sourceRunRoutes, foldedRuntimeEvents: input.source.foldedRuntimeEvents, newlyFoldedRuntimeEvents: input.newlyFoldedRuntimeEvents, + checkpointIntent, }), ); const priorFailure = this.malformedSummaryFailures.get(fingerprint); @@ -1107,6 +1123,15 @@ export class AiSdkCompaction { } const orderedEvents = [...state.priorContentEvents, ...currentTurnEvents]; const memoryDecision = input.memoryCompactionDecision?.(); + const memoryExtractionBoundary = + memoryDecision && orderedEvents.at(-1) + ? { + runId: orderedEvents.at(-1)!.runId, + turnId: orderedEvents.at(-1)!.turnId, + runtimeEventId: orderedEvents.at(-1)!.id, + disposition: memoryDecision.disposition, + } + : undefined; const plan = await planHistoryCompaction({ sessionId: this.sessionId, phase: input.phase ?? 'mid_turn', @@ -1124,30 +1149,29 @@ export class AiSdkCompaction { ? { highWaterName: compactPolicy.highWaterName } : {}), ...(state.previousCheckpoint ? { previousCheckpoint: state.previousCheckpoint } : {}), - ...(memoryDecision && orderedEvents.at(-1) - ? { - memoryExtractionBoundary: { - runId: orderedEvents.at(-1)!.runId, - turnId: orderedEvents.at(-1)!.turnId, - runtimeEventId: orderedEvents.at(-1)!.id, - disposition: memoryDecision.disposition, - }, - } - : {}), + ...(memoryExtractionBoundary ? { memoryExtractionBoundary } : {}), summarize: async ({ coveredRuntimeEvents, newlyFoldedRuntimeEvents, previousCheckpoint }) => { - return await this.summarizeWithFailureCircuit(summarizer, { - sessionId: this.sessionId, - turnId, - ...(input.origin.runId ? { runId: input.origin.runId } : {}), - source: { - foldedRuntimeEvents: [...coveredRuntimeEvents], - invocations: state.priorInvocations, + return await this.summarizeWithFailureCircuit( + summarizer, + { + sessionId: this.sessionId, + turnId, + ...(input.origin.runId ? { runId: input.origin.runId } : {}), + source: { + foldedRuntimeEvents: [...coveredRuntimeEvents], + invocations: state.priorInvocations, + }, + ...(previousCheckpoint ? { previousCheckpoint } : {}), + newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], + ...(abortSignal ? { abortSignal } : {}), + ...(midTurnTracker ? { providerRequestTracker: midTurnTracker } : {}), }, - ...(previousCheckpoint ? { previousCheckpoint } : {}), - newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], - ...(abortSignal ? { abortSignal } : {}), - ...(midTurnTracker ? { providerRequestTracker: midTurnTracker } : {}), - }); + { + phase: input.phase ?? 'mid_turn', + headAnchor: { runtimeEventId: state.headAnchor.id, turnId }, + ...(memoryExtractionBoundary ? { memoryExtractionBoundary } : {}), + }, + ); }, }); From 0300051bd193b0306b71531e1c02eaec70d0f0fc Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 09:41:49 +0800 Subject: [PATCH 05/11] fix(runtime): align compaction dedup with current contract Generated-by: OpenAI Codex --- packages/runtime/src/__tests__/ai-sdk-backend.test.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts index 436eab87fe..ade55a6edc 100644 --- a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts +++ b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts @@ -4547,7 +4547,6 @@ describe('AiSdkBackend model history', () => { now: monotonicClock(), contextBudget: { name: 'in-flight-dedup-test', - maxHistoryEstimatedTokens: 10_000, charsPerToken: 1, }, summarizeHistoryCompact: async () => { From c88ecdb0a95c4ee64eebe0e76ce1979809133ac2 Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 11:23:40 +0800 Subject: [PATCH 06/11] fix(runtime): isolate aborts for shared compaction calls --- packages/runtime/src/ai-sdk-compaction.ts | 84 +++++++++++++++++++---- 1 file changed, 72 insertions(+), 12 deletions(-) diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 3599164b01..4fa99b4bc9 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -155,6 +155,12 @@ export interface AutomaticMemoryCompactionDecision { readonly dispatch: boolean; } +interface InFlightHistorySummary { + readonly abortController: AbortController; + readonly consumers: Set; + readonly promise: Promise; +} + /** Constructor dependencies for AiSdkCompaction. */ export interface AiSdkCompactionDeps { input: AiSdkCompactionCapabilities; @@ -217,10 +223,7 @@ export class AiSdkCompaction { * fingerprint already describes the request that spends provider budget. * Each caller still completes its own checkpoint/turn bookkeeping. */ - private readonly inFlightHistorySummaries = new Map< - string, - Promise - >(); + private readonly inFlightHistorySummaries = new Map(); /** * Session-scoped circuit for exact malformed compaction inputs. A retry or * regeneration on the same backend must not dispatch the same doomed call; @@ -551,13 +554,26 @@ export class AiSdkCompaction { const priorFailure = this.malformedSummaryFailures.get(fingerprint); if (priorFailure) throw new HistoryCompactSummarizerError(priorFailure); - const existing = this.inFlightHistorySummaries.get(fingerprint); - if (existing) return existing; - - const pending = Promise.resolve().then(() => summarizer(input)); - this.inFlightHistorySummaries.set(fingerprint, pending); + let shared = this.inFlightHistorySummaries.get(fingerprint); + if (!shared) { + const abortController = new AbortController(); + const pending = Promise.resolve().then(() => + summarizer({ ...input, abortSignal: abortController.signal }), + ); + shared = { abortController, consumers: new Set(), promise: pending }; + this.inFlightHistorySummaries.set(fingerprint, shared); + void pending.then( + () => this.removeInFlightHistorySummary(fingerprint, shared!), + () => this.removeInFlightHistorySummary(fingerprint, shared!), + ); + } + const consumer = Symbol('history-summary-consumer'); + shared.consumers.add(consumer); try { - return await pending; + // A rider must be able to stop waiting without aborting the physical + // call for other consumers. The shared controller is aborted only when + // every consumer has detached. + return await waitForAbortablePromise(shared.promise, input.abortSignal); } catch (error) { if ( error instanceof HistoryCompactSummarizerError && @@ -573,12 +589,22 @@ export class AiSdkCompaction { } throw error; } finally { - if (this.inFlightHistorySummaries.get(fingerprint) === pending) { - this.inFlightHistorySummaries.delete(fingerprint); + shared.consumers.delete(consumer); + if ( + shared.consumers.size === 0 && + this.inFlightHistorySummaries.get(fingerprint) === shared + ) { + shared.abortController.abort(); } } } + private removeInFlightHistorySummary(fingerprint: string, shared: InFlightHistorySummary): void { + if (this.inFlightHistorySummaries.get(fingerprint) === shared) { + this.inFlightHistorySummaries.delete(fingerprint); + } + } + /** * Fold the durable transition ledger onto any slice of model-visible history. * @@ -1688,6 +1714,40 @@ function waitForQueueProgressOrAbort( }); } +/** Let an individual compaction stop waiting without cancelling shared work. */ +function waitForAbortablePromise( + promise: Promise, + abortSignal?: AbortSignal, +): Promise { + if (!abortSignal) return promise; + if (abortSignal.aborted) return Promise.resolve(undefined); + return new Promise((resolve, reject) => { + let settled = false; + const cleanup = () => abortSignal.removeEventListener('abort', onAbort); + const onAbort = () => { + if (settled) return; + settled = true; + cleanup(); + resolve(undefined); + }; + abortSignal.addEventListener('abort', onAbort, { once: true }); + void promise.then( + (value) => { + if (settled) return; + settled = true; + cleanup(); + resolve(value); + }, + (error: unknown) => { + if (settled) return; + settled = true; + cleanup(); + reject(error); + }, + ); + }); +} + export function hasBlockingReplayDiagnostics(plan: RuntimeEventModelReplayPlan): boolean { // `unmatched_tool_result` is deliberately NOT blocking: the materializer // drops an orphan tool result (its call sliced away or the ledger corrupt) From c596de1035dd10430ac1d27bd1b54cb11de935a2 Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 11:27:07 +0800 Subject: [PATCH 07/11] fix(runtime): preserve shared compaction checkpoint identity --- packages/runtime/src/ai-sdk-compaction.ts | 7 --- .../history-compact-checkpoint-coordinator.ts | 44 +++++++------------ 2 files changed, 16 insertions(+), 35 deletions(-) diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 4fa99b4bc9..069fa081a0 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -301,13 +301,6 @@ export class AiSdkCompaction { public async compactHistory( input: Omit & { runId: string | undefined }, automaticMemoryBoundary?: HistoryCompactMemoryExtractionBoundary, - ): Promise { - return this.compactHistoryOnce(input, automaticMemoryBoundary); - } - - private async compactHistoryOnce( - input: Omit & { runId: string | undefined }, - automaticMemoryBoundary?: HistoryCompactMemoryExtractionBoundary, ): Promise { const historyCompactAbortController = new AbortController(); this.historyCompactAbortControllers.add(historyCompactAbortController); diff --git a/packages/runtime/src/history-compact-checkpoint-coordinator.ts b/packages/runtime/src/history-compact-checkpoint-coordinator.ts index 3478725961..884ab70df9 100644 --- a/packages/runtime/src/history-compact-checkpoint-coordinator.ts +++ b/packages/runtime/src/history-compact-checkpoint-coordinator.ts @@ -101,9 +101,12 @@ export class HistoryCompactCheckpointCoordinator { // from the summarizer coalescer. Keep its run-local ledger event too so // the rider Turn retains provenance even though the session checkpoint // itself is already current. - await run.recordHistoryCompactCheckpoint(checkpoint); - this.checkpoints.set(sessionId, checkpoint); - this.scheduleCleanup(sessionId, checkpoint); + const checkpointToRecord = sameEffectiveCheckpoint ? durableCheckpoint! : checkpoint; + await run.recordHistoryCompactCheckpoint(checkpointToRecord); + if (!sameEffectiveCheckpoint) { + this.checkpoints.set(sessionId, checkpoint); + this.scheduleCleanup(sessionId, checkpoint); + } }) .finally(() => { if (this.writes.get(sessionId) === tracked) { @@ -169,29 +172,14 @@ function hasSameEffectiveCoverage( candidate: HistoryCompactCheckpoint, ): boolean { if (!current) return false; - const currentContent = checkpointContent(current); - const candidateContent = checkpointContent(candidate); - return ( - current.sessionId === candidate.sessionId && - current.version === candidate.version && - current.highWaterName === candidate.highWaterName && - current.phase === candidate.phase && - current.coverage.eventCount === candidate.coverage.eventCount && - current.coverage.turnCount === candidate.coverage.turnCount && - current.coverage.sourceDigest === candidate.coverage.sourceDigest && - current.coverage.through.runId === candidate.coverage.through.runId && - current.coverage.through.turnId === candidate.coverage.through.turnId && - current.coverage.through.runtimeEventId === candidate.coverage.through.runtimeEventId && - JSON.stringify(current.source) === JSON.stringify(candidate.source) && - JSON.stringify(current.headAnchor) === JSON.stringify(candidate.headAnchor) && - JSON.stringify(current.memoryExtractionBoundary) === - JSON.stringify(candidate.memoryExtractionBoundary) && - JSON.stringify(currentContent) === JSON.stringify(candidateContent) - ); -} - -function checkpointContent(checkpoint: HistoryCompactCheckpoint): unknown { - return checkpoint.version === 2 - ? { summary: checkpoint.summary, summaryFormat: checkpoint.summaryFormat } - : { providerState: checkpoint.providerState }; + const stable = (checkpoint: HistoryCompactCheckpoint): string => { + const { + checkpointId: _checkpointId, + createdAt: _createdAt, + highWaterSeq: _highWaterSeq, + ...rest + } = checkpoint; + return JSON.stringify(rest); + }; + return stable(current) === stable(candidate); } From c1981ca43279719266f0332903c05cbad8b9a2ef Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 11:28:43 +0800 Subject: [PATCH 08/11] fix(runtime): align compaction provenance and cancellation --- packages/runtime/src/ai-sdk-compaction.ts | 24 +++++++++++++++-------- 1 file changed, 16 insertions(+), 8 deletions(-) diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 069fa081a0..79ec5dacfb 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -408,10 +408,9 @@ export class AiSdkCompaction { foldedRuntimeEvents: [...coveredRuntimeEvents], ...(input.runtimeContextInvocations ? { - invocations: input.runtimeContextInvocations.filter((invocation) => - coveredRuntimeEvents.some( - (event) => event.invocationId === invocation.invocationId, - ), + invocations: invocationsForFoldedEvents( + input.runtimeContextInvocations, + coveredRuntimeEvents, ), } : {}), @@ -512,10 +511,8 @@ export class AiSdkCompaction { memoryExtractionBoundary?: HistoryCompactMemoryExtractionBoundary; }, ): Promise { - const foldedRunIds = new Set(input.source.foldedRuntimeEvents.map((event) => event.runId)); const sourceRunRoutes = input.source.invocations - ?.filter((invocation) => foldedRunIds.has(invocation.runId)) - .map((invocation) => { + ?.map((invocation) => { const route = invocation.opening.route; return { runId: invocation.runId, @@ -1178,7 +1175,10 @@ export class AiSdkCompaction { ...(input.origin.runId ? { runId: input.origin.runId } : {}), source: { foldedRuntimeEvents: [...coveredRuntimeEvents], - invocations: state.priorInvocations, + invocations: invocationsForFoldedEvents( + state.priorInvocations, + coveredRuntimeEvents, + ), }, ...(previousCheckpoint ? { previousCheckpoint } : {}), newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], @@ -1741,6 +1741,14 @@ function waitForAbortablePromise( }); } +function invocationsForFoldedEvents( + invocations: readonly RuntimeInvocationRecord[], + foldedRuntimeEvents: readonly RuntimeEvent[], +): RuntimeInvocationRecord[] { + const foldedInvocationIds = new Set(foldedRuntimeEvents.map((event) => event.invocationId)); + return invocations.filter((invocation) => foldedInvocationIds.has(invocation.invocationId)); +} + export function hasBlockingReplayDiagnostics(plan: RuntimeEventModelReplayPlan): boolean { // `unmatched_tool_result` is deliberately NOT blocking: the materializer // drops an orphan tool result (its call sliced away or the ledger corrupt) From 02851a33201e806821e227293564c6ef5706f0f6 Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 11:31:24 +0800 Subject: [PATCH 09/11] fix(runtime): align compaction provenance and cancellation --- packages/runtime/src/history-compact-checkpoint-coordinator.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/packages/runtime/src/history-compact-checkpoint-coordinator.ts b/packages/runtime/src/history-compact-checkpoint-coordinator.ts index 884ab70df9..3ee0a9b85a 100644 --- a/packages/runtime/src/history-compact-checkpoint-coordinator.ts +++ b/packages/runtime/src/history-compact-checkpoint-coordinator.ts @@ -177,6 +177,7 @@ function hasSameEffectiveCoverage( checkpointId: _checkpointId, createdAt: _createdAt, highWaterSeq: _highWaterSeq, + previousCheckpointId: _previousCheckpointId, ...rest } = checkpoint; return JSON.stringify(rest); From 155366d7d8e40f220d35e4092e98a5a94f845870 Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 16:29:32 +0800 Subject: [PATCH 10/11] fix(runtime): align invocation provenance after rebase --- packages/runtime/src/ai-sdk-compaction.ts | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 79ec5dacfb..0395bd767c 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -1175,10 +1175,7 @@ export class AiSdkCompaction { ...(input.origin.runId ? { runId: input.origin.runId } : {}), source: { foldedRuntimeEvents: [...coveredRuntimeEvents], - invocations: invocationsForFoldedEvents( - state.priorInvocations, - coveredRuntimeEvents, - ), + invocations: invocationsForFoldedEvents(state.priorInvocations, coveredRuntimeEvents), }, ...(previousCheckpoint ? { previousCheckpoint } : {}), newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], From 12bb2c3e25c41bd55fe96c200d40b61ee70b26a8 Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 16:52:19 +0800 Subject: [PATCH 11/11] fix(runtime): retain run routing in shared compaction --- packages/runtime/src/ai-sdk-compaction.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 0395bd767c..56ee552246 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -1742,8 +1742,8 @@ function invocationsForFoldedEvents( invocations: readonly RuntimeInvocationRecord[], foldedRuntimeEvents: readonly RuntimeEvent[], ): RuntimeInvocationRecord[] { - const foldedInvocationIds = new Set(foldedRuntimeEvents.map((event) => event.invocationId)); - return invocations.filter((invocation) => foldedInvocationIds.has(invocation.invocationId)); + const foldedRunIds = new Set(foldedRuntimeEvents.map((event) => event.runId)); + return invocations.filter((invocation) => foldedRunIds.has(invocation.runId)); } export function hasBlockingReplayDiagnostics(plan: RuntimeEventModelReplayPlan): boolean {