Compare commits

..
Author SHA1 Message Date
arkon d072e773d8 chore: version packages 2026-03-14 18:38:28 +01:00
arkonandClaude Opus 4.6 a649c91b68 fix: WS session lifecycle, reconnection, and CJK session-switch cleanup
- Close WebSocket when session exits (exit event listener) to prevent
  orphaned listeners and stale writes to dead PTY
- Add readyState guard in onTerminal to stop buffering after socket closes
- Simplify heartbeat: remove redundant alive flag, use pongTimeout only
- Add exponential backoff reconnection on unexpected WS close (skip for
  server rejections 4004/4008/4009)
- Clear CJK textarea on session switch to prevent wrong-session input

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 18:37:10 +01:00
arkonandClaude Opus 4.6 3383c23099 fix: address code review findings across WS, CJK input, install.sh, and README
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 <noreply@anthropic.com>
2026-03-14 18:27:58 +01:00
arkonandClaude Opus 4.6 c3e1e731ef fix: use generic placeholders in fork install README example
Replace hardcoded contributor fork URL with <user>/<branch> placeholders
so the documentation is useful for any contributor.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 18:11:43 +01:00
Ark0N 405b711c3a Merge pull request #41 from douchekr/feat/input-cjk-form
feat: add CJK IME input textarea and fork/branch install support
2026-03-14 18:10:09 +01:00
jayparkandClaude Opus 4.6 393a2d9c28 fix: use BRANCH variable in install.sh no-changes update path
Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-03-14 16:52:28 +09:00
jayparkandClaude Opus 4.6 809bf6a614 fix: use BRANCH variable in install.sh update path
The update path was hardcoded to origin/master. Now uses the
CODEMAN_BRANCH variable and updates the remote URL on upgrade.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-03-14 16:51:20 +09:00
jayparkandClaude Opus 4.6 da71d8d01c feat: support custom repo URL and branch in install.sh
Add CODEMAN_REPO_URL and CODEMAN_BRANCH env vars to install.sh
for installing from forks or feature branches. Update README with
fork installation instructions and env var reference table.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-03-14 16:50:06 +09:00
jayparkandClaude Opus 4.6 e5aca6aa4c feat: add CJK IME input textarea with env toggle
Add a dedicated textarea below the terminal for CJK (Korean/Japanese/Chinese)
IME input. xterm.js intercepts IME composition events, preventing composed
characters from displaying correctly. This textarea bypasses xterm entirely
by using native browser IME handling — text accumulates until Enter, then
sends to PTY in one shot.

