Skip to content
Merged
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
12 changes: 6 additions & 6 deletions src/session/assemble-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
/**
Expand Down Expand Up @@ -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 },
),
),
},
Expand Down
11 changes: 8 additions & 3 deletions src/session/compaction-archive.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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;
},
Expand Down
105 changes: 105 additions & 0 deletions src/session/runtime-assembly.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, Uint8Array>();
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<void>((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", () => {
Expand Down
14 changes: 10 additions & 4 deletions src/session/runtime-assembly.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
7 changes: 4 additions & 3 deletions src/tui/runner/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand All @@ -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),
Expand Down
Loading