Compare commits

...
Author SHA1 Message Date
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
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
25 changed files with 1027 additions and 221 deletions
+34
View File
@@ -1,5 +1,39 @@
# aicodeman
## 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
+3 -3
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.3.12 (must match `package.json`)
## Project Overview
@@ -109,8 +109,8 @@ Codeman is a Claude Code session manager with web interface and autonomous Ralph
| **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 |
| **Frontend** | `src/web/public/app.js` ★ (~12.5K lines) + 9 JS modules (incl. `sw.js` service worker) | |
| **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.
-5
View File
@@ -35,11 +35,6 @@ 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>
Binary file not shown.
+18 -2
View File
@@ -1283,7 +1283,16 @@ update() {
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 +1364,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"
+3 -1
View File
@@ -1,6 +1,6 @@
{
"name": "aicodeman",
"version": "0.3.10",
"version": "0.3.12",
"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",
@@ -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",
+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) {
+1 -7
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
// ============================================================================
+121 -10
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;
@@ -1864,6 +1869,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 +1884,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 +2003,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 +2857,72 @@ 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;
}
};
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 = () => {
if (this._ws === ws) {
this._ws = null;
this._wsSessionId = null;
this._wsReady = false;
}
};
ws.onerror = () => {
// onclose will fire after onerror — cleanup happens there
};
}
/** Close the active WebSocket connection (if any). */
_disconnectWs() {
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 +2935,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',
@@ -3632,6 +3728,9 @@ class CodemanApp {
if (selectGen !== this._selectGeneration) return; // newer tab switch won
// Close WebSocket for previous session (new one opens after buffer load)
this._disconnectWs();
// Clean up flicker filter state when switching sessions
if (this.flickerFilterTimeout) {
clearTimeout(this.flickerFilterTimeout);
@@ -3932,6 +4031,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 +5356,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 +6605,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 +7773,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');
+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 */
}
+179
View File
@@ -0,0 +1,179 @@
/**
* @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';
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;
}
// 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) => {
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"}');
}
};
session.on('terminal', onTerminal);
session.on('clearTerminal', onClearTerminal);
session.on('needsRefresh', onNeedsRefresh);
// Heartbeat: detect stale connections (especially through tunnels where
// TCP RST can take minutes to propagate).
let pongTimeout: ReturnType<typeof setTimeout> | null = null;
let alive = true;
socket.on('pong', () => {
alive = true;
if (pongTimeout) {
clearTimeout(pongTimeout);
pongTimeout = null;
}
});
const pingInterval = setInterval(() => {
if (!alive) {
// Previous ping never got a pong — connection is dead
socket.terminate();
return;
}
alive = false;
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);
});
});
}
+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 ==========
+6
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);
}
/**
+362
View File
@@ -0,0 +1,362 @@
/**
* @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.
*
* 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';
const PORT = 3170;
/** Helper: open a WebSocket connection and wait for it to reach OPEN state. */
function connectWs(path: string): Promise<WebSocket> {
return new Promise((resolve, reject) => {
const ws = new WebSocket(`ws://127.0.0.1:${PORT}${path}`);
ws.on('open', () => resolve(ws));
ws.on('error', reject);
});
}
/** 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() });
});
});
}
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();
}
});
});
// ========== 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;
session.writeBuffer = [];
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;
session.writeBuffer = [];
// MAX_INPUT_LENGTH is 64KB
const hugeInput = 'x'.repeat(65 * 1024);
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;
session.writeBuffer = [];
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();
}
});
});
// ========== 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 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);
// Give the server-side close handler time to run
await new Promise((resolve) => setTimeout(resolve, 50));
// Listeners should be cleaned up
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) => {