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
118 changes: 86 additions & 32 deletions agent-chat/adapters/transcript.ts
Original file line number Diff line number Diff line change
Expand Up @@ -360,66 +360,102 @@ export const TRANSCRIPT_INITIAL_WINDOW_BYTES = 8 * 1024 * 1024;
const TRANSCRIPT_POLL_MS = 500;
const TRANSCRIPT_READ_CHUNK = 1024 * 1024;

/** Follows an append-only JSONL file by offset, delivering complete lines. */
/** Follows JSONL appends by offset, resetting when the file is replaced. */
export class TranscriptTail {
private offset = -1;
private identity: { dev: number; ino: number } | null = null;
private pending = "";
private timer: ReturnType<typeof setInterval> | null = null;
private inflight: Promise<void> | null = null;
// Filesystem calls can finish after stop/restart. Only the current lifetime
// may advance decoding state or deliver callbacks; old handles still close.
private generation = 0;
private stopped = false;
private inflight: { generation: number; promise: Promise<void> } | null = null;
private decoder = new TextDecoder();

constructor(
readonly path: string,
private readonly onLines: (lines: string[], mtimeMs: number) => void,
private readonly opts: { pollMs?: number; initialWindowBytes?: number } = {},
private readonly opts: { pollMs?: number; initialWindowBytes?: number; onReset?: () => void } = {},
) {}

start() {
if (this.timer) return;
this.stopped = false;
const generation = this.generation;
// A transcript that stops being readable between stat and open (deleted,
// or a root-owned file) must not reject out of the timer: an unhandled
// rejection ends the whole sidecar. The next poll tries again.
const tick = () => void this.poll().catch(() => {});
const tick = () => {
if (this.ownsRead(generation)) void this.poll().catch(() => {});
};
tick();
this.timer = setInterval(tick, this.opts.pollMs ?? TRANSCRIPT_POLL_MS);
}

stop() {
if (this.timer) clearInterval(this.timer);
this.timer = null;
this.stopped = true;
this.generation++;
}

/** Reads everything appended since the last poll; joins a read in flight. */
/** Reads new appends, joining only a read from the current lifetime. */
poll(): Promise<void> {
if (!this.inflight) {
this.inflight = this.read().finally(() => {
this.inflight = null;
if (this.stopped) return Promise.resolve();
const generation = this.generation;
if (!this.inflight || this.inflight.generation !== generation) {
const flight = { generation, promise: this.read(generation) };
this.inflight = flight;
flight.promise = flight.promise.finally(() => {
if (this.inflight === flight) this.inflight = null;
});
}
return this.inflight;
return this.inflight.promise;
}

private async read(): Promise<void> {
const info = await stat(this.path).catch(() => null);
if (!info) return;
let skipPartialFirstLine = false;
if (this.offset < 0) {
const window = this.opts.initialWindowBytes ?? TRANSCRIPT_INITIAL_WINDOW_BYTES;
this.offset = Math.max(0, info.size - window);
skipPartialFirstLine = this.offset > 0;
} else if (info.size < this.offset) {
// Truncated or replaced: follow the new file from its start.
this.offset = 0;
this.pending = "";
this.decoder = new TextDecoder();
}
if (info.size === this.offset) return;
private ownsRead(generation: number): boolean {
return !this.stopped && this.generation === generation;
}

private async read(generation: number): Promise<void> {
const probe = await stat(this.path).catch(() => null);
if (!this.ownsRead(generation) || !probe) return;
if (this.identity?.dev === probe.dev && this.identity.ino === probe.ino && probe.size === this.offset) return;
const handle = await open(this.path, "r");
try {
if (!this.ownsRead(generation)) return;
// Use the opened file's identity and size: an atomic rename can replace
// the path between the probe and open, even with unchanged size/mtime.
const info = await handle.stat();
if (!this.ownsRead(generation)) return;
const reset = this.identity !== null && (
this.identity.dev !== info.dev || this.identity.ino !== info.ino || info.size < this.offset
);
if (this.identity?.dev !== info.dev || this.identity.ino !== info.ino) {
this.offset = -1;
this.pending = "";
this.decoder = new TextDecoder();
}
this.identity = { dev: info.dev, ino: info.ino };
let skipPartialFirstLine = false;
if (this.offset < 0) {
const window = this.opts.initialWindowBytes ?? TRANSCRIPT_INITIAL_WINDOW_BYTES;
this.offset = Math.max(0, info.size - window);
skipPartialFirstLine = this.offset > 0;
} else if (info.size < this.offset) {
// An in-place truncation keeps its inode but still resets decoding.
this.offset = 0;
this.pending = "";
this.decoder = new TextDecoder();
}
if (reset) this.opts.onReset?.();
if (!this.ownsRead(generation)) return;
if (info.size === this.offset) return;
const buf = new Uint8Array(TRANSCRIPT_READ_CHUNK);
while (this.offset < info.size) {
while (this.ownsRead(generation) && this.offset < info.size) {
const { bytesRead } = await handle.read(buf, 0, Math.min(buf.length, info.size - this.offset), this.offset);
if (bytesRead <= 0) break;
if (!this.ownsRead(generation) || bytesRead <= 0) break;
this.offset += bytesRead;
this.pending += this.decoder.decode(buf.subarray(0, bytesRead), { stream: true });
if (skipPartialFirstLine) {
Expand Down Expand Up @@ -476,25 +512,43 @@ export function attachTranscript(
onTitle?: (title: string) => void,
opts: { pollMs?: number; initialWindowBytes?: number; onTick?: () => void } = {},
): TranscriptTail {
const parser = transcriptParser(agent);
let parser = transcriptParser(agent);
const refreshStatus = () => {
const st = transcriptState(sess);
if (!st) return;
if (!st || st.tail !== tail) return;
sess.setStatus(transcriptLooksRunning(sess.events, st.lastWriteMs) ? "running" : "idle");
opts.onTick?.();
if (transcriptState(sess) === st) opts.onTick?.();
};
const tail = new TranscriptTail(path, (lines, mtimeMs) => {
const st = transcriptState(sess);
if (!st || st.tail !== tail) return;
// Activity comes from the file's own write time, so a transcript that
// went idle long ago does not look busy when its history first loads.
if (st) st.lastWriteMs = mtimeMs;
st.lastWriteMs = mtimeMs;
const title = parser.title;
for (const line of lines) {
for (const evt of parser.parse(line)) sess.emit(evt);
if (transcriptState(sess) !== st) return;
for (const evt of parser.parse(line)) {
if (transcriptState(sess) !== st) return;
sess.emit(evt);
}
}
if (transcriptState(sess) !== st) return;
if (parser.title && parser.title !== title) onTitle?.(parser.title);
refreshStatus();
}, opts);
}, {
...opts,
onReset: () => {
const st = transcriptState(sess);
if (!st || st.tail !== tail) return;
parser = transcriptParser(agent);
st.parser = parser;
st.lastWriteMs = 0;
if (sess.resetHistory) sess.resetHistory();
else sess.events.length = 0;
refreshStatus();
},
});
const state: TranscriptState = {
tail,
parser,
Expand Down
13 changes: 10 additions & 3 deletions agent-chat/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -395,6 +395,9 @@ function createSession(
}
emitSessionEvent(sess, evt);
},
resetHistory() {
resetSessionHistory(sess);
},
setStatus(status: SessionStatus) {
const pendingDone = sess.internal.pendingDoneEmit as Promise<void> | undefined;
if (status === "idle" && pendingDone) {
Expand Down Expand Up @@ -480,6 +483,12 @@ function broadcastSessionHistory(sess: Session) {
for (const ws of sess.sockets) ws.send(payload);
}

export function resetSessionHistory(sess: Session) {
sess.events.length = 0;
delete sess.internal.eventGenerations;
broadcastSessionHistory(sess);
}

function activeAttributionGenerations(sess: Session): number[] {
let generations = sess.internal.activeTurnGenerations as number[] | undefined;
if (!generations) {
Expand Down Expand Up @@ -649,12 +658,10 @@ function ensureTranscriptSession(source: TranscriptSource): Session {
// The agent's transcript moved (for example a resolved fallback path):
// re-point the same session so open pages stay subscribed.
existing.adapter.dispose(existing);
existing.events.length = 0;
delete existing.internal.eventGenerations;
existing.transcript.path = source.path;
existing.internal.transcriptTarget = { agentSessionId: source.sessionId, surfaceId: source.surfaceId };
startTranscriptTail(existing, source);
broadcastSessionHistory(existing);
resetSessionHistory(existing);
return existing;
}
const sess = createSession(source.agent, source.cwd ?? DEFAULT_CWD, false, transcriptTitle(source), {}, {}, {
Expand Down
136 changes: 136 additions & 0 deletions agent-chat/test/transcript-disposal.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
import assert from "node:assert/strict";
import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises";
import { join } from "node:path";
import { attachTranscript, transcriptAdapter, TranscriptTail, type TranscriptAgent } from "../adapters/transcript";
import type { AgentEvent, SessionCtx, SessionStatus } from "../types";

const descriptors = new Map(["setInterval", "clearInterval"].map((key) => [key, Object.getOwnPropertyDescriptor(globalThis, key)]));
const intervals = new Map<number, { callback: () => void; delay: number }>();
let nextInterval = 1;
Object.defineProperty(globalThis, "setInterval", { configurable: true, writable: true, value(callback: () => void, delay: number) {
const id = nextInterval++; intervals.set(id, { callback, delay }); return id;
} });
Object.defineProperty(globalThis, "clearInterval", { configurable: true, writable: true, value(id: number) { intervals.delete(id); } });

const tails: TranscriptTail[] = [];
const reads: Promise<void>[] = [];
let root: string | undefined;
const poll = (tail: TranscriptTail) => { const read = tail.poll(); reads.push(read); return read; };
const remember = (tail: TranscriptTail) => { tails.push(tail); return tail; };
const user = (agent: TranscriptAgent, id: string, text: string) => agent === "claude"
? { type: "user", uuid: id, message: { role: "user", content: text } }
: { type: "event_msg", payload: { type: "user_message", message: text } };
const title = (agent: TranscriptAgent, text: string) => agent === "claude"
? { type: "ai-title", aiTitle: text }
: { type: "event_msg", payload: { type: "thread_name_updated", thread_name: text } };
const jsonl = (...records: unknown[]) => records.map((record) => JSON.stringify(record)).join("\n") + "\n";
function session() {
const statuses: SessionStatus[] = [];
const sess = {
events: [] as AgentEvent[], internal: {} as Record<string, unknown>,
emit(event: AgentEvent) { this.events.push(event); },
setStatus(status: SessionStatus) { statuses.push(status); },
} as unknown as SessionCtx;
return { sess, statuses };
}

try {
const scratch = join(import.meta.dir, "../scratch");
await mkdir(scratch, { recursive: true });
root = await mkdtemp(join(scratch, "transcript-disposal-"));
for (const agent of ["claude", "codex"] as const) {
const path = join(root, `${agent}-disposed.jsonl`);
await writeFile(path, jsonl(user(agent, "old", "disposed content"), title(agent, "disposed title")));
const { sess, statuses } = session();
const titles: string[] = [];
let ticks = 0;
const tail = remember(attachTranscript(sess, agent, path, (next) => titles.push(next), { pollMs: 10_000, onTick: () => ticks++ }));
const activeRead = poll(tail);
transcriptAdapter.dispose(sess);
await activeRead;
assert.deepEqual(sess.events, [], "a disposed transcript read must not emit late events");
assert.deepEqual(titles, [], "a disposed transcript read must not rename the view");
assert.deepEqual(statuses, []);
assert.equal(ticks, 0);
assert.equal(intervals.size, 0);
}

for (const agent of ["claude", "codex"] as const) {
const oldPath = join(root, `${agent}-old.jsonl`);
const newPath = join(root, `${agent}-new.jsonl`);
await writeFile(oldPath, jsonl(user(agent, "old", "old content"), title(agent, "old title")));
await writeFile(newPath, jsonl(user(agent, "new", "new content"), title(agent, "new title")));
const { sess } = session();
const titles: string[] = [];
let oldTicks = 0;
let newTicks = 0;
const oldTail = remember(attachTranscript(sess, agent, oldPath, (next) => titles.push(next), { pollMs: 10_000, onTick: () => oldTicks++ }));
const oldStatus = [...intervals.values()].find((interval) => interval.delay === 2_000)!.callback;
const oldRead = poll(oldTail);
transcriptAdapter.dispose(sess);
const newTail = remember(attachTranscript(sess, agent, newPath, (next) => titles.push(next), { pollMs: 10_000, onTick: () => newTicks++ }));
await Promise.all([oldRead, poll(newTail)]);
assert.deepEqual(sess.events, [{ kind: "user", text: "new content" }], "old reads must not append to a redirected session");
assert.deepEqual(titles, ["new title"]);
assert.equal(oldTicks, 0);
assert.ok(newTicks > 0);
oldStatus();
assert.equal(oldTicks, 0, "a retained old status callback must not inspect the new attachment");
transcriptAdapter.dispose(sess);
assert.equal(intervals.size, 0);
}

const restartPath = join(root, "restart.jsonl");
await writeFile(restartPath, '{"line":1}\n');
const restartedLines: string[] = [];
const restart = remember(new TranscriptTail(restartPath, (lines) => restartedLines.push(...lines), { pollMs: 10_000 }));
const oldRead = poll(restart);
restart.stop();
restart.start();
const newRead = poll(restart);
assert.notEqual(newRead, oldRead, "restart must not join the canceled read from the previous lifetime");
await Promise.all([oldRead, newRead]);
assert.deepEqual(restartedLines, ['{"line":1}']);
restart.stop();
await poll(restart);
assert.deepEqual(restartedLines, ['{"line":1}']);
assert.equal(intervals.size, 0);

const burstPath = join(root, "burst.jsonl");
const rows = Array.from({ length: 256 }, (_, index) => JSON.stringify({ index, text: "x".repeat(16_384) }));
await writeFile(burstPath, rows.join("\n") + "\n");
const received: string[] = [];
let batches = 0;
const burst = remember(new TranscriptTail(burstPath, (lines) => {
batches++; received.push(...lines);
if (batches === 1) burst.stop();
}, { pollMs: 10_000 }));
await poll(burst);
assert.equal(batches, 1, "stop during delivery must prevent reading more chunks");
assert.ok(received.length > 0 && received.length < rows.length);
burst.start();
await poll(burst);
assert.deepEqual(received, rows, "restart resumes at the last delivered byte offset without duplicating rows");
burst.stop();

const batchPath = join(root, "dispose-in-batch.jsonl");
await writeFile(batchPath, jsonl(user("claude", "first", "first"), user("claude", "second", "second"), title("claude", "stale title")));
const { sess, statuses } = session();
sess.emit = (event) => { sess.events.push(event); transcriptAdapter.dispose(sess); };
const titles: string[] = [];
const batch = remember(attachTranscript(sess, "claude", batchPath, (next) => titles.push(next), { pollMs: 10_000 }));
await poll(batch);
assert.deepEqual(sess.events, [{ kind: "user", text: "first" }], "disposal during a batch must stop remaining event delivery");
assert.deepEqual(titles, []);
assert.deepEqual(statuses, []);
assert.equal(intervals.size, 0, "all owned timers are released");
console.log("Transcript disposal, redirected views, retained status callbacks, chunk cancellation, and immediate restart: OK");
} finally {
for (const tail of tails) tail.stop();
await Promise.allSettled(reads);
if (root) await rm(root, { recursive: true, force: true });
for (const [key, descriptor] of descriptors) {
if (descriptor) Object.defineProperty(globalThis, key, descriptor);
else delete (globalThis as any)[key];
}
}
Loading
Loading