diff --git a/CLAUDE.md b/CLAUDE.md index d339d5ea..9b73e6e4 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -292,6 +292,8 @@ Frontend JS modules have `@fileoverview` with `@dependency`/`@loadorder` tags. L **Connection-loss UI** (`computeConnectionLossUi()` in constants.js, writer `_updateConnectionLossUi()` in app.js): the service worker serves the cached app shell, so an unreachable server (phone off the tailnet, VPN down, server stopped) used to render a normal-looking empty dashboard whose only tell was the 8px header dot, which reads as "no sessions", not "no connection". Two surfaces now: a full-screen **overlay** while no server state has loaded this page load (nothing behind it is worth preserving), and a non-blocking **banner** once it has (the terminal scrollback stays readable). ⚠️ A **2.5s grace** is load-bearing: a COM deploy restarts the server and SSE is back in ~200ms, and a banner on every deploy trains the user to ignore it. `navigator.onLine === false` skips the grace, since that is never a blip. Retry re-arms SSE **and** the terminal WS (`planWsReconnect` can 'give-up', and the SSE backoff caps at 30s). +**SSE staleness watchdog** (`computeSseStale()` in constants.js, `_checkSseStale()` + a 5s interval in app.js): an `EventSource` that stops delivering does not always error, so `onerror` never fires, the header dot stays green, and every SSE-driven surface (tab status dots, sessions created on another device, renames) freezes until the user reloads. ⚠️ The 15s server keepalive was an SSE **comment** (`:keepalive`), and comments are **invisible to `EventSource` by spec**, so there was nothing a client could observe: it is now the named `sse:heartbeat` event (`cleanupDeadClients()`, sse-stream-manager.ts), which is exactly why the frame had to change type. ⚠️ Staleness is judged **only while the status is `connected`** and the device is online; that guard is the loop breaker, since a forced `connectSSE()` leaves `connected` immediately and cannot re-fire while a reconnect is in flight. ⚠️ The liveness stamp is applied inside `addListener` itself, so every registered handler (the `_SSE_HANDLER_MAP` wrappers AND the directly-registered ones) feeds it from one place; the heartbeat's own listener is a no-op that exists **only** to be registered, since `EventSource` drops named events nobody listens for. ⚠️ The watchdog interval is cleared at the top of `connectSSE()` and nowhere else (its only teardown path); clearing it elsewhere stacks intervals. Recovery needs no new sync path: the reconnect re-runs `handleInit` → `_resetAllAppState()`. The forced reconnect logs one diagnostic line, because a middlebox that strips heartbeats presents as "silently reconnects every 45s". + **Z-index layers**: subagent windows (1000), plan agents (1100), mobile/tablet fixed header (1200, `mobile.css`), modals on ≤768px (1300 — must beat the fixed header or the modal close button is buried), log viewers (2000), connection-loss overlay (2500, above the fixed header and modals), image popups (3000), local echo overlay (7). **Respawn presets**: `solo-work` (3s/60min), `subagent-workflow` (45s/240min), `team-lead` (90s/480min), `ralph-todo` (8s/480min), `overnight-autonomous` (10s/480min). @@ -320,7 +322,7 @@ Frontend JS modules have `@fileoverview` with `@dependency`/`@loadorder` tags. L ### SSE Event Registry -154 event constants in `src/web/sse-events.ts` (backend) and `SSE_EVENTS` in `constants.js` (frontend). **Both must be kept in sync**, and `test/sse-registry-parity.test.ts` is the guard that pins it (currently exactly in sync, 154 = 154, no drift either direction). The backend file's `@fileoverview` carries the per-category breakdown. +155 event constants in `src/web/sse-events.ts` (backend) and `SSE_EVENTS` in `constants.js` (frontend). **Both must be kept in sync**, and `test/sse-registry-parity.test.ts` is the guard that pins it (currently exactly in sync, 155 = 155, no drift either direction). The backend file's `@fileoverview` carries the per-category breakdown. ### API Routes diff --git a/docs/api-reference.md b/docs/api-reference.md index 25ebd69f..9871865d 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -534,6 +534,29 @@ the stable contract — event names are not renamed without a major bump. An optional `?sessions=` filter suppresses only the high-volume terminal stream; lifecycle/metadata events are delivered to all clients regardless. +### `sse:heartbeat` (liveness) + +Every 15s the server writes a `sse:heartbeat` frame to every connected client: + +``` +event: sse:heartbeat +data: {"t":1755100000000} +``` + +`t` is the server's epoch-ms timestamp at write time. The frame carries no +application state and can be ignored for correctness. It exists so a client can +tell a live stream from a dead one: an `EventSource` whose connection has been +idle-closed by a proxy (or that resumed from sleep on a stale socket) keeps +delivering nothing without ever firing `onerror`. Clients that care should treat +silence longer than about three intervals as a dead stream and reconnect, which +is what the bundled frontend does. + +This replaced a `:keepalive` SSE **comment**, which served the same +proxy-flushing purpose but is invisible to `EventSource` by spec and so could +never be observed by a client. Consumers written against the old behavior are +unaffected: `EventSource` dispatches only events that have a registered +listener, so an unknown event name is dropped. + ## Consuming from JavaScript The bundled frontend reads responses through `_apiJson()` diff --git a/src/web/public/app.js b/src/web/public/app.js index a531a1b1..d80aa730 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -684,6 +684,17 @@ class CodemanApp { this.maxReconnectAttempts = 10; this.isOnline = navigator.onLine; + // SSE staleness watchdog. An EventSource that stops delivering does not + // always error (a proxy that idle-closed it, a resumed laptop), so + // `onerror` never fires and every SSE-driven surface freezes silently. + // The server heartbeats every 15s; going quiet for three of them means the + // stream is a zombie and has to be rebuilt. The decision is pure + // (computeSseStale in constants.js); these are its inputs. The threshold + // is an instance field so a browser test can shrink it. + this._sseLastMessageAt = 0; + this._sseStaleTimeoutMs = window.CodemanSseStale?.TIMEOUT_MS ?? 45000; + this._sseStaleWatchdog = null; + // Connection-loss UI (banner + full-screen overlay). The decision itself is // pure and lives in constants.js (computeConnectionLossUi); these are just // its inputs. `_connDownSince` is the timestamp the transport LEFT the @@ -719,6 +730,11 @@ class CodemanApp { window.addEventListener('pagehide', () => this._persistReliableNow()); document.addEventListener('visibilitychange', () => { if (document.visibilityState === 'hidden') this._persistReliableNow(); + // A background tab's timers are throttled, so the 5s watchdog may not + // have run for minutes, and a wake/unlock is exactly when a stream + // comes back zombie. Checking here is what makes recovery feel instant + // instead of up to a full timeout late. + else this._checkSseStale(); }); // Local echo overlay — DOM overlay positioned at the visible ❯ prompt @@ -1421,6 +1437,14 @@ class CodemanApp { // Clear any pending reconnect timeout to prevent duplicate connections this._clearTimer('sseReconnectTimeout'); + // Same discipline for the staleness watchdog: connectSSE() runs on every + // reconnect and is the only teardown path this page-lifetime interval has, + // so clearing it anywhere else (or not at all) stacks intervals. + if (this._sseStaleWatchdog) { + clearInterval(this._sseStaleWatchdog); + this._sseStaleWatchdog = null; + } + // Clean up existing SSE listeners before creating new connection (prevents listener accumulation) if (this._sseListenerCleanup) { this._sseListenerCleanup(); @@ -1448,11 +1472,20 @@ class CodemanApp { if (this.activeSessionId) _sseParams.set('sessions', this.activeSessionId); this.eventSource = new EventSource(`/api/events?${_sseParams.toString()}`); - // Store all event listeners for cleanup on reconnect + // Store all event listeners for cleanup on reconnect. + // + // Every handler is wrapped so ANY frame that arrives stamps the liveness + // clock the staleness watchdog reads. Doing it here (rather than at the + // three separate registration sites below) is what keeps a future + // addListener() call from silently opting out of it. const listeners = []; const addListener = (event, handler) => { - this.eventSource.addEventListener(event, handler); - listeners.push({ event, handler }); + const stamped = (e) => { + this._sseLastMessageAt = Date.now(); + handler(e); + }; + this.eventSource.addEventListener(event, stamped); + listeners.push({ event, handler: stamped }); }; // Create cleanup function to remove all listeners @@ -1467,6 +1500,10 @@ class CodemanApp { this.eventSource.onopen = () => { this.reconnectAttempts = 0; + // Start the liveness clock here, not at the first frame: the watchdog + // only ever fires while the status is 'connected', and this is the + // moment that becomes true. + this._sseLastMessageAt = Date.now(); this.setConnectionStatus('connected'); }; this.eventSource.onerror = () => { @@ -1614,6 +1651,52 @@ class CodemanApp { } this._onSessionListMaybeChanged(); }); + + // Liveness heartbeat. The handler is deliberately empty: the whole point + // is the stamp inherited from addListener's wrapper. It still has to be + // REGISTERED: EventSource only dispatches named events that have a + // listener, so without this the frame arrives on the wire and is dropped + // before it can prove the stream is alive. + addListener(SSE_EVENTS.HEARTBEAT, () => {}); + + // Watchdog: a stream that goes quiet without erroring is invisible to + // onerror, so poll the pure staleness policy and rebuild the connection + // ourselves. 5s granularity against a 45s threshold: cheap, and it keeps + // the worst-case detection lag well under a heartbeat interval. + this._sseStaleWatchdog = setInterval(() => this._checkSseStale(), 5000); + } + + /** + * Force a reconnect if the SSE stream has gone quiet while still claiming to + * be connected. Called by the 5s watchdog and on tab-visible. + * + * Recovery needs no new sync path: the reconnect re-runs `handleInit`, which + * already calls `_resetAllAppState()` and rebuilds everything from the + * server. The connection-loss UI needs nothing either: `connectSSE()` sets + * status 'connecting' (reconnectAttempts was zeroed by onopen), and the 2.5s + * grace in computeConnectionLossUi means a stream that heals in 200ms shows + * nothing at all. + */ + _checkSseStale() { + const policy = window.CodemanSseStale; + if (!policy) return; + const now = Date.now(); + const stale = policy.compute({ + lastMessageAt: this._sseLastMessageAt, + now, + status: this._connectionStatus, + isOnline: this.isOnline, + timeoutMs: this._sseStaleTimeoutMs, + }); + if (!stale) return; + // If a middlebox ever strips or delays heartbeats, the failure mode is + // "silently reconnects every 45s", and a field report of that would be + // undebuggable without this line. + console.log( + `[SSE] stream stale: no frame for ${now - this._sseLastMessageAt}ms ` + + `(threshold ${this._sseStaleTimeoutMs}ms), forcing reconnect` + ); + this.connectSSE(); } // ═══════════════════════════════════════════════════════════════ diff --git a/src/web/public/constants.js b/src/web/public/constants.js index 3b0678db..1d206fd1 100644 --- a/src/web/public/constants.js +++ b/src/web/public/constants.js @@ -379,6 +379,41 @@ function computeConnectionLossUi(input) { }; } +// SSE staleness policy: is this stream a zombie? +// +// An EventSource that stops delivering does not always error. A proxy that +// idle-closed the connection, a laptop resumed from sleep, a tailnet +// reconnect: `onerror` never fires, the header dot stays green, and every +// SSE-driven surface (tab status dots, sessions created on another device, +// renames) freezes until the user reloads. The server writes a +// `sse:heartbeat` frame every 15s, so silence longer than three of them means +// the stream is dead even though the transport still claims otherwise. +// +// Stale ONLY when the transport believes it is 'connected': the other states +// already have the reconnect/backoff machinery running, and re-firing on top +// of them would stack reconnects. That guard is also the loop breaker: a +// forced reconnect leaves 'connected' immediately, so the watchdog cannot +// fire again while one is in flight. `navigator.onLine === false` is not +// staleness either; there is nothing to reconnect to yet. +// +// Pure: no DOM, no timers, no side effects. `now` is passed in. +const SSE_STALE_TIMEOUT_MS = 45000; // three missed 15s heartbeats + +function computeSseStale(input) { + const { + lastMessageAt = null, + now = 0, + status = 'connected', + isOnline = true, + timeoutMs = SSE_STALE_TIMEOUT_MS, + } = input || {}; + if (!isOnline || status !== 'connected') return false; + // No frame has ever arrived: `init` lands on connect, so this is a stream + // that has not opened yet rather than one that went quiet. + if (typeof lastMessageAt !== 'number' || !(lastMessageAt > 0)) return false; + return now - lastMessageAt >= timeoutMs; +} + if (typeof window !== 'undefined') { window.WEBGL_FALLBACK = WEBGL_FALLBACK; window.evaluateWebGLLongTaskTrip = evaluateWebGLLongTaskTrip; @@ -401,6 +436,10 @@ if (typeof window !== 'undefined') { compute: computeConnectionLossUi, GRACE_MS: CONNECTION_LOSS_GRACE_MS, }; + window.CodemanSseStale = { + compute: computeSseStale, + TIMEOUT_MS: SSE_STALE_TIMEOUT_MS, + }; } // Scheduler API — prioritize terminal writes over background UI updates. @@ -514,6 +553,9 @@ const SSE_EVENTS = { // Core INIT: 'init', + // Transport + HEARTBEAT: 'sse:heartbeat', + // Session lifecycle SESSION_CREATED: 'session:created', SESSION_UPDATED: 'session:updated', diff --git a/src/web/public/session-ui.js b/src/web/public/session-ui.js index 183e8f72..c68e4e11 100644 --- a/src/web/public/session-ui.js +++ b/src/web/public/session-ui.js @@ -1415,9 +1415,54 @@ Object.assign(CodemanApp.prototype, { this.activeFocusTrap.activate(); }, + /** + * Write a name the server has just confirmed into the local session map. + * + * Both rename surfaces re-render the tab strip from `this.sessions` right + * after their PUT, so without this they depended on the `session:updated` SSE + * frame to carry their own write back. On a page whose SSE stream has gone + * quiet without erroring (a proxy that idle-closed it, a laptop resumed from + * sleep) that frame never lands: the PUT stores the new name, the re-render + * repaints the stale one, and the rename looks like it did nothing until a + * full page reload. The response body is authoritative, so apply it directly. + * The SSE frame, when it does arrive, replaces the object with the same name. + */ + _applyLocalSessionName(sessionId, name) { + if (typeof name !== 'string') return; + const session = this.sessions.get(sessionId); + if (!session) return; + session.name = name; + this.sessions.set(sessionId, session); + // Mirrors _onSessionUpdated: subagent windows cache their parent's name. + this.updateSubagentParentNames?.(sessionId); + }, + + /** + * PUT a session name and return the name the server stored, or null if the + * request failed. `_apiPut` swallows network errors into a null Response and + * an API-level failure arrives as a non-ok status or `{success:false}`, so a + * rename that silently did nothing has to be detected here, not thrown. + */ + async _putSessionName(sessionId, name) { + const res = await this._apiPut(`/api/sessions/${sessionId}/name`, { name }); + if (!res || !res.ok) return null; + let payload = null; + try { + payload = await res.json(); + } catch { + return null; + } + if (payload && payload.success === false) return null; + const confirmed = payload?.data?.name; + return typeof confirmed === 'string' ? confirmed : name; + }, + async saveSessionName() { if (!this.editingSessionId) return; - const session = this.sessions.get(this.editingSessionId); + // Captured: the modal can be closed (or switched to another session) while + // the PUT is in flight, and the name belongs to the session that was open. + const sessionId = this.editingSessionId; + const session = this.sessions.get(sessionId); const parsed = session ? parseSessionPrefix(session.name) : null; const inputVal = document.getElementById('modalSessionName').value.trim(); let name; @@ -1426,11 +1471,13 @@ Object.assign(CodemanApp.prototype, { } else { name = inputVal; } - try { - await this._apiPut(`/api/sessions/${this.editingSessionId}/name`, { name }); - } catch (err) { - this.showToast('Failed to save session name: ' + err.message, 'error'); + const confirmed = await this._putSessionName(sessionId, name); + if (confirmed === null) { + this.showToast('Failed to save session name', 'error'); + return; } + this._applyLocalSessionName(sessionId, confirmed); + this.renderSessionTabs(); }, async autoSaveAutoCompact() { @@ -1740,15 +1787,14 @@ Object.assign(CodemanApp.prototype, { // Skip the API call if the session vanished between focus and blur. const stillExists = this.sessions.has(sessionId); if (stillExists && fullName !== session.name) { - try { - await fetch(`/api/sessions/${sessionId}/name`, { - method: 'PUT', - headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify({ name: fullName }) - }); - } catch (err) { + const confirmed = await this._putSessionName(sessionId, fullName); + if (confirmed === null) { tabName.textContent = originalContent; this.showToast('Failed to rename', 'error'); + } else { + // The re-render below repaints from this.sessions, so the new name has + // to be in the map before it runs (see _applyLocalSessionName()). + this._applyLocalSessionName(sessionId, confirmed); } } // Re-render tabs to restore full tab structure diff --git a/src/web/sse-events.ts b/src/web/sse-events.ts index 03eded65..b2feb899 100644 --- a/src/web/sse-events.ts +++ b/src/web/sse-events.ts @@ -5,8 +5,9 @@ * and referenced by the frontend (`SSE_EVENTS` in `constants.js`). * Both files MUST be kept in sync. * - * 154 event constants organized by category: + * 155 event constants organized by category: * - **Core** (1): init + * - **Transport** (1): sse:heartbeat * - **Session lifecycle** (23): created, updated, deleted, terminal, idle, working, ... * - **Session: Ralph** (6): ralphLoopUpdate, todoUpdate, completionDetected, ... * - **Session: Bash tools** (3): bashToolStart, bashToolEnd, bashToolsUpdate @@ -52,6 +53,22 @@ /** Sent to each SSE client on initial connection with full app state. */ export const Init = 'init' as const; +// ─── Transport ─────────────────────────────────────────────────────────────── + +/** + * Liveness frame written to every SSE client every `SSE_HEARTBEAT_INTERVAL`. + * Payload: `{ t: }`. + * + * Carries no application data; its only job is to be *observable*. This was a + * `:keepalive` SSE **comment**, and comments are invisible to `EventSource` by + * spec, so a stream that stopped delivering without erroring (a proxy that + * idle-closed it, a laptop resumed from sleep, a tailnet reconnect) was + * undetectable to the client: `onerror` never fires and the UI freezes until a + * reload. A named event reaches a listener, which is what lets the client's + * staleness watchdog notice the silence and force a reconnect. + */ +export const Heartbeat = 'sse:heartbeat' as const; + // ─── Session Lifecycle ─────────────────────────────────────────────────────── /** New session spawned. */ @@ -443,6 +460,9 @@ export const SseEvent = { // Core Init, + // Transport + Heartbeat, + // Session lifecycle SessionCreated, SessionUpdated, diff --git a/src/web/sse-stream-manager.ts b/src/web/sse-stream-manager.ts index 32c65f71..d1a5c2ea 100644 --- a/src/web/sse-stream-manager.ts +++ b/src/web/sse-stream-manager.ts @@ -470,12 +470,20 @@ export class SseStreamManager { // ========== Client Health ========== /** - * Clean up dead SSE clients and send keep-alive comments. + * Clean up dead SSE clients and send the liveness heartbeat. * Keep-alive prevents proxy/load-balancer timeouts on idle connections. * Dead client cleanup prevents memory leaks from abruptly terminated connections. + * + * The heartbeat is a NAMED event, not the `:keepalive` comment it used to be: + * comments are invisible to `EventSource` by spec, so a stream that stopped + * delivering without erroring was undetectable to the client (see + * `SseEvent.Heartbeat`). Written per-client rather than through `broadcast()` + * deliberately: the frame carries no session data, so it needs no owner + * routing, and this loop is already walking every client to check its socket. */ cleanupDeadClients(): void { const deadClients: FastifyReply[] = []; + const heartbeat = `event: ${SseEvent.Heartbeat}\ndata: ${JSON.stringify({ t: Date.now() })}\n\n`; for (const [client] of this.sseClients) { try { @@ -484,11 +492,9 @@ export class SseStreamManager { if (!socket || socket.destroyed || !socket.writable) { deadClients.push(client); } else { - // Send SSE comment as keep-alive. Only add padding when tunnel is - // active — it flushes Cloudflare proxy buffers but wastes bandwidth - // for direct/Tailscale connections. - const ka = this._isTunnelActive ? ':keepalive\n' + SSE_PADDING : ':keepalive\n\n'; - client.raw.write(ka); + // Only add padding when tunnel is active: it flushes Cloudflare + // proxy buffers but wastes bandwidth for direct/Tailscale connections. + client.raw.write(this._isTunnelActive ? heartbeat + SSE_PADDING : heartbeat); } } catch { // Error accessing socket means client is dead diff --git a/test/inline-rename.test.ts b/test/inline-rename.test.ts index 62155575..76eb6d97 100644 --- a/test/inline-rename.test.ts +++ b/test/inline-rename.test.ts @@ -358,6 +358,83 @@ describe('Inline rename input', () => { expect(result.editingAfter).toBe(null); }); + it('Commit writes the confirmed name into app.sessions WITHOUT any session:updated frame', async () => { + await resetState(); + expect(await startRename('no-sse', 'w9-case')).toBe(true); + + // finishRename() re-renders the tab strip from app.sessions, so the rename + // used to depend on the session:updated SSE frame to carry its own write + // back. On a page whose stream has gone quiet without erroring, the PUT + // stored the new name, the re-render repainted the stale one, and the tab + // only showed it after a full reload. No SSE is dispatched here at all. + const result = await page.evaluate(async () => { + const app = ( + window as unknown as { + app: { sessions: Map }; + } + ).app; + const origFetch = window.fetch; + window.fetch = (async () => + new Response('{"success":true,"data":{"name":"w9-case: fresh"}}', { + status: 200, + headers: { 'Content-Type': 'application/json' }, + })) as typeof window.fetch; + + const inputEl = document.querySelector('input.tab-rename-input') as HTMLInputElement; + inputEl.value = 'fresh'; + inputEl.dispatchEvent(new KeyboardEvent('keydown', { key: 'Enter', bubbles: true })); + await new Promise((r) => setTimeout(r, 60)); + + window.fetch = origFetch; + return { mapName: app.sessions.get('no-sse')?.name ?? null }; + }); + + expect(result.mapName).toBe('w9-case: fresh'); + }); + + it('A rejected rename restores the old label and leaves app.sessions untouched', async () => { + await resetState(); + expect(await startRename('rename-500', 'w9-case')).toBe(true); + + // _apiPut turns a network error into a null Response and an API-level + // failure arrives as a non-ok status, neither of which throws, so a + // rejected rename has to be detected from the response, or it reports + // success and silently discards the user's edit. + const result = await page.evaluate(async () => { + const app = ( + window as unknown as { + app: { sessions: Map; showToast: (m: string, k: string) => void }; + } + ).app; + const toasts: string[] = []; + const origToast = app.showToast; + app.showToast = (msg: string) => void toasts.push(msg); + const origFetch = window.fetch; + window.fetch = (async () => + new Response('{"success":false,"error":"boom","errorCode":"INTERNAL"}', { + status: 500, + headers: { 'Content-Type': 'application/json' }, + })) as typeof window.fetch; + + const inputEl = document.querySelector('input.tab-rename-input') as HTMLInputElement; + inputEl.value = 'never-stored'; + inputEl.dispatchEvent(new KeyboardEvent('keydown', { key: 'Enter', bubbles: true })); + await new Promise((r) => setTimeout(r, 60)); + + window.fetch = origFetch; + app.showToast = origToast; + return { + mapName: app.sessions.get('rename-500')?.name ?? null, + label: document.querySelector('.tab-name[data-session-id="rename-500"]')?.textContent ?? null, + toasts, + }; + }); + + expect(result.mapName).toBe('w9-case'); + expect(result.label).toBe('w9-case'); + expect(result.toasts).toContain('Failed to rename'); + }); + it('Re-entry: starting rename while one is active aborts the previous one', async () => { await resetState(); expect(await startRename('first-id', 'First')).toBe(true); diff --git a/test/sse-heartbeat.test.ts b/test/sse-heartbeat.test.ts new file mode 100644 index 00000000..9dea1dee --- /dev/null +++ b/test/sse-heartbeat.test.ts @@ -0,0 +1,145 @@ +/** + * SSE liveness heartbeat. + * + * `cleanupDeadClients()` runs every SSE_HEARTBEAT_INTERVAL (15s) and does two + * jobs: evict clients whose socket died, and write a liveness frame to the + * ones that are still up. + * + * The regression these guard: that frame used to be an SSE `:keepalive` + * COMMENT, and comments are invisible to `EventSource` by spec. A stream that + * stopped delivering without erroring was therefore undetectable to the + * client: `onerror` never fired, the header dot stayed green, and every + * SSE-driven surface froze until the user reloaded. A named `sse:heartbeat` + * event reaches a listener, which is what lets the client's staleness + * watchdog notice the silence (see test/sse-staleness.test.ts). + * + * No port needed (the manager is driven directly with fake replies). + */ +import type { FastifyReply } from 'fastify'; +import { describe, it, expect } from 'vitest'; +import { SSE_PADDING_SIZE } from '../src/config/server-timing.js'; +import { SseEvent } from '../src/web/sse-events.js'; +import { SseStreamManager } from '../src/web/sse-stream-manager.js'; +import { CleanupManager } from '../src/utils/index.js'; + +/** A FastifyReply stand-in that records every raw write. */ +function fakeClient(opts: { destroyed?: boolean; writable?: boolean; throwOnAccess?: boolean } = {}) { + const writes: string[] = []; + const socket = { destroyed: opts.destroyed ?? false, writable: opts.writable ?? true }; + const raw = { + get socket() { + if (opts.throwOnAccess) throw new Error('socket gone'); + return socket; + }, + write(chunk: string) { + writes.push(chunk); + return true; + }, + }; + return { reply: { raw } as unknown as FastifyReply, writes }; +} + +function makeManager() { + const cleanup = new CleanupManager(); + const manager = new SseStreamManager({ getSessionStateWithRespawn: () => null }, cleanup); + return { manager, cleanup }; +} + +describe('SSE liveness heartbeat', () => { + it('writes a NAMED sse:heartbeat event, not an invisible comment', () => { + const { manager, cleanup } = makeManager(); + const client = fakeClient(); + manager.addClient(client.reply, null, false); + + manager.cleanupDeadClients(); + + expect(client.writes).toHaveLength(1); + const frame = client.writes[0]; + // A comment (`:keepalive`) never reaches an EventSource listener, and that is + // the entire bug. The frame must be a dispatchable named event. + expect(frame.startsWith(':')).toBe(false); + expect(frame).toMatch(/^event: sse:heartbeat\n/); + expect(frame.endsWith('\n\n')).toBe(true); + expect(SseEvent.Heartbeat).toBe('sse:heartbeat'); + cleanup.dispose(); + }); + + it('carries a parseable epoch-ms payload', () => { + const { manager, cleanup } = makeManager(); + const client = fakeClient(); + manager.addClient(client.reply, null, false); + const before = Date.now(); + + manager.cleanupDeadClients(); + + const dataLine = client.writes[0].split('\n').find((l) => l.startsWith('data: ')); + expect(dataLine).toBeDefined(); + const payload = JSON.parse(dataLine!.slice('data: '.length)) as { t: number }; + expect(payload.t).toBeGreaterThanOrEqual(before); + expect(payload.t).toBeLessThanOrEqual(Date.now()); + cleanup.dispose(); + }); + + it('still appends Cloudflare tunnel padding when a tunnel is active', () => { + const { manager, cleanup } = makeManager(); + const client = fakeClient(); + manager.addClient(client.reply, null, false); + manager.setTunnelActive(true); + + manager.cleanupDeadClients(); + + const frame = client.writes[0]; + expect(frame).toMatch(/^event: sse:heartbeat\n/); + // Padding rides AFTER the terminating blank line, so the event still parses. + const [event, padding] = frame.split('\n\n'); + expect(event).toMatch(/^event: sse:heartbeat\ndata: \{/); + expect(padding.startsWith(':')).toBe(true); + expect(padding.length).toBeGreaterThanOrEqual(SSE_PADDING_SIZE); + cleanup.dispose(); + }); + + it('sends no padding without a tunnel', () => { + const { manager, cleanup } = makeManager(); + const client = fakeClient(); + manager.addClient(client.reply, null, false); + + manager.cleanupDeadClients(); + + expect(client.writes[0].length).toBeLessThan(200); + cleanup.dispose(); + }); + + it('still evicts dead clients instead of heartbeating them', () => { + const { manager, cleanup } = makeManager(); + const alive = fakeClient(); + const destroyed = fakeClient({ destroyed: true }); + const unwritable = fakeClient({ writable: false }); + const exploding = fakeClient({ throwOnAccess: true }); + for (const c of [alive, destroyed, unwritable, exploding]) manager.addClient(c.reply, null, false); + expect(manager.clientCount).toBe(4); + + manager.cleanupDeadClients(); + + expect(manager.clientCount).toBe(1); + expect(alive.writes).toHaveLength(1); + for (const c of [destroyed, unwritable, exploding]) expect(c.writes).toHaveLength(0); + cleanup.dispose(); + }); + + it('heartbeats every client on each pass', () => { + const { manager, cleanup } = makeManager(); + const a = fakeClient(); + const b = fakeClient(); + manager.addClient(a.reply, null, false); + manager.addClient(b.reply, null, false); + + manager.cleanupDeadClients(); + manager.cleanupDeadClients(); + + // The frame carries no session data, so it needs no owner routing and is + // written per-client rather than through the scoped broadcast() path. + expect(a.writes).toHaveLength(2); + expect(b.writes).toHaveLength(2); + cleanup.dispose(); + }); +}); diff --git a/test/sse-staleness.test.ts b/test/sse-staleness.test.ts new file mode 100644 index 00000000..fdef7f14 --- /dev/null +++ b/test/sse-staleness.test.ts @@ -0,0 +1,102 @@ +/** + * SSE staleness policy. + * + * `CodemanSseStale.compute(input)` is the pure decision behind app.js's + * watchdog: given when the last SSE frame arrived, the transport status and + * the browser's online flag, it says whether the stream has gone quiet while + * still claiming to be connected: a zombie that has to be rebuilt. + * + * The regression it guards: the server's liveness keepalive used to be an SSE + * `:keepalive` COMMENT, and comments are invisible to `EventSource` by spec. + * A stream that stopped delivering without erroring (a proxy that idle-closed + * it, a laptop resumed from sleep, a tailnet reconnect) never fired `onerror`, + * so the header dot stayed green and tab status dots, sessions created on + * another device, and renames all froze until the user reloaded the page. + * + * Loaded in a plain node VM context (no jsdom), mirroring + * test/connection-loss-ui.test.ts. + */ +import { readFileSync } from 'node:fs'; +import { resolve } from 'node:path'; +import vm from 'node:vm'; +import { describe, expect, it } from 'vitest'; + +type StaleInput = { + lastMessageAt?: number | null; + now?: number; + status?: 'connected' | 'connecting' | 'reconnecting' | 'disconnected' | 'offline'; + isOnline?: boolean; + timeoutMs?: number; +}; + +function loadPolicy() { + const context = vm.createContext({ window: {}, globalThis: {} }); + const source = readFileSync(resolve(import.meta.dirname, '../src/web/public/constants.js'), 'utf8'); + vm.runInContext(source, context, { filename: 'constants.js' }); + return ( + context.window as { + CodemanSseStale: { compute: (input: StaleInput) => boolean; TIMEOUT_MS: number }; + } + ).CodemanSseStale; +} + +const T0 = 1_000_000; + +describe('SSE staleness policy', () => { + it('defaults to three missed 15s heartbeats', () => { + const { TIMEOUT_MS } = loadPolicy(); + expect(TIMEOUT_MS).toBe(45000); + }); + + it('is not stale while frames keep arriving', () => { + const { compute, TIMEOUT_MS } = loadPolicy(); + expect(compute({ lastMessageAt: T0, now: T0 + TIMEOUT_MS - 1, status: 'connected' })).toBe(false); + }); + + it('is stale once the threshold is reached', () => { + const { compute, TIMEOUT_MS } = loadPolicy(); + // Boundary is inclusive: exactly three missed heartbeats already means the + // stream has been silent through a window it was contractually filling. + expect(compute({ lastMessageAt: T0, now: T0 + TIMEOUT_MS, status: 'connected' })).toBe(true); + expect(compute({ lastMessageAt: T0, now: T0 + TIMEOUT_MS * 10, status: 'connected' })).toBe(true); + }); + + it('honours a custom timeoutMs (what a browser test shrinks)', () => { + const { compute } = loadPolicy(); + expect(compute({ lastMessageAt: T0, now: T0 + 999, status: 'connected', timeoutMs: 1000 })).toBe(false); + expect(compute({ lastMessageAt: T0, now: T0 + 1000, status: 'connected', timeoutMs: 1000 })).toBe(true); + }); + + it('is never stale while the transport is already reconnecting', () => { + const { compute, TIMEOUT_MS } = loadPolicy(); + // These states already have the backoff machinery running; firing on top + // of them would stack reconnects. This guard is also the loop breaker: + // a forced reconnect leaves 'connected' immediately, so the watchdog + // cannot re-fire while one is in flight. + for (const status of ['connecting', 'reconnecting', 'disconnected', 'offline'] as const) { + expect(compute({ lastMessageAt: T0, now: T0 + TIMEOUT_MS * 10, status })).toBe(false); + } + }); + + it('is never stale while the device is offline', () => { + const { compute, TIMEOUT_MS } = loadPolicy(); + // Nothing to reconnect to yet; the connection-loss UI already owns this. + expect(compute({ lastMessageAt: T0, now: T0 + TIMEOUT_MS * 10, status: 'connected', isOnline: false })).toBe(false); + }); + + it('is not stale before any frame has ever arrived', () => { + const { compute, TIMEOUT_MS } = loadPolicy(); + // The clock starts at onopen, and `init` lands immediately after. A zero + // stamp means the stream has not opened yet, not that it went quiet. The + // constructor optimistically seeds status 'connected' before the first + // connect, so without this guard the watchdog would fire on page load. + for (const lastMessageAt of [0, null, undefined]) { + expect(compute({ lastMessageAt, now: T0 + TIMEOUT_MS * 10, status: 'connected' })).toBe(false); + } + }); + + it('tolerates a missing input object', () => { + const { compute } = loadPolicy(); + expect(compute(undefined as unknown as StaleInput)).toBe(false); + }); +});