Files
Codeman/test/tile-grid-load-queue.test.ts
T
Codeman maintainer eb5d982c38 fix(tiles): refresh fetches first, then resets in-stream and replays
A tile's refresh (a {t:'r'} or {t:'c'} frame, every reconnect) wiped the
pane with a synchronous xterm clear() at the load's turn, BEFORE its fetch,
and wrote live frames straight through the fetch and the replay. That is
the replay clear CLAUDE.md "Terminal resilience" forbids: bytes still
queued in xterm are parsed after a synchronous clear and fuse into the
snapshot, and clear() keeps the cursor's row, column, SGR and margins, so
the capture (raw rows, no home) started wherever the cursor sat. A failed
or empty fetch left the tile blank.

The refresh now runs in the primary pane's order (_onSessionNeedsRefresh,
_resetTerminalForReplay):
- fetch first, so the tile keeps its last frame through the round trip and
  through a grid tile's wait in the load queue;
- from the response on, live frames are held in _liveQueue with their
  arrival time, as _pullHistory already did, and the body read of a bounded
  window (grid tile, shell) gets the pull's 10 s budget, while Pane B's
  unbounded full=1 keeps the request's own budget;
- then the queued in-stream \x1bc immediately before the replay;
- then the held frames that arrived after the response (_flushLiveQueue,
  now shared with _pullHistory), then the owed marker.
A failed, aborted or empty fetch writes nothing and resets nothing.

The _stampMarkerIfOwed guard for a pending trailing refresh stays (that
refresh settles the marker itself either way); only its rationale changed.
The fake xterm now treats an in-stream RIS like clear() in its row
emulation.

Tests: the ones that counted clear() calls on the refresh path now count
the in-stream reset instead, assert it sits right before the replay and
that clear() is never called (unit single-flight block, the marker
ordering tests, the reconnect test, the grid {t:'r'} and marker tests, and
the scroll test's server-clear overflow case, which now goes through a
refresh). New: the screen is untouched on a failed or empty fetch and on a
failed body read (held frames written in order), frames before the
response are written through and later ones held behind the replay, the
cutoff drops frames the capture covers, the body budgets, and a grid tile
keeps its last frame through its own capture's round trip.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-10-09 09:18:32 +02:00

477 lines
18 KiB
TypeScript

