diff --git a/src/web/public/app.js b/src/web/public/app.js index c63f12be..169bfb19 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -2358,34 +2358,68 @@ class CodemanApp { const text = this.$('connectionText'); if (!indicator || !dot || !text) return; - const status = this._connectionStatus; + const { bytes: totalBytes, count } = this._pendingBytes(); + const hasQueue = count > 0; + // Only surface a backlog once it's more than a few bytes. A single keystroke + // (1B) ACKs in milliseconds, so without this the label flickered "sending 1B" + // on every key press. Above this threshold means input is genuinely backing up. + const BACKLOG_HINT_BYTES = 4; + const showBacklog = totalBytes > BACKLOG_HINT_BYTES; + const formatBytes = (b) => (b < 1024 ? `${b}B` : `${(b / 1024).toFixed(1)}KB`); + const queuedSuffix = showBacklog ? ` · ${formatBytes(totalBytes)} queued` : ''; - // While the connection is healthy, never surface the input queue. With the - // reliable-delivery layer every keystroke is briefly "pending" until its ACK - // lands a few ms later — showing that flashed "Sending 1B…" on every single - // character. The indicator is only meaningful for an actual connection - // problem (reconnecting / offline), where the queued byte count reassures - // the user their typing is safely buffered and will be sent. - if (status === 'connected' || status === 'connecting') { - indicator.style.display = 'none'; + // Hard offline (browser reports no network) dominates everything. + if (!this.isOnline || this._connectionStatus === 'offline') { + indicator.style.display = 'flex'; + dot.className = 'connection-dot offline'; + text.textContent = showBacklog ? `Offline (${formatBytes(totalBytes)} queued)` : 'Offline'; + indicator.title = 'No network connection'; return; } - const { bytes: totalBytes, count } = this._pendingBytes(); - const hasQueue = count > 0; - indicator.style.display = 'flex'; - dot.className = 'connection-dot'; - - const formatBytes = (b) => (b < 1024 ? `${b}B` : `${(b / 1024).toFixed(1)}KB`); - - if (status === 'reconnecting') { - dot.classList.add('reconnecting'); - text.textContent = hasQueue ? `Reconnecting (${formatBytes(totalBytes)} queued)` : 'Reconnecting...'; - } else { - // Offline or disconnected - dot.classList.add('offline'); - text.textContent = hasQueue ? `Offline (${formatBytes(totalBytes)} queued)` : 'Offline'; + // With an active terminal, show its transport (WebSocket vs HTTP fallback). + if (this.activeSessionId) { + let cls, label, detail; + switch (this._wsState) { + case 'connected': + cls = 'connected'; label = 'WS'; detail = 'Terminal connected over WebSocket'; + break; + case 'fallback': + cls = 'fallback'; label = 'HTTP'; detail = 'WebSocket unavailable — input sent over HTTP'; + break; + case 'reconnecting': + cls = 'reconnecting'; label = 'WS…'; detail = 'Reconnecting WebSocket'; + break; + case 'connecting': + default: + cls = 'reconnecting'; label = 'WS…'; detail = 'Connecting WebSocket'; + break; + } + indicator.style.display = 'flex'; + dot.className = `connection-dot ${cls}`; + text.textContent = `${label}${queuedSuffix}`; + indicator.title = detail; + return; } + + // No active terminal — reflect the SSE event stream only when it needs attention. + if (this._connectionStatus === 'reconnecting' || this._connectionStatus === 'disconnected') { + indicator.style.display = 'flex'; + dot.className = 'connection-dot reconnecting'; + text.textContent = showBacklog ? `Reconnecting (${formatBytes(totalBytes)} queued)` : 'Reconnecting...'; + indicator.title = 'Reconnecting to server'; + return; + } + + // Idle dashboard, healthy stream — hide unless input is genuinely queued. + if (!hasQueue) { + indicator.style.display = 'none'; + return; + } + indicator.style.display = 'flex'; + dot.className = 'connection-dot draining'; + text.textContent = showBacklog ? `Sending ${formatBytes(totalBytes)}...` : 'Sending...'; + indicator.title = 'Delivering queued input'; } setupOnlineDetection() { diff --git a/test/connection-indicator.test.ts b/test/connection-indicator.test.ts new file mode 100644 index 00000000..3f8bcc66 --- /dev/null +++ b/test/connection-indicator.test.ts @@ -0,0 +1,204 @@ +/** + * @fileoverview Regression tests for the header connection indicator + * (`CodemanApp._updateConnectionIndicator`) and the invariant that updating it + * never aborts durable input delivery. + * + * Guards two regressions from the upstream v1.1.15 merge (fixed in COD-133): + * 1. The indicator body referenced an undefined `transport` object, throwing + * `ReferenceError: transport is not defined` on every queued state. Because + * `_reliableSend()` calls `_updateConnectionIndicator()` *before* + * `_drainSession()`, the throw skipped immediate delivery on every keystroke + * → input only flushed on the 2s sweep (large typing lag) and the indicator + * never rendered (missing "WS" status). + * 2. The restored body only read the SSE `_connectionStatus`, so it never + * surfaced the terminal WebSocket transport ("WS" / "HTTP"), and it flashed + * "sending 1B" on every single keystroke. + * + * Loaded via `vm` with a stubbed context (no jsdom — see input-send-order.test.ts). + */ +import { readFileSync } from 'node:fs'; +import { performance } from 'node:perf_hooks'; +import { resolve } from 'node:path'; +import vm from 'node:vm'; +import { describe, expect, it, vi } from 'vitest'; + +function loadCodemanAppClass() { + const constants = readFileSync(resolve(import.meta.dirname, '../src/web/public/constants.js'), 'utf8'); + const source = readFileSync(resolve(import.meta.dirname, '../src/web/public/app.js'), 'utf8'); + const context = vm.createContext({ + console, + performance, + setInterval: vi.fn(), + clearInterval: vi.fn(), + setTimeout, + clearTimeout, + requestAnimationFrame: vi.fn(), + HTMLCanvasElement: class HTMLCanvasElement {}, + fetch: (...args: Parameters) => global.fetch(...args), + document: { addEventListener: vi.fn() }, + localStorage: { + length: 0, + key: vi.fn(), + getItem: vi.fn(), + setItem: vi.fn(), + removeItem: vi.fn(), + }, + window: { addEventListener: vi.fn(), removeEventListener: vi.fn() }, + MobileDetection: {}, + }); + vm.runInContext(`${constants}\n${source}\nglobalThis.__CodemanApp = CodemanApp;`, context); + return (context as { __CodemanApp: new () => unknown }).__CodemanApp; +} + +const CodemanApp = loadCodemanAppClass(); + +function fakeElement() { + return { style: { display: '' }, title: '', textContent: '', className: '' }; +} + +type Indicator = { + $: (id: string) => unknown; + _pendingDeliveries: Map>; + _connectionStatus: string; + _wsState: string; + activeSessionId: string | null; + isOnline: boolean; + _updateConnectionIndicator: () => void; +}; + +function makeApp(overrides: Partial & { queuedBytes?: number } = {}) { + const app = Object.create((CodemanApp as { prototype: object }).prototype) as Indicator & { + queuedBytes?: number; + }; + const els: Record> = { + connectionIndicator: fakeElement(), + connectionDot: fakeElement(), + connectionText: fakeElement(), + }; + app.$ = (id: string) => els[id]; + app._pendingDeliveries = new Map(); + app._connectionStatus = 'connected'; + app._wsState = 'disconnected'; + app.activeSessionId = null; + app.isOnline = true; + Object.assign(app, overrides); + const queued = overrides.queuedBytes ?? 0; + if (queued > 0) { + app._pendingDeliveries.set('s1', [{ seq: 1, data: 'x'.repeat(queued) }]); + } + return { app, els }; +} + +describe('connection indicator — transport display', () => { + it('shows "WS" with a connected dot when the terminal WebSocket is open', () => { + const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected' }); + app._updateConnectionIndicator(); + expect(els.connectionIndicator.style.display).toBe('flex'); + expect(els.connectionText.textContent).toBe('WS'); + expect(els.connectionDot.className).toContain('connected'); + }); + + it('shows "HTTP" with a fallback dot when the socket dropped to HTTP POST', () => { + const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'fallback' }); + app._updateConnectionIndicator(); + expect(els.connectionText.textContent).toBe('HTTP'); + expect(els.connectionDot.className).toContain('fallback'); + }); + + it('shows "WS…" while connecting or reconnecting the socket', () => { + for (const state of ['connecting', 'reconnecting']) { + const { app, els } = makeApp({ activeSessionId: 's1', _wsState: state }); + app._updateConnectionIndicator(); + expect(els.connectionText.textContent).toBe('WS…'); + expect(els.connectionDot.className).toContain('reconnecting'); + } + }); + + it('shows "Offline" when the browser reports no network', () => { + const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected', isOnline: false }); + app._updateConnectionIndicator(); + expect(els.connectionText.textContent).toBe('Offline'); + expect(els.connectionDot.className).toContain('offline'); + }); + + it('hides on an idle dashboard (no active session, healthy stream, no queue)', () => { + const { app, els } = makeApp({ activeSessionId: null, _connectionStatus: 'connected' }); + app._updateConnectionIndicator(); + expect(els.connectionIndicator.style.display).toBe('none'); + }); + + it('surfaces SSE reconnecting on the dashboard when there is no active terminal', () => { + const { app, els } = makeApp({ activeSessionId: null, _connectionStatus: 'reconnecting' }); + app._updateConnectionIndicator(); + expect(els.connectionText.textContent).toContain('Reconnecting'); + expect(els.connectionDot.className).toContain('reconnecting'); + }); +}); + +describe('connection indicator — keystroke backlog threshold', () => { + it('does NOT annotate a single-keystroke (1B) queue — no "sending 1B" flicker', () => { + const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected', queuedBytes: 1 }); + app._updateConnectionIndicator(); + expect(els.connectionText.textContent).toBe('WS'); + expect(els.connectionText.textContent).not.toMatch(/queued/); + }); + + it('does NOT annotate at the 4B threshold boundary', () => { + const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected', queuedBytes: 4 }); + app._updateConnectionIndicator(); + expect(els.connectionText.textContent).toBe('WS'); + }); + + it('annotates a genuine backlog (>4B) with a queued byte count', () => { + const { app, els } = makeApp({ activeSessionId: 's1', _wsState: 'connected', queuedBytes: 40 }); + app._updateConnectionIndicator(); + expect(els.connectionText.textContent).toMatch(/^WS · 40B queued$/); + }); +}); + +describe('connection indicator — never throws (the ReferenceError regression)', () => { + it('renders every transport × stream × queue combination without throwing', () => { + const wsStates = ['disconnected', 'connecting', 'connected', 'reconnecting', 'fallback']; + const sseStates = ['connected', 'connecting', 'reconnecting', 'disconnected', 'offline']; + for (const ws of wsStates) { + for (const sse of sseStates) { + for (const active of ['s1', null] as const) { + for (const queuedBytes of [0, 1, 4, 200]) { + for (const isOnline of [true, false]) { + const { app } = makeApp({ + activeSessionId: active, + _wsState: ws, + _connectionStatus: sse, + isOnline, + queuedBytes, + }); + expect(() => app._updateConnectionIndicator()).not.toThrow(); + } + } + } + } + } + }); +}); + +describe('durable input delivery is not aborted by the indicator (typing-lag regression)', () => { + it('_reliableSend reaches _drainSession after updating the indicator', () => { + const { app } = makeApp({ activeSessionId: 's1', _wsState: 'connected' }); + const a = app as unknown as { + _seqCounters: Map; + _persistReliableState: () => void; + _drainSession: (id: string) => void; + _reliableSend: (id: string, data: string, useMux: boolean) => void; + }; + a._seqCounters = new Map(); + a._persistReliableState = vi.fn(); + const drain = vi.fn(); + a._drainSession = drain; + + // A single keystroke. The indicator runs first; if it throws, drain is skipped. + a._reliableSend('s1', 'x', false); + + expect(drain).toHaveBeenCalledWith('s1'); + expect(app._pendingDeliveries.get('s1')).toHaveLength(1); + }); +}); diff --git a/test/input-send-order.test.ts b/test/input-send-order.test.ts new file mode 100644 index 00000000..cc7c6294 --- /dev/null +++ b/test/input-send-order.test.ts @@ -0,0 +1,159 @@ +/** + * @fileoverview Input dispatch ordering for the durable, acknowledged delivery + * layer (`CodemanApp._sendInputAsync` → `_reliableSend` → `_drainSession`). + * + * Local replaced the upstream best-effort "coalescing fallback queue" with the + * durable per-(clientId, seq) layer in commit 1255e28 (docs/reliable-input- + * delivery.md). This suite verifies the client-side ordering guarantees of that + * layer: each input is a distinct seq-tagged frame, delivered in order over the + * WebSocket when open, serialized over HTTP POST when not, and only dropped on a + * server ACK (HTTP 2xx). Exactly-once application is covered server-side in + * test/reliable-input-dedup.test.ts; the header transport indicator ("WS"/"HTTP") + * is covered in test/connection-indicator.test.ts. + * + * Loaded via `vm` with a stubbed context (no jsdom). + */ +import { readFileSync } from 'node:fs'; +import { performance } from 'node:perf_hooks'; +import { resolve } from 'node:path'; +import vm from 'node:vm'; +import { describe, expect, it, vi } from 'vitest'; + +function loadCodemanAppClass() { + const constants = readFileSync(resolve(import.meta.dirname, '../src/web/public/constants.js'), 'utf8'); + const source = readFileSync(resolve(import.meta.dirname, '../src/web/public/app.js'), 'utf8'); + const context = vm.createContext({ + console, + performance, + setInterval: vi.fn(), + clearInterval: vi.fn(), + setTimeout, + clearTimeout, + requestAnimationFrame: vi.fn(), + HTMLCanvasElement: class HTMLCanvasElement {}, + WebSocket: { OPEN: 1 }, + fetch: (...args: Parameters) => global.fetch(...args), + document: { addEventListener: vi.fn() }, + localStorage: { + length: 0, + key: vi.fn(), + getItem: vi.fn(), + setItem: vi.fn(), + removeItem: vi.fn(), + }, + window: { addEventListener: vi.fn(), removeEventListener: vi.fn() }, + MobileDetection: {}, + }); + vm.runInContext(`${constants}\n${source}\nglobalThis.__CodemanApp = CodemanApp;`, context); + return (context as { __CodemanApp: new () => unknown }).__CodemanApp; +} + +const CodemanApp = loadCodemanAppClass(); + +async function waitForCalls(calls: unknown[], count: number) { + for (let i = 0; i < 50; i++) { + if (calls.length >= count) return; + await new Promise((r) => setTimeout(r, 0)); + } +} + +type Frame = { t: string; d: string; seq: number; cid: string }; +type PostBody = { input: string; seq: number; clientId: string }; + +type App = { + _sendInputAsync: (sessionId: string, input: string, opts?: { useMux?: boolean }) => void; + _pendingDeliveries: Map>; + _ws: { readyState: number; send: (data: string) => void } | null; + _wsSessionId: string | null; + activeSessionId: string | null; +}; + +function makeApp(): App { + const app = Object.create((CodemanApp as { prototype: object }).prototype) as App & Record; + app._clientId = 'c-test'; + app._seqCounters = new Map(); + app._pendingDeliveries = new Map(); + app._postDraining = new Set(); + app._persistReliableState = vi.fn(); + app._persistReliableNow = vi.fn(); + app._updateConnectionIndicator = vi.fn(); + app.clearPendingHooks = vi.fn(); + app.activeSessionId = 'session-1'; + app.isOnline = true; + app._connectionStatus = 'connected'; + app._ws = null; + app._wsSessionId = null; + return app as unknown as App; +} + +describe('durable input delivery — send ordering', () => { + it('delivers rapid input as distinct ordered seq frames over an open WebSocket (no coalescing)', () => { + const app = makeApp(); + const frames: Frame[] = []; + app._ws = { readyState: 1, send: (d: string) => frames.push(JSON.parse(d)) }; + app._wsSessionId = 'session-1'; + + app._sendInputAsync('session-1', 'a'); + app._sendInputAsync('session-1', 'b'); + app._sendInputAsync('session-1', 'c'); + + // Each keystroke is its own frame, in seq order — never merged into "abc". + expect(frames.map((f) => f.d)).toEqual(['a', 'b', 'c']); + expect(frames.map((f) => f.seq)).toEqual([1, 2, 3]); + expect(frames.every((f) => f.t === 'i' && f.cid === 'c-test')).toBe(true); + }); + + it('POSTs queued input one frame at a time in seq order when no socket is open', async () => { + const app = makeApp(); + const calls: PostBody[] = []; + const completions: Array<() => void> = []; + global.fetch = vi.fn(async (_url, init) => { + calls.push(JSON.parse(String(init?.body)) as PostBody); + await new Promise((r) => completions.push(r)); + return new Response('{}', { status: 200 }); + }); + + app._sendInputAsync('session-1', 'a'); + app._sendInputAsync('session-1', 'b'); + + // Serialized: only the first frame is in flight until its 2xx ACK lands. + await waitForCalls(calls, 1); + expect(calls.map((c) => c.input)).toEqual(['a']); + + completions.shift()?.(); // ACK 'a' + await waitForCalls(calls, 2); + expect(calls.map((c) => c.input)).toEqual(['a', 'b']); + expect(calls.map((c) => c.seq)).toEqual([1, 2]); + + completions.shift()?.(); + await waitForCalls(calls, 2); + }); + + it('leaves a frame queued (unacked) when HTTP delivery fails', async () => { + const app = makeApp(); + const calls: PostBody[] = []; + global.fetch = vi.fn(async (_url, init) => { + calls.push(JSON.parse(String(init?.body)) as PostBody); + return new Response('busy', { status: 503 }); + }); + + app._sendInputAsync('session-1', 'a'); + await waitForCalls(calls, 1); + await new Promise((r) => setTimeout(r, 0)); + + // 5xx is not an ACK — the frame must survive for the sweep/reconnect to retry. + expect(app._pendingDeliveries.get('session-1')).toHaveLength(1); + expect(app._pendingDeliveries.get('session-1')?.[0].data).toBe('a'); + }); + + it('drops a frame addressed to a vanished session (404) instead of retrying forever', async () => { + const app = makeApp(); + global.fetch = vi.fn(async () => new Response('gone', { status: 404 })); + + app._sendInputAsync('session-1', 'a'); + await new Promise((r) => setTimeout(r, 0)); + await new Promise((r) => setTimeout(r, 0)); + + expect(app._pendingDeliveries.get('session-1')).toBeUndefined(); + }); +});