diff --git a/agent-chat/server.ts b/agent-chat/server.ts index 12711e98d55a..7e42381ddeac 100644 --- a/agent-chat/server.ts +++ b/agent-chat/server.ts @@ -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"; @@ -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. @@ -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; @@ -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, msg: any) { + handleMessage(ws, msg); +} + function handleMessage(ws: Bun.ServerWebSocket, msg: any) { switch (msg.op) { case "start": { @@ -2232,7 +2248,10 @@ function handleMessage(ws: Bun.ServerWebSocket, 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; diff --git a/agent-chat/test/transcript-reconnect.test.ts b/agent-chat/test/transcript-reconnect.test.ts new file mode 100644 index 000000000000..2d3a4ed38891 --- /dev/null +++ b/agent-chat/test/transcript-reconnect.test.ts @@ -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[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[] = []; + 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[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 }); +} diff --git a/agent-chat/test/transcript-terminal-rebind.test.ts b/agent-chat/test/transcript-terminal-rebind.test.ts new file mode 100644 index 000000000000..e13fd256b483 --- /dev/null +++ b/agent-chat/test/transcript-terminal-rebind.test.ts @@ -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 }[] = []; + 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 }); +}