/** * W2 — integration tests for the server-side PTY-inject queue routes and the * idle-drain wiring, against a real startServer. * * Sandbox-runnable (no PTY): Origin/CSRF guard, id/text validation, size + depth * caps, rate limit, QUEUE_ENABLED=0 → 503, unknown-session 404. * * itPty-gated (needs a real shell — auto-skips where posix_spawn is blocked): * - happy path: POST queue → 200 {length:1}; GET /live-sessions shows * queueLength; GET/DELETE .../queue. * - idle-drain: enqueue → POST Stop → the marker echoes through the PTY. * - one-per-idle pacing: two entries, one Stop drains only the first. * - settle guard: a non-idle hook during the window cancels the drain. */ import net from 'node:net' import { afterEach, describe, expect, it } from 'vitest' import WebSocket from 'ws' import * as nodePty from 'node-pty' import { loadConfig } from '../../src/config.js' import { startServer } from '../../src/server.js' const PTY_AVAILABLE = (() => { try { const p = nodePty.spawn(process.env['SHELL'] ?? '/bin/sh', [], { cols: 80, rows: 24 }) p.kill() return true } catch { return false } })() const itPty = PTY_AVAILABLE ? it : it.skip function getFreePort(): Promise { return new Promise((resolve, reject) => { const srv = net.createServer() srv.listen(0, '127.0.0.1', () => { const addr = srv.address() if (addr === null || typeof addr === 'string') { srv.close() reject(new Error('bad addr')) return } const port = addr.port srv.close(() => resolve(port)) }) srv.on('error', reject) }) } const handles: { close(): Promise }[] = [] async function spawnServer( overrides: Record = {}, ): Promise<{ port: number; origin: string }> { const port = await getFreePort() const cfg = loadConfig({ PORT: String(port), BIND_HOST: '127.0.0.1', SHELL_PATH: process.env['SHELL'] ?? '/bin/zsh', ALLOWED_ORIGINS: `http://127.0.0.1:${port}`, USE_TMUX: '0', IDLE_TTL: '86400', ...overrides, }) const handle = startServer(cfg) handles.push(handle) await new Promise((r) => setTimeout(r, 80)) return { port, origin: `http://127.0.0.1:${port}` } } function waitForOpen(ws: WebSocket): Promise { return new Promise((resolve, reject) => { const t = setTimeout(() => reject(new Error('open timeout')), 3000) ws.once('open', () => { clearTimeout(t) resolve() }) ws.once('error', (e) => { clearTimeout(t) reject(e) }) }) } function waitForMessage( ws: WebSocket, pred: (m: Record) => boolean, timeoutMs = 6000, ): Promise> { return new Promise((resolve, reject) => { const t = setTimeout(() => { ws.off('message', onMsg) reject(new Error('message timeout')) }, timeoutMs) const onMsg = (raw: WebSocket.RawData): void => { let m: Record try { m = JSON.parse(raw.toString('utf8')) as Record } catch { return } if (pred(m)) { clearTimeout(t) ws.off('message', onMsg) resolve(m) } } ws.on('message', onMsg) }) } /** Resolve once no WS message has arrived for `quietMs` (the shell prompt settled). */ function waitForQuiet(ws: WebSocket, quietMs = 300, maxMs = 4000): Promise { return new Promise((resolve) => { let quietTimer: ReturnType const done = (): void => { clearTimeout(quietTimer) clearTimeout(hardStop) ws.off('message', onMsg) resolve() } const arm = (): void => { clearTimeout(quietTimer) quietTimer = setTimeout(done, quietMs) } const onMsg = (): void => arm() const hardStop = setTimeout(done, maxMs) ws.on('message', onMsg) arm() }) } /** Collect every output-frame's data string for `windowMs`. */ function collectOutput(ws: WebSocket, windowMs: number): Promise { return new Promise((resolve) => { let acc = '' const onMsg = (raw: WebSocket.RawData): void => { try { const m = JSON.parse(raw.toString('utf8')) as Record if (m['type'] === 'output') acc += String(m['data']) } catch { /* ignore */ } } ws.on('message', onMsg) setTimeout(() => { ws.off('message', onMsg) resolve(acc) }, windowMs) }) } const isAttached = (m: Record): boolean => m['type'] === 'attached' async function attachSession(port: number, origin: string): Promise<{ ws: WebSocket; sessionId: string }> { const ws = new WebSocket(`ws://127.0.0.1:${port}/term`, { headers: { Origin: origin } }) await waitForOpen(ws) ws.send(JSON.stringify({ type: 'attach', sessionId: null })) const attached = await waitForMessage(ws, isAttached) return { ws, sessionId: attached['sessionId'] as string } } function postQueue( port: number, id: string, body: unknown, headers: Record = {}, ): Promise { return fetch(`http://127.0.0.1:${port}/live-sessions/${id}/queue`, { method: 'POST', headers: { 'Content-Type': 'application/json', ...headers }, body: JSON.stringify(body), }) } function postStop(port: number, sessionId: string): Promise { return fetch(`http://127.0.0.1:${port}/hook`, { method: 'POST', headers: { 'Content-Type': 'application/json', 'X-Webterm-Session': sessionId }, body: JSON.stringify({ hook_event_name: 'Stop' }), }) } afterEach(async () => { while (handles.length > 0) await handles.pop()?.close() await new Promise((r) => setTimeout(r, 30)) }) const UNKNOWN = '00000000-0000-4000-8000-000000000000' // ── Route guards + validation (no PTY needed) ───────────────────────────────── describe('POST /live-sessions/:id/queue — guards & validation', () => { it('rejects a foreign Origin with 403 (CSRF)', async () => { const { port } = await spawnServer() const res = await postQueue(port, UNKNOWN, { text: 'x' }, { Origin: 'http://evil.example' }) expect(res.status).toBe(403) }) it('rejects a missing Origin with 403 (default-deny)', async () => { const { port } = await spawnServer() const res = await postQueue(port, UNKNOWN, { text: 'x' }) // no Origin header expect(res.status).toBe(403) }) it('rejects a malformed session id with 400', async () => { const { port, origin } = await spawnServer() const res = await postQueue(port, 'not-a-uuid', { text: 'x' }, { Origin: origin }) expect(res.status).toBe(400) }) it('rejects empty / non-string text with 400', async () => { const { port, origin } = await spawnServer() expect((await postQueue(port, UNKNOWN, { text: '' }, { Origin: origin })).status).toBe(400) expect((await postQueue(port, UNKNOWN, { text: 42 }, { Origin: origin })).status).toBe(400) expect((await postQueue(port, UNKNOWN, {}, { Origin: origin })).status).toBe(400) }) it('rejects oversized text with 413', async () => { const { port, origin } = await spawnServer({ QUEUE_ITEM_MAX_BYTES: '16' }) const res = await postQueue(port, UNKNOWN, { text: 'x'.repeat(17) }, { Origin: origin }) expect(res.status).toBe(413) }) it('returns 404 for a well-formed but unknown session id', async () => { const { port, origin } = await spawnServer() const res = await postQueue(port, UNKNOWN, { text: 'x' }, { Origin: origin }) expect(res.status).toBe(404) }) it('returns 503 when QUEUE_ENABLED=0', async () => { const { port, origin } = await spawnServer({ QUEUE_ENABLED: '0' }) const res = await postQueue(port, UNKNOWN, { text: 'x' }, { Origin: origin }) expect(res.status).toBe(503) }) it('rate-limits after QUEUE_RATE_MAX requests with 429', async () => { const { port, origin } = await spawnServer() // 20 allowed (each 400 on empty text — rate is checked before validation), // the 21st exceeds the per-IP window and returns 429. let sawRate = false for (let i = 0; i < 21; i++) { const res = await postQueue(port, UNKNOWN, { text: '' }, { Origin: origin }) if (res.status === 429) sawRate = true } expect(sawRate).toBe(true) }) }) describe('GET/DELETE /live-sessions/:id/queue — guards', () => { it('GET returns 404 for an unknown session', async () => { const { port } = await spawnServer() const res = await fetch(`http://127.0.0.1:${port}/live-sessions/${UNKNOWN}/queue`) expect(res.status).toBe(404) }) it('DELETE rejects a foreign Origin with 403', async () => { const { port } = await spawnServer() const res = await fetch(`http://127.0.0.1:${port}/live-sessions/${UNKNOWN}/queue`, { method: 'DELETE', headers: { Origin: 'http://evil.example' }, }) expect(res.status).toBe(403) }) it('DELETE returns 404 for an unknown session', async () => { const { port, origin } = await spawnServer() const res = await fetch(`http://127.0.0.1:${port}/live-sessions/${UNKNOWN}/queue`, { method: 'DELETE', headers: { Origin: origin }, }) expect(res.status).toBe(404) }) }) // ── Happy path + drain (real PTY) ───────────────────────────────────────────── describe('inject queue — real session', () => { itPty('enqueues, surfaces queueLength, lists items, and clears', async () => { const { port, origin } = await spawnServer() const { ws, sessionId } = await attachSession(port, origin) await waitForQuiet(ws) const enq = await postQueue(port, sessionId, { text: 'echo hi', appendEnter: true }, { Origin: origin }) expect(enq.status).toBe(200) expect(await enq.json()).toEqual({ length: 1 }) // /live-sessions surfaces the depth. const live = (await (await fetch(`http://127.0.0.1:${port}/live-sessions`)).json()) as Array<{ id: string queueLength?: number }> expect(live.find((s) => s.id === sessionId)?.queueLength).toBe(1) // GET .../queue returns the stored (verbatim) items. const view = (await ( await fetch(`http://127.0.0.1:${port}/live-sessions/${sessionId}/queue`) ).json()) as { length: number; items: string[] } expect(view.length).toBe(1) expect(view.items).toEqual(['echo hi\r']) // DELETE clears the queue. const del = await fetch(`http://127.0.0.1:${port}/live-sessions/${sessionId}/queue`, { method: 'DELETE', headers: { Origin: origin }, }) expect(del.status).toBe(200) const view2 = (await ( await fetch(`http://127.0.0.1:${port}/live-sessions/${sessionId}/queue`) ).json()) as { length: number } expect(view2.length).toBe(0) ws.close() }) itPty('drains the head into the PTY after a Stop hook (idle-drain)', async () => { const { port, origin } = await spawnServer({ QUEUE_SETTLE_MS: '50' }) const { ws, sessionId } = await attachSession(port, origin) await waitForQuiet(ws) const enq = await postQueue( port, sessionId, { text: 'echo QUEUED_MARKER_A', appendEnter: true }, { Origin: origin }, ) expect(enq.status).toBe(200) const marker = waitForMessage( ws, (m) => m['type'] === 'output' && String(m['data']).includes('QUEUED_MARKER_A'), ) const stop = await postStop(port, sessionId) expect(stop.status).toBe(204) await marker // the shell echoed the injected keystrokes → marker present ws.close() }) itPty('fires only one entry per idle (Stop) — the next waits for the next Stop', async () => { const { port, origin } = await spawnServer({ QUEUE_SETTLE_MS: '50' }) const { ws, sessionId } = await attachSession(port, origin) await waitForQuiet(ws) await postQueue(port, sessionId, { text: 'echo MARK_ONE', appendEnter: true }, { Origin: origin }) await postQueue(port, sessionId, { text: 'echo MARK_TWO', appendEnter: true }, { Origin: origin }) // First Stop → only MARK_ONE drains. const one = waitForMessage(ws, (m) => m['type'] === 'output' && String(m['data']).includes('MARK_ONE')) await postStop(port, sessionId) await one // Without another Stop, MARK_TWO must NOT be injected. const idleWindow = await collectOutput(ws, 500) expect(idleWindow.includes('MARK_TWO')).toBe(false) // A second Stop drains MARK_TWO. await waitForQuiet(ws) const two = waitForMessage(ws, (m) => m['type'] === 'output' && String(m['data']).includes('MARK_TWO')) await postStop(port, sessionId) await two ws.close() }) itPty('settle guard: a non-idle hook during the window cancels the drain', async () => { const { port, origin } = await spawnServer({ QUEUE_SETTLE_MS: '300' }) const { ws, sessionId } = await attachSession(port, origin) await waitForQuiet(ws) await postQueue( port, sessionId, { text: 'echo SHOULD_NOT_RUN', appendEnter: true }, { Origin: origin }, ) // Stop schedules a drain at +300ms; a PreToolUse immediately flips the status // back to 'working', so the settle-time idle re-check skips the drain. await postStop(port, sessionId) await fetch(`http://127.0.0.1:${port}/hook`, { method: 'POST', headers: { 'Content-Type': 'application/json', 'X-Webterm-Session': sessionId }, body: JSON.stringify({ hook_event_name: 'PreToolUse', tool_name: 'Bash' }), }) const window = await collectOutput(ws, 600) // past the 300ms settle expect(window.includes('SHOULD_NOT_RUN')).toBe(false) ws.close() }) })