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
14 changes: 13 additions & 1 deletion src/session/assemble-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -556,6 +556,12 @@ export interface ChatAgentWiring {
* side effects (stub notice, onFolded prune) run only after a fold lands.
*/
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.
*/
isCompactionAborted?: () => boolean;
/** Experimental Anthropic prompt shrink. Default off when omitted. */
anthropicCachePrompt?: () => boolean;
/**
Expand Down Expand Up @@ -786,7 +792,13 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent {
primaryArchive === undefined
? undefined
: (pruning) =>
wrapCompactorWithCompletenessGate(pruning, primaryArchive),
wrapCompactorWithCompletenessGate(
pruning,
primaryArchive,
wiring.isCompactionAborted === undefined
? undefined
: { isAborted: wiring.isCompactionAborted },
),
),
},
});
Expand Down
18 changes: 16 additions & 2 deletions src/session/compaction-archive.ts
Original file line number Diff line number Diff line change
Expand Up @@ -922,9 +922,12 @@ async function recordFreshHandoffOutput(
archive: CompactionArchive,
input: readonly ConversationTurn[],
output: readonly ConversationTurn[],
isAborted?: () => boolean,
): Promise<void> {
if (isAborted?.() === true) return;
const fresh = uncoveredContentUnits(output, input);
for (const unit of fresh) {
if (isAborted?.() === true) return;
if (unit.kind !== "text" || unit.role !== "user") continue;
const text = unit.text ?? "";
if (text.length === 0) continue;
Expand All @@ -946,11 +949,14 @@ async function recordFreshHandoffOutput(
* Synthetic handoff spines (and pre-format fat summaries) are adopted into
* the archive as user_message so a later fold may change the live spine.
* After the fold certifies, the new spine is recorded so the next fold can
* drop it even when the summarizer does not echo it verbatim.
* 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.
*/
export function wrapCompactorWithCompletenessGate(
inner: Compactor,
archive: CompactionArchive,
opts?: { isAborted?: () => boolean },
): Compactor {
return {
name: inner.name,
Expand Down Expand Up @@ -982,7 +988,15 @@ export function wrapCompactorWithCompletenessGate(
if (certificate.status !== "complete") {
return incompleteIdentity(inner, turns);
}
await recordFreshHandoffOutput(archive, turns, proposed.output);
if (opts?.isAborted?.() === true) {
return proposed;
}
await recordFreshHandoffOutput(
archive,
turns,
proposed.output,
opts?.isAborted,
);
return proposed;
},
};
Expand Down
108 changes: 108 additions & 0 deletions src/session/runtime-assembly.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,10 @@ import {
createCompactionArchive,
wrapCompactorWithCompletenessGate,
} from "./compaction-archive.js";
import {
COMPACTION_ABORTED_REASON,
createCompactionLifecycle,
} from "./compaction-lifecycle.js";
import { generateSessionId, initSessionDir, sessionDir } from "./index.js";
import type { PluginModule } from "../plugins/loader.js";

Expand Down Expand Up @@ -876,6 +880,110 @@ describe("createSessionPruningCompactor stub fallback", () => {
expect(notices).toHaveLength(1);
expect(notices[0]).toContain("statistics-only stub");
expect(notices[0]).toContain("failed");
const committedHandoffs = (await archive.listOccurrences()).filter(
(occurrence) => occurrence.provenance === "compaction-handoff",
);
expect(committedHandoffs.length).toBeGreaterThan(0);
});

test("TUI abort after a gate-committed stub 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-stub-abort-"));
const blobs = new Map<string, Uint8Array>();
const archive = createCompactionArchive({
sessionId: "sess-gate-stub-abort",
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 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");
return certificate;
},
};
const wrapped = lifecycle.wrapCompactor(
createSessionPruningCompactor({
summarize,
onFolded: (info) => folds.push(info),
onFailure: (text) => notices.push(text),
isAborted,
compactionShape: { tailBudgetTokens: 1 },
wrapPruning: (pruning) => {
const gated = wrapCompactorWithCompletenessGate(
pruning,
abortingArchive,
{ 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([]);
});
});

Expand Down
2 changes: 1 addition & 1 deletion src/session/runtime-assembly.ts
Original file line number Diff line number Diff line change
Expand Up @@ -444,7 +444,7 @@ export interface SessionPruningCompactorArgs {
/**
* Wraps the inner pruning apply before fold-commit side effects.
* assembleChatAgent passes the completeness gate here so a discarded
* fold never notices or prunes.
* fold never notices, prunes, or records a phantom compaction-handoff.
*/
wrapPruning?: (pruning: Compactor) => Compactor;

Expand Down
1 change: 1 addition & 0 deletions src/tui/runner/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -697,6 +697,7 @@ export async function assembleTUISession(
? state.liveDefaultSource
: state.liveSource.id,
anthropicCachePrompt: () => config.anthropicCachePrompt,
isCompactionAborted: () => compactionLifecycle.getSignal().aborted,
getCompactor: (wrapPruning) =>
compactionLifecycle.wrapCompactor(
createSessionPruningCompactor({
Expand Down
Loading