From 3383c23099163a882c7ae7a75164944c141fa66e Mon Sep 17 00:00:00 2001 From: arkon Date: Sat, 14 Mar 2026 18:27:58 +0100 Subject: [PATCH] fix: address code review findings across WS, CJK input, install.sh, and README MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WebSocket route: add socket error handler to prevent process crashes, enforce per-session connection limit (max 5), track/decrement counts on close. CJK input: add destroy() method with proper listener cleanup, guard against double-init, add maxlength/aria-label to textarea, use language-neutral placeholder, explicitly clear cjkActive on hide. install.sh: fix update() to use $BRANCH and $REPO_URL instead of hardcoded origin/master — fork users were silently switched back to master on update. README: fix broken markdown table (paragraph concatenated into last cell), add CODEMAN_NODE_VERSION to env var table. Tests: add 8 new test cases for batch coalescing, flush threshold, unknown message types, connection limit, heartbeat, readyState guards. Import MAX_INPUT_LENGTH from config, add connectWs timeout, replace setTimeout with vi.waitFor in cleanup test. Co-Authored-By: Claude Opus 4.6 --- README.md | 5 +- install.sh | 3 +- src/web/public/app.js | 5 +- src/web/public/index.html | 3 +- src/web/public/input-cjk.js | 40 +++++++-- src/web/routes/ws-routes.ts | 25 ++++++ test/routes/ws-routes.test.ts | 164 +++++++++++++++++++++++++++++++--- 7 files changed, 222 insertions(+), 23 deletions(-) diff --git a/README.md b/README.md index ad414a9e..6806e37e 100644 --- a/README.md +++ b/README.md @@ -45,7 +45,10 @@ The installer supports these environment variables: | `CODEMAN_BRANCH` | `master` | Git branch to install | | `CODEMAN_INSTALL_DIR` | `~/.codeman/app` | Custom install directory | | `CODEMAN_SKIP_SYSTEMD` | `0` | Skip systemd service setup prompt | -| `CODEMAN_NONINTERACTIVE` | `0` | Skip all prompts (for CI/automation) | You'll need at least one AI coding CLI installed — [Claude Code](https://docs.anthropic.com/en/docs/claude-code) or [OpenCode](https://opencode.ai) (or both). After install: +| `CODEMAN_NODE_VERSION` | `22` | Node.js major version to install | +| `CODEMAN_NONINTERACTIVE` | `0` | Skip all prompts (for CI/automation) | + +You'll need at least one AI coding CLI installed — [Claude Code](https://docs.anthropic.com/en/docs/claude-code) or [OpenCode](https://opencode.ai) (or both). After install: ```bash codeman web diff --git a/install.sh b/install.sh index 4c71acbf..2d7d6b8b 100755 --- a/install.sh +++ b/install.sh @@ -1281,8 +1281,9 @@ update() { info "Updating Codeman..." cd "$INSTALL_DIR" + git remote set-url origin "$REPO_URL" 2>/dev/null || true git fetch --quiet origin - git reset --hard origin/master --quiet + git reset --hard "origin/$BRANCH" --quiet npm install --quiet --no-fund --no-audit 2>/dev/null || npm install --no-fund --no-audit npm run build --quiet 2>/dev/null || npm run build success "Updated to $(node -e "console.log(require('./package.json').version)")" diff --git a/src/web/public/app.js b/src/web/public/app.js index 9d158e10..dec17f68 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -3083,7 +3083,10 @@ class CodemanApp { // CJK input form: show/hide based on server env INPUT_CJK_FORM=ON const cjkEl = document.getElementById('cjkInput'); - if (cjkEl) cjkEl.style.display = data.inputCjkForm ? 'block' : 'none'; + if (cjkEl) { + cjkEl.style.display = data.inputCjkForm ? 'block' : 'none'; + if (!data.inputCjkForm) window.cjkActive = false; + } // Update version displays (header and toolbar) if (data.version) { diff --git a/src/web/public/index.html b/src/web/public/index.html index 275285c6..641cd3af 100644 --- a/src/web/public/index.html +++ b/src/web/public/index.html @@ -233,7 +233,8 @@
-
diff --git a/src/web/public/input-cjk.js b/src/web/public/input-cjk.js index a3d0774a..0cba507e 100644 --- a/src/web/public/input-cjk.js +++ b/src/web/public/input-cjk.js @@ -7,14 +7,20 @@ * While this textarea has focus, window.cjkActive = true blocks xterm's onData. * Arrow keys and function keys are forwarded to PTY directly. * - * @globals {object} CjkInput - * @loadorder 5.5 of 9 — loaded after keyboard-accessory.js, before app.js + * @dependency index.html (#cjkInput textarea) + * @globals {object} CjkInput — window.cjkActive (boolean) signals app.js to block xterm onData + * @loadorder 5.5 of 10 — loaded after keyboard-accessory.js, before app.js */ // eslint-disable-next-line no-unused-vars const CjkInput = (() => { let _textarea = null; let _send = null; + let _initialized = false; + let _onMousedown = null; + let _onFocus = null; + let _onBlur = null; + let _onKeydown = null; const PASSTHROUGH_KEYS = { ArrowUp: '\x1b[A', @@ -32,15 +38,21 @@ const CjkInput = (() => { return { init({ send }) { + // Guard against double-init: remove previous listeners + if (_initialized) this.destroy(); + _send = send; _textarea = document.getElementById('cjkInput'); if (!_textarea) return this; - _textarea.addEventListener('mousedown', (e) => { e.stopPropagation(); }); - _textarea.addEventListener('focus', () => { window.cjkActive = true; }); - _textarea.addEventListener('blur', () => { window.cjkActive = false; }); + _onMousedown = (e) => { e.stopPropagation(); }; + _onFocus = () => { window.cjkActive = true; }; + _onBlur = () => { window.cjkActive = false; }; + _textarea.addEventListener('mousedown', _onMousedown); + _textarea.addEventListener('focus', _onFocus); + _textarea.addEventListener('blur', _onBlur); - _textarea.addEventListener('keydown', (e) => { + _onKeydown = (e) => { if (e.isComposing || e.keyCode === 229) return; // Enter: send accumulated text (or bare Enter if empty) @@ -82,11 +94,25 @@ const CjkInput = (() => { _send(PASSTHROUGH_KEYS[e.key]); return; } - }); + }; + _textarea.addEventListener('keydown', _onKeydown); + _initialized = true; return this; }, + destroy() { + if (_textarea) { + if (_onMousedown) _textarea.removeEventListener('mousedown', _onMousedown); + if (_onFocus) _textarea.removeEventListener('focus', _onFocus); + if (_onBlur) _textarea.removeEventListener('blur', _onBlur); + if (_onKeydown) _textarea.removeEventListener('keydown', _onKeydown); + } + window.cjkActive = false; + _onMousedown = _onFocus = _onBlur = _onKeydown = null; + _initialized = false; + }, + get element() { return _textarea; }, }; })(); diff --git a/src/web/routes/ws-routes.ts b/src/web/routes/ws-routes.ts index fe490608..79c55e3d 100644 --- a/src/web/routes/ws-routes.ts +++ b/src/web/routes/ws-routes.ts @@ -52,6 +52,12 @@ const WS_PONG_TIMEOUT_MS = 10_000; const DEC_2026_START = '\x1b[?2026h'; const DEC_2026_END = '\x1b[?2026l'; +/** Max concurrent WS connections per session. Prevents listener/bandwidth multiplication. */ +const MAX_WS_PER_SESSION = 5; + +/** Track active WS connections per session for connection limiting. */ +const sessionWsCount = new Map(); + export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort): void { app.get<{ Params: { id: string } }>('/ws/sessions/:id/terminal', { websocket: true }, (socket: WebSocket, req) => { const { id } = req.params; @@ -62,6 +68,17 @@ export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort): void { return; } + // Enforce per-session connection limit + const currentCount = sessionWsCount.get(id) ?? 0; + if (currentCount >= MAX_WS_PER_SESSION) { + socket.close(4008, 'Too many connections'); + return; + } + sessionWsCount.set(id, currentCount + 1); + + // Swallow socket errors — cleanup happens in 'close' + socket.on('error', () => {}); + // Per-connection micro-batch state let batchChunks: string[] = []; let batchSize = 0; @@ -174,6 +191,14 @@ export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort): void { session.off('terminal', onTerminal); session.off('clearTerminal', onClearTerminal); session.off('needsRefresh', onNeedsRefresh); + + // Decrement per-session connection count + const count = sessionWsCount.get(id) ?? 1; + if (count <= 1) { + sessionWsCount.delete(id); + } else { + sessionWsCount.set(id, count - 1); + } }); }); } diff --git a/test/routes/ws-routes.test.ts b/test/routes/ws-routes.test.ts index a20c8775..1cc6271f 100644 --- a/test/routes/ws-routes.test.ts +++ b/test/routes/ws-routes.test.ts @@ -5,6 +5,8 @@ * a real listening server since inject() doesn't support upgrade requests. * Uses the `ws` package (transitive dep of @fastify/websocket) as the client. * + * @dependency test/mocks/mock-route-context.ts (createMockRouteContext) + * @dependency src/web/routes/ws-routes.ts (registerWsRoutes) * Port: 3170 (ws-routes tests) */ @@ -14,15 +16,23 @@ import fastifyWebsocket from '@fastify/websocket'; import WebSocket from 'ws'; import { createMockRouteContext, type MockRouteContext } from '../mocks/index.js'; import { registerWsRoutes } from '../../src/web/routes/ws-routes.js'; +import { MAX_INPUT_LENGTH } from '../../src/config/terminal-limits.js'; const PORT = 3170; /** Helper: open a WebSocket connection and wait for it to reach OPEN state. */ -function connectWs(path: string): Promise { +function connectWs(path: string, timeoutMs = 5000): Promise { return new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error('WS connection timeout')), timeoutMs); const ws = new WebSocket(`ws://127.0.0.1:${PORT}${path}`); - ws.on('open', () => resolve(ws)); - ws.on('error', reject); + ws.on('open', () => { + clearTimeout(timer); + resolve(ws); + }); + ws.on('error', (err) => { + clearTimeout(timer); + reject(err); + }); }); } @@ -48,6 +58,23 @@ function waitForClose(ws: WebSocket, timeoutMs = 2000): Promise<{ code: number; }); } +/** Helper: collect N messages from a WebSocket. */ +function collectMessages(ws: WebSocket, count: number, timeoutMs = 3000): Promise { + return new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error(`Only received ${msgs.length}/${count} messages`)), timeoutMs); + const msgs: unknown[] = []; + const onMessage = (raw: WebSocket.RawData) => { + msgs.push(JSON.parse(String(raw))); + if (msgs.length >= count) { + clearTimeout(timer); + ws.off('message', onMessage); + resolve(msgs); + } + }; + ws.on('message', onMessage); + }); +} + describe('ws-routes', () => { let app: FastifyInstance; let ctx: MockRouteContext; @@ -123,6 +150,43 @@ describe('ws-routes', () => { ws.close(); } }); + + it('coalesces rapid terminal emissions into a single frame', async () => { + const ws = await connectWs('/ws/sessions/ws-test-session/terminal'); + try { + const session = ctx._session; + + // Emit multiple small chunks in rapid succession (within the 8ms batch window) + session.emit('terminal', 'chunk1'); + session.emit('terminal', 'chunk2'); + session.emit('terminal', 'chunk3'); + + // Should arrive as a single coalesced message + const msg = (await nextMessage(ws)) as { t: string; d: string }; + expect(msg.t).toBe('o'); + expect(msg.d).toContain('chunk1chunk2chunk3'); + } finally { + ws.close(); + } + }); + + it('flushes immediately when batch exceeds size threshold', async () => { + const ws = await connectWs('/ws/sessions/ws-test-session/terminal'); + try { + const session = ctx._session; + + // Emit data larger than WS_BATCH_FLUSH_THRESHOLD (16384) + const largeData = 'X'.repeat(17000); + session.emit('terminal', largeData); + + // Should flush immediately (no 8ms wait) — use a tight timeout + const msg = (await nextMessage(ws, 500)) as { t: string; d: string }; + expect(msg.t).toBe('o'); + expect(msg.d).toContain(largeData); + } finally { + ws.close(); + } + }); }); // ========== Client input ========== @@ -132,7 +196,6 @@ describe('ws-routes', () => { const ws = await connectWs('/ws/sessions/ws-test-session/terminal'); try { const session = ctx._session; - session.writeBuffer = []; ws.send(JSON.stringify({ t: 'i', d: 'ls -la\r' })); @@ -149,10 +212,8 @@ describe('ws-routes', () => { const ws = await connectWs('/ws/sessions/ws-test-session/terminal'); try { const session = ctx._session; - session.writeBuffer = []; - // MAX_INPUT_LENGTH is 64KB - const hugeInput = 'x'.repeat(65 * 1024); + const hugeInput = 'x'.repeat(MAX_INPUT_LENGTH + 1); ws.send(JSON.stringify({ t: 'i', d: hugeInput })); // Send a valid message after to confirm the connection still works @@ -173,7 +234,6 @@ describe('ws-routes', () => { const ws = await connectWs('/ws/sessions/ws-test-session/terminal'); try { const session = ctx._session; - session.writeBuffer = []; ws.send('not-json{{{'); @@ -190,6 +250,28 @@ describe('ws-routes', () => { ws.close(); } }); + + it('ignores unknown message types without breaking the connection', async () => { + const ws = await connectWs('/ws/sessions/ws-test-session/terminal'); + try { + const session = ctx._session; + + // Send unknown type + ws.send(JSON.stringify({ t: 'x', d: 'mystery' })); + + // Connection should still work + ws.send(JSON.stringify({ t: 'i', d: 'still-alive' })); + + await vi.waitFor(() => { + expect(session.writeBuffer).toContain('still-alive'); + }); + + // Unknown type should not have been written + expect(session.writeBuffer).toHaveLength(1); + } finally { + ws.close(); + } + }); }); // ========== Resize validation ========== @@ -332,6 +414,64 @@ describe('ws-routes', () => { }); }); + // ========== Connection limit ========== + + describe('connection limit', () => { + it('closes with 4008 when too many connections per session', async () => { + const connections: WebSocket[] = []; + try { + // Open 5 connections (the max) + for (let i = 0; i < 5; i++) { + connections.push(await connectWs('/ws/sessions/ws-test-session/terminal')); + } + + // 6th connection should be rejected + const ws6 = new WebSocket(`ws://127.0.0.1:${PORT}/ws/sessions/ws-test-session/terminal`); + const { code, reason } = await waitForClose(ws6); + expect(code).toBe(4008); + expect(reason).toBe('Too many connections'); + } finally { + for (const ws of connections) ws.close(); + } + }); + }); + + // ========== Heartbeat ========== + + describe('heartbeat', () => { + it('responds to server ping with pong (connection stays alive)', async () => { + const ws = await connectWs('/ws/sessions/ws-test-session/terminal'); + try { + // The ws library automatically responds to pings with pongs. + // Verify the connection survives by sending data after a brief delay. + const session = ctx._session; + session.emit('terminal', 'heartbeat-test'); + + const msg = (await nextMessage(ws)) as { t: string; d: string }; + expect(msg.t).toBe('o'); + expect(msg.d).toContain('heartbeat-test'); + } finally { + ws.close(); + } + }); + }); + + // ========== readyState guards ========== + + describe('readyState guards', () => { + it('does not throw when clearTerminal fires after close', async () => { + const session = ctx._session; + const ws = await connectWs('/ws/sessions/ws-test-session/terminal'); + + ws.close(); + await waitForClose(ws); + + // These should be no-ops, not throw + expect(() => session.emit('clearTerminal')).not.toThrow(); + expect(() => session.emit('needsRefresh')).not.toThrow(); + }); + }); + // ========== Connection cleanup ========== describe('connection cleanup', () => { @@ -350,11 +490,11 @@ describe('ws-routes', () => { ws.close(); await waitForClose(ws); - // Give the server-side close handler time to run - await new Promise((resolve) => setTimeout(resolve, 50)); + // Wait for server-side close handler + await vi.waitFor(() => { + expect(session.listenerCount('terminal')).toBe(listenersBefore); + }); - // Listeners should be cleaned up - expect(session.listenerCount('terminal')).toBe(listenersBefore); expect(session.listenerCount('clearTerminal')).toBe(0); expect(session.listenerCount('needsRefresh')).toBe(0); });