fix(tiles): cap each tile's live-output backlog and recover dropped output

The server applies no WebSocket backpressure (16 KB / 8 ms batches, no
bufferedAmount check), and a tile wrote every live frame straight into
xterm. A flood a tile could not parse as fast (a shell tile running cat on
a huge log) piled up in xterm's own write queue without bound, on a main
thread up to six tiles share, until xterm's WriteBuffer threw past 50M
code units; onmessage's empty catch then dropped every frame silently and
nothing recaptured the screen. The primary pane caps its queues and drops
then recaptures (_onSessionTerminal, _scheduleDroppedOutputRecovery).

Each tile now writes live output through _writeLive:
- unparsed code units are counted, each write's callback counting its own
  back down; frames held behind a replay (_liveQueue) count too;
- the budget is TerminalTile.LIVE_BACKLOG_BUDGET, 4 MiB, deliberately not
  the primary pane's 128 KB: that caps its own rAF-paced queues, while
  xterm itself paces a tile, and a tight cap would trip on ordinary bursts
  and blank-and-reload the tile over and over;
- past it a frame is dropped, the tile stops writing onto the hole, and one
  refresh is scheduled, debounced and bounded by the primary pane's own
  rule (CodemanDroppedOutput: 2 s, DROP_RECOVERY_MAX_ATTEMPTS, never retried
  after a deadline abort). It is an ordinary refresh, so single-flight,
  bounded by lines=/tail= and paced by the grid's TileLoadQueue. The flag
  clears once a capture taken after the last dropped frame has replayed;
- a write that throws is the same drop, never a "malformed frame";
- past the bound the flag is released, so a tile is never left frozen;
- a reconnect starts the accounting over (an epoch makes callbacks from
  before it count nothing) and drops a pending recovery, since its own
  refresh replaces the screen; destroy() cancels it.
The live-queue flush after a pull or a refresh goes through the same path,
so a throwing write there cannot skip the load's marker and trailing
refresh either.

