diff --git a/src/session/assemble-runtime.ts b/src/session/assemble-runtime.ts index 9998f8f61..8e7960c5c 100644 --- a/src/session/assemble-runtime.ts +++ b/src/session/assemble-runtime.ts @@ -557,11 +557,11 @@ export interface ChatAgentWiring { */ getCompactor: (wrapPruning?: (pruning: Compactor) => Compactor) => Compactor; /** - * Bound to the TUI compaction lifecycle abort signal. The completeness - * gate runs inside wrapCompactor's race, so a discarded certified stub - * must not persist a compaction-handoff for a fold that never landed. + * Bound to the TUI compaction lifecycle abort signal. Captured at + * completeness-gate apply start so onBuilt reset() cannot un-abort an + * in-flight fold that wrapCompactor already discarded. */ - isCompactionAborted?: () => boolean; + getCompactionAbortSignal?: () => AbortSignal; /** Experimental Anthropic prompt shrink. Default off when omitted. */ anthropicCachePrompt?: () => boolean; /** @@ -795,9 +795,9 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent { wrapCompactorWithCompletenessGate( pruning, primaryArchive, - wiring.isCompactionAborted === undefined + wiring.getCompactionAbortSignal === undefined ? undefined - : { isAborted: wiring.isCompactionAborted }, + : { getSignal: wiring.getCompactionAbortSignal }, ), ), }, diff --git a/src/session/compaction-archive.ts b/src/session/compaction-archive.ts index c00a58a2a..40487af4d 100644 --- a/src/session/compaction-archive.ts +++ b/src/session/compaction-archive.ts @@ -952,16 +952,21 @@ async function recordFreshHandoffOutput( * drop it even when the summarizer does not echo it verbatim. `isAborted` * skips that record: the TUI wrapCompactor race can discard a certified * stub, and a handoff for a fold that never landed is a phantom. + * `getSignal` is captured at apply start so onBuilt `reset()` cannot + * un-abort an in-flight fold that wrapCompactor already discarded. */ export function wrapCompactorWithCompletenessGate( inner: Compactor, archive: CompactionArchive, - opts?: { isAborted?: () => boolean }, + opts?: { isAborted?: () => boolean; getSignal?: () => AbortSignal }, ): Compactor { return { name: inner.name, version: inner.version, async apply(turns: ConversationTurn[], ctx: StrategyContext) { + const abortSignal = opts?.getSignal?.(); + const isAborted = () => + abortSignal?.aborted === true || opts?.isAborted?.() === true; await archive.awaitPendingWrites(); const proposed = await inner.apply(turns, ctx); const units = uncoveredContentUnits(turns, proposed.output); @@ -988,14 +993,14 @@ export function wrapCompactorWithCompletenessGate( if (certificate.status !== "complete") { return incompleteIdentity(inner, turns); } - if (opts?.isAborted?.() === true) { + if (isAborted()) { return proposed; } await recordFreshHandoffOutput( archive, turns, proposed.output, - opts?.isAborted, + isAborted, ); return proposed; }, diff --git a/src/session/runtime-assembly.test.ts b/src/session/runtime-assembly.test.ts index c3eeb0942..324a08475 100644 --- a/src/session/runtime-assembly.test.ts +++ b/src/session/runtime-assembly.test.ts @@ -985,6 +985,111 @@ describe("createSessionPruningCompactor stub fallback", () => { ); expect(handoffs).toEqual([]); }); + + test("abort then reset during certifyRange does not record a phantom handoff", async () => { + const notices: string[] = []; + const folds: { stub: boolean }[] = []; + const summarize = createModelSummarizer({ + getSource: () => + ({ + id: "test", + provider: "openai", + model: "test-model", + baseURL: "http://localhost:1", + credentialId: "test", + }) as never, + complete: async () => { + throw new Error("model unreachable"); + }, + }); + const dir = await mkdtemp(join(tmpdir(), "compaction-gate-abort-reset-")); + const blobs = new Map(); + const archive = createCompactionArchive({ + sessionId: "sess-gate-abort-reset", + contextDir: dir, + writeBlob: async (key, bytes) => { + blobs.set(key, bytes); + }, + readBlob: async (key) => { + const bytes = blobs.get(key); + if (bytes === undefined) throw new Error(`missing ${key}`); + return bytes; + }, + }); + const now = Date.now(); + const many = Array.from({ length: 8 }, (_, i) => ({ + role: (i % 2 === 0 ? "user" : "assistant") as "user" | "assistant", + content: [{ type: "text" as const, text: `t${i}` }], + timestamp: now, + })); + for (const turn of many) { + const block = turn.content[0]; + if (block?.type !== "text") continue; + await archive.recordAuthorizedPayload({ + kind: turn.role === "assistant" ? "assistant_text" : "user_message", + payload: block.text, + }); + } + const lifecycle = createCompactionLifecycle(); + const getSignal = () => lifecycle.getSignal(); + const isAborted = () => lifecycle.getSignal().aborted; + let releaseGated: () => void = () => undefined; + const gatedFinished = new Promise((resolve) => { + releaseGated = resolve; + }); + const abortingArchive = { + ...archive, + certifyRange: async (ids: readonly string[]) => { + const certificate = await archive.certifyRange(ids); + lifecycle.abortCompaction("operator interrupt"); + // onBuilt reset() mints a fresh controller; the in-flight apply + // must still treat this compact as aborted. + lifecycle.reset(); + return certificate; + }, + }; + const wrapped = lifecycle.wrapCompactor( + createSessionPruningCompactor({ + summarize, + onFolded: (info) => folds.push(info), + onFailure: (text) => notices.push(text), + getSignal, + isAborted, + compactionShape: { tailBudgetTokens: 1 }, + wrapPruning: (pruning) => { + const gated = wrapCompactorWithCompletenessGate( + pruning, + abortingArchive, + { getSignal, isAborted }, + ); + return { + name: gated.name, + version: gated.version, + apply: async (turns, ctx) => { + try { + return await gated.apply(turns, ctx); + } finally { + releaseGated(); + } + }, + }; + }, + }), + ); + const result = await wrapped.apply(many as never, { + state: {} as never, + trigger: "test", + }); + await gatedFinished; + expect(result.output).toBe(many); + expect(result.record.reason).toBe(COMPACTION_ABORTED_REASON); + expect(folds).toEqual([]); + expect(notices).toEqual([]); + const handoffs = (await archive.listOccurrences()).filter( + (occurrence) => occurrence.provenance === "compaction-handoff", + ); + expect(handoffs).toEqual([]); + }); }); describe("buildCompactionContinuationMessage", () => { diff --git a/src/session/runtime-assembly.ts b/src/session/runtime-assembly.ts index 51e937b86..d0a1e5167 100644 --- a/src/session/runtime-assembly.ts +++ b/src/session/runtime-assembly.ts @@ -441,6 +441,12 @@ export interface SessionPruningCompactorArgs { * no onFolded side effects for work that never landed. */ isAborted?: () => boolean; + /** + * Captured at apply start. Prefer this over a live `isAborted` that + * re-reads getSignal(): onBuilt reset() replaces the controller, and a + * discarded in-flight fold must still look aborted. + */ + getSignal?: () => AbortSignal; /** * Wraps the inner pruning apply before fold-commit side effects. * assembleChatAgent passes the completeness gate here so a discarded @@ -513,15 +519,15 @@ export function createSessionPruningCompactor( pendingStubNotice = undefined; const turnsBefore = turns.length; const startedAt = Date.now(); + const abortSignal = args.getSignal?.(); + const isAborted = () => + abortSignal?.aborted === true || args.isAborted?.() === true; const result = await compactor.apply(turns, ctx); // A discarded compact reports nothing. When the lifecycle abort wins // the outer race, this inner run may still complete with a genuine // fold — but the reactor threw that output away, so emitting telemetry // or onFolded would describe work that never landed (phantom fold). - if ( - args.isAborted?.() === true || - result.record.reason === COMPACTION_ABORTED_REASON - ) { + if (isAborted() || result.record.reason === COMPACTION_ABORTED_REASON) { return result; } // summarizedTurnCount is only set on the branch that actually folded diff --git a/src/tui/runner/session.ts b/src/tui/runner/session.ts index 1b3d21c67..bd01a6219 100644 --- a/src/tui/runner/session.ts +++ b/src/tui/runner/session.ts @@ -697,7 +697,7 @@ export async function assembleTUISession( ? state.liveDefaultSource : state.liveSource.id, anthropicCachePrompt: () => config.anthropicCachePrompt, - isCompactionAborted: () => compactionLifecycle.getSignal().aborted, + getCompactionAbortSignal: () => compactionLifecycle.getSignal(), getCompactor: (wrapPruning) => compactionLifecycle.wrapCompactor( createSessionPruningCompactor({ @@ -710,8 +710,9 @@ export async function assembleTUISession( ), // The outer abort race discards this run's output — a fold that // still completes underneath must not report telemetry or side - // effects for work that never landed. - isAborted: () => compactionLifecycle.getSignal().aborted, + // effects for work that never landed. Capture the signal at apply + // start: onBuilt reset() replaces the live controller. + getSignal: () => compactionLifecycle.getSignal(), // Stub notice waits until the fold commits (verify abort and // completeness-gate discard stay silent). onFailure: (text) => state.systemNotice?.(text),