Compare commits

..
Author SHA1 Message Date
arkon 551461cb31 chore: version packages 2026-03-14 19:26:09 +01:00
arkonandClaude Opus 4.6 c4bae75c59 fix: add onerror handler for lazy-loaded WebGL addon script
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 19:23:15 +01:00
arkonandClaude Opus 4.6 88c415fc37 perf: V8 compile cache, lazy-load WebGL, preload hints, batch tmux reconciliation
- Enable NODE_COMPILE_CACHE in systemd service and npm start for 10-20% faster cold starts
- Lazy-load xterm-addon-webgl.min.js (244KB) only on desktop — mobile never downloads it
- Add <link rel="preload"> hints for critical scripts (xterm, constants, app) in <head>
- Replace per-session tmux subprocess calls with single batch `list-panes -a` call
  (N*2+1+M execSync calls → 1 for reconcileSessions)
- Fix CLAUDE.md frontend module count (10 → 11)

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 19:20:44 +01:00
arkonandClaude Opus 4.6 7175e4b350 docs: update CLAUDE.md and README.md to reflect current codebase
Correct stale counts and add missing entries: route modules 12→13
(ws-routes.ts), frontend modules 9→10 (input-cjk.js), handler count
~111→~114, utilities section expanded, TypeScript badge 5.5→5.9,
frontend extracted modules 8→9.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 18:48:25 +01:00
arkonandClaude Opus 4.6 08a417997f ci: upgrade actions/checkout and actions/setup-node to v6 (Node 24)
Replace v4 (Node 20) with v6 (Node 24 native) to eliminate the
deprecation warning. Remove the FORCE_JAVASCRIPT_ACTIONS_TO_NODE24
workaround since v6 doesn't need it.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 18:42:22 +01:00
arkonandClaude Opus 4.6 e6cb89b0cd ci: use Node.js 24 runtime for actions and bump release node to 22
Opt into Node.js 24 for GitHub Actions runners (actions/checkout@v4,
actions/setup-node@v4) to silence deprecation warnings. Also bump
release.yml from node 20 to 22 to match ci.yml.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 18:40:44 +01:00
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
arkon 4295faefc9 chore: version packages 2026-03-14 18:04:51 +01:00
arkon 93e1ba5110 Merge remote-tracking branch 'origin/feat/ws-terminal-io-upstream' 2026-03-14 18:03:36 +01:00
Ark0N abbbf9e90a Merge pull request #43 from Ark0N/feat/ws-tests
test: add WebSocket terminal I/O route tests
2026-03-14 18:03:04 +01:00
Ark0N 3a41de7b57 Merge pull request #42 from Ark0N/feat/ws-heartbeat
feat: add ping/pong heartbeat to WebSocket connections
2026-03-14 18:03:02 +01:00
Ark0N 8267edc6fe Merge pull request #40 from Spirotot/feat/ws-terminal-io-upstream
feat: WebSocket terminal I/O with server-side DEC 2026 sync
2026-03-14 18:02:55 +01:00
arkonandClaude Opus 4.6 78c568e5f7 test: add automated tests for WebSocket terminal I/O route
16 tests covering session-not-found close code, terminal output with
DEC 2026 sync markers, client input forwarding, resize bounds
validation, malformed message handling, and connection cleanup of
session event listeners.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 18:00:55 +01:00
arkonandClaude Opus 4.6 cc624d2575 feat: add ping/pong heartbeat to WebSocket connections
Detect stale connections that TCP keepalive won't catch for minutes,
especially through tunnels and proxies. Pings every 30s with a 10s
pong timeout — if the client doesn't respond, the socket is terminated
and all timers cleaned up.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 17:58:50 +01:00
arkonandClaude Opus 4.6 5844720525 fix: validate WS resize dimensions to match HTTP route bounds
The HTTP resize route validates via ResizeSchema (cols: 1-500, rows:
1-200, integers only). The WS handler only checked typeof === 'number',
allowing floats, negatives, and extreme values through to ptyProcess.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-14 17:57:20 +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
Aaron FieldsandClaude Opus 4.6 ceaf4624a1 feat: add WebSocket terminal I/O with server-side DEC 2026 sync
Replace per-keystroke HTTP POST + SSE terminal output with a single
bidirectional WebSocket connection for dramatically lower input latency.
The existing SSE+POST paths remain fully functional as fallback.

Server-side: ws-routes.ts provides /ws/sessions/:id/terminal with 8ms
micro-batching and 16KB flush threshold. Each batch is wrapped in
DEC 2026 synchronized update markers so xterm.js renders atomically —
Ink's DA capability negotiation fails through the PTY→server→WS proxy
chain, so without server-injected markers, cursor-up redraws flicker.

Frontend: _connectWs/_disconnectWs manage per-session WS lifecycle.
Input and resize use WS fast path with HTTP POST fallback. SSE terminal
events are suppressed when WS is active to prevent double rendering.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 21:02:28 -04:00
arkonandClaude Opus 4.6 a6597e4a9a fix: patch 5 dependency vulnerabilities (basic-ftp, fastify, minimatch, serialize-javascript)
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-13 00:26:56 +01:00
arkon f869e823af chore: version packages 2026-03-12 23:59:16 +01:00
arkonandClaude Opus 4.6 8d0b179f94 fix: repair 15 pre-existing subagent-watcher test failures
Root causes:
- Mock readline (EventEmitter) lacked .close() method, causing TypeError
  that blocked extractDescriptionFromFile's Promise from ever resolving
- Mock stream lacked .destroy() method (same issue after .close() fix)
- Entry-processing tests shared one readline mock between description
  extraction and tailing — events emitted before tailFile started were lost
- Liveness checker marked agents as 'completed' instead of 'idle' because
  fixed stat timestamps became stale after fake timer advancement