/**
* @fileoverview Every capture a grid tile fetches goes through ONE queue.
*
* `GET /api/sessions/:id/terminal` runs synchronous tmux calls on the server, so
* N tiles loading at once do not load in parallel: they stall every WebSocket
* and SSE stream on the server back to back. The grid therefore hands each
* TerminalTile a `scheduleLoad` (a TileLoadQueue, terminal-tile.js) and the
* tile routes EVERY capture through it: the initial load, the refresh after a
* reconnect, a server `{t:'r'}` refresh and the shell history pull.
*
* Pinned here, with `connect()` and the socket handlers running for real:
* - at most one `/terminal` fetch is in flight at a time, for initial loads and
* for N tiles reconnecting together;
* - the focused tile goes first, then reading order, and a history pull (the
* user is waiting on it) jumps ahead of background refreshes;
* - `{t:'r'}` and `{t:'c'}` go through the same queue, and a tile keeps its
* last frame while it waits its turn and through its own capture's round
* trip (it is reset in-stream only once the capture is in hand);
* - a destroyed tile's queued load is dropped, and destroying the tile whose
* load is running aborts its fetch so the queue moves on;
* - a load that never answers is cut off by its deadline;
* - a close while a load only WAITS writes the disconnected marker at once;
* - grid tiles load a bounded window and keep TILE_SCROLLBACK lines.
*
* Real code under test: constants.js + app.js + terminal-ui.js +
* terminal-tile.js in one `vm` context; xterm, the fit addon and WebSocket are
* fakes (test/mocks/terminal-tile-fakes.ts). Port: N/A.
*/
import { readFileSync } from 'node:fs';
import { performance } from 'node:perf_hooks';
import { resolve } from 'node:path';
import vm from 'node:vm';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { FakeFit, FakeSocket, FakeTerminal } from './mocks/terminal-tile-fakes.js';
/** One `/terminal` fetch the test answers (or lets hang) by hand. */
type Capture = { url: string; settled: boolean; aborted: boolean; answer(body: string): void };
let captures: Capture[] = [];
const fetchMock = vi.fn((url: string, init?: { signal?: AbortSignal }) => {
return new Promise((resolveFetch, rejectFetch) => {
const capture: Capture = {
url,
settled: false,
aborted: false,
answer(body: string) {
capture.settled = true;
resolveFetch({ ok: true, status: 200, json: async () => ({ data: { terminalBuffer: body } }) });
},
};
init?.signal?.addEventListener('abort', () => {
capture.settled = true;
capture.aborted = true;
rejectFetch(new Error('aborted'));
});
captures.push(capture);
});
});
const inFlight = () => captures.filter((c) => !c.settled).length;
const read = (f: string) => readFileSync(resolve(import.meta.dirname, `../src/web/public/${f}`), 'utf8');
const windowStub: Record<string, unknown> = {
addEventListener: vi.fn(),
removeEventListener: vi.fn(),
CodemanBase: { base: '' },
AbortController,
};
const context = vm.createContext({
console: { ...console, log: vi.fn(), debug: vi.fn() },
performance,
setInterval: vi.fn(),
clearInterval: vi.fn(),
// Late-bound so vi.useFakeTimers() reaches code running in this context.
setTimeout: (fn: () => void, ms?: number) => globalThis.setTimeout(fn, ms),
clearTimeout: (id: ReturnType<typeof setTimeout>) => globalThis.clearTimeout(id),
requestAnimationFrame: vi.fn(),
HTMLCanvasElement: class HTMLCanvasElement {},
WebSocket: FakeSocket,
Terminal: FakeTerminal,
FitAddon: { FitAddon: FakeFit },
fetch: (...args: Parameters<typeof fetchMock>) => fetchMock(...args),
location: { protocol: 'http:', host: 'codeman.test' },
document: { addEventListener: vi.fn(), documentElement: { dataset: {} } },
localStorage: { length: 0, key: vi.fn(), getItem: vi.fn(), setItem: vi.fn(), removeItem: vi.fn() },
window: windowStub,
MobileDetection: { isTouchDevice: () => false, isHandheldDevice: () => false, getDeviceType: () => 'desktop' },
});
vm.runInContext(
`${read('constants.js')}\n${read('app.js')}\n${read('terminal-ui.js')}\n${read('terminal-tile.js')}\n` +
'globalThis.__CodemanApp = CodemanApp;',
context
);
const CodemanApp = (context as unknown as { __CodemanApp: { prototype: object } }).__CodemanApp;
const TileGrid = windowStub.CodemanTileGrid as { TILE_SCROLLBACK: number };
const TAIL = 1024 * 1024;
// A grid tile's full captures read no more tmux history than its xterm keeps:
// TILE_SCROLLBACK plus the screen (the fake terminal has 24 rows).
const LINES = `&lines=${TileGrid.TILE_SCROLLBACK + 24}`;
type Tile = {
sessionId: string;
connect(): Promise<void>;
destroy(): void;
reconnectNow(): void;
_maybeLoadMoreHistory(): void;
_destroyed: boolean;
terminal: FakeTerminal | null;
ws: FakeSocket | null;
};
type Queue = {
schedule(tile: object, kind: string, run: () => Promise<void>): Promise<void>;
size: number;
activeTile: object | null;
};
const TerminalTile = windowStub.TerminalTile as new (id: string, mount: unknown, opts?: object) => Tile;
const TileLoadQueue = windowStub.TileLoadQueue as new (opts?: object) => Queue;
function makeApp() {
const app = Object.create(CodemanApp.prototype) as Record<string, unknown>;
app._clientId = 'c-test';
app._wsTabNonce = 'nonce-1';
app._seqCounters = new Map();
app._pendingDeliveries = new Map();
app._postDraining = new Set();
app._extraInputSockets = new Map();
app._persistReliableState = vi.fn();
app._persistReliableNow = vi.fn();
app._updateConnectionIndicator = vi.fn();
app.loadAppSettingsFromStorage = () => ({});
app._estimateReplayRows = (text: string) => text.split('\n').length;
return app;
}
const liveTiles: Tile[] = [];
/** Grid-style tiles sharing one queue: the focused id ranks first, then the order given. */
function makeGrid(ids: string[], { focused = ids[0], modes = {} as Record<string, string> } = {}) {
windowStub.app = makeApp();
const order = [...ids];
const states: Array<[string, string]> = [];
const queue = new TileLoadQueue({
rank: (tile: Tile) => (tile.sessionId === focused ? -1 : order.indexOf(tile.sessionId)),
onChange: (tile: Tile, state: string) => states.push([tile.sessionId, state]),
});
const tiles = ids.map((id) => {
const tile = new TerminalTile(
id,
{ addEventListener: vi.fn(), removeEventListener: vi.fn() },
{
mode: modes[id] ?? 'claude',
scheduleLoad: (t: object, kind: string, run: () => Promise<void>) => queue.schedule(t, kind, run),
scrollback: TileGrid.TILE_SCROLLBACK,
fontSize: 13,
boundedLoad: true,
}
);
liveTiles.push(tile);
return tile;
});
return { tiles, queue, states };
}
/** Let promise chains (fetch, json, chunked write, queue pump) run. */
async function settle() {
for (let i = 0; i < 20; i++) await Promise.resolve();
}
/** Answer every capture as it arrives, one at a time, asserting the queue never overlaps two. */
async function drain(body = 'frame') {
let max = 0;
for (let guard = 0; guard < 50; guard++) {
await settle();
max = Math.max(max, inFlight());
const open = captures.find((c) => !c.settled);
if (!open) break;
open.answer(body);
}
return max;
}
beforeEach(() => {
captures = [];
fetchMock.mockClear();
FakeSocket.instances = [];
});
afterEach(() => {
for (const tile of liveTiles.splice(0)) tile.destroy();
vi.useRealTimers();
});
describe('initial loads', () => {
it('N tiles connecting together fetch one capture at a time', async () => {
const { tiles } = makeGrid(['a', 'b', 'c', 'd']);
const connecting = tiles.map((t) => t.connect());
await settle();
expect(inFlight()).toBe(1);
expect(await drain()).toBe(1);
await Promise.all(connecting);
expect(captures).toHaveLength(4);
});
it('runs the focused tile first, then reading order', async () => {
const { tiles } = makeGrid(['a', 'b', 'c', 'd'], { focused: 'c' });
const connecting = tiles.map((t) => t.connect());
await drain();
await Promise.all(connecting);
// `a` was already running when the others arrived; then focus, then order.
expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a', 'c', 'b', 'd']);
});
it('reports each tile queued, then running, then idle (the quiet loading state)', async () => {
const { tiles, states } = makeGrid(['a', 'b']);
const connecting = tiles.map((t) => t.connect());
await drain();
await Promise.all(connecting);
expect(states.filter(([id]) => id === 'b').map(([, s]) => s)).toEqual(['queued', 'running', 'idle']);
});
it('loads a bounded window: full=1&tail= for a TUI, tail= for a shell', async () => {
const { tiles } = makeGrid(['tui', 'sh'], { modes: { sh: 'shell' } });
const connecting = tiles.map((t) => t.connect());
await drain();
await Promise.all(connecting);
expect(captures.map((c) => c.url)).toEqual([
`/api/sessions/tui/terminal?full=1&tail=${TAIL}${LINES}`,
`/api/sessions/sh/terminal?tail=${TAIL}`,
]);
});
it('keeps TILE_SCROLLBACK lines and the tile font size', async () => {
const { tiles } = makeGrid(['a']);
const connecting = tiles[0].connect();
await drain();
await connecting;
expect(tiles[0].terminal?.options.scrollback).toBe(10000);
expect(tiles[0].terminal?.options.fontSize).toBe(13);
});
});
/** Connects every tile and opens its socket, with the queue drained. */
async function connectAll(tiles: Tile[]) {
const connecting = tiles.map((t) => t.connect());
await drain('first');
await Promise.all(connecting);
for (const tile of tiles) tile.ws?.open();
captures = [];
}
describe('refreshes', () => {
it('N tiles reconnecting together (a deploy restart) refresh one at a time', async () => {
const { tiles } = makeGrid(['a', 'b', 'c', 'd', 'e', 'f']);
await connectAll(tiles);
for (const tile of tiles) tile.ws?.drop(1006);
for (const tile of tiles) {
tile.reconnectNow();
tile.ws?.open();
}
await settle();
expect(inFlight()).toBe(1);
expect(await drain('after')).toBe(1);
expect(captures).toHaveLength(6);
});
it('a server {t:"r"} refresh waits its turn behind another tile, keeping its last frame meanwhile', async () => {
const { tiles } = makeGrid(['a', 'b']);
await connectAll(tiles);
const [a, b] = tiles;
a.ws?.receive({ t: 'r' });
b.ws?.receive({ t: 'r' });
await settle();
expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a']);
// b has not been reset: it shows its last frame until its load runs.
const before = [...(b.terminal?.writes ?? [])];
expect(before).not.toContain('<CLEAR>');
await drain('fresh');
expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a', 'b']);
// Its turn reset it in-stream, right before the replay, never with clear().
expect(b.terminal?.writes).toEqual([...before, '\x1bc', 'fresh']);
});
it("a server {t:'c'} (a Claude pane's first prompt) is a refresh through the same queue", async () => {
const { tiles } = makeGrid(['a', 'b']);
await connectAll(tiles);
const [a, b] = tiles;
a.ws?.receive({ t: 'r' });
b.ws?.receive({ t: 'c' });
await settle();
// b waits behind a, like any refresh: one capture in flight.
expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a']);
expect(await drain('banner')).toBe(1);
expect(captures.map((c) => c.url)).toEqual([
`/api/sessions/a/terminal?full=1&tail=${TAIL}${LINES}`,
`/api/sessions/b/terminal?full=1&tail=${TAIL}${LINES}`,
]);
expect(b.terminal?.writes.at(-1)).toBe('banner');
});
it("a tile keeps its last frame through its own capture's round trip, and one cut off leaves it as it was", async () => {
vi.useFakeTimers();
const { tiles } = makeGrid(['a']);
await connectAll(tiles);
const [a] = tiles;
a.ws?.receive({ t: 'o', d: 'last frame' });
a.ws?.receive({ t: 'r' });
await settle();
// Its turn came and its capture is in flight: nothing reset yet.
expect(inFlight()).toBe(1);
expect(a.terminal?.writes.at(-1)).toBe('last frame');
// The full-capture budget runs out with no answer.
await vi.advanceTimersByTimeAsync(45_000);
await settle();
expect(captures.at(-1)?.aborted).toBe(true);
expect(a.terminal?.writes.at(-1)).toBe('last frame');
expect(a.terminal?.writes).not.toContain('\x1bc');
expect(a.terminal?.writes).not.toContain('<CLEAR>');
});
it('a history pull jumps ahead of background refreshes', async () => {
const { tiles } = makeGrid(['a', 'b', 'sh'], { modes: { sh: 'shell' } });
await connectAll(tiles);
const [a, b, sh] = tiles;
a.ws?.receive({ t: 'r' });
b.ws?.receive({ t: 'r' });
sh._maybeLoadMoreHistory();
await settle();
expect(inFlight()).toBe(1);
await drain();
expect(captures.map((c) => c.url)).toEqual([
`/api/sessions/a/terminal?full=1&tail=${TAIL}${LINES}`,
`/api/sessions/sh/terminal?full=1&tail=${TAIL}${LINES}`,
`/api/sessions/b/terminal?full=1&tail=${TAIL}${LINES}`,
]);
});
it('a close while the refresh only WAITS writes the marker at once, and once', async () => {
const { tiles } = makeGrid(['a', 'b']);
await connectAll(tiles);
const [a, b] = tiles;
a.ws?.receive({ t: 'r' });
b.ws?.receive({ t: 'r' });
await settle();
b.ws?.drop(1006);
const markers = () => (b.terminal?.writes ?? []).filter((w) => w.includes('[disconnected')).length;
expect(markers()).toBe(1);
await drain('fresh');
// Its replay reset the screen, so the marker is written again below the replay: one on screen.
const writes = b.terminal?.writes ?? [];
expect(writes).not.toContain('<CLEAR>');
expect(writes.slice(writes.lastIndexOf('\x1bc'))).toEqual([
'\x1bc',
'fresh',
expect.stringContaining('[disconnected'),
]);
});
});
describe('teardown', () => {
it("drops a destroyed tile's queued load: it never fetches", async () => {
const { tiles, queue } = makeGrid(['a', 'b', 'c']);
const connecting = tiles.map((t) => t.connect());
await settle();
tiles[1].destroy();
await drain();
await Promise.all(connecting);
expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a', 'c']);
expect(queue.size).toBe(0);
expect(tiles[1].ws).toBeNull();
});
it('destroying the tile whose load is running aborts its fetch and the queue moves on', async () => {
const { tiles } = makeGrid(['a', 'b']);
const connecting = tiles.map((t) => t.connect());
await settle();
tiles[0].destroy();
await settle();
expect(captures[0].aborted).toBe(true);
expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a', 'b']);
await drain();
await Promise.all(connecting);
});
it('destroying a tile while xterm still parses its replay settles the load, and the queue moves on', async () => {
// The replay waits for xterm's write callback (writeChunked), which a
// disposed xterm never runs: unsettled, the tile's load would hold the
// grid's one queue, and every other tile behind it, forever.
const { tiles } = makeGrid(['a', 'b']);
const connecting = tiles.map((t) => t.connect());
await settle();
tiles[0].terminal!.holdParse = true;
captures[0].answer('replay of a');
await settle();
// a's replay is still parsing: b waits its turn.
expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a']);
tiles[0].destroy();
await settle();
expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a', 'b']);
captures[1].answer('b');
await settle();
await Promise.all(connecting);
});
it('a capture that never answers is cut off by its deadline, and the next tile loads', async () => {
vi.useFakeTimers();
const { tiles } = makeGrid(['a', 'b']);
const connecting = tiles.map((t) => t.connect());
await settle();
expect(captures).toHaveLength(1);
// The full-capture budget (CodemanFetchDeadline, 45 s for a TUI).
await vi.advanceTimersByTimeAsync(45_000);
await settle();
expect(captures[0].aborted).toBe(true);
expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a', 'b']);
captures[1].answer('b');
await settle();
await Promise.all(connecting);
});
});
describe('TileLoadQueue on its own', () => {
it('never rejects, and moves on when a load throws', async () => {
const queue = new TileLoadQueue();
const order: string[] = [];
const first = queue.schedule({}, 'initial', async () => {
order.push('first');
throw new Error('boom');
});
const second = queue.schedule({}, 'initial', async () => {
order.push('second');
});
await expect(first).resolves.toBeUndefined();
await expect(second).resolves.toBeUndefined();
expect(order).toEqual(['first', 'second']);
expect(queue.activeTile).toBeNull();
});
it('resolves, but never runs, a load whose tile was destroyed while it waited', async () => {
const queue = new TileLoadQueue();
let release: () => void = () => {};
const busy = queue.schedule({}, 'initial', () => new Promise<void>((r) => (release = r)));
const gone = { _destroyed: false };
const run = vi.fn(async () => {});
const waiting = queue.schedule(gone, 'refresh', run);
gone._destroyed = true;
release();
await busy;
await expect(waiting).resolves.toBeUndefined();
expect(run).not.toHaveBeenCalled();
expect(queue.size).toBe(0);
});
it('drop() resolves every load still waiting for a tile', async () => {
const queue = new TileLoadQueue();
let release: () => void = () => {};
const busy = queue.schedule({}, 'initial', () => new Promise<void>((r) => (release = r)));
const tile = {};
const run = vi.fn(async () => {});
const waiting = queue.schedule(tile, 'refresh', run);
(queue as unknown as { drop(t: object): void }).drop(tile);
await expect(waiting).resolves.toBeUndefined();
release();
await busy;
expect(run).not.toHaveBeenCalled();
});
});