From 4c7f50ed4bac8e87b35a8119486404b4db28ba10 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Tue, 29 Sep 2026 19:50:48 -0700 Subject: [PATCH] fix(resume): cap transcript hydration in rows after fold A 600-block mapping window under-fills the retained row tail once call+result pairs fold, so retention never evicts and the dropped-rows notice stays unpainted. Resume now reads a turn window sized for that fold and lets the host cap in rows. --- src/session/optimized-context-store.test.ts | 22 ++- src/session/optimized-context-store.ts | 23 ++- src/tui/history-hydrate.test.ts | 64 +++++++++ src/tui/history-hydrate.ts | 25 ++++ src/tui/product-host.test.ts | 149 ++++++++++++++++++++ src/tui/product-host.ts | 10 +- src/tui/runner/wiring.ts | 28 ++-- src/tui/shell/chrome.ts | 29 ++++ src/tui/turns-to-blocks.test.ts | 8 +- src/tui/turns-to-blocks.ts | 39 ++--- 10 files changed, 341 insertions(+), 56 deletions(-) diff --git a/src/session/optimized-context-store.test.ts b/src/session/optimized-context-store.test.ts index 5318bb981..78adc7ef2 100644 --- a/src/session/optimized-context-store.test.ts +++ b/src/session/optimized-context-store.test.ts @@ -369,7 +369,8 @@ describe("loadRecentTurns", () => { ); const loaded = await loadRecentTurns(dir, 2); - expect(turnTexts(loaded)).toEqual(["d", "e"]); + expect(turnTexts(loaded.turns)).toEqual(["d", "e"]); + expect(loaded.truncated).toBe(true); }); test("walks back into older segments when the window exceeds the newest one", async () => { @@ -388,12 +389,27 @@ describe("loadRecentTurns", () => { // window of 4, so the walk continues into segment 0 (2 turns) whole — // reads are segment-granular, not turn-exact. const loaded = await loadRecentTurns(dir, 4); - expect(turnTexts(loaded)).toEqual(["a", "b", "c", "d", "e"]); + expect(turnTexts(loaded.turns)).toEqual(["a", "b", "c", "d", "e"]); + expect(loaded.truncated).toBe(false); }); test("returns an empty list when no segments exist", async () => { const dir = tempDir(); - expect(await loadRecentTurns(dir, 10)).toEqual([]); + expect(await loadRecentTurns(dir, 10)).toEqual({ + turns: [], + truncated: false, + }); + }); + + test("a single segment that fills the window is not truncated", async () => { + const dir = tempDir(); + fs.writeFileSync( + path.join(dir, TURNS_FILE), + jsonl([turn("a"), turn("b"), turn("c")]), + ); + const loaded = await loadRecentTurns(dir, 2); + expect(turnTexts(loaded.turns)).toEqual(["a", "b", "c"]); + expect(loaded.truncated).toBe(false); }); test("reactor load skips mid-file garbage in an extra segment", async () => { diff --git a/src/session/optimized-context-store.ts b/src/session/optimized-context-store.ts index 45439908c..a9d0092ab 100644 --- a/src/session/optimized-context-store.ts +++ b/src/session/optimized-context-store.ts @@ -287,6 +287,17 @@ async function unlinkExtraSegmentsFrom( return removed; } +/** + * Display-only tail of the turn history. Older segments are unread when + * `truncated` is set — callers that paint a bounded transcript can still + * notice that more history exists even if the loaded window hydrates to + * exactly the retention cap. + */ +export type RecentTurns = { + readonly turns: ConversationTurn[]; + readonly truncated: boolean; +}; + /** * Read only the tail of the turn history needed to satisfy `minTurns`, walking * segments from newest to oldest and stopping as soon as enough turns have @@ -302,12 +313,13 @@ async function unlinkExtraSegmentsFrom( export async function loadRecentTurns( dir: string, minTurns: number, -): Promise { +): Promise { const segments = await listSegmentFiles(dir, TURNS_FILE); - if (segments.length === 0) return []; + if (segments.length === 0) return { turns: [], truncated: false }; const collectedNewestFirst: ConversationTurn[][] = []; let total = 0; + let truncated = false; for (let i = segments.length - 1; i >= 0; i--) { const name = segments[i]; if (name === undefined) continue; @@ -325,7 +337,10 @@ export async function loadRecentTurns( ); collectedNewestFirst.push(turns); total += turns.length; - if (total >= minTurns) break; + if (total >= minTurns) { + truncated = i > 0; + break; + } } const turns: ConversationTurn[] = []; @@ -333,7 +348,7 @@ export async function loadRecentTurns( const chunk = collectedNewestFirst[i]; if (chunk !== undefined) turns.push(...chunk); } - return turns; + return { turns, truncated }; } /** diff --git a/src/tui/history-hydrate.test.ts b/src/tui/history-hydrate.test.ts index c0ea944f1..fc1987f11 100644 --- a/src/tui/history-hydrate.test.ts +++ b/src/tui/history-hydrate.test.ts @@ -7,6 +7,7 @@ import { EMPTY_VIEW_DETAIL, hydrateHistoryRows, MISSING_ERROR_DETAIL, + parseHistoryHydratePayload, rowFromHistoryBlock, type HistoryBlock, } from "./history-hydrate.js"; @@ -410,4 +411,67 @@ describe("resume pipeline end to end (turns-to-blocks into hydrate)", () => { const rows = hydrateHistoryRows(JSON.parse(JSON.stringify(blocks))); expect(rows).toEqual([]); }); + + test("spawn_agent pairs fold to one row each, including a pair at the old 600-block splice", () => { + const turns: ConversationTurn[] = []; + for (let i = 0; i < 400; i++) { + turns.push({ + role: "assistant", + model: "test", + timestamp: i, + content: [ + { + type: "tool_call", + id: `c${i}`, + name: "spawn_agent", + arguments: { description: `job-${i}` }, + }, + ], + } as unknown as ConversationTurn); + turns.push({ + role: "assistant", + model: "test", + timestamp: i, + content: [ + { + type: "tool_result", + callId: `c${i}`, + content: `done c${i}`, + isError: false, + }, + ], + } as unknown as ConversationTurn); + } + turns.push({ + role: "user", + content: [{ type: "text", text: "trailing" }], + timestamp: 400, + } as unknown as ConversationTurn); + const blocks = turnsToContentBlocks(turns); + expect(blocks.length).toBe(801); + const rows = hydrateHistoryRows(blocks); + expect(rows.length).toBe(401); + expect(rows[0]?.pending).not.toBe(true); + expect(rows[0]?.text).toBe("done c0"); + expect(rows[399]?.text).toBe("done c399"); + expect(rows[400]).toEqual({ role: "user", text: "trailing" }); + }); +}); + +describe("parseHistoryHydratePayload", () => { + test("keeps a raw block array and treats a wrapper as truncated only when flagged", () => { + const blocks = [{ type: "user", content: "hi" }]; + expect(parseHistoryHydratePayload(blocks)).toEqual({ + blocks, + truncated: false, + }); + expect(parseHistoryHydratePayload({ blocks, truncated: true })).toEqual({ + blocks, + truncated: true, + }); + expect(parseHistoryHydratePayload({ blocks })).toEqual({ + blocks, + truncated: false, + }); + }); }); diff --git a/src/tui/history-hydrate.ts b/src/tui/history-hydrate.ts index 769e51ea9..30c38bd22 100644 --- a/src/tui/history-hydrate.ts +++ b/src/tui/history-hydrate.ts @@ -203,6 +203,31 @@ export function hydrateHistoryRows(blocks: unknown): StreamRow[] { return rows; } +/** + * Resume hydrate events are either a raw block array (tests, load errors) or + * `{ blocks, truncated }` when loadRecentTurns left older segments unread. + */ +export function parseHistoryHydratePayload(payload: unknown): { + blocks: unknown; + truncated: boolean; +} { + if ( + Array.isArray(payload) || + payload === null || + typeof payload !== "object" + ) { + return { blocks: payload, truncated: false }; + } + const record = payload as Record; + if (!("blocks" in record)) { + return { blocks: payload, truncated: false }; + } + return { + blocks: record.blocks, + truncated: record.truncated === true, + }; +} + /** Argument payload of a tool_call block, wherever the block carries it. */ function callArguments(block: HistoryBlock): string | undefined { if (block.content !== undefined) return block.content; diff --git a/src/tui/product-host.test.ts b/src/tui/product-host.test.ts index bf855b2a6..aef7dd4df 100644 --- a/src/tui/product-host.test.ts +++ b/src/tui/product-host.test.ts @@ -5,6 +5,7 @@ import { EventEmitter } from "node:events"; import { describe, expect, test } from "bun:test"; import type { KeyEvent } from "@opentui/core"; +import type { ConversationTurn } from "@intx/types/runtime"; import { createHarness, type Harness } from "./harness.js"; import { acceptOverlaySelection } from "./shell/overlay-host.js"; import { @@ -16,6 +17,7 @@ import { mountProductHost, type ProductHostConfig } from "./product-host.js"; import { buildModelsFirstCatalog, modelOptionId } from "./model-catalog.js"; import { hydrateHistoryRows } from "./history-hydrate.js"; import { MAX_RETAINED_STREAM_ROWS } from "./long-log.js"; +import { turnsToContentBlocks } from "./turns-to-blocks.js"; import { enterSubagentObserve } from "./shell/observe.js"; import { transcriptMarker } from "./shell/transcript.js"; @@ -97,6 +99,46 @@ function composedKey(glyph: string): KeyEvent { } as KeyEvent; } +function userTextTurn(text: string): ConversationTurn { + return { + role: "user", + content: [{ type: "text", text }], + timestamp: 0, + } as unknown as ConversationTurn; +} + +function spawnCallTurn(id: string): ConversationTurn { + return { + role: "assistant", + model: "test", + timestamp: 0, + content: [ + { + type: "tool_call", + id, + name: "spawn_agent", + arguments: { description: `job-${id}` }, + }, + ], + } as unknown as ConversationTurn; +} + +function spawnResultTurn(id: string): ConversationTurn { + return { + role: "assistant", + model: "test", + timestamp: 0, + content: [ + { + type: "tool_result", + callId: id, + content: `done ${id}`, + isError: false, + }, + ], + } as unknown as ConversationTurn; +} + describe("mountProductHost", () => { test("stream events emitted on the event emitter paint rows into the shell", async () => { const { host, emitter } = await mountHeadless(); @@ -260,6 +302,113 @@ describe("mountProductHost", () => { } }); + test("resume pipeline of long text history paints the eviction marker", async () => { + const { host, emitter } = await mountHeadless(); + try { + const total = MAX_RETAINED_STREAM_ROWS + 200; + const turns = Array.from({ length: total }, (_, i) => + userTextTurn(`row-${i}`), + ); + const blocks = turnsToContentBlocks(turns); + emitter.emit("history.hydrate", blocks); + expect(host.shell.streamLog.length).toBe(MAX_RETAINED_STREAM_ROWS); + expect(host.shell.streamLogBase).toBe(200); + expect(host.shell.streamLog[0]).toEqual({ + role: "user", + text: "row-200", + }); + expect(transcriptMarker(host.shell)).toBeDefined(); + } finally { + host.dispose(); + } + }); + + test("resume pipeline of tool pairs fills the retained row cap", async () => { + const { host, emitter } = await mountHeadless(); + try { + const pairCount = MAX_RETAINED_STREAM_ROWS + 200; + const turns: ConversationTurn[] = []; + for (let i = 0; i < pairCount; i++) { + turns.push(spawnCallTurn(`p${i}`), spawnResultTurn(`p${i}`)); + } + const blocks = turnsToContentBlocks(turns); + emitter.emit("history.hydrate", blocks); + expect(host.shell.streamLog.length).toBe(MAX_RETAINED_STREAM_ROWS); + expect(host.shell.streamLogBase).toBe(200); + const first = host.shell.streamLog[0]; + expect(first?.pending).not.toBe(true); + expect(first?.text).toBe("done p200"); + const last = host.shell.streamLog[host.shell.streamLog.length - 1]; + expect(last?.pending).not.toBe(true); + expect(last?.text).toBe(`done p${pairCount - 1}`); + expect(transcriptMarker(host.shell)).toBeDefined(); + } finally { + host.dispose(); + } + }); + + test("resume pipeline keeps a pair merged when it straddles the old block splice", async () => { + const { host, emitter } = await mountHeadless(); + try { + const turns: ConversationTurn[] = []; + for (let i = 0; i < 400; i++) { + turns.push(spawnCallTurn(`c${i}`), spawnResultTurn(`c${i}`)); + } + turns.push(userTextTurn("trailing")); + const blocks = turnsToContentBlocks(turns); + emitter.emit("history.hydrate", { blocks, truncated: false }); + expect(host.shell.streamLog.length).toBe(401); + expect(host.shell.streamLog[0]?.pending).not.toBe(true); + expect(host.shell.streamLog[0]?.text).toBe("done c0"); + expect(host.shell.streamLog[host.shell.streamLog.length - 1]).toEqual({ + role: "user", + text: "trailing", + }); + expect(host.shell.streamLogBase).toBe(0); + } finally { + host.dispose(); + } + }); + + test("resume pipeline keeps a pair merged across the retained row cap", async () => { + const { host, emitter } = await mountHeadless(); + try { + const turns = [ + ...Array.from({ length: 605 }, (_, i) => userTextTurn(`row-${i}`)), + spawnCallTurn("cut-1"), + spawnResultTurn("cut-1"), + ]; + const blocks = turnsToContentBlocks(turns); + emitter.emit("history.hydrate", blocks); + const allRows = hydrateHistoryRows(blocks); + const expected = allRows.slice(-MAX_RETAINED_STREAM_ROWS); + expect(host.shell.streamLog).toEqual(expected); + const last = host.shell.streamLog[host.shell.streamLog.length - 1]; + expect(last?.pending).not.toBe(true); + expect(last?.text).toBe("done cut-1"); + expect(host.shell.streamLogBase).toBeGreaterThan(0); + expect(transcriptMarker(host.shell)).toBeDefined(); + } finally { + host.dispose(); + } + }); + + test("truncated load of an exact-cap tail still paints the eviction marker", async () => { + const { host, emitter } = await mountHeadless(); + try { + const turns = Array.from({ length: MAX_RETAINED_STREAM_ROWS }, (_, i) => + userTextTurn(`kept-${i}`), + ); + const blocks = turnsToContentBlocks(turns); + emitter.emit("history.hydrate", { blocks, truncated: true }); + expect(host.shell.streamLog.length).toBe(MAX_RETAINED_STREAM_ROWS); + expect(host.shell.streamLogBase).toBeGreaterThan(0); + expect(transcriptMarker(host.shell)).toBeDefined(); + } finally { + host.dispose(); + } + }); + test("session.title updates the shell header", async () => { const { host, emitter } = await mountHeadless(); try { diff --git a/src/tui/product-host.ts b/src/tui/product-host.ts index f29a8e312..927bc0d05 100644 --- a/src/tui/product-host.ts +++ b/src/tui/product-host.ts @@ -43,6 +43,7 @@ import { appendObserveStreamRow, appendStreamRow, clearTranscript, + noteUnloadedHistory, paintChrome, setChromeZones, setHeader, @@ -64,7 +65,10 @@ import { } from "./shell/palette.js"; import { surfaceSystemNotice } from "./shell/prompt.js"; import type { DeliverySettle, QueueKind } from "./delivery-queue.js"; -import { hydrateHistoryRows } from "./history-hydrate.js"; +import { + hydrateHistoryRows, + parseHistoryHydratePayload, +} from "./history-hydrate.js"; import type { StreamRow } from "./stream.js"; import type { PendingImageAttachment } from "./image-attachments.js"; @@ -539,11 +543,13 @@ export async function mountProductHost( throw err; } - function onHistory(blocks: unknown): void { + function onHistory(payload: unknown): void { if (disposed) return; + const { blocks, truncated } = parseHistoryHydratePayload(payload); for (const row of hydrateHistoryRows(blocks)) { appendStreamRow(shell, row); } + if (truncated) noteUnloadedHistory(shell); } function onTitle(title: unknown): void { diff --git a/src/tui/runner/wiring.ts b/src/tui/runner/wiring.ts index 81adc5a34..16835b665 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -44,7 +44,7 @@ import { import { isCodexProviderName } from "../../config/codex-providers.js"; import { RUNTIME_FLASH_MS } from "../runtime-notices.js"; import { - RESUME_TRANSCRIPT_BLOCK_LIMIT, + RESUME_TRANSCRIPT_TURN_LIMIT, turnsToContentBlocks, } from "../turns-to-blocks.js"; import { setPluginNeedsAttention, setStatusFlash } from "../shell/chrome.js"; @@ -501,16 +501,15 @@ export function wirePostStartup( // Hydrate a resumed session's transcript after first paint. Reading history and // mapping it to content blocks is pure I/O with no bearing on the shell, so the // App renders empty immediately and fills in the past turns once they are ready. - // Only the retained transcript tail is read from disk — a long session's - // full history is not needed just to paint a display that itself caps how - // much it keeps. Agent conversation state still loads in full via - // ContextStore.load(); this path is display-only. - void loadRecentTurns(state.workdir, RESUME_TRANSCRIPT_BLOCK_LIMIT) - .then((turns) => { - const blocks = turnsToContentBlocks(turns, { - maxBlocks: RESUME_TRANSCRIPT_BLOCK_LIMIT, - }); - const tasks = hydrateTasksFromTurns(turns); + // The disk window is in turns, sized so a tool-pair-heavy tail can still fill + // the retained row cap after call+result fold to one row. Mapping does not + // splice content-blocks; the host caps in rows via retention eviction. Agent + // conversation state still loads in full via ContextStore.load(); this path + // is display-only. + void loadRecentTurns(state.workdir, RESUME_TRANSCRIPT_TURN_LIMIT) + .then((recent) => { + const blocks = turnsToContentBlocks(recent.turns); + const tasks = hydrateTasksFromTurns(recent.turns); // Restored tasks go to the panel only. They are live state, not something // that happened in the conversation, so putting them in scrollback as well // renders the same list twice on one screen. @@ -520,7 +519,12 @@ export function wirePostStartup( services.directorHolder.instance?.restoreTasks(tasks); services.emitter.emit("tasks", tasks); } - if (blocks.length > 0) services.emitter.emit("history.hydrate", blocks); + if (blocks.length > 0) { + services.emitter.emit("history.hydrate", { + blocks, + truncated: recent.truncated, + }); + } }) .catch((err: unknown) => { // Resume still works without painted history, but a silent empty diff --git a/src/tui/shell/chrome.ts b/src/tui/shell/chrome.ts index cb21cc24f..bbb3ee30e 100644 --- a/src/tui/shell/chrome.ts +++ b/src/tui/shell/chrome.ts @@ -1073,6 +1073,35 @@ export function appendStreamRow(shell: AppShell, row: StreamRow): void { paintAppendStreamRow(shell, row); } +/** + * Paint the dropped-rows notice when older history exists on disk but the + * loaded window hydrated to at most the retention cap, so trim never ran. + */ +export function noteUnloadedHistory(shell: AppShell): void { + if (shell.observe !== null && shell.parentStreamLog !== null) { + if ( + (shell.parentStreamLogBase ?? 0) === 0 && + shell.parentStreamLog.length > 0 + ) { + shell.parentStreamLogBase = 1; + } + return; + } + if (shell.streamLogBase > 0 || shell.streamLog.length === 0) return; + shell.streamLogBase = 1; + const marker = transcriptMarker(shell); + if (marker instanceof TextRenderable) { + marker.content = evictedRowsNotice(shell.streamLogBase); + return; + } + const node = new TextRenderable(shell.renderer as CliRenderer, { + content: evictedRowsNotice(shell.streamLogBase), + fg: UI.textDim, + }); + evictionMarkers.add(node); + shell.transcript.add(node, 1); +} + /** * Append a child stream row while observing a subagent. * Host-pushed live events (not only fixture seed lines). No-op when not observing. diff --git a/src/tui/turns-to-blocks.test.ts b/src/tui/turns-to-blocks.test.ts index 3fd5b0213..f8700aa05 100644 --- a/src/tui/turns-to-blocks.test.ts +++ b/src/tui/turns-to-blocks.test.ts @@ -2,7 +2,7 @@ import { describe, test, expect } from "bun:test"; import type { ConversationTurn } from "@intx/types/runtime"; import { buildMailboxMailPrompt } from "../subagent/mailbox-mail-drive.js"; import { - RESUME_TRANSCRIPT_BLOCK_LIMIT, + RESUME_TRANSCRIPT_TURN_LIMIT, turnsToContentBlocks, } from "./turns-to-blocks.js"; import { MAX_RETAINED_STREAM_ROWS } from "./long-log.js"; @@ -119,9 +119,9 @@ describe("turnsToContentBlocks marks occupancy wakes with system origin", () => }); }); -describe("RESUME_TRANSCRIPT_BLOCK_LIMIT", () => { - test("matches the retained stream tail, not a larger pre-slice", () => { - expect(RESUME_TRANSCRIPT_BLOCK_LIMIT).toBe(MAX_RETAINED_STREAM_ROWS); +describe("RESUME_TRANSCRIPT_TURN_LIMIT", () => { + test("is a turn window large enough to fill the retained row cap after pair fold", () => { + expect(RESUME_TRANSCRIPT_TURN_LIMIT).toBe(MAX_RETAINED_STREAM_ROWS * 2 + 1); }); }); diff --git a/src/tui/turns-to-blocks.ts b/src/tui/turns-to-blocks.ts index bd8b0c184..0583c1c0b 100644 --- a/src/tui/turns-to-blocks.ts +++ b/src/tui/turns-to-blocks.ts @@ -210,13 +210,12 @@ function finalizeResumeToolBlocks( return out; } -// Resume paints into the same retained log live turns use, so the disk/block -// window is the retention cap rather than a larger pre-slice. -export const RESUME_TRANSCRIPT_BLOCK_LIMIT = MAX_RETAINED_STREAM_ROWS; - -interface TurnsToContentBlocksOptions { - maxBlocks?: number; -} +// Resume paints into the same retained log live turns use. The disk window is +// in turns, not content-blocks: a tool call and its result fold to one row, so +// a 1:1 block cap under-fills the 600-row tail. Two turns per row plus one +// extra means a pair-heavy session can still fill the cap, and a full 600-row +// text tail still overflows when older segments exist. +export const RESUME_TRANSCRIPT_TURN_LIMIT = MAX_RETAINED_STREAM_ROWS * 2 + 1; function turnToContentBlocks(turn: ConversationTurn): ContentBlockData[] { const out: ContentBlockData[] = []; @@ -287,33 +286,11 @@ function turnToContentBlocks(turn: ConversationTurn): ContentBlockData[] { /** Best-effort transcript hydration when resuming a TUI session. */ export function turnsToContentBlocks( turns: ConversationTurn[], - options: TurnsToContentBlocksOptions = {}, ): ContentBlockData[] { - const maxBlocks = options.maxBlocks ?? Infinity; - - // Collect turn-blocks newest-first (backward iteration with early exit so a - // deep session only processes recent turns), then flatten oldest-first. - // Building forward with unshift would be O(n²) — each unshift shifts every - // accumulated element. - const collected: ContentBlockData[][] = []; - let total = 0; - for (let i = turns.length - 1; i >= 0; i--) { - const turn = turns[i]; - if (turn == null) continue; - const blocks = turnToContentBlocks(turn); - if (blocks.length === 0) continue; - collected.push(blocks); - total += blocks.length; - if (total >= maxBlocks) break; - } - const out: ContentBlockData[] = []; - for (let i = collected.length - 1; i >= 0; i--) { - const group = collected[i]; - if (group == null) continue; - out.push(...group); + for (const turn of turns) { + out.push(...turnToContentBlocks(turn)); } - if (out.length > maxBlocks) out.splice(0, out.length - maxBlocks); // Collapse present tool calls into view blocks when args are still available. for (let i = 0; i < out.length; i++) {