Fixes:
- Add createMockRl() helper with .close() method
- Use { destroy: vi.fn() } for stream mocks
- Use mockReturnValueOnce() for two-readline pattern in 7 entry tests
- Use mockImplementation() for dynamic stat timestamps

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 16:08:39 +01:00
arkonandClaude Opus 4.6 98fa55b7b2 chore: codebase cleanup — remove dead code, consolidate imports, extract constants
- Remove 3 unused exported constants (TRIM_MESSAGES_TO, MAX_TERMINAL_COLS, MAX_TERMINAL_ROWS)
- Consolidate 8 direct util imports into barrel imports (./utils/index.js)
- Extract magic number 8191 to FILE_PEEK_BYTES constant in buffer-limits.ts
- Add explanatory comments to 9 undocumented .catch(() => {}) handlers

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 15:50:40 +01:00
arkonandClaude Opus 4.6 c46ac30631 fix: hide subagent monitor panel by default
Change showSubagents default from true to false so the subagent
panel doesn't auto-show on page load. Users can still enable it
via Settings.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 15:35:27 +01:00
arkonandClaude Opus 4.6 dfcc14bfd2 fix: one-liner restart command that works for background processes
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 15:34:42 +01:00
arkonandClaude Opus 4.6 a068008409 fix: clarify restart instructions — stop first, then start
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 15:33:34 +01:00
arkonandClaude Opus 4.6 0aa31f100e fix: show restart command when codeman-web is not a systemd service
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 15:31:19 +01:00
arkonandClaude Opus 4.6 314a160458 feat: auto-restart codeman-web service after update if running
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 15:23:55 +01:00
arkonandClaude Opus 4.6 e7ee5595c5 feat: auto-detect existing install and run update instead of fresh install
Re-running the install script now detects ~/.codeman/app/.git and
automatically updates instead of re-installing. Removes the separate
`bash -s update` instructions from README since it's no longer needed.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 15:09:37 +01:00
33 changed files with 1635 additions and 302 deletions
+2 -2
View File
@@ -11,10 +11,10 @@ jobs:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/checkout@v6
- name: Setup Node.js
uses: actions/setup-node@v4
uses: actions/setup-node@v6
with:
node-version: 22
cache: 'npm'
+3 -3
View File
@@ -16,12 +16,12 @@ jobs:
pull-requests: write
steps:
- name: Checkout repo
uses: actions/checkout@v4
uses: actions/checkout@v6
- name: Setup Node.js
uses: actions/setup-node@v4
uses: actions/setup-node@v6
with:
node-version: 20
node-version: 22
cache: npm
registry-url: https://registry.npmjs.org
+54
View File
@@ -1,5 +1,59 @@
# aicodeman
## 0.4.1
### Patch Changes
- Performance optimizations: V8 compile cache for 10-20% faster cold starts, lazy-load WebGL addon (244KB saved on mobile), preload hints for critical scripts, batch tmux reconciliation (N subprocess calls → 1). Also: WebSocket session lifecycle fixes, CJK IME input support, CI upgrade to Node 24/actions v6, install.sh fork support, and CLAUDE.md/README documentation refresh.
## 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
- Add WebSocket terminal I/O with server-side DEC 2026 synchronized update markers. Replaces per-keystroke HTTP POST + SSE terminal output with a single bidirectional WebSocket connection for dramatically lower input latency. Server-side 8ms micro-batching with 16KB flush threshold groups rapid PTY events into single WS frames wrapped in DEC 2026 markers for flicker-free atomic rendering. Includes 30s ping/pong heartbeat with 10s timeout for stale connection detection through tunnels. Existing SSE + HTTP POST paths remain fully functional as transparent fallback. Resize messages validated to match HTTP route bounds (cols 1-500, rows 1-200, integers only). 16 automated route tests added for WS endpoint. Also patches 5 dependency vulnerabilities (basic-ftp, fastify, minimatch, serialize-javascript).
## 0.3.11
### Patch Changes
- ### Session Resume & History
- Add `resumeSessionId` support for conversation resume after reboot
- Add history session resume UI and API with route shell sessions routing fix
- Improve session resume reliability and persist user settings across refresh
- Correct `claudeSessionId` for resumed sessions
### Terminal & Frontend
- Upgrade xterm.js 5.3 → 6.0 with native DEC 2026 synchronized output
- Increase terminal scrollback from 5,000 to 20,000 lines
- Reduce default font size and persist tab state across refresh
- Resolve terminal resize scrollback ghost renders
- Hide subagent monitor panel by default
### Installer
- Auto-detect existing install and run update instead of fresh install
- Auto-restart codeman-web service after update if running
- Show restart command when codeman-web is not a systemd service
- Fix one-liner restart command for background processes
### Codebase Quality
- Remove dead code, consolidate imports, extract constants
- Repair 15 pre-existing subagent-watcher test failures
- Clean up DEC sync dead code
## 0.3.10
### Patch Changes
+7 -7
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.10 (must match `package.json`)
**Version**: 0.4.1 (must match `package.json`)
## Project Overview
@@ -108,9 +108,9 @@ Codeman is a Claude Code session manager with web interface and autonomous Ralph
| **State** | `src/state-store.ts`, `src/run-summary.ts`, `src/session-lifecycle-log.ts` | |
| **Infra** | `src/hooks-config.ts`, `src/push-store.ts`, `src/tunnel-manager.ts`, `src/image-watcher.ts`, `src/file-stream-manager.ts` | |
| **Plan** | `src/plan-orchestrator.ts`, `src/prompts/*.ts`, `src/templates/claude-md.ts` | |
| **Web** | `src/web/server.ts`, `src/web/sse-events.ts`, `src/web/routes/*.ts` (12 route modules + barrel), `src/web/ports/*.ts`, `src/web/middleware/auth.ts`, `src/web/schemas.ts` | |
| **Frontend** | `src/web/public/app.js` ★ (~12.1K lines) + 9 JS modules (incl. `sw.js` service worker) | |
| **Types** | `src/types/index.ts` → 14 domain files | See `@fileoverview` in index.ts |
| **Web** | `src/web/server.ts`, `src/web/sse-events.ts`, `src/web/routes/*.ts` (13 route modules incl. `ws-routes.ts` + barrel), `src/web/ports/*.ts`, `src/web/middleware/auth.ts`, `src/web/schemas.ts` | |
| **Frontend** | `src/web/public/app.js` ★ (~12.5K lines) + 11 JS modules (incl. `sw.js`, `input-cjk.js`) | |
| **Types** | `src/types/index.ts` → 13 domain files | See `@fileoverview` in index.ts |
★ = Large file (>50KB). All files have `@fileoverview` JSDoc — read that before diving in.
@@ -118,7 +118,7 @@ Codeman is a Claude Code session manager with web interface and autonomous Ralph
**Config**: `src/config/` — 9 files. Import from specific files, not barrel.
**Utilities**: `src/utils/` — re-exported via index. Key: `CleanupManager`, `LRUMap`, `StaleExpirationMap`, `BufferAccumulator`, `stripAnsi`, `Debouncer`.
**Utilities**: `src/utils/` — re-exported via index. Key: `CleanupManager`, `LRUMap`, `StaleExpirationMap`, `BufferAccumulator`, `stripAnsi`, `Debouncer`, `KeyedDebouncer`. Also: `claude-cli-resolver`/`opencode-cli-resolver` (CLI path resolution), `string-similarity` (fuzzy matching), `regex-patterns` (ANSI/token/spinner patterns), `assertNever` (exhaustive checks).
### Data Flow
@@ -143,7 +143,7 @@ Codeman is a Claude Code session manager with web interface and autonomous Ralph
### Frontend
Frontend JS modules have `@fileoverview` with `@dependency`/`@loadorder` tags. Load order: `constants.js`(1) → `mobile-handlers.js`(2) → `voice-input.js`(3) → `notification-manager.js`(4) → `keyboard-accessory.js`(5) → `app.js`(6) → `ralph-wizard.js`(7) → `api-client.js`(8) → `subagent-windows.js`(9).
Frontend JS modules have `@fileoverview` with `@dependency`/`@loadorder` tags. Load order: `constants.js`(1) → `mobile-handlers.js`(2) → `voice-input.js`(3) → `notification-manager.js`(4) → `keyboard-accessory.js`(5) → `input-cjk.js`(5.5) → `app.js`(6) → `ralph-wizard.js`(7) → `api-client.js`(8) → `subagent-windows.js`(9). `input-cjk.js` handles CJK IME composition via an always-visible textarea below the terminal (`window.cjkActive` blocks xterm's onData).
**Z-index layers**: subagent windows (1000), plan agents (1100), log viewers (2000), image popups (3000), local echo overlay (7).
@@ -170,7 +170,7 @@ Frontend JS modules have `@fileoverview` with `@dependency`/`@loadorder` tags. L
### API Routes
~111 handlers across 12 route files in `src/web/routes/`: system (35), sessions (24), ralph (9), plan (8), respawn (7), cases (7), files (5), mux (5), scheduled (4), push (4), teams (2), hooks (1). Each file has `@fileoverview` with endpoint details.
~114 handlers across 13 route files in `src/web/routes/`: system (36), sessions (25), ralph (9), plan (8), respawn (7), cases (7), files (5), mux (5), scheduled (4), push (4), teams (2), hooks (1), ws (1 WebSocket). Each file has `@fileoverview` with endpoint details.
## Adding Features
+24 -9
View File
@@ -11,7 +11,7 @@
<p align="center">
<a href="https://opensource.org/licenses/MIT"><img src="https://img.shields.io/badge/License-MIT-1e3a5f?style=flat-square" alt="License: MIT"></a>
<a href="https://nodejs.org/"><img src="https://img.shields.io/badge/Node.js-18%2B-22c55e?style=flat-square&logo=node.js&logoColor=white" alt="Node.js 18+"></a>
<a href="https://www.typescriptlang.org/"><img src="https://img.shields.io/badge/TypeScript-5.5-3b82f6?style=flat-square&logo=typescript&logoColor=white" alt="TypeScript 5.5"></a>
<a href="https://www.typescriptlang.org/"><img src="https://img.shields.io/badge/TypeScript-5.9-3b82f6?style=flat-square&logo=typescript&logoColor=white" alt="TypeScript 5.9"></a>
<a href="https://fastify.dev/"><img src="https://img.shields.io/badge/Fastify-5.x-1e3a5f?style=flat-square&logo=fastify&logoColor=white" alt="Fastify"></a>
<img src="https://img.shields.io/badge/Tests-1435%20total-22c55e?style=flat-square" alt="Tests">
</p>
@@ -28,18 +28,33 @@
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
# Open http://localhost:3000 — press Ctrl+Enter to start your first session
```
**Update to latest version:**
```bash
curl -fsSL https://raw.githubusercontent.com/Ark0N/Codeman/master/install.sh | bash -s update
```
<details>
<summary><strong>Run as a background service</strong></summary>
@@ -484,9 +499,9 @@ The codebase went through a comprehensive 7-phase refactoring that eliminated go
| Phase | What changed | Impact |
|-------|-------------|--------|
| **Performance** | Cached endpoints, SSE adaptive batching, buffer chunking | Sub-16ms terminal latency |
| **Route extraction** | `server.ts` split into 12 domain route modules + auth middleware + port interfaces | **−60%** server.ts LOC (6,736 → 2,697) |
| **Route extraction** | `server.ts` split into 13 domain route modules + auth middleware + port interfaces | **−60%** server.ts LOC (6,736 → 2,697) |
| **Domain splitting** | `types.ts` → 14 domain files, `ralph-tracker` → 7 files, `respawn-controller` → 5 files, `session` → 6 files | No more god files |
| **Frontend modules** | `app.js` → 8 extracted modules (constants, mobile, voice, notifications, keyboard, API, subagent windows) | **−24%** app.js LOC (15.2K → 11.5K) |
| **Frontend modules** | `app.js` → 9 extracted modules (constants, mobile, voice, notifications, keyboard, CJK input, API, Ralph wizard, subagent windows) | **−24%** app.js LOC (15.2K → 11.5K) |
| **Config consolidation** | ~70 scattered magic numbers → 9 domain-focused config files | Zero cross-file duplicates |
| **Test infrastructure** | Shared mock library, 12 route test files, consolidated MockSession | Testable route handlers via `app.inject()` |
Binary file not shown.
+74
View File
@@ -0,0 +1,74 @@
# Codeman Performance Optimization Plan
## Current State
The backend is **already production-grade** — SSE broadcasting, state persistence, terminal batching, buffer management, and memory patterns are all well-optimized. The biggest gains are on the **frontend delivery** side.
## Implemented Optimizations
### 1. V8 Compile Cache (10-20% faster cold start)
**Files:** `scripts/codeman-web.service`, `package.json`
Node.js re-parses and compiles all JS on every cold start. `NODE_COMPILE_CACHE` caches V8 compiled bytecode to disk, reusing it on subsequent starts.
- Added `Environment=NODE_COMPILE_CACHE=/home/arkon/.codeman/compile-cache` to systemd service
- Added to `npm start` script for non-systemd usage
- Zero code changes, immediate win on every restart
### 2. WebGL Addon Lazy-Loading (244KB saved on mobile, non-blocking on desktop)
**Files:** `src/web/public/index.html`, `src/web/public/app.js`
`xterm-addon-webgl.min.js` (244KB) was loaded eagerly for all users via `<script defer>`, but only used on desktop with WebGL2 support.
- Removed `<script defer>` from `index.html`
- Added dynamic script loading in `app.js` — only downloads on desktop when WebGL is needed
- Mobile users never download the file at all (244KB saved)
- Desktop: loads in parallel with page rendering, addon initializes when ready
- Graceful fallback: canvas renderer used if WebGL unavailable or script fails
### 3. Preload Hints (~50-100ms faster perceived load)
**Files:** `src/web/public/index.html`
Browser discovers `<script defer>` tags only when the parser reaches them at the bottom of `<body>`. By then, the HTML parse has blocked for hundreds of lines.
- Added `<link rel="preload" as="script">` in `<head>` for `vendor/xterm.min.js`, `constants.js`, `app.js`
- Browser starts fetching critical scripts immediately during HTML parse (before reaching `<body>`)
- Zero runtime overhead — just hints for the browser's preload scanner
### 4. Batch Tmux Reconciliation (N subprocess calls → 1)
**Files:** `src/tmux-manager.ts`
`reconcileSessions()` previously called `tmux has-session` + `tmux display-message` per known session, plus `tmux list-sessions` for discovery, plus `tmux display-message` per discovered session. With 20 sessions: 41+ subprocess calls.
- Replaced with single `tmux list-panes -a -F '#{session_name}\t#{pane_pid}'` call
- Builds a Map from the result, then does O(1) lookups for both known and discovered sessions
- Also replaced inner O(n) `isKnown` scan with a Set lookup
- 20 sessions: 41 subprocess calls → 1, with faster lookups
### 5. Asset Hashing / Cache Busting (already implemented)
**Files:** `scripts/build.mjs` (pre-existing)
Content-hash cache busting was already implemented in the build script:
- All app JS/CSS files get content hashes (`app.abc123.js`)
- `index.html` rewritten to reference hashed filenames
- Pre-compressed with gzip + Brotli
- 1-year immutable cache works correctly — new deploys get new filenames
## Already Optimized (No Action Needed)
| Area | Why It's Fine |
|------|---------------|
| **SSE Broadcasting** | Single serialization per broadcast, preformatted frames, backpressure handling, session subscription filtering |
| **State Persistence** | 500ms debounce, incremental per-session JSON caching, async atomic writes, circuit breaker on failures |
| **Terminal Batching** | Adaptive intervals (16-50ms), per-session queues, immediate flush at 32KB, array-based accumulation |
| **Buffer Management** | BufferAccumulator (array-push, lazy join), auto-trim at 2MB/1MB, no string concatenation in hot paths |
| **ANSI Stripping** | Pre-compiled regex via factory functions, single-pass processing |
| **Static File Serving** | @fastify/static with 1-year cache, pre-compressed Brotli/gzip, no-cache for HTML |
| **Memory Management** | CleanupManager, LRUMap, StaleExpirationMap, bounded buffers, explicit listener cleanup |
| **Import Patterns** | Pure ESM, lazy web server import, no circular deps, no dynamic imports in hot paths |
| **Config Loading** | Small constant files, no I/O at import time, specific imports (no barrel) |
+28 -7
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,13 +1281,23 @@ 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)")"
echo ""
echo -e " ${DIM}Restart codeman web to use the new version.${NC}"
# Auto-restart systemd service if it's running, otherwise tell the user
if systemctl --user is-active codeman-web.service &>/dev/null; then
info "Restarting codeman-web service..."
systemctl --user restart codeman-web.service
success "codeman-web service restarted"
else
echo -e " ${DIM}Restart codeman web to use the new version:${NC}"
echo -e " ${CYAN}pkill -f 'codeman.*web'; codeman web &${NC}"
fi
echo ""
}
@@ -1355,5 +1369,12 @@ uninstall() {
case "${1:-}" in
update) update ;;
uninstall) uninstall ;;
*) main "$@" ;;
*)
if [[ -z "${1:-}" && -d "$INSTALL_DIR/.git" ]]; then
print_banner
update
else
main "$@"
fi
;;
esac
+103 -25
View File
@@ -1,12 +1,12 @@
{
"name": "aicodeman",
"version": "0.3.8",
"version": "0.3.11",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "aicodeman",
"version": "0.3.8",
"version": "0.3.11",
"hasInstallScript": true,
"license": "MIT",
"workspaces": [
@@ -17,6 +17,7 @@
"@fastify/compress": "^8.3.1",
"@fastify/cookie": "^11.0.2",
"@fastify/static": "^8.0.0",
"@fastify/websocket": "^11.2.0",
"@xterm/addon-fit": "^0.11.0",
"@xterm/addon-unicode11": "^0.9.0",
"@xterm/addon-webgl": "^0.19.0",
@@ -45,6 +46,7 @@
"@types/react": "^19.2.14",
"@types/uuid": "^10.0.0",
"@types/web-push": "^3.6.4",
"@types/ws": "^8.18.1",
"@vitest/coverage-v8": "^4.0.18",
"agent-browser": "^0.6.0",
"esbuild": "^0.27.3",
@@ -945,6 +947,53 @@
"glob": "^11.0.0"
}
},
"node_modules/@fastify/websocket": {
"version": "11.2.0",
"resolved": "https://registry.npmjs.org/@fastify/websocket/-/websocket-11.2.0.tgz",
"integrity": "sha512-3HrDPbAG1CzUCqnslgJxppvzaAZffieOVbLp1DAy1huCSynUWPifSvfdEDUR8HlJLp3sp1A36uOM2tJogADS8w==",
"funding": [
{
"type": "github",
"url": "https://github.com/sponsors/fastify"
},
{
"type": "opencollective",
"url": "https://opencollective.com/fastify"
}
],
"license": "MIT",
"dependencies": {
"duplexify": "^4.1.3",
"fastify-plugin": "^5.0.0",
"ws": "^8.16.0"
}
},
"node_modules/@fastify/websocket/node_modules/duplexify": {
"version": "4.1.3",
"resolved": "https://registry.npmjs.org/duplexify/-/duplexify-4.1.3.tgz",
"integrity": "sha512-M3BmBhwJRZsSx38lZyhE53Csddgzl5R7xGJNk7CVddZD6CcmwMCH8J+7AprIrQKH7TonKxaCjcv27Qmf+sQ+oA==",
"license": "MIT",
"dependencies": {
"end-of-stream": "^1.4.1",
"inherits": "^2.0.3",
"readable-stream": "^3.1.1",
"stream-shift": "^1.0.2"
}
},
"node_modules/@fastify/websocket/node_modules/readable-stream": {
"version": "3.6.2",
"resolved": "https://registry.npmjs.org/readable-stream/-/readable-stream-3.6.2.tgz",
"integrity": "sha512-9u/sniCrY3D5WdsERHzHE4G2YCXqoG5FTHUiCC4SIbr6XcLZBY05ya9EKjYek9O5xOAwjGq+1JdGBAS7Q9ScoA==",
"license": "MIT",
"dependencies": {
"inherits": "^2.0.3",
"string_decoder": "^1.1.1",
"util-deprecate": "^1.0.1"
},
"engines": {
"node": ">= 6"
}
},
"node_modules/@humanfs/core": {
"version": "0.19.1",
"dev": true,
@@ -2118,6 +2167,16 @@
"@types/node": "*"
}
},
"node_modules/@types/ws": {
"version": "8.18.1",
"resolved": "https://registry.npmjs.org/@types/ws/-/ws-8.18.1.tgz",
"integrity": "sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==",
"dev": true,
"license": "MIT",
"dependencies": {
"@types/node": "*"
}
},
"node_modules/@types/yauzl": {
"version": "2.10.3",
"dev": true,
@@ -3019,7 +3078,9 @@
}
},
"node_modules/basic-ftp": {
"version": "5.1.0",
"version": "5.2.0",
"resolved": "https://registry.npmjs.org/basic-ftp/-/basic-ftp-5.2.0.tgz",
"integrity": "sha512-VoMINM2rqJwJgfdHq6RiUudKt2BV+FY5ZFezP/ypmwayk68+NzzAQy4XXLlqsGD4MCzq3DrmNFD/uUmBJuGoXw==",
"dev": true,
"license": "MIT",
"engines": {
@@ -4416,7 +4477,9 @@
"license": "BSD-3-Clause"
},
"node_modules/fastify": {
"version": "5.7.4",
"version": "5.8.2",
"resolved": "https://registry.npmjs.org/fastify/-/fastify-5.8.2.tgz",
"integrity": "sha512-lZmt3navvZG915IE+f7/TIVamxIwmBd+OMB+O9WBzcpIwOo6F0LTh0sluoMFk5VkrKTvvrwIaoJPkir4Z+jtAg==",
"funding": [
{
"type": "github",
@@ -4438,7 +4501,7 @@
"fast-json-stringify": "^6.0.0",
"find-my-way": "^9.0.0",
"light-my-request": "^6.0.0",
"pino": "^10.1.0",
"pino": "^9.14.0 || ^10.1.0",
"process-warning": "^5.0.0",
"rfdc": "^1.3.1",
"secure-json-parse": "^4.0.0",
@@ -4601,6 +4664,20 @@
"dev": true,
"license": "Unlicense"
},
"node_modules/fsevents": {
"version": "2.3.3",
"resolved": "https://registry.npmjs.org/fsevents/-/fsevents-2.3.3.tgz",
"integrity": "sha512-5xoDfX+fL7faATnagmWPpbFtwh/R77WmMMqqHGS65C3vvB0YHrgF+B1YmZ3441tMj5n63k0212XNoJwzlhffQw==",
"hasInstallScript": true,
"license": "MIT",
"optional": true,
"os": [
"darwin"
],
"engines": {
"node": "^8.16.0 || ^10.6.0 || >=11.0.0"
}
},
"node_modules/function-bind": {
"version": "1.1.2",
"dev": true,
@@ -5641,7 +5718,9 @@
"license": "ISC"
},
"node_modules/minimatch": {
"version": "10.2.2",
"version": "10.2.4",
"resolved": "https://registry.npmjs.org/minimatch/-/minimatch-10.2.4.tgz",
"integrity": "sha512-oRjTw/97aTBN0RHbYCdtF1MQfvusSIBQM0IZEgzl6426+8jSC0nF1a/GmnVLpfB9yyr6g6FTqWqiZVbxrtaCIg==",
"license": "BlueOak-1.0.0",
"dependencies": {
"brace-expansion": "^5.0.2"
@@ -6163,6 +6242,21 @@
"node": ">=18"
}
},
"node_modules/playwright/node_modules/fsevents": {
"version": "2.3.2",
"resolved": "https://registry.npmjs.org/fsevents/-/fsevents-2.3.2.tgz",
"integrity": "sha512-xiqMQR4xAeHTuB9uWm+fFRcIOgKBMiOBP+eXiyT7jsgVCq1bkVygt00oASowB7EdtpOHaaPgKt812P9ab+DDKA==",
"dev": true,
"hasInstallScript": true,
"license": "MIT",
"optional": true,
"os": [
"darwin"
],
"engines": {
"node": "^8.16.0 || ^10.6.0 || >=11.0.0"
}
},
"node_modules/pngjs": {
"version": "7.0.0",
"dev": true,
@@ -6625,14 +6719,6 @@
"version": "4.0.4",
"license": "MIT"
},
"node_modules/randombytes": {
"version": "2.1.0",
"dev": true,
"license": "MIT",
"dependencies": {
"safe-buffer": "^5.1.0"
}
},
"node_modules/react": {
"version": "19.2.4",
"dev": true,
@@ -7023,14 +7109,6 @@
"node": ">=10"
}
},
"node_modules/serialize-javascript": {
"version": "6.0.2",
"dev": true,
"license": "BSD-3-Clause",
"dependencies": {
"randombytes": "^2.1.0"
}
},
"node_modules/set-blocking": {
"version": "2.0.0",
"license": "ISC"
@@ -7401,14 +7479,15 @@
}
},
"node_modules/terser-webpack-plugin": {
"version": "5.3.16",
"version": "5.4.0",
"resolved": "https://registry.npmjs.org/terser-webpack-plugin/-/terser-webpack-plugin-5.4.0.tgz",
"integrity": "sha512-Bn5vxm48flOIfkdl5CaD2+1CiUVbonWQ3KQPyP7/EuIl9Gbzq/gQFOzaMFUEgVjB1396tcK0SG8XcNJ/2kDH8g==",
"dev": true,
"license": "MIT",
"dependencies": {
"@jridgewell/trace-mapping": "^0.3.25",
"jest-worker": "^27.4.5",
"schema-utils": "^4.3.0",
"serialize-javascript": "^6.0.2",
"terser": "^5.31.1"
},
"engines": {
@@ -8488,7 +8567,6 @@
},
"node_modules/ws": {
"version": "8.19.0",
"dev": true,
"license": "MIT",
"engines": {
"node": ">=10.0.0"
+4 -2
View File
@@ -1,6 +1,6 @@
{
"name": "aicodeman",
"version": "0.3.10",
"version": "0.4.1",
"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",
@@ -11,7 +11,7 @@
"scripts": {
"postinstall": "node scripts/postinstall.js",
"build": "node scripts/build.mjs",
"start": "node dist/index.js",
"start": "NODE_COMPILE_CACHE=${HOME}/.codeman/compile-cache node dist/index.js",
"dev": "tsx src/index.ts web",
"web": "node dist/index.js web",
"clean": "rm -rf dist",
@@ -51,6 +51,7 @@
"@fastify/compress": "^8.3.1",
"@fastify/cookie": "^11.0.2",
"@fastify/static": "^8.0.0",
"@fastify/websocket": "^11.2.0",
"@xterm/addon-fit": "^0.11.0",
"@xterm/addon-unicode11": "^0.9.0",
"@xterm/addon-webgl": "^0.19.0",
@@ -76,6 +77,7 @@
"@types/react": "^19.2.14",
"@types/uuid": "^10.0.0",
"@types/web-push": "^3.6.4",
"@types/ws": "^8.18.1",
"@vitest/coverage-v8": "^4.0.18",
"agent-browser": "^0.6.0",
"esbuild": "^0.27.3",
+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',
+1
View File
@@ -11,6 +11,7 @@ RestartSec=5
KillMode=process
Environment=NODE_ENV=production
Environment=HOME=/home/arkon
Environment=NODE_COMPILE_CACHE=/home/arkon/.codeman/compile-cache
# Logging
StandardOutput=journal
+1 -2
View File
@@ -28,8 +28,7 @@ import { existsSync, readFileSync, unlinkSync, writeFileSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { EventEmitter } from 'node:events';
import { getAugmentedPath } from './utils/claude-cli-resolver.js';
import { ANSI_ESCAPE_PATTERN_SIMPLE } from './utils/index.js';
import { getAugmentedPath, ANSI_ESCAPE_PATTERN_SIMPLE } from './utils/index.js';
import { AI_CHECK_MAX_BACKOFF_MS } from './config/ai-defaults.js';
// ========== Security Validation ==========
+11 -5
View File
@@ -56,11 +56,6 @@ export const TRIM_TEXT_TO = 768 * 1024; // 768KB
*/
export const MAX_MESSAGES = 1000;
/**
* Number of messages to keep when trimming (80% of max).
*/
export const TRIM_MESSAGES_TO = 800;
// ============================================================================
// Line Buffer Limits
// ============================================================================
@@ -85,3 +80,14 @@ export const MAX_RESPAWN_BUFFER_SIZE = 1 * 1024 * 1024; // 1MB
* Size to trim respawn buffer to when max is exceeded.
*/
export const TRIM_RESPAWN_BUFFER_TO = 512 * 1024; // 512KB
// ============================================================================
// File Peek Limits
// ============================================================================
/**
* Maximum bytes to read when peeking at the beginning of a file.
* Used with `createReadStream({ end })` (inclusive) to read the first 8KB,
* which is enough to extract metadata from the first few JSONL lines.
*/
export const FILE_PEEK_BYTES = 8 * 1024 - 1; // 8KB (inclusive end offset)
-6
View File
@@ -11,11 +11,5 @@
/** Max input length per API request (bytes) */
export const MAX_INPUT_LENGTH = 64 * 1024;
/** Max terminal columns for resize requests */
export const MAX_TERMINAL_COLS = 500;
/** Max terminal rows for resize requests */
export const MAX_TERMINAL_ROWS = 200;
/** Max session name length (chars) */
export const MAX_SESSION_NAME_LENGTH = 128;
+2 -2
View File
@@ -521,7 +521,7 @@ export class PlanOrchestrator {
} finally {
// Always clean up session and progress interval — centralizing here
// prevents the race where cancel() and catch both try to manage the set
await session.stop().catch(() => {});
await session.stop().catch(() => {}); // Ignore - session cleanup is best-effort in finally block
this.runningSessions.delete(session);
clearInterval(progressInterval);
}
@@ -651,7 +651,7 @@ export class PlanOrchestrator {
} finally {
// Always clean up session and progress interval — centralizing here
// prevents the race where cancel() and catch both try to manage the set
await session.stop().catch(() => {});
await session.stop().catch(() => {}); // Ignore - session cleanup is best-effort in finally block
this.runningSessions.delete(session);
clearInterval(progressInterval);
}
+1 -2
View File
@@ -49,8 +49,7 @@ import { Session } from './session.js';
import { AiIdleChecker, type AiCheckResult, type AiCheckState } from './ai-idle-checker.js';
import { AiPlanChecker, type AiPlanCheckResult } from './ai-plan-checker.js';
import type { TeamWatcher } from './team-watcher.js';
import { BufferAccumulator } from './utils/buffer-accumulator.js';
import { ANSI_ESCAPE_PATTERN_SIMPLE, assertNever, CleanupManager } from './utils/index.js';
import { BufferAccumulator, ANSI_ESCAPE_PATTERN_SIMPLE, assertNever, CleanupManager } from './utils/index.js';
import { MAX_RESPAWN_BUFFER_SIZE, TRIM_RESPAWN_BUFFER_TO as RESPAWN_BUFFER_TRIM_SIZE } from './config/buffer-limits.js';
import {
isCompletionMessage,
+1 -1
View File
@@ -9,7 +9,7 @@
*/
import type { ClaudeMode } from './types.js';
import { getAugmentedPath } from './utils/claude-cli-resolver.js';
import { getAugmentedPath } from './utils/index.js';
/**
* Build Claude CLI permission flags based on the configured mode.
+1 -1
View File
@@ -48,8 +48,8 @@ import type { TerminalMultiplexer, MuxSession } from './mux-interface.js';
import { TaskTracker, type BackgroundTask } from './task-tracker.js';
import { RalphTracker } from './ralph-tracker.js';
import { BashToolParser } from './bash-tool-parser.js';
import { BufferAccumulator } from './utils/buffer-accumulator.js';
import {
BufferAccumulator,
ANSI_ESCAPE_PATTERN_FULL,
TOKEN_PATTERN,
SPINNER_PATTERN,
+5 -4
View File
@@ -17,7 +17,7 @@
* Tracks per-agent: status, token counts, model, description, tool call count, liveness (PID).
*
* @dependencies config/map-limits (MAX_TRACKED_AGENTS, PENDING_TOOL_CALL_TTL_MS),
* utils (CleanupManager, KeyedDebouncer)
* config/buffer-limits (FILE_PEEK_BYTES), utils (CleanupManager, KeyedDebouncer)
* @consumedby web/server (SSE broadcast), session (subagent-session correlation)
* @emits subagent:discovered, subagent:updated, subagent:tool_call, subagent:tool_result,
* subagent:progress, subagent:message, subagent:completed
@@ -35,6 +35,7 @@ import { execFile } from 'node:child_process';
import { readFile, readdir, stat as statAsync } from 'node:fs/promises';
import { PENDING_TOOL_CALL_TTL_MS, MAX_PENDING_TOOL_CALLS, MAX_TRACKED_AGENTS } from './config/map-limits.js';
import { STALE_DATA_MAX_AGE_MS } from './config/server-timing.js';
import { FILE_PEEK_BYTES } from './config/buffer-limits.js';
import { CleanupManager, KeyedDebouncer } from './utils/index.js';
// ========== Types ==========
@@ -1009,7 +1010,7 @@ export class SubagentWatcher extends EventEmitter {
private async extractDescriptionFromFile(filePath: string): Promise<string | undefined> {
try {
// Only read the first 8KB — more than enough for 5 JSONL lines
const stream = createReadStream(filePath, { end: 8191 });
const stream = createReadStream(filePath, { end: FILE_PEEK_BYTES });
const rl = createInterface({ input: stream });
return await new Promise<string | undefined>((resolve) => {
@@ -1141,10 +1142,10 @@ export class SubagentWatcher extends EventEmitter {
if (this.fileAgentContext.has(filePath)) {
// Known file — handle content change
this.handleFileChange(filePath).catch(() => {});
this.handleFileChange(filePath).catch(() => {}); // Ignore - errors logged internally, don't crash watcher callback
} else {
// New file — register it
this.registerAgentFile(filePath, projectHash, sessionId).catch(() => {});
this.registerAgentFile(filePath, projectHash, sessionId).catch(() => {}); // Ignore - errors logged internally, don't crash watcher callback
}
});
});
+5 -5
View File
@@ -61,7 +61,7 @@ export class TeamWatcher extends EventEmitter {
persistent: false,
});
const teamsHandler = () => this.pollAsync().catch(() => {});
const teamsHandler = () => this.pollAsync().catch(() => {}); // Ignore - poll errors are non-fatal, next poll will retry
this.teamsWatcher.on('add', teamsHandler);
this.teamsWatcher.on('change', teamsHandler);
this.teamsWatcher.on('unlink', teamsHandler);
@@ -82,8 +82,8 @@ export class TeamWatcher extends EventEmitter {
persistent: false,
});
this.tasksWatcher.on('add', () => this.pollTasks().catch(() => {}));
this.tasksWatcher.on('change', () => this.pollTasks().catch(() => {}));
this.tasksWatcher.on('add', () => this.pollTasks().catch(() => {})); // Ignore - poll errors are non-fatal, next poll will retry
this.tasksWatcher.on('change', () => this.pollTasks().catch(() => {})); // Ignore - poll errors are non-fatal, next poll will retry
this.tasksWatcher.on('error', (err) => {
console.warn('[TeamWatcher] chokidar tasks watcher error:', err);
});
@@ -95,11 +95,11 @@ export class TeamWatcher extends EventEmitter {
stop(): void {
// Close chokidar watchers
if (this.teamsWatcher) {
this.teamsWatcher.close().catch(() => {});
this.teamsWatcher.close().catch(() => {}); // Ignore - watcher cleanup is best-effort during shutdown
this.teamsWatcher = null;
}
if (this.tasksWatcher) {
this.tasksWatcher.close().catch(() => {});
this.tasksWatcher.close().catch(() => {}); // Ignore - watcher cleanup is best-effort during shutdown
this.tasksWatcher = null;
}
if (this.pollTimer) {
+50 -55
View File
@@ -40,8 +40,7 @@ import {
type SessionMode,
type OpenCodeConfig,
} from './types.js';
import { wrapWithNice } from './utils/nice-wrapper.js';
import { SAFE_PATH_PATTERN } from './utils/regex-patterns.js';
import { wrapWithNice, SAFE_PATH_PATTERN, findClaudeDir, resolveOpenCodeDir } from './utils/index.js';
import type {
TerminalMultiplexer,
MuxSession,
@@ -50,11 +49,6 @@ import type {
RespawnPaneOptions,
} from './mux-interface.js';
// Claude CLI PATH resolution — shared utility
import { findClaudeDir } from './utils/claude-cli-resolver.js';
// OpenCode CLI PATH resolution
import { resolveOpenCodeDir } from './utils/opencode-cli-resolver.js';
// ============================================================================
// Timing Constants
// ============================================================================
@@ -911,13 +905,34 @@ export class TmuxManager extends EventEmitter implements TerminalMultiplexer {
const dead: string[] = [];
const discovered: string[] = [];
// Check known sessions
// Batch: single tmux call to get all session names + pane PIDs (replaces N per-session subprocess calls)
const activeSessions = new Map<string, number>();
try {
const output = execSync("tmux list-panes -a -F '#{session_name}\t#{pane_pid}' 2>/dev/null || true", {
encoding: 'utf-8',
timeout: EXEC_TIMEOUT_MS,
}).trim();
for (const line of output.split('\n')) {
if (!line) continue;
const sep = line.indexOf('\t');
if (sep === -1) continue;
const name = line.slice(0, sep);
const pid = parseInt(line.slice(sep + 1), 10);
if (name && !Number.isNaN(pid)) {
activeSessions.set(name, pid);
}
}
} catch (err) {
console.error('[TmuxManager] Failed to list tmux panes:', err);
}
// Check known sessions against the batch result (O(1) map lookup instead of subprocess per session)
for (const [sessionId, session] of this.sessions) {
if (this.sessionExists(session.muxName)) {
const pid = activeSessions.get(session.muxName);
if (pid !== undefined) {
alive.push(sessionId);
// Update PID if it changed
const pid = this.getPanePid(session.muxName);
if (pid && pid !== session.pid) {
if (pid !== session.pid) {
session.pid = pid;
}
} else {
@@ -927,51 +942,31 @@ export class TmuxManager extends EventEmitter implements TerminalMultiplexer {
}
}
// Discover unknown codeman sessions
try {
const output = execSync("tmux list-sessions -F '#{session_name}' 2>/dev/null || true", {
encoding: 'utf-8',
timeout: EXEC_TIMEOUT_MS,
}).trim();
// Discover unknown codeman/claudeman sessions from the same batch result
const knownMuxNames = new Set<string>();
for (const session of this.sessions.values()) {
knownMuxNames.add(session.muxName);
}
for (const line of output.split('\n')) {
const sessionName = line.trim();
if (!sessionName || (!sessionName.startsWith('codeman-') && !sessionName.startsWith('claudeman-'))) continue;
for (const [sessionName, pid] of activeSessions) {
if (!sessionName.startsWith('codeman-') && !sessionName.startsWith('claudeman-')) continue;
if (knownMuxNames.has(sessionName)) continue;
// Check if this session is already known
let isKnown = false;
for (const session of this.sessions.values()) {
if (session.muxName === sessionName) {
isKnown = true;
break;
}
}
if (!isKnown) {
// Extract session ID fragment from name
const fragment = sessionName.replace(/^(?:codeman|claudeman)-/, '');
const sessionId = `restored-${fragment}`;
const pid = this.getPanePid(sessionName);
if (pid) {
const session: MuxSession = {
sessionId,
muxName: sessionName,
pid,
createdAt: Date.now(),
workingDir: process.cwd(),
mode: 'claude',
attached: false,
name: `Restored: ${sessionName}`,
};
this.sessions.set(sessionId, session);
discovered.push(sessionId);
console.log(`[TmuxManager] Discovered unknown tmux session: ${sessionName} (PID ${pid})`);
}
}
}
} catch (err) {
console.error('[TmuxManager] Failed to discover sessions:', err);
const fragment = sessionName.replace(/^(?:codeman|claudeman)-/, '');
const sessionId = `restored-${fragment}`;
const session: MuxSession = {
sessionId,
muxName: sessionName,
pid,
createdAt: Date.now(),
workingDir: process.cwd(),
mode: 'claude',
attached: false,
name: `Restored: ${sessionName}`,
};
this.sessions.set(sessionId, session);
discovered.push(sessionId);
console.log(`[TmuxManager] Discovered unknown tmux session: ${sessionName} (PID ${pid})`);
}
if (dead.length > 0 || discovered.length > 0) {
+202 -22
View File
@@ -172,9 +172,9 @@ const _SSE_HANDLER_MAP = [
[SSE_EVENTS.SESSION_CREATED, '_onSessionCreated'],
[SSE_EVENTS.SESSION_UPDATED, '_onSessionUpdated'],
[SSE_EVENTS.SESSION_DELETED, '_onSessionDeleted'],
[SSE_EVENTS.SESSION_TERMINAL, '_onSessionTerminal'],
[SSE_EVENTS.SESSION_NEEDS_REFRESH, '_onSessionNeedsRefresh'],
[SSE_EVENTS.SESSION_CLEAR_TERMINAL, '_onSessionClearTerminal'],
[SSE_EVENTS.SESSION_TERMINAL, '_onSSETerminal'],
[SSE_EVENTS.SESSION_NEEDS_REFRESH, '_onSSENeedsRefresh'],
[SSE_EVENTS.SESSION_CLEAR_TERMINAL, '_onSSEClearTerminal'],
[SSE_EVENTS.SESSION_COMPLETION, '_onSessionCompletion'],
[SSE_EVENTS.SESSION_ERROR, '_onSessionError'],
[SSE_EVENTS.SESSION_EXIT, '_onSessionExit'],
@@ -361,6 +361,11 @@ class CodemanApp {
// Tracks pending hook events that need resolution (permission_prompt, elicitation_dialog, idle_prompt)
this.pendingHooks = new Map();
// WebSocket terminal I/O (low-latency bypass of HTTP POST + SSE)
this._ws = null; // WebSocket instance for active session
this._wsSessionId = null; // Session ID the WS is connected to
this._wsReady = false; // True when WS is open and ready for I/O
// Terminal write batching with DEC 2026 sync support
this.pendingWrites = [];
this.writeFrameScheduled = false;
@@ -566,6 +571,20 @@ class CodemanApp {
document.body.classList.add('app-loaded');
}
_initWebGL() {
if (typeof WebglAddon === 'undefined') return;
try {
this._webglAddon = new WebglAddon.WebglAddon();
this._webglAddon.onContextLoss(() => {
console.error('[CRASH-DIAG] WebGL context LOST — falling back to canvas renderer');
this._webglAddon.dispose();
this._webglAddon = null;
});
this.terminal.loadAddon(this._webglAddon);
console.log('[CRASH-DIAG] WebGL renderer enabled');
} catch (_e) { /* WebGL2 unavailable — canvas renderer used */ }
}
// ═══════════════════════════════════════════════════════════════
// Terminal Setup — xterm.js config and input handling
// ═══════════════════════════════════════════════════════════════
@@ -623,28 +642,49 @@ 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
// oversized terminal.write() calls that triggered the stalls.
// Disable with ?nowebgl URL param if GPU issues return.
// Lazy-loaded: script downloaded only on desktop (saves 244KB on mobile).
this._webglAddon = null;
const skipWebGL = MobileDetection.getDeviceType() !== 'desktop';
if (!skipWebGL && !new URLSearchParams(location.search).has('nowebgl') && typeof WebglAddon !== 'undefined') {
try {
this._webglAddon = new WebglAddon.WebglAddon();
this._webglAddon.onContextLoss(() => {
console.error('[CRASH-DIAG] WebGL context LOST — falling back to canvas renderer');
this._webglAddon.dispose();
this._webglAddon = null;
});
this.terminal.loadAddon(this._webglAddon);
console.log('[CRASH-DIAG] WebGL renderer enabled via ?webgl param');
} catch (_e) { /* WebGL2 unavailable — canvas renderer used */ }
if (!skipWebGL && !new URLSearchParams(location.search).has('nowebgl')) {
if (typeof WebglAddon !== 'undefined') {
this._initWebGL();
} else {
// Lazy-load WebGL addon — not bundled in <head> to avoid blocking mobile
const wglScript = document.createElement('script');
wglScript.src = 'vendor/xterm-addon-webgl.min.js';
wglScript.onload = () => this._initWebGL();
wglScript.onerror = () => console.warn('[CRASH-DIAG] Failed to load WebGL addon — using canvas renderer');
document.head.appendChild(wglScript);
}
}
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' &&
@@ -857,6 +897,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.
@@ -1632,6 +1674,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();
@@ -1864,6 +1909,7 @@ class CodemanApp {
}
_onSessionDeleted(data) {
if (this._wsSessionId === data.id) this._disconnectWs();
this._cleanupSessionData(data.id);
if (this.activeSessionId === data.id) {
this.activeSessionId = null;
@@ -1878,6 +1924,21 @@ class CodemanApp {
if (this.sessions.size === 0) this.stopSystemStatsPolling();
}
// SSE wrappers — skip terminal events when WebSocket is delivering for this session.
// WS handler calls the underlying _onSession* methods directly.
_onSSETerminal(data) {
if (this._wsReady && this._wsSessionId === data.id) return;
this._onSessionTerminal(data);
}
_onSSENeedsRefresh(data) {
if (this._wsReady && this._wsSessionId === data?.id) return;
this._onSessionNeedsRefresh(data);
}
_onSSEClearTerminal(data) {
if (this._wsReady && this._wsSessionId === data?.id) return;
this._onSessionClearTerminal(data);
}
_onSessionTerminal(data) {
if (data.id === this.activeSessionId) {
if (data.data.length > 32768) _crashDiag.log(`TERMINAL: ${(data.data.length/1024).toFixed(0)}KB`);
@@ -1982,6 +2043,7 @@ class CodemanApp {
}
_onSessionExit(data) {
if (this._wsSessionId === data.id) this._disconnectWs();
const session = this.sessions.get(data.id);
if (session) {
session.status = 'stopped';
@@ -2835,6 +2897,91 @@ class CodemanApp {
}
}
// ═══════════════════════════════════════════════════════════════
// WebSocket Terminal I/O
// ═══════════════════════════════════════════════════════════════
/**
* Open a WebSocket for terminal I/O on the given session.
* Replaces HTTP POST input and SSE terminal output with a single
* bidirectional connection. Falls back to SSE+POST if WS fails.
*/
_connectWs(sessionId) {
this._disconnectWs();
const proto = location.protocol === 'https:' ? 'wss:' : 'ws:';
const url = `${proto}//${location.host}/ws/sessions/${sessionId}/terminal`;
const ws = new WebSocket(url);
this._ws = ws;
this._wsSessionId = sessionId;
ws.onopen = () => {
// Only mark ready if this is still the intended session
if (this._ws === ws) {
this._wsReady = true;
this._wsReconnectAttempts = 0;
}
};
ws.onmessage = (event) => {
if (this._ws !== ws) return;
try {
const msg = JSON.parse(event.data);
if (msg.t === 'o') {
// Terminal output — route through the same batching pipeline as SSE
this._onSessionTerminal({ id: sessionId, data: msg.d });
} else if (msg.t === 'c') {
this._onSessionClearTerminal({ id: sessionId });
} else if (msg.t === 'r') {
this._onSessionNeedsRefresh({ id: sessionId });
}
} catch {
// Ignore malformed messages
}
};
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);
}
};
ws.onerror = () => {
// onclose will fire after onerror — cleanup happens there
};
}
/** 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();
this._ws = null;
this._wsSessionId = null;
this._wsReady = false;
}
}
/**
* Send input to server without blocking the keystroke flush cycle.
* Uses a sequential promise chain to preserve character ordering
@@ -2847,11 +2994,19 @@ class CodemanApp {
return;
}
// Chain on dispatch only — wait for the previous request to be sent before
// dispatching the next one (preserves keystroke ordering), but don't wait
// for the server's response. The server handles writeViaMux as
// fire-and-forget anyway, so the HTTP response carries no useful data
// beyond success/failure for retry purposes.
// Fast path: WebSocket — fire-and-forget, inherently ordered (single TCP stream).
if (this._wsReady && this._wsSessionId === sessionId) {
try {
this._ws.send(JSON.stringify({ t: 'i', d: input }));
this.clearPendingHooks(sessionId);
return;
} catch {
// WS send failed — fall through to HTTP POST
}
}
// Slow path: HTTP POST — chain on dispatch only, don't wait for response.
// The server handles writeViaMux as fire-and-forget anyway.
this._inputSendChain = this._inputSendChain.then(() => {
const fetchPromise = fetch(`/api/sessions/${sessionId}/input`, {
method: 'POST',
@@ -2960,6 +3115,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');
@@ -3632,6 +3794,13 @@ class CodemanApp {
if (selectGen !== this._selectGeneration) return; // newer tab switch won
// 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);
@@ -3669,7 +3838,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.
@@ -3932,6 +4100,9 @@ class CodemanApp {
}
});
// Open WebSocket for low-latency terminal I/O (after buffer load completes)
this._connectWs(sessionId);
_crashDiag.log('FOCUS');
this.terminal.focus();
this.terminal.scrollToBottom();
@@ -5254,6 +5425,15 @@ class CodemanApp {
async sendResize(sessionId) {
const dims = this.getTerminalDimensions();
if (!dims) return;
// Fast path: WebSocket resize
if (this._wsReady && this._wsSessionId === sessionId) {
try {
this._ws.send(JSON.stringify({ t: 'z', c: dims.cols, r: dims.rows }));
return;
} catch {
// Fall through to HTTP POST
}
}
await fetch(`/api/sessions/${sessionId}/resize`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
@@ -6494,7 +6674,7 @@ class CodemanApp {
document.getElementById('appSettingsShowMonitor').checked = settings.showMonitor ?? defaults.showMonitor ?? true;
document.getElementById('appSettingsShowProjectInsights').checked = settings.showProjectInsights ?? defaults.showProjectInsights ?? false;
document.getElementById('appSettingsShowFileBrowser').checked = settings.showFileBrowser ?? defaults.showFileBrowser ?? false;
document.getElementById('appSettingsShowSubagents').checked = settings.showSubagents ?? defaults.showSubagents ?? true;
document.getElementById('appSettingsShowSubagents').checked = settings.showSubagents ?? defaults.showSubagents ?? false;
document.getElementById('appSettingsSubagentTracking').checked = settings.subagentTrackingEnabled ?? defaults.subagentTrackingEnabled ?? true;
document.getElementById('appSettingsSubagentActiveTabOnly').checked = settings.subagentActiveTabOnly ?? defaults.subagentActiveTabOnly ?? true;
document.getElementById('appSettingsImageWatcherEnabled').checked = settings.imageWatcherEnabled ?? defaults.imageWatcherEnabled ?? false;
@@ -7662,7 +7842,7 @@ class CodemanApp {
const settings = this.loadAppSettingsFromStorage();
const defaults = this.getDefaultSettings();
const showMonitor = settings.showMonitor ?? defaults.showMonitor ?? true;
const showSubagents = settings.showSubagents ?? defaults.showSubagents ?? true;
const showSubagents = settings.showSubagents ?? defaults.showSubagents ?? false;
const showFileBrowser = settings.showFileBrowser ?? defaults.showFileBrowser ?? false;
const monitorPanel = document.getElementById('monitorPanel');
+13 -2
View File
@@ -14,11 +14,16 @@
<!-- xterm.css loaded async — terminal won't display until xterm.js runs anyway -->
<link rel="preload" href="vendor/xterm.css" as="style" onload="this.onload=null;this.rel='stylesheet'">
<noscript><link rel="stylesheet" href="vendor/xterm.css"></noscript>
<!-- Preload critical resources — lets browser discover these during HTML parse
instead of waiting until <script> tags at bottom-of-body are reached. -->
<link rel="preload" href="vendor/xterm.min.js" as="script">
<link rel="preload" href="constants.js" as="script">
<link rel="preload" href="app.js" as="script">
<!-- Self-hosted xterm.js — eliminates CDN DNS/TLS latency (~100ms).
'defer' preserves execution order (xterm loads before fit addon). -->
<script defer src="vendor/xterm.min.js"></script>
<script defer src="vendor/xterm-addon-fit.min.js"></script>
<script defer src="vendor/xterm-addon-webgl.min.js"></script>
<!-- WebGL addon lazy-loaded by app.js on desktop only (skipped on mobile, saving 244KB) -->
<script defer src="vendor/xterm-addon-unicode11.min.js"></script>
<script defer src="vendor/xterm-zerolag-input.js"></script>
<!-- Synchronous mobile detection — runs before first paint to prevent panel flash -->
@@ -231,7 +236,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 +1699,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;
}
+1
View File
@@ -14,3 +14,4 @@ export { registerSessionRoutes } from './session-routes.js';
export { registerRespawnRoutes } from './respawn-routes.js';
export { registerRalphRoutes } from './ralph-routes.js';
export { registerPlanRoutes } from './plan-routes.js';
export { registerWsRoutes } from './ws-routes.js';
+1 -1
View File
@@ -505,7 +505,7 @@ export function registerRalphRoutes(
settings.lastUsedCase = caseName;
const dir = dirname(SETTINGS_PATH);
if (!existsSync(dir)) mkdirSync(dir, { recursive: true });
fs.writeFile(SETTINGS_PATH, JSON.stringify(settings, null, 2)).catch(() => {});
fs.writeFile(SETTINGS_PATH, JSON.stringify(settings, null, 2)).catch(() => {}); // Ignore - persisting lastUsedCase is non-critical
} catch {
/* non-critical */
}
+206
View File
@@ -0,0 +1,206 @@
/**
* @fileoverview WebSocket terminal I/O route.
*
* Provides a low-latency bidirectional channel for terminal input/output,
* bypassing the HTTP POST + SSE path that adds per-request middleware overhead.
* Auth is checked once on the WebSocket upgrade handshake (cookies are included
* automatically by the browser). After upgrade, the connection is raw — no
* per-message middleware processing.
*
* Additive: the existing HTTP POST /api/sessions/:id/input and SSE session:terminal
* paths remain fully functional. The frontend opts into WS when available and
* falls back transparently.
*
* Terminal output is micro-batched at 8ms to group Ink's rapid cursor-up redraws
* into single frames, preventing flicker from split ANSI sequences. This matches
* the SSE path's server-side batching (16-50ms) but at a shorter interval since
* WS has no Traefik buffering overhead.
*
* Protocol (all JSON text frames):
* Server -> Client:
* {"t":"o","d":"..."} — terminal output
* {"t":"c"} — clear terminal
* {"t":"r"} — needs refresh (reload buffer)
* Client -> Server:
* {"t":"i","d":"..."} — input (keystroke or paste)
* {"t":"z","c":N,"r":N} — resize terminal
*/
import { FastifyInstance } from 'fastify';
import type { WebSocket } from 'ws';
import type { SessionPort } from '../ports/session-port.js';
import { MAX_INPUT_LENGTH } from '../../config/terminal-limits.js';
/** Micro-batch interval for terminal output (ms). Short enough for low latency,
* long enough to group Ink's rapid cursor-up redraw sequences into single frames. */
const WS_BATCH_INTERVAL_MS = 8;
/** Flush immediately when batch exceeds this size (bytes) for responsiveness. */
const WS_BATCH_FLUSH_THRESHOLD = 16384;
/** How often to ping each WebSocket client (ms). Detects stale connections that
* TCP keepalive won't catch for minutes, especially through tunnels/proxies. */
const WS_PING_INTERVAL_MS = 30_000;
/** If pong isn't received within this window after a ping, terminate the socket. */
const WS_PONG_TIMEOUT_MS = 10_000;
/** DEC 2026 synchronized update markers. Wrapping output in these tells xterm.js
* to buffer all content and render atomically in a single frame — eliminates
* flicker from cursor-up redraws that Ink sends without its own sync markers
* (DA capability negotiation fails through the PTY→server→WS proxy chain). */
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;
const session = ctx.sessions.get(id);
if (!session) {
socket.close(4004, 'Session not found');
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;
let batchTimer: ReturnType<typeof setTimeout> | null = null;
const flushBatch = () => {
batchTimer = null;
if (batchChunks.length === 0 || socket.readyState !== 1) {
batchChunks = [];
batchSize = 0;
return;
}
const data = batchChunks.join('');
batchChunks = [];
batchSize = 0;
socket.send(`{"t":"o","d":${JSON.stringify(DEC_2026_START + data + DEC_2026_END)}}`);
};
// Attach message handler synchronously BEFORE any async work
// (@fastify/websocket requirement to avoid dropped messages).
socket.on('message', (raw) => {
try {
const msg = JSON.parse(String(raw));
if (msg.t === 'i' && typeof msg.d === 'string') {
if (msg.d.length > MAX_INPUT_LENGTH) return;
session.write(msg.d);
} else if (
msg.t === 'z' &&
Number.isInteger(msg.c) &&
Number.isInteger(msg.r) &&
msg.c >= 1 &&
msg.c <= 500 &&
msg.r >= 1 &&
msg.r <= 200
) {
session.resize(msg.c, msg.r);
}
} catch {
// Ignore malformed messages
}
});
// Terminal output -> micro-batched WS send
const onTerminal = (data: string) => {
if (socket.readyState !== 1) return;
batchChunks.push(data);
batchSize += data.length;
// Flush immediately for large batches (responsiveness during bulk output)
if (batchSize > WS_BATCH_FLUSH_THRESHOLD) {
if (batchTimer) {
clearTimeout(batchTimer);
}
flushBatch();
return;
}
// Start timer if not already running
if (!batchTimer) {
batchTimer = setTimeout(flushBatch, WS_BATCH_INTERVAL_MS);
}
};
const onClearTerminal = () => {
if (socket.readyState === 1) {
socket.send('{"t":"c"}');
}
};
const onNeedsRefresh = () => {
if (socket.readyState === 1) {
socket.send('{"t":"r"}');
}
};
// 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;
socket.on('pong', () => {
if (pongTimeout) {
clearTimeout(pongTimeout);
pongTimeout = null;
}
});
const pingInterval = setInterval(() => {
if (socket.readyState !== 1) return;
socket.ping();
pongTimeout = setTimeout(() => {
socket.terminate();
}, WS_PONG_TIMEOUT_MS);
}, WS_PING_INTERVAL_MS);
socket.on('close', () => {
clearInterval(pingInterval);
if (pongTimeout) clearTimeout(pongTimeout);
if (batchTimer) clearTimeout(batchTimer);
batchChunks = [];
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 -1
View File
@@ -8,7 +8,7 @@
*/
import { z } from 'zod';
import { SAFE_PATH_PATTERN } from '../utils/regex-patterns.js';
import { SAFE_PATH_PATTERN } from '../utils/index.js';
// ========== Path Validation ==========
+7
View File
@@ -31,6 +31,7 @@ import Fastify, { FastifyInstance, FastifyReply } from 'fastify';
import fastifyCompress from '@fastify/compress';
import fastifyCookie from '@fastify/cookie';
import fastifyStatic from '@fastify/static';
import fastifyWebsocket from '@fastify/websocket';
import { join, dirname } from 'node:path';
import { fileURLToPath } from 'node:url';
import { existsSync, mkdirSync, readFileSync, chmodSync } from 'node:fs';
@@ -103,6 +104,7 @@ import {
registerRespawnRoutes,
registerRalphRoutes,
registerPlanRoutes,
registerWsRoutes,
} from './routes/index.js';
const __dirname = dirname(fileURLToPath(import.meta.url));
@@ -563,6 +565,9 @@ export class WebServer extends EventEmitter {
this.qrAuthFailures = authState.qrAuthFailures;
}
// WebSocket support (terminal I/O — low-latency bidirectional channel)
await this.app.register(fastifyWebsocket);
// Security headers + CORS
registerSecurityHeaders(this.app, this.https);
// Service worker must never be cached — browsers check for SW updates on navigation
@@ -700,6 +705,7 @@ export class WebServer extends EventEmitter {
registerRespawnRoutes(this.app, ctx);
registerRalphRoutes(this.app, ctx);
registerPlanRoutes(this.app, ctx);
registerWsRoutes(this.app, ctx);
}
/**
@@ -1947,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 };
+502
View File
@@ -0,0 +1,502 @@
/**
* @fileoverview Tests for WebSocket terminal I/O route.
*
* Unlike other route tests that use app.inject(), WebSocket testing requires
* 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)
*/
import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest';
import Fastify, { type FastifyInstance } from 'fastify';
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, 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', () => {
clearTimeout(timer);
resolve(ws);
});
ws.on('error', (err) => {
clearTimeout(timer);
reject(err);
});
});
}
/** Helper: wait for the next WS message, parsed as JSON. */
function nextMessage(ws: WebSocket, timeoutMs = 2000): Promise<unknown> {
return new Promise((resolve, reject) => {
const timer = setTimeout(() => reject(new Error('WS message timeout')), timeoutMs);
ws.once('message', (raw) => {
clearTimeout(timer);
resolve(JSON.parse(String(raw)));
});
});
}
/** Helper: wait for WS close event and return { code, reason }. */
function waitForClose(ws: WebSocket, timeoutMs = 2000): Promise<{ code: number; reason: string }> {
return new Promise((resolve, reject) => {
const timer = setTimeout(() => reject(new Error('WS close timeout')), timeoutMs);
ws.on('close', (code, reason) => {
clearTimeout(timer);
resolve({ code, reason: reason.toString() });
});
});
}
/** 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;
beforeEach(async () => {
app = Fastify({ logger: false });
await app.register(fastifyWebsocket);
ctx = createMockRouteContext({ sessionId: 'ws-test-session' });
registerWsRoutes(app, ctx as never);
await app.listen({ port: PORT, host: '127.0.0.1' });
});
afterEach(async () => {
await app.close();
});
// ========== Session not found ==========
describe('session not found', () => {
it('closes with 4004 when session does not exist', async () => {
const ws = new WebSocket(`ws://127.0.0.1:${PORT}/ws/sessions/nonexistent/terminal`);
const { code, reason } = await waitForClose(ws);
expect(code).toBe(4004);
expect(reason).toBe('Session not found');
});
});
// ========== Terminal output ==========
describe('terminal output', () => {
it('receives terminal output via WS with DEC 2026 sync markers', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
// Emit terminal data from mock session
session.emit('terminal', 'hello world');
// Wait for the micro-batched message (8ms batch interval + margin)
const msg = (await nextMessage(ws)) as { t: string; d: string };
expect(msg.t).toBe('o');
// Should contain DEC 2026 sync markers wrapping the data
expect(msg.d).toContain('hello world');
expect(msg.d).toMatch(/^\x1b\[\?2026h/); // starts with DEC 2026 start
expect(msg.d).toMatch(/\x1b\[\?2026l$/); // ends with DEC 2026 end
} finally {
ws.close();
}
});
it('sends clearTerminal event as {"t":"c"}', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
ctx._session.emit('clearTerminal');
const msg = (await nextMessage(ws)) as { t: string };
expect(msg.t).toBe('c');
} finally {
ws.close();
}
});
it('sends needsRefresh event as {"t":"r"}', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
ctx._session.emit('needsRefresh');
const msg = (await nextMessage(ws)) as { t: string };
expect(msg.t).toBe('r');
} finally {
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 ==========
describe('client input', () => {
it('forwards input messages to session.write()', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send(JSON.stringify({ t: 'i', d: 'ls -la\r' }));
// Give the message handler time to process
await vi.waitFor(() => {
expect(session.writeBuffer).toContain('ls -la\r');
});
} finally {
ws.close();
}
});
it('ignores input exceeding MAX_INPUT_LENGTH', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
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
ws.send(JSON.stringify({ t: 'i', d: 'ok' }));
await vi.waitFor(() => {
expect(session.writeBuffer).toContain('ok');
});
// The oversized input should not have been written
expect(session.writeBuffer).not.toContain(hugeInput);
} finally {
ws.close();
}
});
it('ignores malformed JSON messages', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send('not-json{{{');
// Send valid input to verify connection still alive
ws.send(JSON.stringify({ t: 'i', d: 'after-bad' }));
await vi.waitFor(() => {
expect(session.writeBuffer).toContain('after-bad');
});
// Only 'after-bad' should be in the buffer
expect(session.writeBuffer).toHaveLength(1);
} finally {
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 ==========
describe('resize validation', () => {
it('accepts valid resize within bounds', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send(JSON.stringify({ t: 'z', c: 120, r: 40 }));
await vi.waitFor(() => {
expect(session.resize).toHaveBeenCalledWith(120, 40);
});
} finally {
ws.close();
}
});
it('accepts resize at minimum bounds (1x1)', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send(JSON.stringify({ t: 'z', c: 1, r: 1 }));
await vi.waitFor(() => {
expect(session.resize).toHaveBeenCalledWith(1, 1);
});
} finally {
ws.close();
}
});
it('accepts resize at maximum bounds (500x200)', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send(JSON.stringify({ t: 'z', c: 500, r: 200 }));
await vi.waitFor(() => {
expect(session.resize).toHaveBeenCalledWith(500, 200);
});
} finally {
ws.close();
}
});
it('rejects resize with cols out of bounds (0 cols)', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send(JSON.stringify({ t: 'z', c: 0, r: 40 }));
// Send a valid message to confirm processing continues
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
await vi.waitFor(() => {
expect(session.writeBuffer).toContain('sentinel');
});
expect(session.resize).not.toHaveBeenCalled();
} finally {
ws.close();
}
});
it('rejects resize with cols exceeding 500', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send(JSON.stringify({ t: 'z', c: 501, r: 40 }));
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
await vi.waitFor(() => {
expect(session.writeBuffer).toContain('sentinel');
});
expect(session.resize).not.toHaveBeenCalled();
} finally {
ws.close();
}
});
it('rejects resize with rows exceeding 200', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send(JSON.stringify({ t: 'z', c: 80, r: 201 }));
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
await vi.waitFor(() => {
expect(session.writeBuffer).toContain('sentinel');
});
expect(session.resize).not.toHaveBeenCalled();
} finally {
ws.close();
}
});
it('rejects resize with non-integer values', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send(JSON.stringify({ t: 'z', c: 80.5, r: 24 }));
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
await vi.waitFor(() => {
expect(session.writeBuffer).toContain('sentinel');
});
expect(session.resize).not.toHaveBeenCalled();
} finally {
ws.close();
}
});
it('rejects resize with negative values', async () => {
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
try {
const session = ctx._session;
ws.send(JSON.stringify({ t: 'z', c: -1, r: 24 }));
ws.send(JSON.stringify({ t: 'i', d: 'sentinel' }));
await vi.waitFor(() => {
expect(session.writeBuffer).toContain('sentinel');
});
expect(session.resize).not.toHaveBeenCalled();
} finally {
ws.close();
}
});
});
// ========== 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', () => {
it('removes session event listeners on close', async () => {
const session = ctx._session;
const listenersBefore = session.listenerCount('terminal');
const ws = await connectWs('/ws/sessions/ws-test-session/terminal');
// A listener was added for 'terminal'
expect(session.listenerCount('terminal')).toBe(listenersBefore + 1);
expect(session.listenerCount('clearTerminal')).toBeGreaterThanOrEqual(1);
expect(session.listenerCount('needsRefresh')).toBeGreaterThanOrEqual(1);
// Close the WS connection
ws.close();
await waitForClose(ws);
// Wait for server-side close handler
await vi.waitFor(() => {
expect(session.listenerCount('terminal')).toBe(listenersBefore);
});
expect(session.listenerCount('clearTerminal')).toBe(0);
expect(session.listenerCount('needsRefresh')).toBe(0);
});
});
});
+167 -138
View File
@@ -100,11 +100,7 @@ function createAssistantTextEntry(text: string, timestamp?: string): string {
});
}
function createToolUseEntry(
toolName: string,
input: Record<string, unknown>,
timestamp?: string
): string {
function createToolUseEntry(toolName: string, input: Record<string, unknown>, timestamp?: string): string {
return JSON.stringify({
type: 'assistant',
timestamp: timestamp || new Date().toISOString(),
@@ -134,6 +130,13 @@ function createToolResultEntry(content: string, timestamp?: string): string {
});
}
/** Create a mock readline interface with .close() method */
function createMockRl() {
const rl = new EventEmitter() as EventEmitter & { close: ReturnType<typeof vi.fn> };
rl.close = vi.fn();
return rl;
}
describe('SubagentWatcher', () => {
let watcher: SubagentWatcher;
let mockExistsSync: Mock;
@@ -166,12 +169,12 @@ describe('SubagentWatcher', () => {
// Default mocks - no projects exist
mockExistsSync.mockReturnValue(false);
mockStatSync.mockReturnValue({
mockStatSync.mockImplementation(() => ({
isDirectory: () => true,
birthtime: new Date(),
mtime: new Date(),
size: 0,
});
}));
mockReaddirSync.mockReturnValue([]);
mockReadFileSync.mockReturnValue('');
mockWatch.mockReturnValue({ close: vi.fn(), on: vi.fn(), off: vi.fn() });
@@ -239,9 +242,9 @@ describe('SubagentWatcher', () => {
const lines = [validEntry];
// Setup mock readline interface
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
// Setup file discovery
mockExistsSync.mockReturnValue(true);
@@ -281,16 +284,11 @@ describe('SubagentWatcher', () => {
it('should skip malformed JSON lines', async () => {
const validEntry = createUserEntry('Valid entry');
const malformedLines = [
'not json at all',
'{"incomplete": true',
validEntry,
'}{bad json}{',
];
const malformedLines = ['not json at all', '{"incomplete": true', validEntry, '}{bad json}{'];
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -330,9 +328,9 @@ describe('SubagentWatcher', () => {
// Simulate partial write where line is incomplete
const partialContent = '{"type": "user", "timestamp": "2024-01-01T00:00:00Z"';
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -369,9 +367,9 @@ describe('SubagentWatcher', () => {
const validEntry = createUserEntry('Valid');
const contentWithEmptyLines = ['', validEntry, ' ', '', validEntry].join('\n');
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -408,9 +406,9 @@ describe('SubagentWatcher', () => {
describe('Status Lifecycle', () => {
it('should start agents as active', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -442,9 +440,9 @@ describe('SubagentWatcher', () => {
});
it('should transition to idle after IDLE_TIMEOUT_MS (30s)', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -453,12 +451,12 @@ describe('SubagentWatcher', () => {
if (path.includes('project1')) return ['session1'];
return ['project1'];
});
mockStatSync.mockReturnValue({
mockStatSync.mockImplementation(() => ({
isDirectory: () => true,
birthtime: new Date(),
mtime: new Date(),
size: 100,
});
}));
mockReadFileSync.mockReturnValue(createUserEntry('Test subagent task'));
watcher.start();
@@ -481,9 +479,9 @@ describe('SubagentWatcher', () => {
});
it('should reset to active on new activity', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
const mockWatcher = { close: vi.fn(), on: vi.fn(), off: vi.fn() };
mockWatch.mockReturnValue(mockWatcher);
@@ -495,12 +493,12 @@ describe('SubagentWatcher', () => {
if (path.includes('project1')) return ['session1'];
return ['project1'];
});
mockStatSync.mockReturnValue({
mockStatSync.mockImplementation(() => ({
isDirectory: () => true,
birthtime: new Date(),
mtime: new Date(),
size: 100,
});
}));
mockReadFileSync.mockReturnValue(createUserEntry('Test subagent task'));
watcher.start();
@@ -514,19 +512,20 @@ describe('SubagentWatcher', () => {
expect(watcher.getSubagents()[0].status).toBe('idle');
// Simulate file change event - get the callback from mockWatch
const watchCallback = mockWatch.mock.calls.find(
(call: unknown[]) => typeof call[1] === 'function'
)?.[1];
const watchCallback = mockWatch.mock.calls.find((call: unknown[]) => typeof call[1] === 'function')?.[1];
if (watchCallback) {
// Need to reset the readline mock for the new read
const newMockRl = new EventEmitter();
const newMockRl = createMockRl();
mockCreateInterface.mockReturnValue(newMockRl);
// Trigger file change
watchCallback('change', 'agent-reactive.jsonl');
// Complete the new readline
// Advance past the fileDeb debounce (100ms) so handleFileChange runs
await vi.advanceTimersByTimeAsync(150);
// Complete the new readline (tailFile)
newMockRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
@@ -536,9 +535,9 @@ describe('SubagentWatcher', () => {
});
it('should transition to completed when file becomes stale', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -588,9 +587,10 @@ describe('SubagentWatcher', () => {
it('should extract tool_use entries', async () => {
const toolEntry = createToolUseEntry('WebSearch', { query: 'test query' });
const mockRl = new EventEmitter();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
const descRl = createMockRl();
const tailRl = createMockRl();
mockCreateInterface.mockReturnValueOnce(descRl).mockReturnValue(tailRl);
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -612,9 +612,14 @@ describe('SubagentWatcher', () => {
watcher.start();
await flushAsyncScan();
mockRl.emit('line', toolEntry);
mockRl.emit('close');
// Resolve extractDescriptionFromFile
descRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
// Now tailFile is set up — emit entries on tailRl
tailRl.emit('line', toolEntry);
tailRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
expect(toolCallHandler).toHaveBeenCalled();
@@ -626,9 +631,10 @@ describe('SubagentWatcher', () => {
it('should extract text messages', async () => {
const textEntry = createAssistantTextEntry('This is the assistant response');
const mockRl = new EventEmitter();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
const descRl = createMockRl();
const tailRl = createMockRl();
mockCreateInterface.mockReturnValueOnce(descRl).mockReturnValue(tailRl);
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -650,9 +656,14 @@ describe('SubagentWatcher', () => {
watcher.start();
await flushAsyncScan();
mockRl.emit('line', textEntry);
mockRl.emit('close');
// Resolve extractDescriptionFromFile
descRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
// Now tailFile is set up — emit entries on tailRl
tailRl.emit('line', textEntry);
tailRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
expect(messageHandler).toHaveBeenCalled();
@@ -682,9 +693,9 @@ describe('SubagentWatcher', () => {
},
});
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -720,9 +731,10 @@ describe('SubagentWatcher', () => {
const longText = 'x'.repeat(1000);
const textEntry = createAssistantTextEntry(longText);
const mockRl = new EventEmitter();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
const descRl = createMockRl();
const tailRl = createMockRl();
mockCreateInterface.mockReturnValueOnce(descRl).mockReturnValue(tailRl);
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -744,9 +756,14 @@ describe('SubagentWatcher', () => {
watcher.start();
await flushAsyncScan();
mockRl.emit('line', textEntry);
mockRl.emit('close');
// Resolve extractDescriptionFromFile
descRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
// Now tailFile is set up — emit entries on tailRl
tailRl.emit('line', textEntry);
tailRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
expect(messageHandler).toHaveBeenCalled();
@@ -757,9 +774,10 @@ describe('SubagentWatcher', () => {
it('should extract progress events', async () => {
const progressEntry = createProgressEntry('query_update', { query: 'searching for files' });
const mockRl = new EventEmitter();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
const descRl = createMockRl();
const tailRl = createMockRl();
mockCreateInterface.mockReturnValueOnce(descRl).mockReturnValue(tailRl);
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -781,9 +799,14 @@ describe('SubagentWatcher', () => {
watcher.start();
await flushAsyncScan();
mockRl.emit('line', progressEntry);
mockRl.emit('close');
// Resolve extractDescriptionFromFile
descRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
// Now tailFile is set up — emit entries on tailRl
tailRl.emit('line', progressEntry);
tailRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
expect(progressHandler).toHaveBeenCalled();
@@ -795,9 +818,9 @@ describe('SubagentWatcher', () => {
describe('Memory Management', () => {
it('should track agents in agentInfo map', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -832,9 +855,10 @@ describe('SubagentWatcher', () => {
const toolEntry1 = createToolUseEntry('Read', { file_path: '/test1.ts' });
const toolEntry2 = createToolUseEntry('Write', { file_path: '/test2.ts' });
const mockRl = new EventEmitter();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
const descRl = createMockRl();
const tailRl = createMockRl();
mockCreateInterface.mockReturnValueOnce(descRl).mockReturnValue(tailRl);
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -853,10 +877,15 @@ describe('SubagentWatcher', () => {
watcher.start();
await flushAsyncScan();
mockRl.emit('line', toolEntry1);
mockRl.emit('line', toolEntry2);
mockRl.emit('close');
// Resolve extractDescriptionFromFile
descRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
// Now tailFile is set up — emit entries on tailRl
tailRl.emit('line', toolEntry1);
tailRl.emit('line', toolEntry2);
tailRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
const agent = watcher.getSubagent('toolcount');
@@ -871,9 +900,10 @@ describe('SubagentWatcher', () => {
createToolUseEntry('Read', { file_path: '/test.ts' }),
];
const mockRl = new EventEmitter();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
const descRl = createMockRl();
const tailRl = createMockRl();
mockCreateInterface.mockReturnValueOnce(descRl).mockReturnValue(tailRl);
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -892,11 +922,16 @@ describe('SubagentWatcher', () => {
watcher.start();
await flushAsyncScan();
for (const entry of entries) {
mockRl.emit('line', entry);
}
mockRl.emit('close');
// Resolve extractDescriptionFromFile
descRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
// Now tailFile is set up — emit entries on tailRl
for (const entry of entries) {
tailRl.emit('line', entry);
}
tailRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
const agent = watcher.getSubagent('entrycount');
@@ -907,9 +942,9 @@ describe('SubagentWatcher', () => {
// Note: Current implementation has no cleanup/eviction policy
// This documents the behavior as a known issue
it('should retain all agents indefinitely (no cleanup policy)', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
@@ -949,9 +984,9 @@ describe('SubagentWatcher', () => {
createToolUseEntry('Read', { file_path: '/test.ts' }),
];
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -979,13 +1014,11 @@ describe('SubagentWatcher', () => {
});
it('should limit transcript entries when limit is specified', async () => {
const entries = Array.from({ length: 10 }, (_, i) =>
createUserEntry(`Message ${i}`)
);
const entries = Array.from({ length: 10 }, (_, i) => createUserEntry(`Message ${i}`));
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1013,9 +1046,9 @@ describe('SubagentWatcher', () => {
});
it('should return empty array for unknown agent', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockReturnValue([]);
@@ -1037,9 +1070,7 @@ describe('SubagentWatcher', () => {
sessionId: 'sess1',
message: {
role: 'assistant',
content: [
{ type: 'tool_use', name: 'WebSearch', input: { query: 'test query' } },
],
content: [{ type: 'tool_use', name: 'WebSearch', input: { query: 'test query' } }],
},
},
];
@@ -1092,9 +1123,9 @@ describe('SubagentWatcher', () => {
it('should extract description from first user message', async () => {
const userEntry = createUserEntry('Create comprehensive tests for the module');
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1116,6 +1147,7 @@ describe('SubagentWatcher', () => {
watcher.start();
await flushAsyncScan();
mockRl.emit('line', userEntry);
mockRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
@@ -1134,9 +1166,9 @@ describe('SubagentWatcher', () => {
const userEntry = createUserEntry(longPrompt);
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1158,6 +1190,7 @@ describe('SubagentWatcher', () => {
watcher.start();
await flushAsyncScan();
mockRl.emit('line', userEntry);
mockRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
@@ -1169,9 +1202,9 @@ describe('SubagentWatcher', () => {
});
it('should emit subagent:updated when description is extracted from processEntry', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1234,7 +1267,7 @@ describe('SubagentWatcher', () => {
},
});
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
// createReadStream now used for parent transcript reading (stream tail)
@@ -1287,9 +1320,9 @@ describe('SubagentWatcher', () => {
describe('getRecentSubagents', () => {
it('should return only recent subagents', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1317,9 +1350,9 @@ describe('SubagentWatcher', () => {
});
it('should sort by lastActivityAt descending', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1354,9 +1387,9 @@ describe('SubagentWatcher', () => {
describe('getSubagentsForSession', () => {
it('should filter subagents by working directory', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1392,9 +1425,9 @@ describe('SubagentWatcher', () => {
});
it('should return false for already completed agent', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1428,9 +1461,9 @@ describe('SubagentWatcher', () => {
});
it('should emit completed event when killing active agent', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1485,9 +1518,9 @@ describe('SubagentWatcher', () => {
});
it('should handle readline errors gracefully', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1524,9 +1557,9 @@ describe('SubagentWatcher', () => {
// Only directory watchers are created (no per-file watchers)
mockWatch.mockReturnValue(mockDirWatcher);
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1556,9 +1589,9 @@ describe('SubagentWatcher', () => {
});
it('should clear idle timers on stop', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1596,9 +1629,9 @@ describe('SubagentWatcher', () => {
describe('Project Hash Conversion', () => {
it('should convert working directory to project hash format', async () => {
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
@@ -1640,9 +1673,7 @@ describe('SubagentWatcher', () => {
sessionId: 'sess1',
message: {
role: 'assistant',
content: [
{ type: 'tool_use', name: 'WebSearch', input: { query: 'nodejs best practices' } },
],
content: [{ type: 'tool_use', name: 'WebSearch', input: { query: 'nodejs best practices' } }],
},
},
];
@@ -1661,9 +1692,7 @@ describe('SubagentWatcher', () => {
sessionId: 'sess1',
message: {
role: 'assistant',
content: [
{ type: 'tool_use', name: 'Read', input: { file_path: '/src/index.ts' } },
],
content: [{ type: 'tool_use', name: 'Read', input: { file_path: '/src/index.ts' } }],
},
},
];
@@ -1682,9 +1711,7 @@ describe('SubagentWatcher', () => {
sessionId: 'sess1',
message: {
role: 'assistant',
content: [
{ type: 'tool_use', name: 'Bash', input: { command: 'npm test' } },
],
content: [{ type: 'tool_use', name: 'Bash', input: { command: 'npm test' } }],
},
},
];
@@ -1704,9 +1731,7 @@ describe('SubagentWatcher', () => {
sessionId: 'sess1',
message: {
role: 'assistant',
content: [
{ type: 'tool_use', name: 'Bash', input: { command: longCommand } },
],
content: [{ type: 'tool_use', name: 'Bash', input: { command: longCommand } }],
},
},
];
@@ -1753,9 +1778,10 @@ describe('SubagentWatcher', () => {
it('should emit user messages under 500 chars', async () => {
const userEntry = createUserEntry('Short user message');
const mockRl = new EventEmitter();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
const descRl = createMockRl();
const tailRl = createMockRl();
mockCreateInterface.mockReturnValueOnce(descRl).mockReturnValue(tailRl);
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1777,9 +1803,14 @@ describe('SubagentWatcher', () => {
watcher.start();
await flushAsyncScan();
mockRl.emit('line', userEntry);
mockRl.emit('close');
// Resolve extractDescriptionFromFile
descRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
// Now tailFile is set up — emit entries on tailRl
tailRl.emit('line', userEntry);
tailRl.emit('close');
await vi.advanceTimersByTimeAsync(100);
expect(messageHandler).toHaveBeenCalled();
@@ -1790,9 +1821,9 @@ describe('SubagentWatcher', () => {
it('should not emit long user messages (over 500 chars)', async () => {
const longUserEntry = createUserEntry('x'.repeat(600));
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {
@@ -1820,9 +1851,7 @@ describe('SubagentWatcher', () => {
await vi.advanceTimersByTimeAsync(100);
// Long user messages are filtered out
const userMessages = messageHandler.mock.calls.filter(
(call) => (call[0] as SubagentMessage).role === 'user'
);
const userMessages = messageHandler.mock.calls.filter((call) => (call[0] as SubagentMessage).role === 'user');
expect(userMessages.length).toBe(0);
});
});
@@ -1831,9 +1860,9 @@ describe('SubagentWatcher', () => {
it('should not emit message for empty text content', async () => {
const emptyTextEntry = createAssistantTextEntry(' ');
const mockRl = new EventEmitter();
const mockRl = createMockRl();
mockCreateInterface.mockReturnValue(mockRl);
mockCreateReadStream.mockReturnValue({});
mockCreateReadStream.mockReturnValue({ destroy: vi.fn() });
mockExistsSync.mockReturnValue(true);
mockReaddirSync.mockImplementation((path: string) => {