From 3479dcc3a7e414377dbdeed96cf3bb039ec6beef Mon Sep 17 00:00:00 2001 From: testikun Date: Thu, 3 Sep 2026 14:42:19 +0800 Subject: [PATCH 1/9] 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 926ba7c698..58678f55fb 100644 --- a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts +++ b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts @@ -4525,6 +4525,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 fec8bcbfbc..14184f498a 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 { @@ -1377,6 +1407,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 b4e70cf59baa9aeca1537986408475bbf99e5a71 Mon Sep 17 00:00:00 2001 From: testikun Date: Thu, 3 Sep 2026 15:04:39 +0800 Subject: [PATCH 2/9] 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 58678f55fb..ae2b8eb4e4 100644 --- a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts +++ b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts @@ -4582,7 +4582,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 14184f498a..91fc41a0c9 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); @@ -1409,11 +1409,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 c4d7bbabcf99d5a97dc21753c2376f9c1af48f3f Mon Sep 17 00:00:00 2001 From: testikun Date: Thu, 3 Sep 2026 22:33:24 +0800 Subject: [PATCH 3/9] 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 | 75 +++++++------------ .../history-compact-checkpoint-coordinator.ts | 42 ++++++++++- 3 files changed, 82 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 ae2b8eb4e4..8eb28f8003 100644 --- a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts +++ b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts @@ -4594,11 +4594,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 91fc41a0c9..08b0ae0081 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; @@ -419,7 +406,11 @@ export class AiSdkCompaction { source: { foldedRuntimeEvents: [...coveredRuntimeEvents], ...(input.runtimeContextRunHeaders - ? { runHeaders: input.runtimeContextRunHeaders } + ? { + runHeaders: input.runtimeContextRunHeaders.filter((run) => + coveredRuntimeEvents.some((event) => event.runId === run.runId), + ), + } : {}), }, newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], @@ -494,9 +485,7 @@ export class AiSdkCompaction { }), }; } finally { - if (this.historyCompactAbortController === historyCompactAbortController) { - this.historyCompactAbortController = null; - } + this.historyCompactAbortControllers.delete(historyCompactAbortController); } } @@ -539,8 +528,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 && @@ -555,6 +549,10 @@ export class AiSdkCompaction { } } throw error; + } finally { + if (this.inFlightHistorySummaries.get(fingerprint) === pending) { + this.inFlightHistorySummaries.delete(fingerprint); + } } } @@ -1407,25 +1405,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 684e71c7b5..a498bf3b24 100644 --- a/packages/runtime/src/history-compact-checkpoint-coordinator.ts +++ b/packages/runtime/src/history-compact-checkpoint-coordinator.ts @@ -86,9 +86,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); @@ -145,3 +153,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 76bf43a48fc75e7191df8676dcb95319a486b0ed Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 09:33:42 +0800 Subject: [PATCH 4/9] fix(runtime): scope compaction coalescing by checkpoint intent Generated-by: OpenAI Codex --- packages/runtime/src/ai-sdk-compaction.ts | 112 ++++++++++++++-------- 1 file changed, 73 insertions(+), 39 deletions(-) diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 08b0ae0081..5b990aa610 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'; @@ -399,25 +400,38 @@ 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.runtimeContextRunHeaders - ? { - runHeaders: input.runtimeContextRunHeaders.filter((run) => - coveredRuntimeEvents.some((event) => event.runId === run.runId), - ), - } + await this.summarizeWithFailureCircuit( + summarizer, + { + sessionId: this.sessionId, + turnId: input.turnId, + runId: input.runId, + source: { + foldedRuntimeEvents: [...coveredRuntimeEvents], + ...(input.runtimeContextRunHeaders + ? { + runHeaders: input.runtimeContextRunHeaders.filter((run) => + coveredRuntimeEvents.some((event) => event.runId === run.runId), + ), + } + : {}), + }, + newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], + ...(previousCheckpoint ? { previousCheckpoint } : {}), + inputBudget: { + maxEstimatedTokens: policy.maxHistoryEstimatedTokens ?? estimatedTokensBefore, + charsPerToken, + }, + 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' } }; @@ -496,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.runHeaders @@ -523,6 +542,7 @@ export class AiSdkCompaction { sourceRunRoutes, foldedRuntimeEvents: input.source.foldedRuntimeEvents, newlyFoldedRuntimeEvents: input.newlyFoldedRuntimeEvents, + checkpointIntent, }), ); const priorFailure = this.malformedSummaryFailures.get(fingerprint); @@ -1100,6 +1120,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', @@ -1117,30 +1146,35 @@ 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], - runHeaders: state.priorRunHeaders, + return await this.summarizeWithFailureCircuit( + summarizer, + { + sessionId: this.sessionId, + turnId, + ...(input.origin.runId ? { runId: input.origin.runId } : {}), + source: { + foldedRuntimeEvents: [...coveredRuntimeEvents], + runHeaders: state.priorRunHeaders, + }, + ...(previousCheckpoint ? { previousCheckpoint } : {}), + newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], + inputBudget: { + ...(state.capacity !== undefined + ? { maxEstimatedTokens: Math.max(1, state.capacity - reserveTokens) } + : {}), + charsPerToken, + }, + ...(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 32c0d9500c7fdc91011c62d57aa89ea7e1de949d Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 09:41:49 +0800 Subject: [PATCH 5/9] fix(runtime): align compaction dedup with current contract Generated-by: OpenAI Codex --- packages/runtime/src/__tests__/ai-sdk-backend.test.ts | 1 - packages/runtime/src/ai-sdk-compaction.ts | 10 ---------- 2 files changed, 11 deletions(-) diff --git a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts index 8eb28f8003..2c591c0dc4 100644 --- a/packages/runtime/src/__tests__/ai-sdk-backend.test.ts +++ b/packages/runtime/src/__tests__/ai-sdk-backend.test.ts @@ -4545,7 +4545,6 @@ describe('AiSdkBackend model history', () => { now: monotonicClock(), contextBudget: { name: 'in-flight-dedup-test', - maxHistoryEstimatedTokens: 10_000, charsPerToken: 1, }, summarizeHistoryCompact: async () => { diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 5b990aa610..c358179513 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -418,10 +418,6 @@ export class AiSdkCompaction { }, newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], ...(previousCheckpoint ? { previousCheckpoint } : {}), - inputBudget: { - maxEstimatedTokens: policy.maxHistoryEstimatedTokens ?? estimatedTokensBefore, - charsPerToken, - }, abortSignal: historyCompactAbortController.signal, ...(tracker ? { providerRequestTracker: tracker } : {}), }, @@ -1160,12 +1156,6 @@ export class AiSdkCompaction { }, ...(previousCheckpoint ? { previousCheckpoint } : {}), newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], - inputBudget: { - ...(state.capacity !== undefined - ? { maxEstimatedTokens: Math.max(1, state.capacity - reserveTokens) } - : {}), - charsPerToken, - }, ...(abortSignal ? { abortSignal } : {}), ...(midTurnTracker ? { providerRequestTracker: midTurnTracker } : {}), }, From 6d0e2e2a04f44d0c5136ab2b8528aa23cd69c0b5 Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 11:23:40 +0800 Subject: [PATCH 6/9] 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 c358179513..977fb2ed25 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; @@ -544,13 +547,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 && @@ -566,12 +582,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. * @@ -1677,6 +1703,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 02042ea172e25a5f9e05f6916a9fb8bd0e060d41 Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 11:27:07 +0800 Subject: [PATCH 7/9] 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 977fb2ed25..30d14a35c0 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 a498bf3b24..7f168cde1f 100644 --- a/packages/runtime/src/history-compact-checkpoint-coordinator.ts +++ b/packages/runtime/src/history-compact-checkpoint-coordinator.ts @@ -97,9 +97,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) { @@ -159,29 +162,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 f0a24f3d09a0da37bdfbfc39d01ff17dcd42c537 Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 11:28:43 +0800 Subject: [PATCH 8/9] fix(runtime): align compaction provenance and cancellation --- packages/runtime/src/ai-sdk-compaction.ts | 31 +++++++++++++++-------- 1 file changed, 20 insertions(+), 11 deletions(-) diff --git a/packages/runtime/src/ai-sdk-compaction.ts b/packages/runtime/src/ai-sdk-compaction.ts index 30d14a35c0..7890505126 100644 --- a/packages/runtime/src/ai-sdk-compaction.ts +++ b/packages/runtime/src/ai-sdk-compaction.ts @@ -406,8 +406,9 @@ export class AiSdkCompaction { foldedRuntimeEvents: [...coveredRuntimeEvents], ...(input.runtimeContextRunHeaders ? { - runHeaders: input.runtimeContextRunHeaders.filter((run) => - coveredRuntimeEvents.some((event) => event.runId === run.runId), + runHeaders: runHeadersForFoldedEvents( + input.runtimeContextRunHeaders, + coveredRuntimeEvents, ), } : {}), @@ -508,15 +509,15 @@ export class AiSdkCompaction { memoryExtractionBoundary?: HistoryCompactMemoryExtractionBoundary; }, ): Promise { - const foldedRunIds = new Set(input.source.foldedRuntimeEvents.map((event) => event.runId)); const sourceRunRoutes = input.source.runHeaders - ?.filter((run) => foldedRunIds.has(run.runId)) - .map((run) => ({ - runId: run.runId, - connectionId: run.llmConnectionId, - modelId: run.modelId, - })) - .sort((left, right) => left.runId.localeCompare(right.runId)); + ? input.source.runHeaders + .map((run) => ({ + runId: run.runId, + connectionId: run.llmConnectionId, + modelId: run.modelId, + })) + .sort((left, right) => left.runId.localeCompare(right.runId)) + : undefined; const fingerprint = sha256( stableStringifyForSignature({ version: 2, @@ -1171,7 +1172,7 @@ export class AiSdkCompaction { ...(input.origin.runId ? { runId: input.origin.runId } : {}), source: { foldedRuntimeEvents: [...coveredRuntimeEvents], - runHeaders: state.priorRunHeaders, + runHeaders: runHeadersForFoldedEvents(state.priorRunHeaders, coveredRuntimeEvents), }, ...(previousCheckpoint ? { previousCheckpoint } : {}), newlyFoldedRuntimeEvents: [...newlyFoldedRuntimeEvents], @@ -1730,6 +1731,14 @@ function waitForAbortablePromise( }); } +function runHeadersForFoldedEvents( + runHeaders: readonly AgentRunHeader[], + foldedRuntimeEvents: readonly RuntimeEvent[], +): AgentRunHeader[] { + const foldedRunIds = new Set(foldedRuntimeEvents.map((event) => event.runId)); + return runHeaders.filter((run) => foldedRunIds.has(run.runId)); +} + 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 cb8ff3daa484552d75ff9f314b18822d2a20a796 Mon Sep 17 00:00:00 2001 From: testikun Date: Fri, 4 Sep 2026 11:31:24 +0800 Subject: [PATCH 9/9] 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 7f168cde1f..6ce4c0e200 100644 --- a/packages/runtime/src/history-compact-checkpoint-coordinator.ts +++ b/packages/runtime/src/history-compact-checkpoint-coordinator.ts @@ -167,6 +167,7 @@ function hasSameEffectiveCoverage( checkpointId: _checkpointId, createdAt: _createdAt, highWaterSeq: _highWaterSeq, + previousCheckpointId: _previousCheckpointId, ...rest } = checkpoint; return JSON.stringify(rest);