- Always-visible textarea below terminal (inside .terminal-wrap flex column)
- focus/blur sets window.cjkActive flag to block xterm onData
- Enter sends textarea.value + \r to PTY, Escape clears
- Arrow keys, Ctrl+C/D/L/Z, Tab, Backspace pass through to PTY when empty
- attachCustomKeyEventHandler suppresses xterm key handling during composition
- INPUT_CJK_FORM=ON|OFF env var toggle (default: off, passed via SSE init)

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-03-14 16:23:25 +09:00
13 changed files with 460 additions and 35 deletions
+14
View File
@@ -1,5 +1,19 @@
# aicodeman
## 0.4.0
### Minor Changes
- Add CJK IME input textarea for xterm.js terminal (env toggle INPUT_CJK_FORM=ON). Always-visible textarea below terminal handles native browser IME composition, forwarding completed text to PTY on Enter. Supports arrow keys, Ctrl combos, backspace passthrough, and Escape to clear.
Add fork installation support to install.sh with CODEMAN_REPO_URL and CODEMAN_BRANCH env vars, allowing custom repository and branch for git clone/update operations. README updated with fork installation instructions.
Fix WebSocket session lifecycle: close WS connections when session exits (prevents orphaned listeners and stale writes to dead PTY), add readyState guard in onTerminal to stop buffering after socket closes, simplify heartbeat by removing redundant alive flag.
Add WebSocket reconnection with exponential backoff (1s-10s) on unexpected close, skipping server rejection codes (4004/4008/4009). Falls back gracefully to SSE+POST during reconnection.
Clear CJK textarea on session switch to prevent sending stale text to wrong session.
## 0.3.12
### Patch Changes
+1 -1
View File
@@ -52,7 +52,7 @@ When user says "COM":
4. **Sync CLAUDE.md version**: Update the `**Version**` line below to match the new version from `package.json`
5. **Commit and deploy**: `git add -A && git commit -m "chore: version packages" && git push && npm run build && systemctl --user restart codeman-web`
**Version**: 0.3.12 (must match `package.json`)
**Version**: 0.4.0 (must match `package.json`)
## Project Overview
+21 -1
View File
@@ -28,7 +28,27 @@
curl -fsSL https://raw.githubusercontent.com/Ark0N/Codeman/master/install.sh | bash
```
This installs Node.js and tmux if missing, clones Codeman to `~/.codeman/app`, and builds it. 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:
This installs Node.js and tmux if missing, clones Codeman to `~/.codeman/app`, and builds it.
**Install from a fork or specific branch:**
```bash
curl -fsSL https://raw.githubusercontent.com/<user>/Codeman/<branch>/install.sh | \
CODEMAN_REPO_URL=https://github.com/<user>/Codeman.git \
CODEMAN_BRANCH=<branch> bash
```
The installer supports these environment variables:
| Variable | Default | Description |
|----------|---------|-------------|
| `CODEMAN_REPO_URL` | upstream Codeman | Custom git repository URL |
| `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_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
+10 -5
View File
@@ -9,6 +9,8 @@
# CODEMAN_INSTALL_DIR - Custom install directory (default: ~/.codeman/app)
# CODEMAN_SKIP_SYSTEMD=1 - Skip systemd service setup prompt
# CODEMAN_NODE_VERSION - Node.js major version to install (default: 22)
# CODEMAN_REPO_URL - Custom git repository URL (default: upstream Codeman)
# CODEMAN_BRANCH - Git branch to install (default: master)
set -euo pipefail
@@ -17,7 +19,8 @@ set -euo pipefail
# ============================================================================
INSTALL_DIR="${CODEMAN_INSTALL_DIR:-$HOME/.codeman/app}"
REPO_URL="https://github.com/Ark0N/Codeman.git"
REPO_URL="${CODEMAN_REPO_URL:-https://github.com/Ark0N/Codeman.git}"
BRANCH="${CODEMAN_BRANCH:-master}"
MIN_NODE_VERSION=18
TARGET_NODE_VERSION="${CODEMAN_NODE_VERSION:-22}"
NONINTERACTIVE="${CODEMAN_NONINTERACTIVE:-0}"
@@ -1062,26 +1065,27 @@ main() {
if [[ -d "$INSTALL_DIR/.git" ]]; then
info "Existing installation found, updating..."
cd "$INSTALL_DIR"
git remote set-url origin "$REPO_URL" 2>/dev/null || true
# Check for local changes
if ! git diff --quiet 2>/dev/null || ! git diff --staged --quiet 2>/dev/null; then
warn "Local changes detected in $INSTALL_DIR"
if prompt_yes_no "Discard local changes and update?" "n"; then
git fetch --quiet origin
git reset --hard origin/master --quiet
git reset --hard "origin/$BRANCH" --quiet
else
info "Keeping existing installation, skipping update"
fi
else
git fetch --quiet origin
git reset --hard origin/master --quiet
git reset --hard "origin/$BRANCH" --quiet
fi
else
# Create parent directory
mkdir -p "$(dirname "$INSTALL_DIR")"
# Clone repository (shallow for speed)
git clone --quiet --depth 1 "$REPO_URL" "$INSTALL_DIR"
git clone --quiet --depth 1 --branch "$BRANCH" "$REPO_URL" "$INSTALL_DIR"
cd "$INSTALL_DIR"
fi
@@ -1277,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)")"
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "aicodeman",
"version": "0.3.12",
"version": "0.4.0",
"description": "The missing control plane for AI coding agents - run 20 autonomous agents with real-time monitoring and session persistence",
"type": "module",
"main": "dist/index.js",
+2
View File
@@ -59,6 +59,7 @@ appendFileSync(
);
// 4. Minify frontend assets
run('minify input-cjk.js', 'npx esbuild dist/web/public/input-cjk.js --minify --outfile=dist/web/public/input-cjk.js --allow-overwrite');
run('minify app.js', 'npx esbuild dist/web/public/app.js --minify --outfile=dist/web/public/app.js --allow-overwrite');
run('minify styles.css', 'npx esbuild dist/web/public/styles.css --minify --outfile=dist/web/public/styles.css --allow-overwrite');
run('minify mobile.css', 'npx esbuild dist/web/public/mobile.css --minify --outfile=dist/web/public/mobile.css --allow-overwrite');
@@ -75,6 +76,7 @@ console.log('\n[build] content-hash cache busting');
'voice-input.js',
'notification-manager.js',
'keyboard-accessory.js',
'input-cjk.js',
'app.js',
'ralph-wizard.js',
'api-client.js',
+60 -6
View File
@@ -628,6 +628,14 @@ class CodemanApp {
const container = document.getElementById('terminalContainer');
this.terminal.open(container);
// Suppress xterm key handling during CJK IME composition.
// Without this, xterm processes raw keyDown events (e.g., "Process" key)
// during composition, causing duplicate or garbled input.
this.terminal.attachCustomKeyEventHandler((ev) => {
if (ev.isComposing || ev.keyCode === 229) return false;
return true;
});
// WebGL renderer for GPU-accelerated terminal rendering.
// Previously caused "page unresponsive" crashes from synchronous GPU stalls,
// but the 48KB/frame flush cap in flushPendingWrites() now prevents
@@ -650,6 +658,18 @@ class CodemanApp {
this._localEchoOverlay = new LocalEchoOverlay(this.terminal);
// CJK IME input — textarea in index.html, just wire up send
this._cjkInput = null;
if (typeof CjkInput !== 'undefined') {
this._cjkInput = CjkInput.init({
send: (text) => {
if (this.activeSessionId) {
this._sendInputAsync(this.activeSessionId, text);
}
},
});
}
// On mobile Safari, delay initial fit() to allow layout to settle
// This prevents 0-column terminals caused by fit() running before container is sized
const isMobileSafari = MobileDetection.getDeviceType() === 'mobile' &&
@@ -862,6 +882,8 @@ class CodemanApp {
// survives tab switches and reconnects.
this.terminal.onData((data) => {
// CJK input has focus — block xterm from sending to PTY
if (window.cjkActive || document.activeElement?.id === 'cjkInput') return;
if (this.activeSessionId) {
// Filter out terminal query responses that xterm.js generates automatically.
// These are responses to DA (Device Attributes), DSR (Device Status Report), etc.
@@ -1637,6 +1659,9 @@ class CodemanApp {
setupEventListeners() {
// Use capture to handle before terminal
document.addEventListener('keydown', (e) => {
// Don't intercept keys during CJK IME composition
if (e.isComposing || e.keyCode === 229) return;
// Escape - close panels and modals
if (e.key === 'Escape') {
this.closeAllPanels();
@@ -2879,6 +2904,7 @@ class CodemanApp {
// Only mark ready if this is still the intended session
if (this._ws === ws) {
this._wsReady = true;
this._wsReconnectAttempts = 0;
}
};
@@ -2899,11 +2925,24 @@ class CodemanApp {
}
};
ws.onclose = () => {
if (this._ws === ws) {
this._ws = null;
this._wsSessionId = null;
this._wsReady = false;
ws.onclose = (event) => {
if (this._ws !== ws) return;
this._ws = null;
this._wsSessionId = null;
this._wsReady = false;
// Reconnect on unexpected close (server restart, network blip, ping timeout).
// Don't reconnect if we intentionally disconnected (_disconnectWs nulls onclose)
// or if the server rejected the session (4004=not found, 4008=too many, 4009=terminated).
if (event.code < 4004 && this.activeSessionId === sessionId) {
const delay = Math.min(1000 * Math.pow(2, this._wsReconnectAttempts || 0), 10000);
this._wsReconnectAttempts = (this._wsReconnectAttempts || 0) + 1;
this._wsReconnectTimer = setTimeout(() => {
this._wsReconnectTimer = null;
if (this.activeSessionId === sessionId) {
this._connectWs(sessionId);
}
}, delay);
}
};
@@ -2914,6 +2953,11 @@ class CodemanApp {
/** Close the active WebSocket connection (if any). */
_disconnectWs() {
if (this._wsReconnectTimer) {
clearTimeout(this._wsReconnectTimer);
this._wsReconnectTimer = null;
}
this._wsReconnectAttempts = 0;
if (this._ws) {
this._ws.onclose = null; // Prevent re-entrant cleanup
this._ws.close();
@@ -3056,6 +3100,13 @@ class CodemanApp {
}
const gen = ++this._initGeneration;
// 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 (!data.inputCjkForm) window.cjkActive = false;
}
// Update version displays (header and toolbar)
if (data.version) {
const versionEl = this.$('versionDisplay');
@@ -3731,6 +3782,10 @@ class CodemanApp {
// Close WebSocket for previous session (new one opens after buffer load)
this._disconnectWs();
// Clear CJK textarea to prevent sending stale text to the wrong session
const cjkEl = document.getElementById('cjkInput');
if (cjkEl) cjkEl.value = '';
// Clean up flicker filter state when switching sessions
if (this.flickerFilterTimeout) {
clearTimeout(this.flickerFilterTimeout);
@@ -3768,7 +3823,6 @@ class CodemanApp {
if (ta) ta.dispatchEvent(new CompositionEvent('compositionend', { data: '' }));
}
} catch {}
// Flush local echo text to PTY before switching tabs.
// Send as a single batch (no Enter) so it lands in the session's readline
// input buffer — avoids "old text resent on Enter" and overlay render bugs.
+7 -1
View File
@@ -231,7 +231,12 @@
<!-- Main Terminal Area -->
<main class="main">
<div class="terminal-container" id="terminalContainer"></div>
<div class="terminal-wrap">
<div class="terminal-container" id="terminalContainer"></div>
<textarea id="cjkInput" rows="1" placeholder="CJK input (Enter = send, Esc = clear)"
maxlength="65536" aria-label="CJK IME input field"
autocomplete="off" autocorrect="off" autocapitalize="off" spellcheck="false"></textarea>
</div>
<!-- Welcome Overlay (shown when no session active) -->
<div class="welcome-overlay" id="welcomeOverlay">
@@ -1689,6 +1694,7 @@
<script defer src="voice-input.js"></script>
<script defer src="notification-manager.js"></script>
<script defer src="keyboard-accessory.js"></script>
<script defer src="input-cjk.js"></script>
<script defer src="app.js"></script>
<script defer src="ralph-wizard.js"></script>
<script defer src="api-client.js"></script>
+118
View File
@@ -0,0 +1,118 @@
/**
* @fileoverview CJK IME input for xterm.js terminal.
*
* Always-visible textarea below the terminal (in index.html).
* The browser handles IME composition natively — we just read
* textarea.value on Enter and send it to PTY.
* While this textarea has focus, window.cjkActive = true blocks xterm's onData.
* Arrow keys and function keys are forwarded to PTY directly.
*
* @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',
ArrowDown: '\x1b[B',
ArrowLeft: '\x1b[D',
ArrowRight: '\x1b[C',
Home: '\x1b[H',
End: '\x1b[F',
Tab: '\t',
};
const CTRL_KEYS = {
c: '\x03', d: '\x04', l: '\x0c', z: '\x1a', a: '\x01', e: '\x05',
};
return {
init({ send }) {
// Guard against double-init: remove previous listeners
if (_initialized) this.destroy();
_send = send;
_textarea = document.getElementById('cjkInput');
if (!_textarea) return this;
_onMousedown = (e) => { e.stopPropagation(); };
_onFocus = () => { window.cjkActive = true; };
_onBlur = () => { window.cjkActive = false; };
_textarea.addEventListener('mousedown', _onMousedown);
_textarea.addEventListener('focus', _onFocus);
_textarea.addEventListener('blur', _onBlur);
_onKeydown = (e) => {
if (e.isComposing || e.keyCode === 229) return;
// Enter: send accumulated text (or bare Enter if empty)
if (e.key === 'Enter') {
e.preventDefault();
if (_textarea.value) {
_send(_textarea.value + '\r');
_textarea.value = '';
} else {
_send('\r');
}
return;
}
// Escape: clear textarea
if (e.key === 'Escape') {
e.preventDefault();
_textarea.value = '';
return;
}
// Ctrl combos: forward to PTY
if (e.ctrlKey && CTRL_KEYS[e.key]) {
e.preventDefault();
_send(CTRL_KEYS[e.key]);
return;
}
// Backspace: delete from textarea if has text, else forward to PTY
if (e.key === 'Backspace' && !_textarea.value) {
e.preventDefault();
_send('\x7f');
return;
}
// Arrow/function keys: forward to PTY when textarea is empty
if (PASSTHROUGH_KEYS[e.key] && !_textarea.value) {
e.preventDefault();
_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; },
};
})();
+38
View File
@@ -1882,6 +1882,13 @@ body {
position: relative;
}
.terminal-wrap {
flex: 1;
display: flex;
flex-direction: column;
overflow: hidden;
}
.terminal-container {
flex: 1;
background: #0d0d0d;
@@ -7502,3 +7509,34 @@ kbd {
.advanced-options-content {
padding-left: 0.5rem;
}
/* ═══════════════════════════════════════════════════════════════
CJK IME Input
═══════════════════════════════════════════════════════════════ */
#cjkInput {
display: none;
flex-shrink: 0;
width: 100%;
font-family: 'Fira Code', 'Cascadia Code', 'JetBrains Mono', 'SF Mono', Monaco, monospace;
font-size: 14px;
background: #1a1a2e;
color: #e0e0e0;
border: 1px solid #333;
border-top: none;
padding: 6px 10px;
outline: none;
resize: none;
line-height: 1.4;
box-sizing: border-box;
}
#cjkInput:focus {
border-color: #339af0;
background: #111;
}
#cjkInput::placeholder {
color: #495057;
font-size: 12px;
}
+35 -8
View File
@@ -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<string, number>();
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;
@@ -106,6 +123,7 @@ export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort): void {
// Terminal output -> micro-batched WS send
const onTerminal = (data: string) => {
if (socket.readyState !== 1) return;
batchChunks.push(data);
batchSize += data.length;
@@ -136,17 +154,22 @@ export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort): void {
}
};
// Close WS when session exits (deleted, respawned, or crashed) — prevents
// orphaned listeners and stale writes to a dead PTY.
const onSessionExit = () => {
socket.close(4009, 'Session terminated');
};
session.on('terminal', onTerminal);
session.on('clearTerminal', onClearTerminal);
session.on('needsRefresh', onNeedsRefresh);
session.on('exit', onSessionExit);
// Heartbeat: detect stale connections (especially through tunnels where
// TCP RST can take minutes to propagate).
let pongTimeout: ReturnType<typeof setTimeout> | null = null;
let alive = true;
socket.on('pong', () => {
alive = true;
if (pongTimeout) {
clearTimeout(pongTimeout);
pongTimeout = null;
@@ -154,12 +177,7 @@ export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort): void {
});
const pingInterval = setInterval(() => {
if (!alive) {
// Previous ping never got a pong — connection is dead
socket.terminate();
return;
}
alive = false;
if (socket.readyState !== 1) return;
socket.ping();
pongTimeout = setTimeout(() => {
socket.terminate();
@@ -174,6 +192,15 @@ export function registerWsRoutes(app: FastifyInstance, ctx: SessionPort): void {
session.off('terminal', onTerminal);
session.off('clearTerminal', onClearTerminal);
session.off('needsRefresh', onNeedsRefresh);
session.off('exit', onSessionExit);
// Decrement per-session connection count
const count = sessionWsCount.get(id) ?? 1;
if (count <= 1) {
sessionWsCount.delete(id);
} else {
sessionWsCount.set(id, count - 1);
}
});
});
}
+1
View File
@@ -1953,6 +1953,7 @@ export class WebServer extends EventEmitter {
globalStats: this.store.getAggregateStats(activeSessionTokens),
subagents: subagentWatcher.getRecentSubagents(15), // 15 min to avoid stale agents
timestamp: now,
inputCjkForm: process.env.INPUT_CJK_FORM?.toUpperCase() === 'ON',
};
this.cachedLightState = { data: result, timestamp: now };
+152 -12
View File
@@ -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<WebSocket> {
function connectWs(path: string, timeoutMs = 5000): Promise<WebSocket> {
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<unknown[]> {
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);
});