fix(split-pane): open Pane B's live queue after the response, keep the disconnected marker last

Follow-ups from the #506 review.

The live-frame queue opened before the fetch, freezing Pane B for the
whole round trip. It now opens beside capturedAt; the request uses the
shared terminal fetch deadline and the body read a 10 s one.

A {t:'r'} refresh queued behind a pull ran its clear() after the
disconnected marker was written and wiped it, and a close during a
refresh load wrote the marker above the replay. The marker is now an
owed flag (_markerOwed) that each load settles in its own finally.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
timkjr
2026-10-02 10:00:11 -05:00
co-authored by Claude Sonnet 5.5
parent 9240493c43
commit df398c5c68
3 changed files with 252 additions and 57 deletions
File diff suppressed because one or more lines are too long
+53 -20
View File
@@ -77,11 +77,14 @@
this._bufferLoading = false;
this._bufferRefreshPending = false;
// Scroll-to-top history pull (shell panes only), see _maybeLoadMoreHistory().
// `_liveQueue` is non-null exactly while a pull is replaying: live frames
// are held there with their arrival time instead of written under it.
// `_liveQueue` is non-null from the pull's response until its finally
// block: live frames are held there with their arrival time instead of
// written under the replay. `_markerOwed` is the "disconnected" marker a
// load still has to write (see _onSocketClosed()/_stampMarkerIfOwed()).
this._historyPullAt = 0;
this._historyPullUseless = false;
this._liveQueue = null;
this._markerOwed = false;
this._onWheel = null;
}
@@ -308,7 +311,18 @@
_onSocketClosed() {
this._wsReady = false;
this._wsClosed = true;
if (!this._liveQueue) this._writeDisconnectedMarker();
if (this._bufferLoading) this._markerOwed = true;
else this._writeDisconnectedMarker();
}
// Settles a marker the pane owes: set when a close lands during a load (the
// replay would otherwise sit below it) or when a load wipes the terminal on
// a closed socket. Called from each load's own finally, before a trailing
// refresh starts, so a nested refresh stamps its own.
_stampMarkerIfOwed() {
const owed = this._markerOwed;
this._markerOwed = false;
if (owed && this._wsClosed && !this._destroyed) this._writeDisconnectedMarker();
}
// Extracted so both _onSocketClosed() and a history pull that ends on a
@@ -351,6 +365,7 @@
} catch {
/* Best-effort — live output still arrives once the socket connects. */
} finally {
this._stampMarkerIfOwed();
this._endBufferLoad();
}
}
@@ -430,27 +445,41 @@
// main thread) and replays it under the reader's current place. Holds the
// single-flight flag across the fetch AND the replay, like _loadBuffer().
async _pullHistory() {
// A close before the pull already wrote its marker; one during it did not.
const closedBefore = this._wsClosed;
this._bufferLoading = true;
this._liveQueue = [];
let replayed = false;
let capturedAt = 0;
// Two budgets on one signal. The request itself gets the primary pane's
// (CodemanFetchDeadline, constants.js): nothing is held while it runs. Once
// the headers land live output IS held, so the body read gets the short one
// instead: a body that hangs would otherwise freeze the pane for the long
// budget. Aborting lands in the catch below, which releases the flag and
// the queue. AbortSignal.timeout() alone cannot be re-armed, hence the
// controller; without AbortController the pull simply has no deadline.
const controller = global.AbortController ? new global.AbortController() : null;
let abortTimer = null;
const armDeadline = (ms) => {
if (!controller) return;
clearTimeout(abortTimer);
abortTimer = setTimeout(() => controller.abort(), ms);
};
try {
// A deadline, because live output is held for as long as this runs: a
// request that hangs would otherwise freeze the whole pane. Aborting
// lands in the catch below, which releases the flag and the queue. It
// covers the body read too, not just the headers.
armDeadline(global.CodemanFetchDeadline?.terminalFetchDeadlineMs?.({ full: true }) ?? HISTORY_PULL_TIMEOUT_MS);
const res = await fetch(`/api/sessions/${this.sessionId}/terminal?full=1&tail=${TERMINAL_TAIL_SIZE}`, {
signal: global.AbortSignal?.timeout?.(HISTORY_PULL_TIMEOUT_MS),
signal: controller?.signal,
});
armDeadline(HISTORY_PULL_TIMEOUT_MS);
// The cutoff below is the response's arrival, the same `since` rule the
// primary pane uses (_finishBufferLoad). It is a client clock standing in
// for the instant tmux took the capture, which lies somewhere in the
// round trip, so a frame in that window can be lost or doubled. Bounded
// by one round trip and not closable without a server-side capture time.
capturedAt = performance.now();
// Opened only now: a frame from before the response is either replaced by
// the capture or written unchanged, so holding it for the round trip
// bought nothing and froze the pane for as long as the fetch took.
this._liveQueue = [];
const payload = (await res.json())?.data;
clearTimeout(abortTimer);
const buffer = payload?.terminalBuffer;
const term = this.terminal;
if (!buffer || !term || this._destroyed) return;
@@ -476,6 +505,7 @@
this._historyPullUseless = false;
term.write('\x1bc');
replayed = true;
if (this._wsClosed) this._markerOwed = true;
await writeChunked(term, buffer, () => this._destroyed);
if (this._destroyed || !this.terminal) return;
// xterm parses asynchronously: an empty write's callback fires only
@@ -490,6 +520,7 @@
} catch {
/* Best-effort — live output keeps arriving whatever happens here. */
} finally {
clearTimeout(abortTimer);
const queued = this._liveQueue ?? [];
this._liveQueue = null;
// After a replay, only frames that arrived after the capture are news;
@@ -500,15 +531,13 @@
if (entry.clear) this.terminal?.clear();
else this.terminal?.write(entry.data);
}
// A replay's own `\x1bc` wipes a marker written before the pull,
// painting a fresh, current-looking history while onData keeps
// silently dropping every keystroke on the dead socket, so re-stamp it
// after a replay. A close DURING the pull wrote no marker at all
// (_onSocketClosed() defers it while the queue is live), so write it
// whether or not this pull replayed. Checked after the queue flush so
// it is the last thing on screen, matching what the close would have
// left had the pull never run.
if (this._wsClosed && (replayed || !closedBefore)) this._writeDisconnectedMarker();
// Settled after the queue flush so the marker is the last thing on
// screen: a close during the pull wrote nothing (_onSocketClosed() defers
// it while a load runs), and a replay's own `\x1bc` (flagged above) wipes
// one written before it, which would paint a fresh, current-looking
// history while onData keeps silently dropping every keystroke on the
// dead socket. A trailing refresh (_endBufferLoad) settles its own.
this._stampMarkerIfOwed();
this._endBufferLoad();
}
}
@@ -525,6 +554,10 @@
return;
}
this.terminal?.clear();
// The clear wipes a "disconnected" marker (a `{t:'r'}` frame can queue a
// trailing refresh behind a pull that the socket's close then interrupts),
// so a refresh on a closed socket owes it back once its replay is written.
if (this._wsClosed) this._markerOwed = true;
void this._loadBuffer();
}
+198 -36
View File
@@ -71,6 +71,8 @@ type PaneUnderTest = {
const fetchMock = vi.fn();
/** requestAnimationFrame stand-in: chunked writes queue here and are drained by hand. */
const rafQueue: Array<() => void> = [];
/** Recorded deadline timers (see the context's setTimeout); `fn` aborts the request. */
const deadlines: Array<{ fn: () => void; ms: number; cleared: boolean }> = [];
const SOURCE = readFileSync(resolve(import.meta.dirname, '../src/web/public/terminal-split.js'), 'utf8');
function loadSplitTerminalPane() {
@@ -80,9 +82,27 @@ function loadSplitTerminalPane() {
// compares it with the pane's own row count.
window: {
app: { _estimateReplayRows: (text: string) => text.split('\n').length },
AbortSignal: { timeout: (ms: number) => ({ timeoutMs: ms }) },
// The primary pane's capture budget (constants.js): the full-capture default.
CodemanFetchDeadline: { terminalFetchDeadlineMs: () => 45_000 },
AbortController: class {
signal = { aborted: false };
abort() {
this.signal.aborted = true;
}
},
},
performance: { now: () => clock },
// Deadline timers (>= 1 s) are recorded, never run: a test fires one by hand
// and reads what it aborted. Anything shorter (xterm chunk pacing) is real.
setTimeout: (fn: () => void, ms?: number) => {
if ((ms ?? 0) < 1000) return setTimeout(fn, ms);
deadlines.push({ fn, ms: ms as number, cleared: false });
return -deadlines.length; // negative: never collides with a real timer id
},
clearTimeout: (id: unknown) => {
if (typeof id === 'number' && id < 0) deadlines[-id - 1].cleared = true;
else clearTimeout(id as Parameters<typeof clearTimeout>[0]);
},
fetch: (...args: unknown[]) => fetchMock(...args),
requestAnimationFrame: (fn: () => void) => rafQueue.push(fn),
// The constants.js globals the module reads at call time.
@@ -126,6 +146,27 @@ function jsonResponse(terminalBuffer: string, extra: Record<string, unknown> = {
return { json: async () => ({ data: { terminalBuffer, ...extra } }) };
}
/**
* A response whose headers have landed but whose body has not: the window in
* which the pull's live queue is open and nothing else has happened yet.
* `release(buffer)` delivers the body; `fail()` errors the body read.
*/
function headersOnly() {
let release!: (body: ReturnType<typeof jsonResponse> | Error) => void;
const body = new Promise<{ data: Record<string, unknown> }>((resolve, reject) => {
release = (value) => {
if (value instanceof Error) reject(value);
else void value.json().then(resolve);
};
});
return {
response: { json: () => body },
release: (terminalBuffer: string, extra: Record<string, unknown> = {}) =>
release(jsonResponse(terminalBuffer, extra)),
fail: () => release(new Error('body read failed')),
};
}
function deferred<T>() {
let resolve!: (value: T) => void;
const promise = new Promise<T>((r) => {
@@ -142,6 +183,7 @@ const settle = () => new Promise((r) => setTimeout(r, 0));
beforeEach(() => {
fetchMock.mockReset();
rafQueue.length = 0;
deadlines.length = 0;
clock = 0;
});
@@ -301,10 +343,9 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
pane._maybeLoadMoreHistory();
await settle();
// With a deadline: live output is held for as long as the pull runs, so a
// request that never answers would freeze the pane.
// With a deadline, so a request that never answers cannot pin the pane.
expect(fetchMock).toHaveBeenCalledWith(`/api/sessions/s1/terminal?full=1&tail=${TERMINAL_TAIL_SIZE}`, {
signal: { timeoutMs: 10_000 },
signal: expect.objectContaining({ aborted: false }),
});
expect(term.write).toHaveBeenCalledWith('\x1bc');
expect(term.write).toHaveBeenCalledWith(rowsOf(100));
@@ -463,12 +504,14 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
fetchMock.mockReturnValueOnce(response.promise);
pane._maybeLoadMoreHistory();
expect(pane._liveQueue).toEqual([]);
// The queue opens with the response, not the request.
expect(pane._liveQueue).toBeNull();
// Arrives before the response does: it is IN the capture already.
// Arrives before the response does: written straight through (the pane
// keeps painting during the round trip), and the replay then replaces it.
clock = 1;
pane._onLiveOutput('early');
expect(term.write).not.toHaveBeenCalledWith('early');
expect(term.write).toHaveBeenCalledWith('early');
await settle();
// 200 rows (more than the pane holds, so it replays) of 400 columns each:
@@ -490,7 +533,9 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
await settle();
const written = term.write.mock.calls.map((call) => call[0]);
expect(written).not.toContain('early');
// 'early' went out before the reset, so the replay wiped it and it is not repeated.
expect(written.indexOf('early')).toBeLessThan(written.indexOf('\x1bc'));
expect(written.filter((w) => w === 'early')).toHaveLength(1);
expect(written.at(-1)).toBe('late');
expect(pane._liveQueue).toBeNull();
expect(pane._bufferLoading).toBe(false);
@@ -498,26 +543,29 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
it('writes every held frame when the pull ends without replaying', async () => {
const pane = makePane('shell');
const response = deferred<ReturnType<typeof jsonResponse>>();
fetchMock.mockReturnValueOnce(response.promise);
const held = headersOnly();
fetchMock.mockResolvedValueOnce(held.response);
pane._maybeLoadMoreHistory();
pane._onLiveOutput('held');
await settle();
response.resolve(jsonResponse(rowsOf(30))); // nothing to gain: no replay
pane._onLiveOutput('held');
expect(pane.terminal.write).not.toHaveBeenCalledWith('held');
held.release(rowsOf(30)); // nothing to gain: no replay
await settle();
// Nothing replaced the terminal, so the frame is news even though it
// arrived before the response did.
// Nothing replaced the terminal, so the held frame is news.
expect(pane.terminal.write).toHaveBeenCalledWith('held');
});
it('a failed fetch releases the flag and the queue, so live output flows again', async () => {
const pane = makePane('shell');
fetchMock.mockRejectedValueOnce(new Error('offline'));
const held = headersOnly();
fetchMock.mockResolvedValueOnce(held.response);
pane._maybeLoadMoreHistory();
await settle();
pane._onLiveOutput('held');
held.fail(); // the body read dies with the queue open
await settle();
expect(pane._bufferLoading).toBe(false);
@@ -554,17 +602,18 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
done?.();
});
term.clear.mockImplementation(() => order.push('clear'));
const response = deferred<ReturnType<typeof jsonResponse>>();
fetchMock.mockReturnValueOnce(response.promise);
const held = headersOnly();
fetchMock.mockResolvedValueOnce(held.response);
pane._maybeLoadMoreHistory();
await settle();
pane._onLiveOutput('before');
pane._onLiveClear();
pane._onLiveOutput('after');
// Held: clearing now would wipe a half-written snapshot.
expect(order).toEqual([]);
response.resolve(jsonResponse(rowsOf(30))); // nothing to gain: no replay
held.release(rowsOf(30)); // nothing to gain: no replay
await settle();
expect(order).toEqual(['write:before', 'clear', 'write:after']);
@@ -583,25 +632,27 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
pane._maybeLoadMoreHistory();
clock = 1;
pane._onLiveClear(); // already reflected in the capture
pane._onLiveClear(); // before the response: applied now, already in the capture
expect(term.clear).toHaveBeenCalledTimes(1);
clock = 2;
response.resolve(jsonResponse(rowsOf(100)));
await settle();
expect(term.write).toHaveBeenCalledWith('\x1bc');
expect(term.clear).not.toHaveBeenCalled();
expect(term.clear).toHaveBeenCalledTimes(1); // not replayed after the capture
});
it('destroy() mid-pull leaves nothing running and nothing written to the dead terminal', async () => {
const pane = makePane('shell');
const term = pane.terminal;
const response = deferred<ReturnType<typeof jsonResponse>>();
fetchMock.mockReturnValueOnce(response.promise);
const held = headersOnly();
fetchMock.mockResolvedValueOnce(held.response);
pane._maybeLoadMoreHistory();
await settle();
pane._onLiveOutput('held');
pane.destroy();
response.resolve(jsonResponse(rowsOf(100)));
held.release(rowsOf(100));
await settle();
expect(pane._bufferLoading).toBe(false);
@@ -613,10 +664,13 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
it('a pull whose request is aborted (the deadline) frees the pane', async () => {
const pane = makePane('shell');
fetchMock.mockRejectedValueOnce(new Error('The operation timed out'));
const held = headersOnly();
fetchMock.mockResolvedValueOnce(held.response);
pane._maybeLoadMoreHistory();
await settle();
pane._onLiveOutput('held');
held.fail(); // the deadline aborts the body read
await settle();
expect(pane._bufferLoading).toBe(false);
@@ -698,18 +752,18 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
expect(pane.terminal.write).toHaveBeenCalledWith(marker);
});
it('writes the disconnected marker once, after the replay, if the socket closes mid-fetch', async () => {
// The close lands while the capture is in flight, so the HTTP pull still
it('ends on the disconnected marker if the socket closes mid-fetch, before the queue opens', async () => {
// The close lands while the request is in flight, so the HTTP pull still
// succeeds (a Codeman restart drops the WS while the tmux session, and so
// the pull, survives) and the replay that follows is what the marker must
// end up below.
// the pull, survives). The marker waits for the load to finish, then goes
// below the replay.
const pane = makePane('shell');
const response = deferred<ReturnType<typeof jsonResponse>>();
fetchMock.mockReturnValueOnce(response.promise);
const pull = pane._pullHistory();
pane._onSocketClosed(); // the close arrives mid-fetch, before the response
expect(pane.terminal.write).not.toHaveBeenCalled();
pane._onSocketClosed();
expect(pane.terminal.write.mock.calls.map((c) => c[0]).filter(isMarker)).toHaveLength(0);
response.resolve(jsonResponse(rowsOf(100)));
await pull;
@@ -718,25 +772,101 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
expect(isMarker(writes.at(-1))).toBe(true);
});
it('a close mid-fetch that ends without a replay still writes exactly one marker', async () => {
const pane = makePane('shell');
const response = deferred<ReturnType<typeof jsonResponse>>();
fetchMock.mockReturnValueOnce(response.promise);
const pull = pane._pullHistory();
pane._onSocketClosed();
response.resolve(jsonResponse(rowsOf(30))); // already held in full: no replay
await pull;
const writes = pane.terminal.write.mock.calls.map((c) => c[0]);
expect(writes.filter(isMarker)).toHaveLength(1);
});
it('a close during a refresh load lands the marker below the replay, not above it', async () => {
// A {t:'r'} refresh clears and refetches; a close mid-fetch used to write
// the marker at once, and the replay then landed underneath it.
const pane = makePane('shell');
const response = deferred<ReturnType<typeof jsonResponse>>();
fetchMock.mockReturnValueOnce(response.promise);
pane._refreshBuffer();
pane._onSocketClosed();
response.resolve(jsonResponse('refreshed'));
await settle();
const writes = pane.terminal.write.mock.calls.map((c) => c[0]);
expect(writes.indexOf('refreshed')).toBeLessThan(writes.findIndex(isMarker));
expect(isMarker(writes.at(-1))).toBe(true);
expect(writes.filter(isMarker)).toHaveLength(1);
});
it('back-to-back refreshes on a closed socket leave exactly one marker, at the end', async () => {
// R1's finally runs the trailing refresh R2; each load settles its own
// marker, so R1 never stamps onto R2's freshly cleared terminal.
const pane = makePane('shell');
pane._wsClosed = true;
const first = deferred<ReturnType<typeof jsonResponse>>();
const second = deferred<ReturnType<typeof jsonResponse>>();
fetchMock.mockReturnValueOnce(first.promise).mockReturnValueOnce(second.promise);
pane._refreshBuffer(); // R1
pane._refreshBuffer(); // coalesced into the trailing R2
first.resolve(jsonResponse('first'));
await settle();
// R1 settled its marker, then R2 cleared and is still fetching.
expect(pane.terminal.clear).toHaveBeenCalledTimes(2);
second.resolve(jsonResponse('second'));
await settle();
const writes = pane.terminal.write.mock.calls.map((c) => c[0]);
expect(writes.at(-1)).toSatisfy(isMarker);
expect(writes.lastIndexOf('second')).toBe(writes.length - 2);
});
it('the pull gives the request the long budget and the body read the short one', async () => {
const pane = makePane('shell');
const held = headersOnly();
fetchMock.mockResolvedValueOnce(held.response);
const pull = pane._pullHistory();
expect(deadlines).toHaveLength(1);
expect(deadlines[0].ms).toBe(45_000);
await settle(); // headers landed
expect(deadlines).toHaveLength(2);
expect(deadlines[0].cleared).toBe(true);
expect(deadlines[1].ms).toBe(10_000);
held.release(rowsOf(30));
await pull;
expect(deadlines[1].cleared).toBe(true); // nothing left to abort a settled pull
});
it.each([
['a skip', 40, () => fetchMock.mockResolvedValueOnce(jsonResponse(rowsOf(30)))],
['a downgrade', 500, () => fetchMock.mockResolvedValueOnce(jsonResponse(rowsOf(5)))],
['a failed fetch', 40, () => fetchMock.mockRejectedValueOnce(new Error('offline'))],
['a skip', 40, (h: ReturnType<typeof headersOnly>) => h.release(rowsOf(30))],
['a downgrade', 500, (h: ReturnType<typeof headersOnly>) => h.release(rowsOf(5))],
['a failed body read', 40, (h: ReturnType<typeof headersOnly>) => h.fail()],
])(
'a close mid-fetch that ends in %s writes the marker last, after the held frames',
async (_label, rowsHeld, mockFetch) => {
'a close with the queue open that ends in %s writes the marker last, after the held frames',
async (_label, rowsHeld, finish) => {
// No replay ever runs here, so nothing would wipe a marker written at the
// close; written straight away it sat ABOVE the output the pull was still
// holding, which the finally block then flushed underneath it.
const pane = makePane('shell');
pane.terminal.buffer.active.length = rowsHeld;
mockFetch();
const held = headersOnly();
fetchMock.mockResolvedValueOnce(held.response);
const pull = pane._pullHistory();
await settle(); // the response landed: the queue is open
pane._onLiveOutput('frame-A');
pane._onLiveOutput('frame-B');
pane._onSocketClosed();
expect(pane.terminal.write).not.toHaveBeenCalled();
finish(held);
await pull;
const writes = pane.terminal.write.mock.calls.map((c) => c[0]);
@@ -771,6 +901,38 @@ describe('SplitTerminalPane scroll-to-top history pull', () => {
expect(isMarker(writes.at(-1))).toBe(true);
});
it('a refresh queued behind a pull on a closed socket does not wipe the marker', async () => {
// The refresh's clear() runs after the pull's finally block has written the
// marker, so without a re-stamp the dead pane would look current again.
const pane = makePane('shell');
const held = headersOnly();
fetchMock.mockResolvedValueOnce(held.response).mockResolvedValueOnce(jsonResponse('refreshed'));
const pull = pane._pullHistory();
await settle();
pane._refreshBuffer(); // coalesced into one trailing re-run
pane._onSocketClosed(); // deferred: the queue is open
held.release(rowsOf(30)); // no replay
await pull;
await settle();
const writes = pane.terminal.write.mock.calls.map((c) => c[0]);
expect(pane.terminal.clear).toHaveBeenCalledTimes(1);
expect(writes).toContain('refreshed');
expect(isMarker(writes.at(-1))).toBe(true);
expect(writes.lastIndexOf('refreshed')).toBeLessThan(writes.length - 1);
});
it('a refresh on an open socket does not stamp a marker', async () => {
const pane = makePane('shell');
fetchMock.mockResolvedValueOnce(jsonResponse('refreshed'));
pane._refreshBuffer();
await settle();
expect(pane.terminal.write.mock.calls.map((c) => c[0]).filter(isMarker)).toHaveLength(0);
});
it('does not re-stamp the marker when the socket is still open', async () => {
const pane = makePane('shell');
fetchMock.mockResolvedValueOnce(jsonResponse(rowsOf(100)));