From fca7acd05f06182e32c2d16f64eb26f033c5c4e2 Mon Sep 17 00:00:00 2001 From: Codeman maintainer Date: Tue, 6 Oct 2026 16:05:56 +0200 Subject: [PATCH] feat(tiles): one load queue for every capture a grid tile fetches GET /api/sessions/:id/terminal runs synchronous tmux calls on the server, so N tiles loading at once would stall every WebSocket and SSE stream back to back (and after a deploy restart all N reopen within the same second). TerminalTile takes the options PR 1 deferred to the grid: - scheduleLoad(tile, kind, run): every capture the tile fetches (initial load, reconnect refresh, server {t:'r'} refresh, shell history pull) runs when its owner says so. Absent (the split's Pane B), a load runs at once. - scrollback (the grid passes TILE_SCROLLBACK) and fontSize. - boundedLoad: a TUI tile loads the bounded full=1&tail= window, never its whole history. TileLoadQueue (terminal-tile.js, DOM-free) is that one queue: concurrency 1, a history pull ahead of background refreshes, then the owner's rank (the grid ranks the focused tile first, then reading order). A destroyed tile's waiting loads are dropped unrun, and destroy() aborts the running fetch so the queue moves on. Also, for Pane B as well: the load now has a deadline covering the body (CodemanFetchDeadline), so a capture that never answers cannot hold the single-flight flag (or the queue) forever; a refresh clears the screen at its turn rather than when it is asked for; and a close while a load only waits in the queue writes the disconnected marker at once. Co-Authored-By: Claude Opus 5.5 (1M context) --- src/web/public/terminal-tile.js | 249 ++++++++++++++-- test/terminal-tile-unit.test.ts | 5 +- test/tile-grid-load-queue.test.ts | 478 ++++++++++++++++++++++++++++++ 3 files changed, 706 insertions(+), 26 deletions(-) create mode 100644 test/tile-grid-load-queue.test.ts diff --git a/src/web/public/terminal-tile.js b/src/web/public/terminal-tile.js index f42adea9..6061a3f7 100644 --- a/src/web/public/terminal-tile.js +++ b/src/web/public/terminal-tile.js @@ -4,8 +4,10 @@ * @fileoverview TerminalTile: one independent live terminal pane bound to one * session, with its own xterm instance and its own * `/ws/sessions/:id/terminal` WebSocket. The split pane (terminal-split.js) - * uses one as its second pane ("Pane B"); the tile grid planned in - * docs/tile-grid-plan.md reuses the same class for every tile. + * uses one as its second pane ("Pane B"); the tile grid (tile-grid.js, + * docs/tile-grid-plan.md) uses one per tile, and feeds every capture they + * fetch through ONE TileLoadQueue (below), because each capture is a + * synchronous tmux call that blocks the server's event loop. * * Deliberately plainer than the primary pane (this.terminal/this._ws in * terminal-ui.js): no local-echo overlay, no CJK IME, no touch/mobile @@ -13,7 +15,7 @@ * docs/split-pane-sessions-plan.md. * * @dependency vendor/xterm.js, vendor/xterm-addon-fit.js - * @dependency constants.js (window.CodemanTerminalFont, DEFAULT_SCROLLBACK, TERMINAL_TAIL_SIZE, TERMINAL_CHUNK_SIZE) + * @dependency constants.js (window.CodemanTerminalFont, window.CodemanFetchDeadline, DEFAULT_SCROLLBACK, TERMINAL_TAIL_SIZE, TERMINAL_CHUNK_SIZE) * @dependency terminal-ui.js (codemanCurrentXtermTheme, codemanCurrentSkinIsLight) * @loadorder 7.4 of 16, loaded after terminal-ui.js and before terminal-split.js */ @@ -72,6 +74,20 @@ // seen by _sendResize() below, or it re-creates the exact PTY-size // fight the split picker already refuses to open at pick time. this.detachedSessions = opts.detachedSessions; + // Lines of scrollback this pane's xterm keeps (the grid passes its smaller + // TILE_SCROLLBACK) and its font size (the grid's own tile font); absent, + // the primary pane's values. + this.scrollback = Number.isFinite(opts.scrollback) ? opts.scrollback : null; + this.fontSize = Number.isFinite(opts.fontSize) ? opts.fontSize : null; + // `scheduleLoad(tile, kind, run)` runs every capture this pane fetches + // (`kind`: 'initial', 'refresh' or 'history') when its owner says so, and + // resolves once `run` has finished or was dropped. The grid passes its one + // queue so N tiles never fetch at once; absent, a load runs straight away. + this._scheduleLoad = typeof opts.scheduleLoad === 'function' ? opts.scheduleLoad : null; + // Loads a BOUNDED window (`full=1&tail=` for a TUI, `tail=` for a shell) + // instead of a TUI's whole history. Grid tiles do; full history is one + // "leave the grid" away in the primary pane. + this.boundedLoad = opts.boundedLoad === true; this.terminal = null; this.fitAddon = null; this.ws = null; @@ -81,6 +97,12 @@ // Single-flight state for _loadBuffer()/_refreshBuffer() below. this._bufferLoading = false; this._bufferRefreshPending = false; + // True only while a load's work runs, not while it waits in the owner's + // queue (see _runLoad): a close during the wait writes its marker at once. + this._loadRunning = false; + // Aborts the running load's fetch; destroy() uses it so a removed tile + // does not hold the owner's queue for a whole deadline. + this._loadAbort = null; // Scroll-to-top history pull (shell panes only), see _maybeLoadMoreHistory(). // `_liveQueue` is non-null from the pull's response until its finally // block: live frames are held there with their arrival time instead of @@ -119,7 +141,7 @@ } async connect() { - const savedFontSize = parseInt(localStorage.getItem('codeman-font-size'), 10); + const savedFontSize = this.fontSize ?? parseInt(localStorage.getItem('codeman-font-size'), 10); this.terminal = new Terminal({ theme: { ...global.codemanCurrentXtermTheme() }, fontFamily: global.CodemanTerminalFont.resolve(this.fontSettings.terminalFontFamily), @@ -129,7 +151,7 @@ cursorBlink: false, cursorStyle: 'block', minimumContrastRatio: global.codemanCurrentSkinIsLight() ? 4.5 : 1, - scrollback: DEFAULT_SCROLLBACK, + scrollback: this.scrollback ?? DEFAULT_SCROLLBACK, allowTransparency: true, allowProposedApi: true, }); @@ -453,7 +475,7 @@ const code = event?.code; const permanent = TerminalTile.STOP_MARKERS[code]; this._markerText = permanent || TerminalTile.MARKER_RECONNECTING; - if (this._bufferLoading) this._markerOwed = true; + if (this._loadRunning) this._markerOwed = true; else this._writeDisconnectedMarker(); if (this._destroyed) return; if (permanent) { @@ -571,21 +593,82 @@ // interleave their chunks into one terminal. A second call while one is // in flight is dropped here; _refreshBuffer() is the caller that queues // a trailing re-run instead. - async _loadBuffer() { + async _loadBuffer({ refresh = false } = {}) { if (this._bufferLoading) return; this._bufferLoading = true; - try { - const query = this.sessionMode === 'shell' ? `tail=${TERMINAL_TAIL_SIZE}` : 'full=1'; - const res = await fetch(`/api/sessions/${this.sessionId}/terminal?${query}`); - const payload = (await res.json())?.data ?? {}; - if (payload.terminalBuffer && this.terminal) { - await writeChunked(this.terminal, payload.terminalBuffer, () => this._destroyed); + await this._runLoad(refresh ? 'refresh' : 'initial', async () => { + this._loadRunning = true; + try { + if (this._destroyed) return; + if (refresh) { + // Cleared at the load's turn, not when it was asked for: a grid tile + // waiting in the queue keeps its last frame instead of sitting blank. + 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; + } + const shell = this.sessionMode === 'shell'; + let query = shell ? `tail=${TERMINAL_TAIL_SIZE}` : 'full=1'; + if (this.boundedLoad && !shell) query = `full=1&tail=${TERMINAL_TAIL_SIZE}`; + // A deadline covering the body as well as the headers (the primary + // pane's budgets, CodemanFetchDeadline): a capture that never answers + // would otherwise hold this pane's single-flight flag, and in the grid + // the one load queue every tile waits behind, forever. + const controller = global.AbortController ? new global.AbortController() : null; + this._loadAbort = controller; + const budget = global.CodemanFetchDeadline?.terminalFetchDeadlineMs?.({ full: !shell }) ?? 45000; + const timer = controller ? setTimeout(() => controller.abort(), budget) : null; + let payload; + try { + const res = await fetch( + `/api/sessions/${this.sessionId}/terminal?${query}`, + controller ? { signal: controller.signal } : undefined + ); + payload = (await res.json())?.data ?? {}; + } finally { + clearTimeout(timer); + this._loadAbort = null; + } + if (payload.terminalBuffer && this.terminal) { + await writeChunked(this.terminal, payload.terminalBuffer, () => this._destroyed); + } + } catch { + /* Best-effort: live output still arrives once the socket connects. */ + } finally { + this._loadRunning = false; + this._stampMarkerIfOwed(); + this._endBufferLoad(); } + }); + } + + // Runs a load's work now, or when the owner's queue gives this pane its turn + // (`scheduleLoad`). The single-flight flag is already set by the caller, so a + // load waiting in the queue still coalesces refreshes and blocks a second + // pull; the work itself sets `_loadRunning`. A load the queue drops (this + // pane was destroyed while it waited) never runs, so its flags are released + // here. + async _runLoad(kind, work) { + let ran = false; + const run = () => { + ran = true; + return work(); + }; + if (!this._scheduleLoad) { + await run(); + return; + } + try { + await this._scheduleLoad(this, kind, run); } catch { - /* Best-effort — live output still arrives once the socket connects. */ - } finally { - this._stampMarkerIfOwed(); - this._endBufferLoad(); + /* The queue never rejects; a load that failed already settled itself. */ + } + if (!ran) { + this._bufferLoading = false; + this._bufferRefreshPending = false; } } @@ -656,7 +739,8 @@ const now = Date.now(); if (now - this._historyPullAt < cooldown) return; this._historyPullAt = now; - void this._pullHistory(); + this._bufferLoading = true; + void this._runLoad('history', () => this._pullHistory()); } // Pulls a BOUNDED window of tmux's full history (the same TERMINAL_TAIL_SIZE @@ -665,6 +749,11 @@ // single-flight flag across the fetch AND the replay, like _loadBuffer(). async _pullHistory() { this._bufferLoading = true; + if (this._destroyed) { + this._endBufferLoad(); + return; + } + this._loadRunning = true; let replayed = false; let capturedAt = 0; // Two budgets on one signal. The request itself gets the primary pane's @@ -678,6 +767,7 @@ // be re-armed, hence the controller; without AbortController the pull // simply has no deadline. const controller = global.AbortController ? new global.AbortController() : null; + this._loadAbort = controller; let abortTimer = null; const armDeadline = (ms) => { if (!controller) return; @@ -743,6 +833,8 @@ /* Best-effort — live output keeps arriving whatever happens here. */ } finally { clearTimeout(abortTimer); + this._loadAbort = null; + this._loadRunning = false; const queued = this._liveQueue ?? []; this._liveQueue = null; // After a replay, only frames that arrived after the capture are news; @@ -776,12 +868,7 @@ this._bufferRefreshPending = true; 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(); + void this._loadBuffer({ refresh: true }); } // Local reflow only — no PTY resize frame. Split out so a divider drag @@ -851,6 +938,14 @@ // keystroke loses nothing. clearTimeout(this._reconnectTimer); this._reconnectTimer = null; + // A load still fetching would otherwise hold the owner's queue (and the + // server's attention) for a pane nobody can see any more. + try { + this._loadAbort?.abort(); + } catch { + /* Already settled. */ + } + this._loadAbort = null; if (this._onWheel) { this.mountEl?.removeEventListener('wheel', this._onWheel, { capture: true }); this._onWheel = null; @@ -883,5 +978,111 @@ 4010: '[disconnected: another connection took over this pane]', }; + /** + * ONE queue for every capture a set of tiles fetches (`GET + * /api/sessions/:id/terminal`): the initial load, the refresh after a + * reconnect, a server `{t:'r'}` refresh and the shell history pull. Each + * capture runs synchronous tmux calls on the server, so N of them at once do + * not run in parallel, they stall every WebSocket and SSE stream on it back to + * back. After a deploy restart all N tiles reopen within the same second; this + * drains their refreshes one at a time. + * + * Concurrency 1. Next up is a history pull (the user is waiting on it), then + * the lowest `rank(tile)` (the grid ranks the focused tile first, then reading + * order), then arrival order. A destroyed tile's entries are dropped, never run. + * DOM-free, so the grid owns the policy and tests drive it directly. + */ + class TileLoadQueue { + /** + * @param {{rank?: (tile: object) => number, onChange?: (tile: object, state: 'queued'|'running'|'idle') => void}} [opts] + */ + constructor(opts = {}) { + this._rank = typeof opts.rank === 'function' ? opts.rank : () => 0; + this._onChange = typeof opts.onChange === 'function' ? opts.onChange : null; + this._pending = []; + this._active = null; + this._seq = 0; + } + + /** The `scheduleLoad` a TerminalTile takes. Resolves once `run` finished or was dropped; never rejects. */ + schedule(tile, kind, run) { + return new Promise((resolve) => { + this._pending.push({ tile, kind, run, resolve, seq: this._seq++ }); + this._notify(tile, 'queued'); + this._pump(); + }); + } + + /** Drops every load still waiting for `tile` (the running one, if any, finishes on its own). */ + drop(tile) { + const keep = []; + for (const entry of this._pending) { + if (entry.tile === tile) entry.resolve(); + else keep.push(entry); + } + this._pending = keep; + if (this._active?.tile !== tile) this._notify(tile, 'idle'); + } + + /** How many loads are waiting (not counting the running one). */ + get size() { + return this._pending.length; + } + + /** The tile whose load is running, or null. */ + get activeTile() { + return this._active?.tile ?? null; + } + + _notify(tile, state) { + try { + this._onChange?.(tile, state); + } catch { + /* A display callback never stops the queue. */ + } + } + + _takeNext() { + let best = -1; + let bestKey = null; + for (let i = 0; i < this._pending.length; i++) { + const entry = this._pending[i]; + if (entry.tile?._destroyed) continue; + const key = [entry.kind === 'history' ? 0 : 1, this._rank(entry.tile), entry.seq]; + const order = bestKey ? key[0] - bestKey[0] || key[1] - bestKey[1] || key[2] - bestKey[2] : -1; + if (order < 0) { + best = i; + bestKey = key; + } + } + // Destroyed tiles' entries go now, resolved but never run. + const dropped = this._pending.filter((entry) => entry.tile?._destroyed); + const next = best === -1 ? null : this._pending[best]; + this._pending = this._pending.filter((entry) => entry !== next && !entry.tile?._destroyed); + for (const entry of dropped) entry.resolve(); + return next; + } + + async _pump() { + if (this._active) return; + const entry = this._takeNext(); + if (!entry) return; + this._active = entry; + this._notify(entry.tile, 'running'); + try { + await entry.run(); + } catch { + /* A load settles its own failure; the queue only moves on. */ + } finally { + this._active = null; + const stillQueued = this._pending.some((e) => e.tile === entry.tile); + this._notify(entry.tile, stillQueued ? 'queued' : 'idle'); + entry.resolve(); + void this._pump(); + } + } + } + global.TerminalTile = TerminalTile; + global.TileLoadQueue = TileLoadQueue; })(window); diff --git a/test/terminal-tile-unit.test.ts b/test/terminal-tile-unit.test.ts index 8a77470e..17bb55de 100644 --- a/test/terminal-tile-unit.test.ts +++ b/test/terminal-tile-unit.test.ts @@ -221,7 +221,8 @@ describe('TerminalTile server-refresh single-flight', () => { await settle(); expect(pane.terminal.clear).toHaveBeenCalledTimes(1); - expect(fetchMock).toHaveBeenCalledWith('/api/sessions/s1/terminal?full=1'); + // The second argument carries the load's deadline (an AbortSignal). + expect(fetchMock).toHaveBeenCalledWith('/api/sessions/s1/terminal?full=1', expect.anything()); expect(pane.terminal.write).toHaveBeenCalledWith('one'); expect(pane._bufferLoading).toBe(false); }); @@ -233,7 +234,7 @@ describe('TerminalTile server-refresh single-flight', () => { pane._refreshBuffer(); await settle(); - expect(fetchMock).toHaveBeenCalledWith(`/api/sessions/s1/terminal?tail=${1024 * 1024}`); + expect(fetchMock).toHaveBeenCalledWith(`/api/sessions/s1/terminal?tail=${1024 * 1024}`, expect.anything()); }); it('refreshes arriving mid-fetch neither clear nor fetch again, and run ONCE after the replay lands', async () => { diff --git a/test/tile-grid-load-queue.test.ts b/test/tile-grid-load-queue.test.ts new file mode 100644 index 00000000..ddb84fb1 --- /dev/null +++ b/test/tile-grid-load-queue.test.ts @@ -0,0 +1,478 @@ +/** + * @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'}` goes through the same queue, and a tile waiting its turn keeps its + * last frame (the clear happens at its turn); + * - 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. 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'; + +class FakeSocket { + static instances: FakeSocket[] = []; + readyState = 0; + sent: Array> = []; + onopen: (() => void) | null = null; + onmessage: ((ev: { data: string }) => void) | null = null; + onclose: ((ev?: { code: number }) => void) | null = null; + onerror: (() => void) | null = null; + constructor(public url: string) { + FakeSocket.instances.push(this); + } + send(data: string) { + this.sent.push(JSON.parse(data)); + } + close = vi.fn(() => { + this.readyState = 3; + }); + open() { + this.readyState = 1; + this.onopen?.(); + } + receive(msg: object) { + this.onmessage?.({ data: JSON.stringify(msg) }); + } + drop(code = 1006) { + this.readyState = 3; + this.onclose?.({ code }); + } +} + +class FakeFit { + term: FakeTerminal | null = null; + fit() {} + proposeDimensions() { + return { cols: 80, rows: 24 }; + } +} + +class FakeTerminal { + options: Record; + cols = 80; + rows = 24; + buffer = { active: { type: 'normal', viewportY: 0, length: 24 } }; + writes: string[] = []; + constructor(options: Record) { + this.options = { ...options }; + } + loadAddon(addon: FakeFit) { + addon.term = this; + } + open() {} + onData() {} + attachCustomKeyEventHandler() {} + registerLinkProvider() {} + textarea = { addEventListener() {}, removeEventListener() {} }; + write(data: string, cb?: () => void) { + this.writes.push(data); + cb?.(); + } + clear() { + this.writes.push(''); + } + resize(cols: number, rows: number) { + this.cols = cols; + this.rows = rows; + } + scrollToLine() {} + scrollToTop() {} + dispose() {} +} + +/** 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 = { + 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) => globalThis.clearTimeout(id), + requestAnimationFrame: vi.fn(), + HTMLCanvasElement: class HTMLCanvasElement {}, + WebSocket: FakeSocket, + Terminal: FakeTerminal, + FitAddon: { FitAddon: FakeFit }, + fetch: (...args: Parameters) => 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; + +type Tile = { + sessionId: string; + connect(): Promise; + destroy(): void; + reconnectNow(): void; + _maybeLoadMoreHistory(): void; + _destroyed: boolean; + terminal: FakeTerminal | null; + ws: FakeSocket | null; +}; +type Queue = { + schedule(tile: object, kind: string, run: () => Promise): Promise; + 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; + 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 } = {}) { + 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) => 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}`, + `/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 cleared: it shows its last frame until its load runs. + expect(b.terminal?.writes).not.toContain(''); + + await drain('fresh'); + expect(captures.map((c) => c.url.split('/')[3])).toEqual(['a', 'b']); + expect(b.terminal?.writes.slice(-2)).toEqual(['', 'fresh']); + }); + + 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}`, + `/api/sessions/sh/terminal?full=1&tail=${TAIL}`, + `/api/sessions/b/terminal?full=1&tail=${TAIL}`, + ]); + }); + + 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 turn cleared the screen, so the marker is written again below the replay: one on screen. + const writes = b.terminal?.writes ?? []; + expect(writes.slice(writes.lastIndexOf(''))).toEqual([ + '', + '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('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((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((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(); + }); +});