Tests (input harness, real constants and fake timers): the default budget
lets a 1 MiB unparsed burst through, parsed bytes stop counting, a trip
stops writing and ONE debounced refresh recaptures, a write throw takes the
same recovery, a hole in the held queue is recovered by another refresh,
bounded retries then release, no retry after a deadline, a reconnect resets
the count, and destroy cancels. The fake xterm can now hold and release
parses and throw on a write. Live writes now carry a callback, so the unit
tests match them on the data argument (a `.not.toHaveBeenCalledWith(data)`
would otherwise pass for nothing).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Codeman maintainer
2026-10-09 09:26:51 +02:00
parent cdb3c34ae0
commit 1def7de146
5 changed files with 444 additions and 19 deletions
+18 -2
View File
@@ -149,9 +149,23 @@ export class FakeTerminal {
}
registerLinkProvider() {}
writes: string[] = [];
/** Set by a test: write callbacks never run, as on a disposed xterm. */
/**
* Set by a test: write callbacks do not run, as while xterm is still parsing
* (or never, on a disposed xterm). They wait in `heldParses` until `parse()`.
*/
holdParse = false;
heldParses: Array<() => void> = [];
/** xterm catches up: runs every write callback held so far, in order. */
parse() {
for (const cb of this.heldParses.splice(0)) cb();
}
/** Set by a test: the next write of exactly this data throws, as xterm's WriteBuffer does past 50M. */
throwOnWrite: string | null = null;
write(data: string, cb?: () => void) {
if (this.throwOnWrite !== null && data === this.throwOnWrite) {
this.throwOnWrite = null;
throw new Error('write data discarded, use flow control to avoid losing data');
}
// An empty write puts nothing on screen; the replay queues one only to hear
// (its callback) that everything before it has been parsed.
if (data) this.writes.push(data);
@@ -162,7 +176,9 @@ export class FakeTerminal {
this.lineCount += data.slice(reset === -1 ? 0 : reset + 2).split('\n').length - 1;
this.settleRows();
}
if (!this.holdParse) cb?.();
if (!cb) return;
if (this.holdParse) this.heldParses.push(cb);
else cb();
}
clear() {
this.writes.push('<CLEAR>');
+246
View File
@@ -802,6 +802,252 @@ describe('TerminalTile claims the keyboard for the app-level shortcuts', () => {
});
});
describe('TerminalTile live-output flow control: a flood cannot pile up in xterm', () => {
// The server applies no backpressure, and a tile used to write every live
// frame straight into xterm: a flood it could not parse as fast piled up in
// xterm's own queue without bound, until xterm's WriteBuffer threw past 50M
// code units and onmessage's catch silently lost the frames. Now unparsed
// (and held) output is counted per tile; past the budget a frame is dropped,
// the tile stops writing onto the hole, and one debounced, bounded refresh
// recaptures the screen (the primary pane's _scheduleDroppedOutputRecovery).
const TileStatics = TerminalTile as unknown as { LIVE_BACKLOG_BUDGET: number };
const defaultBudget = TileStatics.LIVE_BACKLOG_BUDGET;
afterEach(() => {
TileStatics.LIVE_BACKLOG_BUDGET = defaultBudget;
delete windowStub.AbortController;
});
const out = (ws: FakeSocket, d: string) => ws.receive({ t: 'o', d });
const inFlight = (tile: Tile) => (tile as unknown as { _liveInFlight: number })._liveInFlight;
function serve(terminalBuffer: string) {
fetchMock.mockClear();
fetchMock.mockImplementation(async () => ({
ok: true,
status: 200,
json: async () => ({ data: { terminalBuffer } }),
}));
}
it("the budget is a few MB, not the primary pane's 128 KB: a 1 MiB burst xterm has not parsed yet is still written", async () => {
expect(defaultBudget).toBeGreaterThanOrEqual(2 * 1024 * 1024);
expect(defaultBudget).toBeLessThanOrEqual(16 * 1024 * 1024);
vi.useFakeTimers();
const { ws, term } = await connectTile(makeApp());
ws.open();
serve('');
term.holdParse = true;
out(ws, 'z'.repeat(1024 * 1024));
expect(term.writes.at(-1)).toHaveLength(1024 * 1024);
await vi.advanceTimersByTimeAsync(10_000);
expect(fetchMock).not.toHaveBeenCalled();
});
it('counts output until xterm has parsed it, so parsed bytes no longer hold the budget', async () => {
TileStatics.LIVE_BACKLOG_BUDGET = 100;
const { tile, ws, term } = await connectTile(makeApp());
ws.open();
term.holdParse = true;
out(ws, 'a'.repeat(40));
out(ws, 'b'.repeat(40));
expect(inFlight(tile)).toBe(80);
term.parse();
expect(inFlight(tile)).toBe(0);
out(ws, 'c'.repeat(90));
expect(term.writes.slice(-3)).toEqual(['a'.repeat(40), 'b'.repeat(40), 'c'.repeat(90)]);
});
it('stops writing past the budget, and ONE debounced refresh recaptures the screen', async () => {
TileStatics.LIVE_BACKLOG_BUDGET = 100;
vi.useFakeTimers();
const { ws, term } = await connectTile(makeApp());
ws.open();
serve('recovered screen');
term.holdParse = true; // xterm falls behind
out(ws, 'a'.repeat(60));
out(ws, 'b'.repeat(60)); // 120 unparsed: past the budget, dropped
out(ws, 'c');
term.parse(); // xterm catches up, but the stream already has a hole:
out(ws, 'd'); // nothing more is written onto it
expect(term.writes.filter((w) => /^[abcd]/.test(w))).toEqual(['a'.repeat(60)]);
expect(fetchMock).not.toHaveBeenCalled();
term.holdParse = false;
await vi.advanceTimersByTimeAsync(1999);
expect(fetchMock).not.toHaveBeenCalled(); // debounced, like the primary pane's
await vi.advanceTimersByTimeAsync(1);
await settle();
expect(fetchMock).toHaveBeenCalledTimes(1);
expect(fetchMock.mock.calls[0][0]).toBe('/api/sessions/s-tile/terminal?full=1');
expect(term.writes.slice(-2)).toEqual(['\x1bc', 'recovered screen']);
out(ws, 'live again');
expect(term.writes.at(-1)).toBe('live again');
await vi.advanceTimersByTimeAsync(10_000);
expect(fetchMock).toHaveBeenCalledTimes(1); // one recovery for the whole burst
});
it("a write xterm refuses (its 50M throw) is a drop that schedules the same recovery, never a swallowed 'malformed frame'", async () => {
vi.useFakeTimers();
const { ws, term } = await connectTile(makeApp());
ws.open();
serve('recovered screen');
term.throwOnWrite = 'boom';
out(ws, 'boom');
out(ws, 'after the hole');
expect(term.writes).not.toContain('after the hole');
await vi.advanceTimersByTimeAsync(2000);
await settle();
expect(fetchMock).toHaveBeenCalledTimes(1);
expect(term.writes.slice(-2)).toEqual(['\x1bc', 'recovered screen']);
});
it('frames held behind a replay count against the budget too, and a hole there is recovered by another refresh', async () => {
TileStatics.LIVE_BACKLOG_BUDGET = 100;
vi.useFakeTimers();
const { ws, term } = await connectTile(makeApp());
ws.open();
fetchMock.mockClear();
let answer!: (body: string) => void;
fetchMock.mockImplementationOnce(async () => ({
ok: true,
status: 200,
json: () =>
new Promise((resolveBody) => {
answer = (terminalBuffer) => resolveBody({ data: { terminalBuffer } });
}),
}));
fetchMock.mockImplementation(async () => ({
ok: true,
status: 200,
json: async () => ({ data: { terminalBuffer: 'second capture' } }),
}));
ws.receive({ t: 'r' });
await settle(); // the headers landed: from here frames are held
out(ws, 'x'.repeat(60));
out(ws, 'y'.repeat(60)); // the held queue would pass the budget: dropped
answer('first capture');
await settle();
// The capture predates the hole, so nothing after it is written onto it.
expect(term.writes.slice(-2)).toEqual(['\x1bc', 'first capture']);
expect(term.writes).not.toContain('y'.repeat(60));
await vi.advanceTimersByTimeAsync(2000);
await settle();
expect(fetchMock).toHaveBeenCalledTimes(2);
expect(term.writes.slice(-2)).toEqual(['\x1bc', 'second capture']);
out(ws, 'live again');
expect(term.writes.at(-1)).toBe('live again');
});
it('a recovery that keeps failing is retried a bounded number of times, then lets live output through again', async () => {
vi.useFakeTimers();
const { ws, term } = await connectTile(makeApp());
ws.open();
fetchMock.mockClear();
fetchMock.mockImplementation(async () => {
throw new Error('offline');
});
term.throwOnWrite = 'boom';
out(ws, 'boom');
out(ws, 'held back');
for (let i = 0; i < 10; i++) {
await vi.advanceTimersByTimeAsync(2000);
await settle();
}
const { DROP_RECOVERY_MAX_ATTEMPTS } = windowStub.CodemanDroppedOutput as { DROP_RECOVERY_MAX_ATTEMPTS: number };
expect(fetchMock).toHaveBeenCalledTimes(DROP_RECOVERY_MAX_ATTEMPTS);
expect(term.writes).not.toContain('held back');
out(ws, 'flowing');
expect(term.writes.at(-1)).toBe('flowing'); // never left frozen
});
it('a recovery cut off at its deadline is not retried (a stalled link), and lets live output through again', async () => {
windowStub.AbortController = AbortController;
vi.useFakeTimers();
const { ws, term } = await connectTile(makeApp());
ws.open();
fetchMock.mockClear();
fetchMock.mockImplementation(
(_url: string, init?: { signal?: AbortSignal }) =>
new Promise((_resolveFetch, rejectFetch) => {
init?.signal?.addEventListener('abort', () => rejectFetch(new Error('aborted')));
})
);
term.throwOnWrite = 'boom';
out(ws, 'boom');
await vi.advanceTimersByTimeAsync(2000);
expect(fetchMock).toHaveBeenCalledTimes(1);
await vi.advanceTimersByTimeAsync(45_000); // the full-capture budget runs out
await settle();
out(ws, 'flowing');
expect(term.writes.at(-1)).toBe('flowing');
await vi.advanceTimersByTimeAsync(20_000);
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it('a reconnect starts the count over, and callbacks from before it count nothing', async () => {
TileStatics.LIVE_BACKLOG_BUDGET = 100;
vi.useFakeTimers();
const { tile, ws, term } = await connectTile(makeApp());
ws.open();
term.holdParse = true;
out(ws, 'a'.repeat(60));
expect(inFlight(tile)).toBe(60);
ws.drop(1006);
await vi.advanceTimersByTimeAsync(300);
FakeSocket.instances.at(-1)!.open(); // reconnected: its refresh replaces the screen
expect(inFlight(tile)).toBe(0);
term.parse(); // the old write's callback lands after the reset
expect(inFlight(tile)).toBe(0);
});
it('a reconnect drops a pending recovery: its own refresh replaces the screen', async () => {
vi.useFakeTimers();
const { ws, term } = await connectTile(makeApp());
ws.open();
serve('');
term.throwOnWrite = 'boom';
out(ws, 'boom');
ws.drop(1006);
await vi.advanceTimersByTimeAsync(300);
const ws2 = FakeSocket.instances.at(-1)!;
ws2.open();
await settle();
expect(fetchMock).toHaveBeenCalledTimes(1); // the reconnect's refresh
out(ws2, 'flowing');
expect(term.writes.at(-1)).toBe('flowing');
await vi.advanceTimersByTimeAsync(5000);
expect(fetchMock).toHaveBeenCalledTimes(1);
});
it('destroy() cancels a pending recovery', async () => {
vi.useFakeTimers();
const { tile, ws, term } = await connectTile(makeApp());
ws.open();
fetchMock.mockClear();
term.throwOnWrite = 'boom';
out(ws, 'boom');
tile.destroy();
await vi.advanceTimersByTimeAsync(5000);
expect(fetchMock).not.toHaveBeenCalled();
});
});
describe('the server coming back kicks Pane B', () => {
it("handleInit's reconnect branch asks the split pane's tile to reconnect without waiting out its backoff", () => {
// handleInit needs a whole app to run, so the wiring is pinned by source;
+15 -9
View File
@@ -185,6 +185,10 @@ const isMarker = (data: unknown) => typeof data === 'string' && data.includes('[
*/
const screenWrites = (pane: { terminal: FakeTerminal }) =>
pane.terminal.write.mock.calls.map((call) => call[0]).filter((data) => data !== '');
// Live frames are written with a callback (the tile counts them until xterm
// has parsed them, _writeLive), so they are matched on the data argument
// alone, never with toHaveBeenCalledWith(data): that would also miss, and a
// `.not` on it would then pass for nothing.
/**
* Holds xterm's write callbacks, as a real xterm still parsing a replay does:
@@ -738,7 +742,7 @@ describe('TerminalTile scroll-to-top history pull', () => {
// keeps painting during the round trip), and the replay then replaces it.
clock = 1;
pane._onLiveOutput('early');
expect(term.write).toHaveBeenCalledWith('early');
expect(screenWrites(pane)).toContain('early');
await settle();
// 200 rows (more than the pane holds, so it replays) of 400 columns each:
@@ -748,12 +752,14 @@ describe('TerminalTile scroll-to-top history pull', () => {
clock = 2; // the response arrives: this is the cutoff
response.resolve(jsonResponse(bigReplay));
await settle();
expect(xterm.held).toHaveLength(1);
// Two parses pending: 'early' (a live write, counted until xterm parses it)
// and the replay's end marker.
expect(xterm.held).toHaveLength(2);
// Arrives while the snapshot is still being parsed: must not land under it.
clock = 3;
pane._onLiveOutput('late');
expect(term.write).not.toHaveBeenCalledWith('late');
expect(screenWrites(pane)).not.toContain('late');
// The replay parsed, then the pull's own settle write before it scrolls.
xterm.parse();
@@ -778,12 +784,12 @@ describe('TerminalTile scroll-to-top history pull', () => {
pane._maybeLoadMoreHistory();
await settle();
pane._onLiveOutput('held');
expect(pane.terminal.write).not.toHaveBeenCalledWith('held');
expect(screenWrites(pane)).not.toContain('held');
held.release(rowsOf(30)); // nothing to gain: no replay
await settle();
// Nothing replaced the terminal, so the held frame is news.
expect(pane.terminal.write).toHaveBeenCalledWith('held');
expect(screenWrites(pane)).toContain('held');
});
it('a failed fetch releases the flag and the queue, so live output flows again', async () => {
@@ -799,9 +805,9 @@ describe('TerminalTile scroll-to-top history pull', () => {
expect(pane._bufferLoading).toBe(false);
expect(pane._liveQueue).toBeNull();
expect(pane.terminal.write).toHaveBeenCalledWith('held');
expect(screenWrites(pane)).toContain('held');
pane._onLiveOutput('after');
expect(pane.terminal.write).toHaveBeenLastCalledWith('after');
expect(screenWrites(pane).at(-1)).toBe('after');
});
it('a refresh frame during the pull runs once behind it', async () => {
@@ -921,7 +927,7 @@ describe('TerminalTile scroll-to-top history pull', () => {
expect(pane._liveQueue).toBeNull();
expect(pane.terminal).toBeNull();
expect(term.write).not.toHaveBeenCalledWith('\x1bc');
expect(term.write).not.toHaveBeenCalledWith('held');
expect(term.write.mock.calls.map((call) => call[0])).not.toContain('held');
});
it('a pull whose request is aborted (the deadline) frees the pane', async () => {
@@ -937,7 +943,7 @@ describe('TerminalTile scroll-to-top history pull', () => {
expect(pane._bufferLoading).toBe(false);
expect(pane._liveQueue).toBeNull();
expect(pane.terminal.write).toHaveBeenCalledWith('held');
expect(screenWrites(pane)).toContain('held');
});
it('the wheel listener is capture-phase, and only a wheel UP can trigger a pull', async () => {