diff --git a/e2e/integration-harness.ts b/e2e/integration-harness.ts index 664cf4b11..34d883d39 100644 --- a/e2e/integration-harness.ts +++ b/e2e/integration-harness.ts @@ -348,26 +348,23 @@ export async function openIntegrationSession( ...(opts.compactionCompletion !== undefined && primaryArchive !== undefined ? { compactors: { - "pruning-compactor": wrapCompactorWithCompletenessGate( - createSessionPruningCompactor({ - // Absent compactionShape falls back to the shared production - // default inside createSessionPruningCompactor. - ...(opts.compactionShape !== undefined - ? { compactionShape: opts.compactionShape } - : {}), - summarize: createModelSummarizer({ - getSource: () => INTEGRATION_SOURCE, - deps: harness.deps, - complete: opts.compactionCompletion, - getArchive: () => evidenceArchiveHolder.current, - }), - readPriorHandoff: () => - tryReadPriorHandoffFile((key) => - storageForAgent.readBlob(key), - ), + "pruning-compactor": createSessionPruningCompactor({ + // Absent compactionShape falls back to the shared production + // default inside createSessionPruningCompactor. + ...(opts.compactionShape !== undefined + ? { compactionShape: opts.compactionShape } + : {}), + summarize: createModelSummarizer({ + getSource: () => INTEGRATION_SOURCE, + deps: harness.deps, + complete: opts.compactionCompletion, + getArchive: () => evidenceArchiveHolder.current, }), - primaryArchive, - ), + readPriorHandoff: () => + tryReadPriorHandoffFile((key) => storageForAgent.readBlob(key)), + wrapPruning: (pruning) => + wrapCompactorWithCompletenessGate(pruning, primaryArchive), + }), }, } : {}), diff --git a/src/exec/runner.ts b/src/exec/runner.ts index ec524d443..d4854dbd7 100644 --- a/src/exec/runner.ts +++ b/src/exec/runner.ts @@ -931,7 +931,7 @@ export async function runExec(config: Config): Promise { getDefaultSource: () => liveDefaultSource.length > 0 ? liveDefaultSource : liveSource.id, anthropicCachePrompt: () => config.anthropicCachePrompt, - getCompactor: () => + getCompactor: (wrapPruning) => createSessionPruningCompactor({ summarize: summarizeForCompaction, summaryContext: () => { @@ -965,6 +965,7 @@ export async function runExec(config: Config): Promise { }, }); }, + ...(wrapPruning !== undefined ? { wrapPruning } : {}), }), getCacheWriteSeed: () => resumeCacheWriteSeed({ diff --git a/src/session/assemble-runtime.ts b/src/session/assemble-runtime.ts index 1e5c4007c..abb223ed3 100644 --- a/src/session/assemble-runtime.ts +++ b/src/session/assemble-runtime.ts @@ -550,8 +550,12 @@ export interface ChatAgentWiring { inferenceDeps: Awaited>; getSources: () => InferenceSource[]; getDefaultSource: () => string; - /** Read at each build so a compaction-mode toggle is visible on rebuild. */ - getCompactor: () => Compactor; + /** + * Read at each build so a compaction-mode toggle is visible on rebuild. + * assemble passes the completeness gate as `wrapPruning` so fold-commit + * side effects (stub notice, onFolded prune) run only after a fold lands. + */ + getCompactor: (wrapPruning?: (pruning: Compactor) => Compactor) => Compactor; /** Experimental Anthropic prompt shrink. Default off when omitted. */ anthropicCachePrompt?: () => boolean; /** @@ -778,13 +782,12 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent { defaultId: `${ID_PREFIX}/chat`, }), compactors: { - "pruning-compactor": + "pruning-compactor": wiring.getCompactor( primaryArchive === undefined - ? wiring.getCompactor() - : wrapCompactorWithCompletenessGate( - wiring.getCompactor(), - primaryArchive, - ), + ? undefined + : (pruning) => + wrapCompactorWithCompletenessGate(pruning, primaryArchive), + ), }, }); const seed = wiring.getCacheWriteSeed?.(); diff --git a/src/session/runtime-assembly.test.ts b/src/session/runtime-assembly.test.ts index 459fe096d..3e6e0c682 100644 --- a/src/session/runtime-assembly.test.ts +++ b/src/session/runtime-assembly.test.ts @@ -11,7 +11,7 @@ import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { getLogger } from "@intx/log"; -import type { ToolCall } from "@intx/types/runtime"; +import type { ConversationTurn, ToolCall } from "@intx/types/runtime"; import { LOG_NAMESPACE_ROOT } from "../branding.js"; import * as permissionStore from "../permission/store.js"; @@ -33,6 +33,10 @@ import type { SubAgentSourcesConfig } from "./runtime-assembly.js"; import type { Settings } from "../config/settings.js"; import type { Telemetry } from "../telemetry/index.js"; import { createModelSummarizer } from "./summarizer.js"; +import { + createCompactionArchive, + wrapCompactorWithCompletenessGate, +} from "./compaction-archive.js"; import { generateSessionId, initSessionDir, sessionDir } from "./index.js"; import type { PluginModule } from "../plugins/loader.js"; @@ -732,6 +736,147 @@ describe("createSessionPruningCompactor stub fallback", () => { expect(folds).toEqual([]); expect(notices).toEqual([]); }); + + test("completeness-gate discard after a stub fold fires no notice or onFolded", async () => { + const notices: string[] = []; + const folds: { stub: boolean }[] = []; + const captured: { event: string }[] = []; + const telemetry: Telemetry = { + enabled: true, + installationId: "test", + capture: (event) => { + captured.push({ event }); + }, + captureIntentional: () => false, + flush: async () => undefined, + discard: () => undefined, + }; + 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-")); + const blobs = new Map(); + const archive = createCompactionArchive({ + sessionId: "sess-gate-stub", + 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: ConversationTurn[] = Array.from({ length: 8 }, (_, i) => ({ + role: i % 2 === 0 ? "user" : "assistant", + content: [{ type: "text", text: `t${i}` }], + timestamp: now, + })); + const result = await createSessionPruningCompactor({ + summarize, + telemetry, + onFolded: (info) => folds.push(info), + onFailure: (text) => notices.push(text), + compactionShape: { tailBudgetTokens: 1 }, + wrapPruning: (pruning) => + wrapCompactorWithCompletenessGate(pruning, archive), + }).apply(many, { state: {} as never, trigger: "test" }); + expect(result.output).toBe(many); + expect(result.record.reason).toBe("incomplete-evidence-archive"); + expect(result.record.decisions.summarizedTurnCount).toBeUndefined(); + expect(folds).toEqual([]); + expect(notices).toEqual([]); + expect(captured).toEqual([]); + }); + + test("a committed stub fold still notices and prunes", async () => { + const notices: string[] = []; + const folds: { turnsBefore: number; turnsAfter: number; stub: boolean }[] = + []; + const captured: { event: string }[] = []; + const telemetry: Telemetry = { + enabled: true, + installationId: "test", + capture: (event) => { + captured.push({ event }); + }, + captureIntentional: () => false, + flush: async () => undefined, + discard: () => undefined, + }; + 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-ok-")); + const blobs = new Map(); + const archive = createCompactionArchive({ + sessionId: "sess-gate-stub-ok", + 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 result = await createSessionPruningCompactor({ + summarize, + telemetry, + onFolded: (info) => folds.push(info), + onFailure: (text) => notices.push(text), + compactionShape: { tailBudgetTokens: 1 }, + wrapPruning: (pruning) => + wrapCompactorWithCompletenessGate(pruning, archive), + }).apply(many as never, { state: {} as never, trigger: "test" }); + expect(result.record.decisions.summarizeFailed).toBe(1); + expect(result.record.reason).toContain("statistics-only stub"); + expect(result.output).not.toBe(many); + expect(folds).toEqual([ + { turnsBefore: 8, turnsAfter: result.output.length, stub: true }, + ]); + expect(captured).toEqual([]); + expect(notices).toHaveLength(1); + expect(notices[0]).toContain("statistics-only stub"); + expect(notices[0]).toContain("failed"); + }); }); describe("buildCompactionContinuationMessage", () => { diff --git a/src/session/runtime-assembly.ts b/src/session/runtime-assembly.ts index 0269ad8f1..58ab8e139 100644 --- a/src/session/runtime-assembly.ts +++ b/src/session/runtime-assembly.ts @@ -430,7 +430,8 @@ export interface SessionPruningCompactorArgs { }) => void; /** * Operator-visible notice for a statistics-only stub that actually replaced - * turns. Verify abort (keeping prior context) does not fire this. + * turns. Verify abort and completeness-gate discard (keeping prior context) + * do not fire this. */ onFailure?: (text: string) => void; /** @@ -440,6 +441,13 @@ export interface SessionPruningCompactorArgs { * no onFolded side effects for work that never landed. */ isAborted?: () => boolean; + /** + * Wraps the inner pruning apply before fold-commit side effects. + * assembleChatAgent passes the completeness gate here so a discarded + * fold never notices or prunes. + */ + wrapPruning?: (pruning: Compactor) => Compactor; + /** * CL-9489 budgeted-tail shape override. Absent means the shared production * default (DEFAULT_TAIL_COMPACTION_SHAPE); tests pin a tiny budget so small @@ -481,7 +489,7 @@ export function createSessionPruningCompactor( throw error; } }; - const compactor = createPruningCompactor({ + const pruning = createPruningCompactor({ summaryMaxChars: SESSION_COMPACTOR_SUMMARY_MAX_CHARS, // CL-9489 budgeted-tail shape: explicit defaults (same object the record // carries under parameters.compactionShape). Zero recent turns stay whole @@ -496,6 +504,8 @@ export function createSessionPruningCompactor( ? { readPriorHandoff: args.readPriorHandoff } : {}), }); + const compactor = args.wrapPruning?.(pruning) ?? pruning; + const telemetry = args.telemetry ?? NOOP_TELEMETRY; return { ...compactor, diff --git a/src/tui/runner/session.ts b/src/tui/runner/session.ts index c3d1b34e4..7683f59b4 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, - getCompactor: () => + getCompactor: (wrapPruning) => compactionLifecycle.wrapCompactor( createSessionPruningCompactor({ summarize: compactionSummarize, @@ -711,7 +711,8 @@ export async function assembleTUISession( // still completes underneath must not report telemetry or side // effects for work that never landed. isAborted: () => compactionLifecycle.getSignal().aborted, - // Stub notice waits until the fold commits (verify abort stays silent). + // Stub notice waits until the fold commits (verify abort and + // completeness-gate discard stay silent). onFailure: (text) => state.systemNotice?.(text), // Main-session folds only — exec runner and subagents stay silent. onFolded: (info) => { @@ -731,6 +732,7 @@ export async function assembleTUISession( }); if (!info.stub) emitter.emit("compaction", info); }, + ...(wrapPruning !== undefined ? { wrapPruning } : {}), }), ), getCacheWriteSeed: () =>