mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 20:49:41 +02:00
Compare commits
9
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d072e773d8 | ||
|
|
a649c91b68 | ||
|
|
3383c23099 | ||
|
|
c3e1e731ef | ||
|
|
405b711c3a | ||
|
|
393a2d9c28 | ||
|
|
809bf6a614 | ||
|
|
da71d8d01c | ||
|
|
e5aca6aa4c |
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
@@ -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
@@ -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",
|
||||
|
||||
@@ -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
@@ -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.
|
||||
|
||||
@@ -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>
|
||||
|
||||
@@ -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; },
|
||||
};
|
||||
})();
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
@@ -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
@@ -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);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user