diff --git a/src/web/public/app.js b/src/web/public/app.js index 5e9ccfc7..8acdae7e 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -1483,11 +1483,24 @@ class CodemanApp { * (just at higher cost), and the next reconnect carries the filter via * the SSE query string. */ + /** + * The session id the SSE filter names for `sessionId`: itself, or while the + * tile grid owns the terminal the grid's fixed filter (TILE_GRID_SSE_FILTER, + * constants.js), which no session matches. Both places that set the filter + * ask here: the live re-subscribe below and the connect URL (connectSSE), + * which an SSE reconnect rebuilds with the grid still open. + */ + _sseFilterSessionId(sessionId) { + if (this._tilesOwnTerminal?.()) return window.CodemanTileGrid?.TILE_GRID_SSE_FILTER || sessionId; + return sessionId; + } + _updateSseSubscription(sessionId) { try { + const filterId = this._sseFilterSessionId(sessionId); const body = JSON.stringify({ clientId: this._clientId, - sessions: sessionId ? [sessionId] : null, + sessions: filterId ? [filterId] : null, }); fetch('/api/events/subscribe', { method: 'POST', @@ -1883,7 +1896,8 @@ class CodemanApp { // session we're rendering. Lifecycle/metadata events are sent globally // regardless of filter (server side). const _sseParams = new URLSearchParams({ clientId: this._clientId }); - if (this.activeSessionId) _sseParams.set('sessions', this.activeSessionId); + const _sseFilterId = this._sseFilterSessionId(this.activeSessionId); + if (_sseFilterId) _sseParams.set('sessions', _sseFilterId); this.eventSource = new EventSource(CodemanBase.url(`/api/events?${_sseParams.toString()}`)); // Store all event listeners for cleanup on reconnect. diff --git a/src/web/public/constants.js b/src/web/public/constants.js index 58802118..bdcbd790 100644 --- a/src/web/public/constants.js +++ b/src/web/public/constants.js @@ -1662,6 +1662,13 @@ const TILE_GRID_WIDE_3X1 = 1800; const TILE_SCROLLBACK = 10000; // Tiles have their own per-device font size (a tile is a fraction of the screen). const TILE_FONT_SIZE_DEFAULT = 13; +// What the page's SSE filter names while tiles own the terminal: a value no +// session id takes (ids are UUIDs), so the server, whose filter gates only +// session:terminal batches, sends none. The tiles carry their own output over +// their own sockets, and the parked main terminal only parsed those frames to +// drop them (16 to 18 a second for one busy shell). Every other event still +// arrives (test/sse-tile-grid-filter.test.ts pins the server's side of this). +const TILE_GRID_SSE_FILTER = 'tile-grid'; /** * Columns and rows for `count` tiles, by count (the spec's table), and whether @@ -2156,6 +2163,7 @@ if (typeof window !== 'undefined') { TILE_MIN_H, TILE_SCROLLBACK, TILE_FONT_SIZE_DEFAULT, + TILE_GRID_SSE_FILTER, }; window.CodemanRenderLiveness = { shouldKickRenderer, RENDER_STALL_MS, RENDER_LIVENESS_POLL_MS }; window.CodemanFetchDeadline = { diff --git a/test/sse-tile-grid-filter.test.ts b/test/sse-tile-grid-filter.test.ts new file mode 100644 index 00000000..b9ea7dac --- /dev/null +++ b/test/sse-tile-grid-filter.test.ts @@ -0,0 +1,212 @@ +/** + * @fileoverview The server's side of the tile grid's SSE filter (live server, + * multi-user mode, port 3287). + * + * While tiles own the terminal the page subscribes with TILE_GRID_SSE_FILTER + * (constants.js), an id that names no session, so the server sends it no + * session:terminal batches (test/tile-grid-sse-filter.test.ts covers the page). + * That holds only while: + * - the server takes an id it does not know, on the connect query and on + * POST /api/events/subscribe (a validation change that refused it would + * leave the grid's stream unfiltered, or with no live filter updates at all); + * - the filter gates nothing but terminal batches: in multi-user mode SSE + * routing is fail-closed, so this also checks that session:updated and hook + * events still reach their owner, and only their owner, through the filter. + * + * Every server read of the filter, for the record: the connect route parses + * `?sessions=` (server.ts, GET /api/events), POST /api/events/subscribe + * replaces it (SseStreamManager.updateClientFilter), and its ONE use is + * flushSessionTerminalBatch. broadcast(), the ownership check (canDeliver), the + * heartbeat, the order and tab-layout frames and the shutdown notice never read it. + */ +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'; +import fs from 'node:fs/promises'; +import { readFileSync } from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import vm from 'node:vm'; +import { WebServer } from '../src/web/server.js'; +import { TmuxManager } from '../src/tmux-manager.js'; +import { createUser, invalidateUsersCache } from '../src/user-store.js'; + +vi.spyOn(TmuxManager, 'isTmuxAvailable').mockReturnValue(true); + +const PORT = 3287; +const url = (p: string) => `http://localhost:${PORT}${p}`; +const basic = (u: string, p: string) => 'Basic ' + Buffer.from(`${u}:${p}`).toString('base64'); +const alice = { Authorization: basic('alice', 'alicepass1') }; + +/** The page's own constant, read from constants.js, so the two can never drift. */ +function gridFilter(): string { + const window: Record = {}; + const context = vm.createContext({ window, globalThis: {} }); + vm.runInContext(readFileSync(path.resolve(import.meta.dirname, '../src/web/public/constants.js'), 'utf8'), context); + return window.CodemanTileGrid.TILE_GRID_SSE_FILTER as string; +} +const FILTER = gridFilter(); + +type Received = { event: string; data: unknown }; +/** An open SSE stream for `headers`, collecting every event until close(). */ +async function openStream(query: string, headers: Record) { + const controller = new AbortController(); + const received: Received[] = []; + const res = await fetch(url(`/api/events${query}`), { headers, signal: controller.signal }); + let text = ''; + const reading = (async () => { + const reader = res.body!.getReader(); + try { + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + text += new TextDecoder().decode(value); + let cut: number; + while ((cut = text.indexOf('\n\n')) !== -1) { + const frame = text.slice(0, cut); + text = text.slice(cut + 2); + const event = /^event: (.*)$/m.exec(frame)?.[1]; + const data = /^data: (.*)$/m.exec(frame)?.[1]; + if (event && data !== undefined) { + try { + received.push({ event, data: JSON.parse(data) }); + } catch { + received.push({ event, data }); + } + } + } + } + } catch { + /* aborted */ + } + })(); + return { + status: res.status, + received, + close: async () => { + controller.abort(); + await reading; + }, + }; +} +const wait = (ms: number) => new Promise((r) => setTimeout(r, ms)); +const until = async (fn: () => boolean, ms = 3000) => { + const end = Date.now() + ms; + while (!fn() && Date.now() < end) await wait(20); +}; + +let server: WebServer; +let dataDir: string; +let spacesDir: string; +const saved: Record = {}; +type Internals = { + sessions: Map; + broadcast(event: string, data: unknown): void; + batchTerminalData(sessionId: string, data: string): void; +}; + +beforeAll(async () => { + dataDir = await fs.mkdtemp(path.join(os.tmpdir(), 'sse-grid-data-')); + spacesDir = await fs.mkdtemp(path.join(os.tmpdir(), 'sse-grid-spaces-')); + for (const k of [ + 'CODEMAN_DATA_DIR', + 'CODEMAN_USER_SPACES_DIR', + 'CODEMAN_MULTIUSER', + 'CODEMAN_PASSWORD', + 'CODEMAN_USERNAME', + ]) { + saved[k] = process.env[k]; + } + process.env.CODEMAN_DATA_DIR = dataDir; + process.env.CODEMAN_USER_SPACES_DIR = spacesDir; + process.env.CODEMAN_MULTIUSER = '1'; + delete process.env.CODEMAN_PASSWORD; + delete process.env.CODEMAN_USERNAME; + invalidateUsersCache(); + await createUser({ username: 'root', role: 'admin', password: 'rootpass123' }); + await createUser({ username: 'alice', role: 'user', password: 'alicepass1' }); + await createUser({ username: 'bob', role: 'user', password: 'bobpass1234' }); + server = new WebServer(PORT, false, true); + await server.start(); +}); + +afterAll(async () => { + await server?.stop(); + for (const [k, v] of Object.entries(saved)) { + if (v === undefined) delete process.env[k]; + else process.env[k] = v; + } + invalidateUsersCache(); + await fs.rm(dataDir, { recursive: true, force: true }).catch(() => {}); + await fs.rm(spacesDir, { recursive: true, force: true }).catch(() => {}); +}); + +describe('the tile grid SSE filter, multi-user', () => { + it('is taken on a live re-subscribe though it names no session, and then withholds terminal output', async () => { + // A tile focus re-subscribes the open stream (POST /api/events/subscribe). + // A 204 alone would not do: a validation that dropped the id would still + // answer 204 with the filter cleared, and the stream would carry every + // session's output. So the filter is checked by what it withholds. + expect(FILTER).toBe('tile-grid'); + const internals = server as unknown as Internals; + const grid = await openStream('?clientId=grid-tolerance-1', alice); + const control = await openStream('?clientId=control-1', alice); + expect(grid.status).toBe(200); + await until(() => [grid, control].every((s) => s.received.some((r) => r.event === 'init'))); + const res = await fetch(url('/api/events/subscribe'), { + method: 'POST', + headers: { ...alice, 'Content-Type': 'application/json' }, + body: JSON.stringify({ clientId: 'grid-tolerance-1', sessions: [FILTER] }), + }); + expect(res.status).toBe(204); + internals.sessions.set('alice-2', { id: 'alice-2', owner: 'alice', toState: () => ({ id: 'alice-2' }) }); + try { + internals.batchTerminalData('alice-2', 'output of alice-2'); + await until(() => control.received.some((r) => r.event === 'session:terminal')); + await wait(150); + } finally { + internals.sessions.delete('alice-2'); + } + expect(control.received.some((r) => r.event === 'session:terminal')).toBe(true); + expect(grid.received.some((r) => r.event === 'session:terminal')).toBe(false); + await grid.close(); + await control.close(); + }); + + it("withholds only terminal output: the subscriber still gets its own sessions' updates and hooks", async () => { + const internals = server as unknown as Internals; + const grid = await openStream(`?clientId=grid-1&sessions=${FILTER}`, alice); + const plain = await openStream('?clientId=plain-1', alice); + await until(() => [grid, plain].every((s) => s.received.some((r) => r.event === 'init'))); + + // One session of alice's and one of bob's (spawning is a no-op under test). + const fake = (id: string, owner: string) => ({ id, owner, toState: () => ({ id, owner }) }); + internals.sessions.set('alice-1', fake('alice-1', 'alice')); + internals.sessions.set('bob-1', fake('bob-1', 'bob')); + try { + internals.broadcast('session:updated', { id: 'alice-1', session: { id: 'alice-1', status: 'busy' } }); + internals.broadcast('session:updated', { id: 'bob-1', session: { id: 'bob-1', status: 'busy' } }); + internals.broadcast('hook:idle_prompt', { sessionId: 'alice-1' }); + internals.batchTerminalData('alice-1', 'output of alice-1'); + await until(() => plain.received.some((r) => r.event === 'session:terminal')); + await wait(150); + } finally { + internals.sessions.delete('alice-1'); + internals.sessions.delete('bob-1'); + } + // Only what names the two sessions above (heartbeats and anything else the + // server sends meanwhile are not this test's business). + const kinds = (s: { received: Received[] }) => + s.received + .map((r) => { + const d = (r.data ?? {}) as { id?: string; sessionId?: string }; + return `${r.event} ${d.id ?? d.sessionId}`; + }) + .filter((k) => k.endsWith(' alice-1') || k.endsWith(' bob-1')); + + expect(kinds(grid)).toEqual(['session:updated alice-1', 'hook:idle_prompt alice-1']); + // The same user without the filter got the terminal batch: it was sent. + expect(kinds(plain)).toContain('session:terminal alice-1'); + expect(kinds(plain)).not.toContain('session:updated bob-1'); + await grid.close(); + await plain.close(); + }); +}); diff --git a/test/tile-grid-sse-filter.test.ts b/test/tile-grid-sse-filter.test.ts new file mode 100644 index 00000000..e444f3ee --- /dev/null +++ b/test/tile-grid-sse-filter.test.ts @@ -0,0 +1,87 @@ +/** + * @fileoverview While the tile grid owns the terminal, the page's SSE filter + * names TILE_GRID_SSE_FILTER (constants.js), which no session matches. + * + * The SSE filter gates only session:terminal batches (server side, pinned by + * test/sse-tile-grid-filter.test.ts). With the grid open those frames were for + * the focused tile alone and were only parsed to be dropped: the main terminal + * is parked and the tiles carry their own output over their own sockets. Both + * places that set the filter ask `_sseFilterSessionId()`: the live re-subscribe + * (`_updateSseSubscription`, run by every tile focus) and the connect URL + * (`connectSSE`, rebuilt by every SSE reconnect with the grid still open). + * Leaving the grid gives the filter back to the session shown. + * + * Real code: the shared vm harness (test/mocks/tile-grid-vm.ts). Port: N/A. + */ +import { readFileSync } from 'node:fs'; +import { resolve } from 'node:path'; +import { beforeEach, describe, expect, it } from 'vitest'; +import { fetchSpy, makeGridApp, resetGridHarness, windowStub, type GridApp } from './mocks/tile-grid-vm.js'; + +const IDS = ['s-a', 's-b', 's-c']; +const FILTER = (windowStub.CodemanTileGrid as { TILE_GRID_SSE_FILTER: string }).TILE_GRID_SSE_FILTER; + +/** The app with the REAL _updateSseSubscription, and the filters it posted, in order. */ +function makeApp(): GridApp { + const app = makeGridApp(IDS); + delete app._updateSseSubscription; + app._clientId = 'client-1'; + return app; +} +const posted = () => + fetchSpy.mock.calls + .filter(([url]) => url === '/api/events/subscribe') + .map(([, init]) => JSON.parse((init as { body: string }).body).sessions); + +beforeEach(() => { + resetGridHarness(); + fetchSpy.mockClear(); +}); + +describe('the SSE filter while tiles own the terminal', () => { + it('is a fixed id no session can take', () => { + expect(FILTER).toBe('tile-grid'); + // Session ids are UUIDs. + expect(FILTER).not.toMatch(/^[0-9a-f-]{36}$/); + }); + + it('opening the grid and every tile focus subscribe with it, never a session id', async () => { + const app = makeApp(); + app.openTileGrid(IDS); + delete app.selectSession; + await app.selectSession('s-b'); + await app.selectSession('s-c'); + expect(posted().length).toBeGreaterThanOrEqual(3); + expect(new Set(posted().map((s: string[]) => s.join()))).toEqual(new Set([FILTER])); + }); + + it('leaving the grid gives the filter back to the session shown', () => { + const app = makeApp(); + app.openTileGrid(IDS); + expect(app._sseFilterSessionId('s-a')).toBe(FILTER); + app.closeTileGrid({ keepStored: true, reselect: false }); + // What the single view's selectSession (app.js) then posts. + app._updateSseSubscription('s-a'); + expect(posted().at(-1)).toEqual(['s-a']); + expect(app._sseFilterSessionId('s-b')).toBe('s-b'); + }); + + it('a page with no grid open (a reload into the single view) filters on the session as before', () => { + const app = makeApp(); + expect(app._sseFilterSessionId('s-a')).toBe('s-a'); + expect(app._sseFilterSessionId(null)).toBe(null); + app._updateSseSubscription('s-a'); + expect(posted()).toEqual([['s-a']]); + }); + + it('the connect URL asks the same question, so an SSE reconnect with the grid open keeps the filter', () => { + // connectSSE builds an EventSource, which this harness has none of, so its + // URL building is read from source: one place, through the helper. + const app = readFileSync(resolve(import.meta.dirname, '../src/web/public/app.js'), 'utf8'); + const start = app.indexOf('const _sseParams = new URLSearchParams('); + expect(start).toBeGreaterThan(-1); + const block = app.slice(start, app.indexOf('this.eventSource = new EventSource(', start)); + expect(block).toContain('this._sseFilterSessionId(this.activeSessionId)'); + expect(block).not.toMatch(/set\('sessions', this\.activeSessionId\)/); + }); +});