Skip to content
Closed
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
25 changes: 22 additions & 3 deletions agent-chat/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ import { claudeAdapter } from "./adapters/claude";
import { codexAdapter } from "./adapters/codex";
import { piAdapter } from "./adapters/pi";
import { makeAcpAdapter } from "./adapters/acp";
import { attachTranscript, focusTranscriptTerminal, transcriptAdapter, type TranscriptAgent } from "./adapters/transcript";
import { attachTranscript, focusTranscriptTerminal, transcriptAdapter, transcriptTarget, type TranscriptAgent } from "./adapters/transcript";
import { resolveSessionTranscript, resolveSurfaceTranscript, transcriptAttention, type TranscriptSource } from "./transcript-sources";
import { pickAccentColor, resolveGhosttyTheme, resolveGhosttyThemeAsync, type GhosttyTheme } from "./theme";
import { agentModelCatalog, type AgentModelProviderCatalog } from "./catalog";
Expand Down Expand Up @@ -644,7 +644,15 @@ function transcriptTitle(source: TranscriptSource): string {
function ensureTranscriptSession(source: TranscriptSource): Session {
const id = transcriptSessionId(source);
const existing = sessions.get(id);
if (existing?.transcript?.path === source.path) return existing;
if (existing?.transcript?.path === source.path) {
const target = transcriptTarget(existing);
if (target?.agentSessionId !== source.sessionId || target?.surfaceId !== source.surfaceId) {
// A resumed agent can keep its transcript while moving to another
// terminal. Replace changed bindings to also fence old RPC replies.
existing.internal.transcriptTarget = { agentSessionId: source.sessionId, surfaceId: source.surfaceId };
}
return existing;
}
if (existing?.transcript) {
// The agent's transcript moved (for example a resolved fallback path):
// re-point the same session so open pages stay subscribed.
Expand All @@ -667,6 +675,10 @@ function ensureTranscriptSession(source: TranscriptSource): Session {
return sess;
}

export function ensureTranscriptSessionForTest(source: TranscriptSource): Session {
return ensureTranscriptSession(source);
}

function startTranscriptTail(sess: Session, source: TranscriptSource) {
attachTranscript(sess, source.agent, source.path, (title) => {
if (sess.title === title) return;
Expand Down Expand Up @@ -2156,6 +2168,10 @@ function sendWsErrorDetails(
ws.send(JSON.stringify({ kind: "error", op, message: safeErrorMessage(op, err, { provider }), ...publicDetails }));
}

export function handleMessageForTest(ws: Bun.ServerWebSocket<WsData>, msg: any) {
handleMessage(ws, msg);
}

function handleMessage(ws: Bun.ServerWebSocket<WsData>, msg: any) {
switch (msg.op) {
case "start": {
Expand Down Expand Up @@ -2232,7 +2248,10 @@ function handleMessage(ws: Bun.ServerWebSocket<WsData>, msg: any) {
}
case "subscribe": {
const sessionId = String(msg.sessionId);
const sess = sessions.get(sessionId) ?? resolveTranscriptSessionById(sessionId);
// The agent may have resumed elsewhere while this page was offline.
// Refresh terminal transcript bindings even when the session is cached;
// retain cached history if the hook record is temporarily unavailable.
const sess = resolveTranscriptSessionById(sessionId) ?? sessions.get(sessionId);
if (!sess) {
ws.send(JSON.stringify({ kind: "no-session", sessionId: msg.sessionId }));
return;
Expand Down
84 changes: 84 additions & 0 deletions agent-chat/test/transcript-reconnect.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
import assert from "node:assert/strict";
import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises";
import { join } from "node:path";
import { ensureTranscriptSessionForTest, handleMessageForTest } from "../server";
import { setTranscriptRpcForTest, transcriptAdapter, type TranscriptTail } from "../adapters/transcript";
import type { TranscriptSource } from "../transcript-sources";
import type { SessionCtx } from "../types";

const priorHooks = process.env.CMUX_AGENT_HOOK_STATE_DIR;
const priorClaude = process.env.CMUX_CLAUDE_HOOK_STATE_PATH;
const sessions: SessionCtx[] = [];
let root: string | undefined;
try {
const scratch = join(import.meta.dir, "../scratch");
await mkdir(scratch, { recursive: true });
root = await mkdtemp(join(scratch, "transcript-reconnect-"));
process.env.CMUX_AGENT_HOOK_STATE_DIR = root;
delete process.env.CMUX_CLAUDE_HOOK_STATE_PATH;
await Promise.all(["claude", "codex"].map((agent) => writeFile(join(root!, `${agent}-hook-sessions.json`), "{}")));
for (const agent of ["claude", "codex"] as const) {
const oldPath = join(root, `${agent}-old.jsonl`);
const newPath = join(root, `${agent}-new.jsonl`);
await Promise.all([oldPath, newPath].map((path) => writeFile(path, "")));
const source: TranscriptSource = { agent, path: oldPath, sessionId: crypto.randomUUID(), surfaceId: "OLD-SURFACE", updatedAt: 1 };
const sess = ensureTranscriptSessionForTest(source);
sessions.push(sess);
await (sess.internal.transcript as { tail: TranscriptTail }).tail.poll();
sess.events.push({ kind: "user", text: "cached history" });
const attachment = sess.internal.transcript;
const messages: any[] = [];
const ws = {
data: { subscribed: null as string | null },
send(payload: string) { messages.push(JSON.parse(payload)); return 1; },
} as unknown as Parameters<typeof handleMessageForTest>[0];
const subscribe = () => { messages.length = 0; handleMessageForTest(ws, { op: "subscribe", sessionId: sess.id }); };
const store = join(root, `${agent}-hook-sessions.json`);
const record = async (path: string) => writeFile(store, JSON.stringify({ sessions: {
[source.sessionId]: { surfaceId: "NEW-SURFACE", transcriptPath: path, updatedAt: 2 },
} }));
await record(oldPath);
subscribe();
assert.equal(ws.data.subscribed, sess.id);
assert.deepEqual(messages.find((message) => message.kind === "history")?.events, [{ kind: "user", text: "cached history" }]);
assert.equal(sess.internal.transcript, attachment, "same-path reconnect preserves the existing reader");
const calls: Record<string, unknown>[] = [];
setTranscriptRpcForTest(async (method, params) => { calls.push({ method, ...params }); return { ok: true }; });
handleMessageForTest(ws, { op: "focus-terminal", sessionId: sess.id });
await Promise.resolve();
assert.equal(calls[0]?.surface_id, "NEW-SURFACE", "reconnecting a cached transcript must refresh its terminal binding");

// Reconnection must also reattach if the hook now points at another file.
await record(newPath);
subscribe();
await (sess.internal.transcript as { tail: TranscriptTail }).tail.poll();
assert.deepEqual(messages.find((message) => message.kind === "history")?.events, [], "the old file's history must not replay after transcript redirection");
const line = agent === "claude"
? { type: "user", uuid: "fresh", message: { role: "user", content: "new transcript input" } }
: { type: "event_msg", payload: { type: "user_message", message: "new transcript input" } };
await writeFile(newPath, JSON.stringify(line) + "\n");
await (sess.internal.transcript as { tail: TranscriptTail }).tail.poll();
assert.equal(messages.some((message) => message.kind === "event" && message.evt?.text === "new transcript input"), true, "the connected page receives events from the new transcript");

// A temporarily unavailable hook record must not discard usable history.
await writeFile(store, "{}");
subscribe();
assert.equal(messages.some((message) => message.kind === "no-session"), false);
assert.equal(messages.find((message) => message.kind === "history")?.events.some((event: any) => event.text === "new transcript input"), true);
handleMessageForTest(ws, { op: "delete", sessionId: sess.id });
transcriptAdapter.dispose(sess);
}
const missingMessages: any[] = [];
const missingWs = { data: { subscribed: null }, send(payload: string) { missingMessages.push(JSON.parse(payload)); return 1; } } as unknown as Parameters<typeof handleMessageForTest>[0];
handleMessageForTest(missingWs, { op: "subscribe", sessionId: "t-missing-session-id" });
assert.deepEqual(missingMessages, [{ kind: "no-session", sessionId: "t-missing-session-id" }]);
console.log("Claude and Codex reconnects refresh cached bindings and transcript paths, stream new history, and retain cached history on hook misses: OK");
} finally {
setTranscriptRpcForTest(null);
for (const sess of sessions) transcriptAdapter.dispose(sess);
if (priorHooks === undefined) delete process.env.CMUX_AGENT_HOOK_STATE_DIR;
else process.env.CMUX_AGENT_HOOK_STATE_DIR = priorHooks;
if (priorClaude === undefined) delete process.env.CMUX_CLAUDE_HOOK_STATE_PATH;
else process.env.CMUX_CLAUDE_HOOK_STATE_PATH = priorClaude;
if (root) await rm(root, { recursive: true, force: true });
}
63 changes: 63 additions & 0 deletions agent-chat/test/transcript-terminal-rebind.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
import assert from "node:assert/strict";
import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises";
import { join } from "node:path";
import { ensureTranscriptSessionForTest } from "../server";
import { focusTranscriptTerminal, setTranscriptRpcForTest, transcriptAdapter } from "../adapters/transcript";
import type { TranscriptSource } from "../transcript-sources";
import type { CmuxRpcResult } from "../cmux-rpc";
import type { SessionCtx } from "../types";

const sessions: SessionCtx[] = [];
let root: string | undefined;
try {
const scratch = join(import.meta.dir, "../scratch");
await mkdir(scratch, { recursive: true });
root = await mkdtemp(join(scratch, "terminal-rebind-"));
for (const agent of ["claude", "codex"] as const) {
const path = join(root, `${agent}.jsonl`);
await writeFile(path, "");
const source: TranscriptSource = { agent, path, sessionId: crypto.randomUUID(), surfaceId: "old-terminal", updatedAt: 1 };
const sess = ensureTranscriptSessionForTest(source);
sessions.push(sess);
sess.events.push({ kind: "user", text: "existing history" });
const attachment = sess.internal.transcript;

let finishOldFocus!: (result: CmuxRpcResult) => void;
setTranscriptRpcForTest(() => new Promise((resolve) => { finishOldFocus = resolve; }));
const oldFocus = focusTranscriptTerminal(sess);

const reopened = ensureTranscriptSessionForTest({ ...source, surfaceId: "new-terminal", updatedAt: 2 });
assert.equal(reopened, sess, "reopening preserves the session and its subscribers");
assert.equal(sess.internal.transcript, attachment, "rebinding does not restart transcript history");
const calls: { method: string; params: Record<string, unknown> }[] = [];
setTranscriptRpcForTest(async (method, params) => { calls.push({ method, params }); return { ok: true }; });
await focusTranscriptTerminal(reopened);
assert.equal(calls.length, 1);
assert.equal(calls[0].method, "surface.focus");
assert.equal(calls[0].params.surface_id, "new-terminal", "reopening the same transcript must focus its current terminal");
finishOldFocus({ ok: false, error: "old terminal disappeared" });
await oldFocus;
assert.deepEqual(sess.events, [{ kind: "user", text: "existing history" }], "the previous terminal's reply must not contaminate the reopened view");

// Repeated resolution with identical IDs must not discard a valid reply.
setTranscriptRpcForTest(() => new Promise((resolve) => { finishOldFocus = resolve; }));
const currentFocus = focusTranscriptTerminal(sess);
ensureTranscriptSessionForTest({ ...source, surfaceId: "new-terminal", updatedAt: 3 });
finishOldFocus({ ok: false, error: "current terminal failed" });
await currentFocus;
assert.equal(sess.events.length, 2, "unchanged bindings still report their own failures");

const detached = ensureTranscriptSessionForTest({ ...source, surfaceId: undefined, updatedAt: 4 });
calls.length = 0;
setTranscriptRpcForTest(async (method, params) => { calls.push({ method, params }); return { ok: true }; });
const detachedFocus = await focusTranscriptTerminal(detached);
assert.equal(detachedFocus.ok, false, "a removed surface must not retain a focus destination");
assert.deepEqual(calls, [], "a removed binding must never focus the stale terminal");
transcriptAdapter.dispose(sess);
}
console.log("Reopening Claude and Codex transcripts refreshes terminal focus without resetting history or accepting stale failures: OK");
} finally {
setTranscriptRpcForTest(null);
for (const sess of sessions) transcriptAdapter.dispose(sess);
if (root) await rm(root, { recursive: true, force: true });
}
Loading