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
22 changes: 19 additions & 3 deletions src/session/optimized-context-store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand All @@ -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 () => {
Expand Down
23 changes: 19 additions & 4 deletions src/session/optimized-context-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -302,12 +313,13 @@ async function unlinkExtraSegmentsFrom(
export async function loadRecentTurns(
dir: string,
minTurns: number,
): Promise<ConversationTurn[]> {
): Promise<RecentTurns> {
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;
Expand All @@ -325,15 +337,18 @@ export async function loadRecentTurns(
);
collectedNewestFirst.push(turns);
total += turns.length;
if (total >= minTurns) break;
if (total >= minTurns) {
truncated = i > 0;
break;
}
}

const turns: ConversationTurn[] = [];
for (let i = collectedNewestFirst.length - 1; i >= 0; i--) {
const chunk = collectedNewestFirst[i];
if (chunk !== undefined) turns.push(...chunk);
}
return turns;
return { turns, truncated };
}

/**
Expand Down
64 changes: 64 additions & 0 deletions src/tui/history-hydrate.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
EMPTY_VIEW_DETAIL,
hydrateHistoryRows,
MISSING_ERROR_DETAIL,
parseHistoryHydratePayload,
rowFromHistoryBlock,
type HistoryBlock,
} from "./history-hydrate.js";
Expand Down Expand Up @@ -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,
});
});
});
25 changes: 25 additions & 0 deletions src/tui/history-hydrate.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>;
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;
Expand Down
149 changes: 149 additions & 0 deletions src/tui/product-host.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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";

Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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 {
Expand Down
Loading
Loading