diff --git a/CLAUDE.md b/CLAUDE.md index fa484486..09564094 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -26,7 +26,7 @@ When user says "COM": 1. Increment version in BOTH `package.json` AND `CLAUDE.md` (verify they match with `grep version package.json && grep Version CLAUDE.md`) 2. Run: `git add -A && git commit -m "chore: bump version to X.XXXX" && git push && npm run build && systemctl --user restart claudeman-web` -**Version**: 0.1476 (must match `package.json` for npm publish) +**Version**: 0.1477 (must match `package.json` for npm publish) ## Project Overview @@ -174,31 +174,9 @@ journalctl --user -u claudeman-web -f ## Default Settings -UI defaults are optimized for minimal distraction. Set in `src/web/public/app.js` (using `??` operator). +UI defaults are set in `src/web/public/app.js` using `??` fallbacks. To change defaults, edit `openAppSettings()` and `apply*Visibility()` functions. -**Display Settings** (default values): -| Setting | Default | Description | -|---------|---------|-------------| -| `showFontControls` | `false` | Font size controls in header | -| `showSystemStats` | `true` | CPU/memory stats in header | -| `showTokenCount` | `true` | Token counter in header | -| `showCost` | `false` | Cost display | -| `showMonitor` | `true` | Monitor panel | -| `showProjectInsights` | `false` | Project insights panel | -| `showFileBrowser` | `false` | File browser panel | -| `showSubagents` | `false` | Subagent windows panel | - -**Tracking Settings**: -| Setting | Default | Description | -|---------|---------|-------------| -| `ralphTrackerEnabled` | `false` | Ralph/Todo loop tracking | -| `subagentTrackingEnabled` | `true` | Background agent monitoring | -| `subagentActiveTabOnly` | `true` | Show subagents only for active session | -| `imageWatcherEnabled` | `false` | Watch for image file creation | - -**Notification Defaults**: Browser notifications enabled, audio alerts disabled. Critical events (permission prompts, questions) notify by default; info events (respawn cycles, token milestones) are silent. - -To change defaults, edit the `??` fallback values in `openAppSettings()` and `apply*Visibility()` functions. +**Key defaults:** Most panels hidden (monitor, subagents shown), notifications enabled (audio disabled), subagent tracking on, Ralph tracking off. ## Testing @@ -231,6 +209,18 @@ curl localhost:3000/api/sessions/:id/run-summary | jq # Session timeline **Avoid port 3000 during E2E tests** — tests use ports 3183-3193 (see `test/e2e/e2e.config.ts`). +## Troubleshooting + +| Problem | Check | Fix | +|---------|-------|-----| +| Session won't start | `screen -ls` for orphans | Kill orphaned screens, check Claude CLI installed | +| Port 3000 in use | `lsof -i :3000` | Kill conflicting process or use `--port` flag | +| SSE not connecting | Browser console for errors | Check CORS, ensure server running | +| Respawn not triggering | Session settings → Respawn enabled? | Enable respawn, check idle timeout config | +| Terminal blank on tab switch | Network tab for `/api/sessions/:id/buffer` | Check session exists, restart server | +| Tests failing on screen limits | `screen -ls \| wc -l` | Clean up test screens: `screen -ls \| grep test \| awk '{print $1}' \| xargs -I{} screen -X -S {} quit` | +| State not persisting | `cat ~/.claudeman/state.json` | Check file permissions, disk space | + ## Performance Constraints The app must stay fast with 20 sessions and 50 agent windows: @@ -241,132 +231,17 @@ The app must stay fast with 20 sessions and 50 agent windows: ## Terminal Anti-Flicker System -Claude Code uses [Ink](https://github.com/vadimdemedes/ink) (React for terminals), which redraws the entire screen on every state change. Without special handling, users see constant flickering. Claudeman implements a 6-layer anti-flicker pipeline: +Claude Code uses Ink (React for terminals), which redraws the screen on every state change. Claudeman implements a 6-layer anti-flicker pipeline for smooth 60fps output: ``` -PTY Output → Server Batching → DEC 2026 Wrap → SSE → Client rAF → Sync Parser → xterm.js +PTY Output → Server Batching (16-50ms) → DEC 2026 Wrap → SSE → Client rAF → xterm.js ``` -### Layer Details +**Key functions:** `server.ts:batchTerminalData()`, `server.ts:flushTerminalBatches()`, `app.js:batchTerminalWrite()`, `app.js:extractSyncSegments()` -| Layer | Location | Technique | Latency | -|-------|----------|-----------|---------| -| **1. Server Batching** | `server.ts:batchTerminalData()` | Adaptive 16-50ms collection window | 16-50ms | -| **2. DEC Mode 2026** | `server.ts:flushTerminalBatches()` | Wraps with `\x1b[?2026h`...`\x1b[?2026l` | 0ms | -| **3. SSE Broadcast** | `server.ts:broadcast()` | JSON serialize once, send to all clients | 0ms | -| **4. Client rAF** | `app.js:batchTerminalWrite()` | `requestAnimationFrame` batching | 0-16ms | -| **5. Sync Block Parser** | `app.js:extractSyncSegments()` | Strips DEC 2026 markers, waits for complete blocks | 0-50ms | -| **6. Chunked Loading** | `app.js:chunkedTerminalWrite()` | 64KB/frame for large buffers | variable | +**Typical latency:** 16-32ms. Optional per-session flicker filter adds ~50ms for problematic terminals. -### Server-Side Implementation (`server.ts`) - -**Constants:** -```typescript -const TERMINAL_BATCH_INTERVAL = 16; // Base: 60fps -const BATCH_FLUSH_THRESHOLD = 32 * 1024; // Flush immediately if >32KB -const DEC_SYNC_START = '\x1b[?2026h'; // Begin synchronized update -const DEC_SYNC_END = '\x1b[?2026l'; // End synchronized update -``` - -**Adaptive Batching** (`batchTerminalData()`): -- Tracks event frequency per session via `lastTerminalEventTime` Map -- Event gap <10ms → 50ms batch window (rapid-fire Ink redraws) -- Event gap <20ms → 32ms batch window -- Otherwise → 16ms (60fps) -- Flushes immediately if batch exceeds 32KB for responsiveness - -**Flush Logic** (`flushTerminalBatches()`): -```typescript -const syncData = DEC_SYNC_START + data + DEC_SYNC_END; -this.broadcast('session:terminal', { id: sessionId, data: syncData }); -``` - -### Client-Side Implementation (`app.js`) - -**batchTerminalWrite(data):** -1. Checks if flicker filter is enabled (optional, per-session) -2. If flicker filter active: buffers screen-clear patterns (`ESC[2J`, `ESC[H ESC[J`, `ESC[nA`) -3. Accumulates data in `pendingWrites` -4. Schedules `requestAnimationFrame` if not already scheduled -5. On rAF callback: checks for incomplete sync blocks (start without end) -6. If incomplete: waits up to 50ms via `syncWaitTimeout` -7. Calls `flushPendingWrites()` when complete - -**extractSyncSegments(data):** -- Parses DEC 2026 markers, returns array of content segments -- Content before sync blocks returned as-is -- Content inside sync blocks returned without markers -- Incomplete blocks (start without end) returned with marker for next chunk - -**flushPendingWrites():** -```javascript -const segments = extractSyncSegments(this.pendingWrites); -this.pendingWrites = ''; // Clear before writing -for (const segment of segments) { - if (segment && !segment.startsWith(DEC_SYNC_START)) { - this.terminal.write(segment); // Skip incomplete blocks (start with marker) - } -} -``` -Note: Segments starting with `DEC_SYNC_START` are incomplete blocks awaiting more data. These are skipped (discarded if timeout forces flush). - -**chunkedTerminalWrite(buffer, chunkSize=128KB):** -- For large buffer restoration (session switch, reconnect) -- Writes 128KB per `requestAnimationFrame` to avoid UI jank -- Strips any embedded DEC 2026 markers from historical data - -**selectSession() optimizations:** -- Starts buffer fetch immediately before other setup -- Shows "Loading session..." indicator while fetching -- Parallelizes session attach with buffer fetch -- Fire-and-forget resize (doesn't block tab switch) - -### Optional Flicker Filter - -Per-session toggle via Session Settings. Adds ~50ms latency but eliminates remaining flicker on problematic terminals. - -**Detection patterns:** -- `ESC[2J` — Clear entire screen -- `ESC[H ESC[J` — Cursor home + clear to end -- `ESC[?25l ESC[H` — Hide cursor + home (Ink pattern) -- `ESC[nA` (n≥1) — Cursor up (Ink line redraw) - -When detected, buffers 50ms of subsequent output before flushing atomically. - -### Responsiveness Considerations - -**Latency sources:** -| Source | Best Case | Worst Case | Notes | -|--------|-----------|------------|-------| -| Server batching | 0ms (flush) | 50ms (rapid events) | Immediate flush if >32KB | -| Sync block wait | 0ms | 50ms | Only if marker split across packets | -| Flicker filter | 0ms (disabled) | 50ms (enabled) | Optional per-session | -| rAF scheduling | 0ms | 16ms | Display refresh sync | -| **Total** | **0ms** | **~115ms** | Worst case rare in practice | - -**Typical latency:** 16-32ms (server batch + rAF) - -**Edge cases handled:** -- Incomplete sync blocks: 50ms timeout forces flush (content discarded to prevent freeze) -- Large buffers: Chunked writing prevents UI freeze -- Server shutdown: Skips batching via `_isStopping` flag -- Session switch: Clears flicker filter state, pending writes, and sync timeout (prevents cross-session data bleed) -- SSE reconnect: `handleInit()` clears all pending write state - -**Trade-off:** If a sync block is split across SSE packets and the end marker doesn't arrive within 50ms, the incomplete content is discarded. This prioritizes responsiveness over completeness. In practice this is rare since the server always sends complete `SYNC_START...SYNC_END` pairs and SSE typically delivers them atomically. - -### Files Involved - -| File | Key Functions | -|------|---------------| -| `src/web/server.ts` | `batchTerminalData()`, `flushTerminalBatches()`, `broadcast()` | -| `src/web/public/app.js` | `batchTerminalWrite()`, `extractSyncSegments()`, `flushPendingWrites()`, `flushFlickerBuffer()`, `chunkedTerminalWrite()` | - -### DEC Mode 2026 Compatibility - -Terminals that natively support DEC 2026 will buffer and render atomically. Terminals that don't support it ignore the escape sequences harmlessly. xterm.js doesn't support DEC 2026 natively, so the client implements its own buffering by parsing the markers. - -**Supporting terminals:** WezTerm, Kitty, Ghostty, iTerm2 3.5+, Windows Terminal, VSCode terminal +See `docs/terminal-anti-flicker.md` for full implementation details (adaptive batching, DEC 2026 markers, edge cases). ## Resource Limits @@ -396,6 +271,7 @@ Use `LRUMap` for bounded caches with eviction, `StaleExpirationMap` for TTL-base | **Respawn state machine** | `docs/respawn-state-machine.md` | | **Ralph Loop guide** | `docs/ralph-wiggum-guide.md` | | **Claude Code hooks** | `docs/claude-code-hooks-reference.md` | +| **Terminal anti-flicker** | `docs/terminal-anti-flicker.md` | | **Browser/E2E testing** | `docs/browser-testing-guide.md` | | **API routes** | `src/web/server.ts:buildServer()` or README.md (full endpoint tables) | | **SSE events** | Search `broadcast(` in `server.ts` | diff --git a/docs/terminal-anti-flicker.md b/docs/terminal-anti-flicker.md new file mode 100644 index 00000000..0d69781d --- /dev/null +++ b/docs/terminal-anti-flicker.md @@ -0,0 +1,138 @@ +# Terminal Anti-Flicker System + +Claude Code uses [Ink](https://github.com/vadimdemedes/ink) (React for terminals), which redraws the entire screen on every state change. Without special handling, users see constant flickering. Claudeman implements a 6-layer anti-flicker pipeline. + +## Pipeline Overview + +``` +PTY Output → Server Batching → DEC 2026 Wrap → SSE → Client rAF → Sync Parser → xterm.js +``` + +| Layer | Location | Technique | Latency | +|-------|----------|-----------|---------| +| **1. Server Batching** | `server.ts:batchTerminalData()` | Adaptive 16-50ms collection window | 16-50ms | +| **2. DEC Mode 2026** | `server.ts:flushTerminalBatches()` | Wraps with `\x1b[?2026h`...`\x1b[?2026l` | 0ms | +| **3. SSE Broadcast** | `server.ts:broadcast()` | JSON serialize once, send to all clients | 0ms | +| **4. Client rAF** | `app.js:batchTerminalWrite()` | `requestAnimationFrame` batching | 0-16ms | +| **5. Sync Block Parser** | `app.js:extractSyncSegments()` | Strips DEC 2026 markers, waits for complete blocks | 0-50ms | +| **6. Chunked Loading** | `app.js:chunkedTerminalWrite()` | 64KB/frame for large buffers | variable | + +## Server-Side Implementation (`server.ts`) + +### Constants + +```typescript +const TERMINAL_BATCH_INTERVAL = 16; // Base: 60fps +const BATCH_FLUSH_THRESHOLD = 32 * 1024; // Flush immediately if >32KB +const DEC_SYNC_START = '\x1b[?2026h'; // Begin synchronized update +const DEC_SYNC_END = '\x1b[?2026l'; // End synchronized update +``` + +### Adaptive Batching (`batchTerminalData()`) + +- Tracks event frequency per session via `lastTerminalEventTime` Map +- Event gap <10ms → 50ms batch window (rapid-fire Ink redraws) +- Event gap <20ms → 32ms batch window +- Otherwise → 16ms (60fps) +- Flushes immediately if batch exceeds 32KB for responsiveness + +### Flush Logic (`flushTerminalBatches()`) + +```typescript +const syncData = DEC_SYNC_START + data + DEC_SYNC_END; +this.broadcast('session:terminal', { id: sessionId, data: syncData }); +``` + +## Client-Side Implementation (`app.js`) + +### `batchTerminalWrite(data)` + +1. Checks if flicker filter is enabled (optional, per-session) +2. If flicker filter active: buffers screen-clear patterns (`ESC[2J`, `ESC[H ESC[J`, `ESC[nA`) +3. Accumulates data in `pendingWrites` +4. Schedules `requestAnimationFrame` if not already scheduled +5. On rAF callback: checks for incomplete sync blocks (start without end) +6. If incomplete: waits up to 50ms via `syncWaitTimeout` +7. Calls `flushPendingWrites()` when complete + +### `extractSyncSegments(data)` + +- Parses DEC 2026 markers, returns array of content segments +- Content before sync blocks returned as-is +- Content inside sync blocks returned without markers +- Incomplete blocks (start without end) returned with marker for next chunk + +### `flushPendingWrites()` + +```javascript +const segments = extractSyncSegments(this.pendingWrites); +this.pendingWrites = ''; // Clear before writing +for (const segment of segments) { + if (segment && !segment.startsWith(DEC_SYNC_START)) { + terminal.write(segment); // Skip incomplete blocks (start with marker) + } +} +``` + +Note: Segments starting with `DEC_SYNC_START` are incomplete blocks awaiting more data. These are skipped (discarded if timeout forces flush). + +### `chunkedTerminalWrite(buffer, chunkSize=128KB)` + +- For large buffer restoration (session switch, reconnect) +- Writes 128KB per `requestAnimationFrame` to avoid UI jank +- Strips any embedded DEC 2026 markers from historical data + +### `selectSession()` Optimizations + +- Starts buffer fetch immediately before other setup +- Shows "Loading session..." indicator while fetching +- Parallelizes session attach with buffer fetch +- Fire-and-forget resize (doesn't block tab switch) + +## Optional Flicker Filter + +Per-session toggle via Session Settings. Adds ~50ms latency but eliminates remaining flicker on problematic terminals. + +### Detection Patterns + +- `ESC[2J` — Clear entire screen +- `ESC[H ESC[J` — Cursor home + clear to end +- `ESC[?25l ESC[H` — Hide cursor + home (Ink pattern) +- `ESC[nA` (n≥1) — Cursor up (Ink line redraw) + +When detected, buffers 50ms of subsequent output before flushing atomically. + +## Latency Analysis + +| Source | Best Case | Worst Case | Notes | +|--------|-----------|------------|-------| +| Server batching | 0ms (flush) | 50ms (rapid events) | Immediate flush if >32KB | +| Sync block wait | 0ms | 50ms | Only if marker split across packets | +| Flicker filter | 0ms (disabled) | 50ms (enabled) | Optional per-session | +| rAF scheduling | 0ms | 16ms | Display refresh sync | +| **Total** | **0ms** | **~115ms** | Worst case rare in practice | + +**Typical latency:** 16-32ms (server batch + rAF) + +## Edge Cases + +- **Incomplete sync blocks**: 50ms timeout forces flush (content discarded to prevent freeze) +- **Large buffers**: Chunked writing prevents UI freeze +- **Server shutdown**: Skips batching via `_isStopping` flag +- **Session switch**: Clears flicker filter state, pending writes, and sync timeout (prevents cross-session data bleed) +- **SSE reconnect**: `handleInit()` clears all pending write state + +**Trade-off:** If a sync block is split across SSE packets and the end marker doesn't arrive within 50ms, the incomplete content is discarded. This prioritizes responsiveness over completeness. In practice this is rare since the server always sends complete `SYNC_START...SYNC_END` pairs and SSE typically delivers them atomically. + +## DEC Mode 2026 Compatibility + +Terminals that natively support DEC 2026 will buffer and render atomically. Terminals that don't support it ignore the escape sequences harmlessly. xterm.js doesn't support DEC 2026 natively, so the client implements its own buffering by parsing the markers. + +**Supporting terminals:** WezTerm, Kitty, Ghostty, iTerm2 3.5+, Windows Terminal, VSCode terminal + +## Files Involved + +| File | Key Functions | +|------|---------------| +| `src/web/server.ts` | `batchTerminalData()`, `flushTerminalBatches()`, `broadcast()` | +| `src/web/public/app.js` | `batchTerminalWrite()`, `extractSyncSegments()`, `flushPendingWrites()`, `flushFlickerBuffer()`, `chunkedTerminalWrite()` | diff --git a/package.json b/package.json index fc819427..5b1de108 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "claudeman", - "version": "0.1476", + "version": "0.1477", "description": "The missing control plane for Claude Code - run 20 autonomous agents with real-time monitoring and session persistence", "type": "module", "main": "dist/index.js", diff --git a/src/ai-checker-base.ts b/src/ai-checker-base.ts index a4db9153..f26b6a53 100644 --- a/src/ai-checker-base.ts +++ b/src/ai-checker-base.ts @@ -31,6 +31,28 @@ import { EventEmitter } from 'node:events'; import { getAugmentedPath } from './session.js'; import { ANSI_ESCAPE_PATTERN_SIMPLE } from './utils/index.js'; +// ========== Security Validation ========== + +/** + * Validates that a model name is safe for shell use. + * Model names should only contain alphanumeric characters, hyphens, underscores, and dots. + */ +function isValidModelName(model: string): boolean { + if (!model || typeof model !== 'string') return false; + // Allow: alphanumeric, hyphens, underscores, dots, slashes (for model paths like claude/opus-4.5) + // Max length 100 to prevent abuse + return /^[a-zA-Z0-9._/-]+$/.test(model) && model.length <= 100; +} + +/** + * Validates that a screen name is safe for shell use. + * Screen names should only contain alphanumeric characters, hyphens, and underscores. + */ +function isValidScreenName(screenName: string): boolean { + if (!screenName || typeof screenName !== 'string') return false; + return /^[a-zA-Z0-9_-]+$/.test(screenName) && screenName.length <= 100; +} + // ========== Types ========== /** Base configuration shared by all AI checkers */ @@ -351,6 +373,11 @@ export abstract class AiCheckerBase< // ========== Private Methods ========== private async runCheck(terminalBuffer: string): Promise { + // Security: Validate model name before use in shell commands + if (!isValidModelName(this.config.model)) { + throw new Error(`Invalid model name: ${this.config.model.substring(0, 50)}`); + } + // Prepare the terminal buffer (strip ANSI, trim to maxContextChars) const stripped = terminalBuffer.replace(ANSI_ESCAPE_PATTERN_SIMPLE, ''); const trimmed = stripped.length > this.config.maxContextChars @@ -367,6 +394,11 @@ export abstract class AiCheckerBase< this.checkPromptFile = join(tmpdir(), `${this.tempFilePrefix}-prompt-${shortId}-${timestamp}.txt`); this.checkScreenName = `${this.screenNamePrefix}${shortId}`; + // Security: Validate screen name before use in shell commands + if (!isValidScreenName(this.checkScreenName)) { + throw new Error(`Invalid screen name generated: ${this.checkScreenName.substring(0, 50)}`); + } + // Ensure output temp file exists (empty) so we can poll it writeFileSync(this.checkTempFile, ''); @@ -375,16 +407,17 @@ export abstract class AiCheckerBase< writeFileSync(this.checkPromptFile, prompt); // Build the command - read prompt from file via stdin to avoid argument size limits - const modelArg = `--model ${this.config.model}`; + // Quote model name to prevent command injection (model names should be simple alphanumeric but be safe) + const modelArg = `--model "${this.config.model.replace(/"/g, '\\"')}"`; const augmentedPath = getAugmentedPath(); const claudeCmd = `cat "${this.checkPromptFile}" | claude -p ${modelArg} --output-format text`; const fullCmd = `export PATH="${augmentedPath}"; ${claudeCmd} > "${this.checkTempFile}" 2>&1; echo "${this.doneMarker}" >> "${this.checkTempFile}"; rm -f "${this.checkPromptFile}"`; // Spawn screen try { - // Kill any leftover screen with this name first + // Kill any leftover screen with this name first (screen name already validated above) try { - execSync(`screen -X -S ${this.checkScreenName} quit 2>/dev/null`, { timeout: 3000 }); + execSync(`screen -X -S "${this.checkScreenName}" quit 2>/dev/null`, { timeout: 3000 }); } catch { // No existing screen, that's fine } @@ -407,9 +440,12 @@ export abstract class AiCheckerBase< const startTime = this.checkStartTime; this.checkResolve = resolve; + // Guard flag to prevent both poll and timeout from resolving (race condition) + let resolved = false; + this.checkPollTimer = setInterval(() => { - if (this.checkCancelled) { - // Cancel was already handled by cancel() calling resolve + if (this.checkCancelled || resolved) { + // Cancel or already resolved - stop polling return; } @@ -417,19 +453,21 @@ export abstract class AiCheckerBase< if (!this.checkTempFile || !existsSync(this.checkTempFile)) return; const content = readFileSync(this.checkTempFile, 'utf-8'); if (content.includes(this.doneMarker)) { + resolved = true; // Mark as resolved first to prevent timeout race const durationMs = Date.now() - startTime; const result = this.parseOutput(content, durationMs); this.checkResolve = null; resolve(result); } } catch { - // File might not be ready yet, keep polling + // File might not be ready yet or was deleted during cleanup, keep polling } }, POLL_INTERVAL_MS); // Set timeout this.checkTimeoutTimer = setTimeout(() => { - if (this._status === 'checking' && !this.checkCancelled) { + if (this._status === 'checking' && !this.checkCancelled && !resolved) { + resolved = true; // Mark as resolved first to prevent poll race this.checkResolve = null; reject(new Error(`${this.checkDescription} timed out after ${this.config.checkTimeoutMs}ms`)); } @@ -467,12 +505,21 @@ export abstract class AiCheckerBase< this.checkTimeoutTimer = null; } - // Kill the screen + // Kill the screen with fallback for stubborn processes if (this.checkScreenName) { + const screenName = this.checkScreenName; + // Screen name was validated before use, but still quote for defense-in-depth try { - execSync(`screen -X -S ${this.checkScreenName} quit 2>/dev/null`, { timeout: 3000 }); + execSync(`screen -X -S "${screenName}" quit 2>/dev/null`, { timeout: 2000 }); } catch { - // Screen may already be dead + // First attempt failed - try force kill via pkill as fallback + // Escape regex metacharacters in screen name to prevent pattern injection + const escapedName = screenName.replace(/[.*+?^${}()|[\]\\]/g, '\\$&'); + try { + execSync(`pkill -f "SCREEN.*${escapedName}" 2>/dev/null`, { timeout: 1000 }); + } catch { + // Screen may already be dead or no matching process + } } this.checkScreenName = null; } diff --git a/src/file-stream-manager.ts b/src/file-stream-manager.ts index d1a9bf13..187e1136 100644 --- a/src/file-stream-manager.ts +++ b/src/file-stream-manager.ts @@ -169,7 +169,9 @@ export class FileStreamManager extends EventEmitter { error: `File too large (${Math.round(stats.size / 1024 / 1024)}MB > ${MAX_FILE_SIZE / 1024 / 1024}MB limit)`, }; } - } catch { + } catch (err) { + const errorCode = err instanceof Error && 'code' in err ? (err as NodeJS.ErrnoException).code : 'UNKNOWN'; + console.warn(`[FileStreamManager] Failed to stat file "${absolutePath}" (${errorCode}):`, err instanceof Error ? err.message : String(err)); return { success: false, error: 'File not found or not accessible' }; } @@ -250,17 +252,26 @@ export class FileStreamManager extends EventEmitter { stream.active = false; + // Remove all event listeners to prevent memory leaks (closures hold references) + stream.process.stdout?.removeAllListeners(); + stream.process.stderr?.removeAllListeners(); + stream.process.removeAllListeners(); + // Kill the tail process try { stream.process.kill('SIGTERM'); // Force kill after 1 second if still running - setTimeout(() => { + const forceKillTimer = setTimeout(() => { try { stream.process.kill('SIGKILL'); } catch { // Already dead } }, 1000); + // Clear the force kill timer if process exits naturally + stream.process.once('exit', () => { + clearTimeout(forceKillTimer); + }); } catch { // Process may have already exited } diff --git a/src/plan-orchestrator.ts b/src/plan-orchestrator.ts index 3ddfb83c..1ec8b3a1 100644 --- a/src/plan-orchestrator.ts +++ b/src/plan-orchestrator.ts @@ -253,15 +253,18 @@ export class PlanOrchestrator { return md; } - cancel(): void { + async cancel(): Promise { this.cancelled = true; + // Stop all running sessions and await cleanup to prevent PTY process leaks + const stopPromises: Promise[] = []; for (const session of this.runningSessions) { - try { - session.stop(); - } catch (err) { - console.error('[PlanOrchestrator] Failed to stop session during cancel:', err); - } + stopPromises.push( + session.stop().catch((err) => { + console.error('[PlanOrchestrator] Failed to stop session during cancel:', err); + }) + ); } + await Promise.all(stopPromises); this.runningSessions.clear(); } diff --git a/src/ralph-loop.ts b/src/ralph-loop.ts index a996886f..23e6baca 100644 --- a/src/ralph-loop.ts +++ b/src/ralph-loop.ts @@ -184,6 +184,11 @@ export class RalphLoop extends EventEmitter { this.tasksCompleted = 0; this.tasksGenerated = 0; + // Re-setup event handlers if they were cleaned up during stop() + if (!this.sessionEventHandlers) { + this.setupEventHandlers(); + } + this.store.setRalphLoopState({ status: 'running', startedAt: this.startedAt, @@ -209,6 +214,9 @@ export class RalphLoop extends EventEmitter { this.loopTimer = null; } + // Clean up event handlers to prevent memory leaks + this.cleanupEventHandlers(); + this.store.setRalphLoopState({ status: 'stopped', lastCheckAt: Date.now(), @@ -251,10 +259,15 @@ export class RalphLoop extends EventEmitter { this.tick() .catch((err) => { - this.emit('error', err); + // Only emit if still running (stop() may have been called during tick()) + if (this._status === 'running') { + this.emit('error', err); + } }) .finally(() => { - if (this._status === 'running') { + // Guard: only reschedule if still running AND no timer is pending + // (prevents race where stop() clears timer between our check and setTimeout) + if (this._status === 'running' && this.loopTimer === null) { this.loopTimer = setTimeout(() => this.runLoop(), this.pollIntervalMs); } }); @@ -263,13 +276,13 @@ export class RalphLoop extends EventEmitter { private async tick(): Promise { this.store.setRalphLoopState({ lastCheckAt: Date.now() }); - // Check for timed out tasks - await this.checkTimeouts(); + // Run independent checks in parallel for better performance + await Promise.all([ + this.checkTimeouts(), + this.assignTasks(), + ]); - // Assign tasks to idle sessions - await this.assignTasks(); - - // Check if we should auto-generate tasks + // Check if we should auto-generate tasks (depends on assignment results) if (this.autoGenerateTasks && this.shouldGenerateTasks()) { await this.generateFollowUpTasks(); } @@ -480,3 +493,8 @@ export function destroyRalphLoop(): void { loopInstance = null; } } + +// Ensure cleanup on process exit (prevents orphaned event handlers and circular references) +process.on('exit', () => { + destroyRalphLoop(); +}); diff --git a/src/ralph-tracker.ts b/src/ralph-tracker.ts index f8aa61ce..9adb99dd 100644 --- a/src/ralph-tracker.ts +++ b/src/ralph-tracker.ts @@ -398,6 +398,55 @@ const COMPLETION_INDICATOR_PATTERNS = [ /project\s+(?:is\s+)?(?:completed?|done|finished)/i, ]; +// ---------- Priority Detection Patterns ---------- +// Pre-compiled for performance; avoids repeated allocation in parsePriority() + +/** P0 (Critical) priority patterns - highest severity issues */ +const P0_PRIORITY_PATTERNS = [ + /\bP0\b|\(P0\)|:?\s*P0\s*:/, // Explicit P0 + /\bCRITICAL\b/, // Critical keyword + /\bBLOCKER\b/, // Blocker + /\bURGENT\b/, // Urgent + /\bSECURITY\b/, // Security issues + /\bCRASH(?:ES|ING)?\b/, // Crash, crashes, crashing + /\bBROKEN\b/, // Broken + /\bDATA\s*LOSS\b/, // Data loss + /\bPRODUCTION\s*(?:DOWN|ISSUE|BUG)\b/, // Production issues + /\bHOTFIX\b/, // Hotfix + /\bSEVERITY\s*1\b/, // Severity 1 +]; + +/** P1 (High) priority patterns - important issues requiring attention */ +const P1_PRIORITY_PATTERNS = [ + /\bP1\b|\(P1\)|:?\s*P1\s*:/, // Explicit P1 + /\bHIGH\s*PRIORITY\b/, // High priority + /\bIMPORTANT\b/, // Important + /\bBUG\b/, // Bug + /\bFIX\b/, // Fix (as task type) + /\bERROR\b/, // Error + /\bFAIL(?:S|ED|ING|URE)?\b/, // Fail variants + /\bREGRESSION\b/, // Regression + /\bMUST\s*(?:HAVE|FIX|DO)\b/, // Must have/fix/do + /\bSEVERITY\s*2\b/, // Severity 2 + /\bREQUIRED\b/, // Required +]; + +/** P2 (Medium) priority patterns - lower priority improvements */ +const P2_PRIORITY_PATTERNS = [ + /\bP2\b|\(P2\)|:?\s*P2\s*:/, // Explicit P2 + /\bNICE\s*TO\s*HAVE\b/, // Nice to have + /\bLOW\s*PRIORITY\b/, // Low priority + /\bREFACTOR\b/, // Refactor + /\bCLEANUP\b/, // Cleanup + /\bIMPROVE(?:MENT)?\b/, // Improve/Improvement + /\bOPTIMIZ(?:E|ATION)\b/, // Optimize/Optimization + /\bCONSIDER\b/, // Consider + /\bWOULD\s*BE\s*NICE\b/, // Would be nice + /\bENHANCE(?:MENT)?\b/, // Enhance/Enhancement + /\bTECH(?:NICAL)?\s*DEBT\b/, // Tech debt + /\bDOCUMENT(?:ATION)?\b/, // Documentation +]; + // ========== Event Types ========== /** @@ -553,6 +602,8 @@ export class RalphTracker extends EventEmitter { /** File watcher for @fix_plan.md */ private _fixPlanWatcher: FSWatcher | null = null; + /** Error handler for FSWatcher (stored for cleanup to prevent memory leak) */ + private _fixPlanWatcherErrorHandler: ((err: Error) => void) | null = null; /** Debounce timer for file change events */ private _fixPlanReloadTimer: NodeJS.Timeout | null = null; @@ -610,8 +661,8 @@ export class RalphTracker extends EventEmitter { /** Whether stall warning has been emitted */ private _iterationStallWarned: boolean = false; - /** Alternate completion phrases (P1-003: multi-phrase support) */ - private _alternateCompletionPhrases: string[] = []; + /** Alternate completion phrases (P1-003: multi-phrase support) - Set for O(1) lookup */ + private _alternateCompletionPhrases: Set = new Set(); // ========== P1-009: Progress Estimation ========== @@ -644,9 +695,9 @@ export class RalphTracker extends EventEmitter { * @param phrase - Additional phrase that can trigger completion */ addAlternateCompletionPhrase(phrase: string): void { - if (!this._alternateCompletionPhrases.includes(phrase)) { - this._alternateCompletionPhrases.push(phrase); - this._loopState.alternateCompletionPhrases = [...this._alternateCompletionPhrases]; + if (!this._alternateCompletionPhrases.has(phrase)) { + this._alternateCompletionPhrases.add(phrase); + this._loopState.alternateCompletionPhrases = Array.from(this._alternateCompletionPhrases); this.emit('loopUpdate', this.loopState); } } @@ -656,10 +707,8 @@ export class RalphTracker extends EventEmitter { * @param phrase - Phrase to remove */ removeAlternateCompletionPhrase(phrase: string): void { - const index = this._alternateCompletionPhrases.indexOf(phrase); - if (index !== -1) { - this._alternateCompletionPhrases.splice(index, 1); - this._loopState.alternateCompletionPhrases = [...this._alternateCompletionPhrases]; + if (this._alternateCompletionPhrases.delete(phrase)) { + this._loopState.alternateCompletionPhrases = Array.from(this._alternateCompletionPhrases); this.emit('loopUpdate', this.loopState); } } @@ -785,6 +834,15 @@ export class RalphTracker extends EventEmitter { this.handleFixPlanChange(); }); } + // Add error handler to prevent unhandled errors and clean up on failure + // Store handler reference for proper cleanup in stopWatchingFixPlan() + if (this._fixPlanWatcher) { + this._fixPlanWatcherErrorHandler = (err: Error) => { + console.log(`[RalphTracker] FSWatcher error for @fix_plan.md: ${err.message}`); + this.stopWatchingFixPlan(); + }; + this._fixPlanWatcher.on('error', this._fixPlanWatcherErrorHandler); + } } catch (err) { console.log(`[RalphTracker] Could not watch @fix_plan.md: ${err}`); } @@ -810,6 +868,11 @@ export class RalphTracker extends EventEmitter { */ stopWatchingFixPlan(): void { if (this._fixPlanWatcher) { + // Remove error handler before closing to prevent memory leak + if (this._fixPlanWatcherErrorHandler) { + this._fixPlanWatcher.off('error', this._fixPlanWatcherErrorHandler); + this._fixPlanWatcherErrorHandler = null; + } this._fixPlanWatcher.close(); this._fixPlanWatcher = null; } @@ -1448,7 +1511,8 @@ export class RalphTracker extends EventEmitter { // Check if the count matches our todo count (e.g., "All 8 files created") const countMatch = line.match(ALL_COUNT_PATTERN); - const mentionedCount = countMatch ? parseInt(countMatch[1]) : null; + const parsedCount = countMatch ? parseInt(countMatch[1], 10) : NaN; + const mentionedCount = Number.isNaN(parsedCount) ? null : parsedCount; const todoCount = this._todos.size; // If a count is mentioned, it should match our todo count (within reason) @@ -1494,7 +1558,8 @@ export class RalphTracker extends EventEmitter { // Only act on explicit task number references like "Task 8 is done" const taskNumMatch = line.match(/task\s*#?(\d+)/i); if (taskNumMatch) { - const taskNum = parseInt(taskNumMatch[1]); + const taskNum = parseInt(taskNumMatch[1], 10); + if (Number.isNaN(taskNum)) return; // Find the nth todo (by order) and mark it complete let count = 0; for (const [_id, todo] of this._todos) { @@ -1850,8 +1915,8 @@ export class RalphTracker extends EventEmitter { // Check for max iterations setting const maxIterMatch = line.match(MAX_ITERATIONS_PATTERN); if (maxIterMatch) { - const maxIter = parseInt(maxIterMatch[1]); - if (!isNaN(maxIter) && maxIter > 0) { + const maxIter = parseInt(maxIterMatch[1], 10); + if (!Number.isNaN(maxIter) && maxIter > 0) { this._loopState.maxIterations = maxIter; this._loopState.lastActivity = Date.now(); // Use debounced emit for settings changes @@ -1863,10 +1928,11 @@ export class RalphTracker extends EventEmitter { const iterMatch = line.match(ITERATION_PATTERN); if (iterMatch) { // Pattern captures: group 1&2 for "Iteration X/Y", group 3&4 for "[X/Y]" - const currentIter = parseInt(iterMatch[1] || iterMatch[3]); - const maxIter = iterMatch[2] || iterMatch[4] ? parseInt(iterMatch[2] || iterMatch[4]) : null; + const currentIter = parseInt(iterMatch[1] || iterMatch[3], 10); + const maxIterStr = iterMatch[2] || iterMatch[4]; + const maxIter = maxIterStr ? parseInt(maxIterStr, 10) : null; - if (!isNaN(currentIter)) { + if (!Number.isNaN(currentIter)) { this.activateLoopIfNeeded(); // Track iteration changes for stall detection if (currentIter !== this._lastObservedIteration) { @@ -1892,7 +1958,7 @@ export class RalphTracker extends EventEmitter { } } this._loopState.cycleCount = currentIter; - if (maxIter !== null && !isNaN(maxIter)) { + if (maxIter !== null && !Number.isNaN(maxIter)) { this._loopState.maxIterations = maxIter; } this._loopState.lastActivity = Date.now(); @@ -1913,8 +1979,8 @@ export class RalphTracker extends EventEmitter { // Check for cycle count (legacy pattern) const cycleMatch = line.match(CYCLE_PATTERN); if (cycleMatch) { - const cycleNum = parseInt(cycleMatch[1] || cycleMatch[2]); - if (!isNaN(cycleNum) && cycleNum > this._loopState.cycleCount) { + const cycleNum = parseInt(cycleMatch[1] || cycleMatch[2], 10); + if (!Number.isNaN(cycleNum) && cycleNum > this._loopState.cycleCount) { this._loopState.cycleCount = cycleNum; this._loopState.lastActivity = Date.now(); // Use debounced emit for cycle updates @@ -2119,68 +2185,23 @@ export class RalphTracker extends EventEmitter { private parsePriority(content: string): RalphTodoPriority { const upper = content.toUpperCase(); - // P0 patterns - Critical issues - const p0Patterns = [ - /\bP0\b|\(P0\)|:?\s*P0\s*:/, // Explicit P0 - /\bCRITICAL\b/, // Critical keyword - /\bBLOCKER\b/, // Blocker - /\bURGENT\b/, // Urgent - /\bSECURITY\b/, // Security issues - /\bCRASH(?:ES|ING)?\b/, // Crash, crashes, crashing - /\bBROKEN\b/, // Broken - /\bDATA\s*LOSS\b/, // Data loss - /\bPRODUCTION\s*(?:DOWN|ISSUE|BUG)\b/, // Production issues - /\bHOTFIX\b/, // Hotfix - /\bSEVERITY\s*1\b/, // Severity 1 - ]; - - // P1 patterns - High priority issues - const p1Patterns = [ - /\bP1\b|\(P1\)|:?\s*P1\s*:/, // Explicit P1 - /\bHIGH\s*PRIORITY\b/, // High priority - /\bIMPORTANT\b/, // Important - /\bBUG\b/, // Bug - /\bFIX\b/, // Fix (as task type) - /\bERROR\b/, // Error - /\bFAIL(?:S|ED|ING|URE)?\b/, // Fail variants - /\bREGRESSION\b/, // Regression - /\bMUST\s*(?:HAVE|FIX|DO)\b/, // Must have/fix/do - /\bSEVERITY\s*2\b/, // Severity 2 - /\bREQUIRED\b/, // Required - ]; - - // P2 patterns - Lower priority - const p2Patterns = [ - /\bP2\b|\(P2\)|:?\s*P2\s*:/, // Explicit P2 - /\bNICE\s*TO\s*HAVE\b/, // Nice to have - /\bLOW\s*PRIORITY\b/, // Low priority - /\bREFACTOR\b/, // Refactor - /\bCLEANUP\b/, // Cleanup - /\bIMPROVE(?:MENT)?\b/, // Improve/Improvement - /\bOPTIMIZ(?:E|ATION)\b/, // Optimize/Optimization - /\bCONSIDER\b/, // Consider - /\bWOULD\s*BE\s*NICE\b/, // Would be nice - /\bENHANCE(?:MENT)?\b/, // Enhance/Enhancement - /\bTECH(?:NICAL)?\s*DEBT\b/, // Tech debt - /\bDOCUMENT(?:ATION)?\b/, // Documentation - ]; - // Check P0 first (highest priority wins) - for (const pattern of p0Patterns) { + // Uses pre-compiled module-level patterns for performance + for (const pattern of P0_PRIORITY_PATTERNS) { if (pattern.test(upper)) { return 'P0'; } } // Check P1 - for (const pattern of p1Patterns) { + for (const pattern of P1_PRIORITY_PATTERNS) { if (pattern.test(upper)) { return 'P1'; } } // Check P2 - for (const pattern of p2Patterns) { + for (const pattern of P2_PRIORITY_PATTERNS) { if (pattern.test(upper)) { return 'P2'; } @@ -2281,12 +2302,17 @@ export class RalphTracker extends EventEmitter { return; } - // Add new todo - if (this._todos.size >= MAX_TODOS_PER_SESSION) { - // Remove oldest todo to make room + // Add new todo with guaranteed eviction if at capacity + // Use while loop to ensure we always have room (handles edge case where findOldestTodo returns undefined) + while (this._todos.size >= MAX_TODOS_PER_SESSION) { const oldest = this.findOldestTodo(); if (oldest) { this._todos.delete(oldest.id); + } else { + // Safety valve: if somehow no oldest found, clear a random entry + const firstKey = this._todos.keys().next().value; + if (firstKey) this._todos.delete(firstKey); + else break; // Map is empty somehow, exit loop } } @@ -2985,7 +3011,7 @@ export class RalphTracker extends EventEmitter { const tasksMatch = trimmedLine.match(RALPH_TASKS_COMPLETED_PATTERN); if (tasksMatch) { const value = parseInt(tasksMatch[1], 10); - if (!isNaN(value) && value >= 0) { + if (!Number.isNaN(value) && value >= 0) { block.tasksCompletedThisLoop = value; } else { parseErrors.push(`Invalid TASKS_COMPLETED_THIS_LOOP value: "${tasksMatch[1]}". Expected: non-negative integer`); @@ -2997,7 +3023,7 @@ export class RalphTracker extends EventEmitter { const filesMatch = trimmedLine.match(RALPH_FILES_MODIFIED_PATTERN); if (filesMatch) { const value = parseInt(filesMatch[1], 10); - if (!isNaN(value) && value >= 0) { + if (!Number.isNaN(value) && value >= 0) { block.filesModified = value; } else { parseErrors.push(`Invalid FILES_MODIFIED value: "${filesMatch[1]}". Expected: non-negative integer`); diff --git a/src/respawn-controller.ts b/src/respawn-controller.ts index 14b6aaa7..20bfbebe 100644 --- a/src/respawn-controller.ts +++ b/src/respawn-controller.ts @@ -43,6 +43,7 @@ import { BufferAccumulator } from './utils/buffer-accumulator.js'; import { ANSI_ESCAPE_PATTERN_SIMPLE, TOKEN_PATTERN, + assertNever, } from './utils/index.js'; import { MAX_RESPAWN_BUFFER_SIZE, @@ -1323,9 +1324,14 @@ export class RespawnController extends EventEmitter { * * If in 'watching' state, immediately checks for idle condition. * Otherwise, continues from current state. + * Re-setups terminal listener if it was removed (e.g., after stop()). */ resume(): void { this.log('Resuming respawn'); + // Re-setup terminal listener if it was removed (robustness for stop->resume case) + if (!this.terminalHandler) { + this.setupTerminalListener(); + } if (this._state === 'watching') { this.checkIdleAndMaybeStart(); } @@ -1410,6 +1416,7 @@ export class RespawnController extends EventEmitter { // In waiting states, also use confirmation timer (same detection logic) // This ensures we wait for Claude to finish before proceeding + // Note: 'watching' is already handled above and returns early switch (this._state) { case 'waiting_update': this.startStepConfirmTimer('update'); @@ -1423,6 +1430,19 @@ export class RespawnController extends EventEmitter { case 'waiting_kickstart': this.startStepConfirmTimer('kickstart'); break; + // Non-waiting states: completion message is ignored + case 'confirming_idle': + case 'ai_checking': + case 'sending_update': + case 'sending_clear': + case 'sending_init': + case 'monitoring_init': + case 'sending_kickstart': + case 'stopped': + // Completion message during these states is ignored + break; + default: + assertNever(this._state, `Unhandled RespawnState in completion detection: ${this._state}`); } return; } @@ -1512,6 +1532,19 @@ export class RespawnController extends EventEmitter { case 'waiting_kickstart': this.startStepConfirmTimer('kickstart'); break; + // Non-waiting states: prompt detection is informational only + case 'watching': + case 'confirming_idle': + case 'ai_checking': + case 'sending_update': + case 'sending_clear': + case 'sending_init': + case 'sending_kickstart': + case 'stopped': + // Prompt detection during these states doesn't trigger action + break; + default: + assertNever(this._state, `Unhandled RespawnState in prompt detection: ${this._state}`); } } } @@ -3243,6 +3276,12 @@ export class RespawnController extends EventEmitter { case 'error': agg.errorCycles++; break; + case 'cancelled': + // Cancelled cycles don't count towards any specific category + // but are still counted in totalCycles + break; + default: + assertNever(metrics.outcome, `Unhandled CycleOutcome: ${metrics.outcome}`); } // Recalculate averages using all recent metrics diff --git a/src/screen-manager.ts b/src/screen-manager.ts index 1e1b5fe7..275302fc 100644 --- a/src/screen-manager.ts +++ b/src/screen-manager.ts @@ -136,6 +136,10 @@ function isValidPath(path: string): boolean { path.includes('"') || path.includes('\n') || path.includes('\r')) { return false; } + // Check for path traversal attempts (prevents escaping workingDir) + if (path.includes('..')) { + return false; + } return SAFE_PATH_PATTERN.test(path); } @@ -268,8 +272,9 @@ export class ScreenManager extends EventEmitter { // Base command for the mode (just the executable, env vars are exported separately) // For Claude sessions: pass --session-id to use the SAME ID as the Claudeman session // This ensures subagents (discovered via file path) can be directly matched to the correct tab + // Note: sessionId is quoted for safety even though it should be a valid UUID const baseCmd = mode === 'claude' - ? `claude --dangerously-skip-permissions --session-id ${sessionId}` + ? `claude --dangerously-skip-permissions --session-id "${sessionId}"` : '$SHELL'; // Apply nice priority if configured @@ -355,7 +360,7 @@ export class ScreenManager extends EventEmitter { timeout: EXEC_TIMEOUT_MS }).trim(); if (output) { - for (const childPid of output.split('\n').map(p => parseInt(p, 10)).filter(p => !isNaN(p))) { + for (const childPid of output.split('\n').map(p => parseInt(p, 10)).filter(p => !Number.isNaN(p))) { pids.push(childPid); // Recursively get grandchildren pids.push(...this.getChildPids(childPid)); @@ -381,7 +386,8 @@ export class ScreenManager extends EventEmitter { // Verify all PIDs are dead, with retry private async verifyProcessesDead(pids: number[], maxWaitMs: number = 1000): Promise { const startTime = Date.now(); - const checkInterval = 50; + // 100ms interval balances responsiveness with CPU usage (was 50ms, too aggressive) + const checkInterval = 100; while (Date.now() - startTime < maxWaitMs) { const aliveCount = pids.filter(pid => this.isProcessAlive(pid)).length; @@ -463,7 +469,7 @@ export class ScreenManager extends EventEmitter { // Strategy 3: Kill screen session by name try { - execSync(`screen -S ${screen.screenName} -X quit`, { + execSync(`screen -S "${screen.screenName}" -X quit`, { timeout: EXEC_TIMEOUT_MS }); } catch { @@ -652,11 +658,11 @@ export class ScreenManager extends EventEmitter { for (const line of pgrepOutput.split('\n')) { const [pidStr, childrenStr] = line.split(':'); const screenPid = parseInt(pidStr, 10); - if (!isNaN(screenPid)) { + if (!Number.isNaN(screenPid)) { const children = (childrenStr || '') .split(',') .map(s => parseInt(s.trim(), 10)) - .filter(n => !isNaN(n) && n > 0); + .filter(n => !Number.isNaN(n) && n > 0); descendantMap.set(screenPid, children); } } @@ -685,7 +691,7 @@ export class ScreenManager extends EventEmitter { const pid = parseInt(parts[0], 10); const rss = parseFloat(parts[1]) || 0; const cpu = parseFloat(parts[2]) || 0; - if (!isNaN(pid)) { + if (!Number.isNaN(pid)) { processStats.set(pid, { rss, cpu }); } } @@ -848,21 +854,24 @@ export class ScreenManager extends EventEmitter { // Send text first (if any) if (escapedText) { - const textCmd = `screen -S ${screen.screenName} -p 0 -X stuff "${escapedText}"`; + const textCmd = `screen -S "${screen.screenName}" -p 0 -X stuff "${escapedText}"`; execSync(textCmd, { encoding: 'utf-8', timeout: EXEC_TIMEOUT_MS }); } // Send carriage return separately (Enter key for Ink) - // Use a synchronous sleep to ensure screen processes the text first + // Use synchronous delay to ensure screen processes the text first if (hasCarriageReturn) { // Delay to let screen process the text before sending Enter // This prevents race conditions where Enter arrives before the text is processed // 100ms is needed for reliability - screen's internal buffering can be slow if (escapedText) { - execSync('sleep 0.1', { timeout: 1000 }); + // Use Atomics.wait for synchronous delay instead of shell command + const sharedBuffer = new SharedArrayBuffer(4); + const sharedArray = new Int32Array(sharedBuffer); + Atomics.wait(sharedArray, 0, 0, 100); } - const crCmd = `screen -S ${screen.screenName} -p 0 -X stuff "$(printf '\\015')"`; + const crCmd = `screen -S "${screen.screenName}" -p 0 -X stuff "$(printf '\\015')"`; // Try up to CR_MAX_ATTEMPTS times with increasing delays let success = false; @@ -873,7 +882,11 @@ export class ScreenManager extends EventEmitter { } catch (crErr) { console.warn(`[ScreenManager] Carriage return attempt ${attempt}/${CR_MAX_ATTEMPTS} failed`); if (attempt < CR_MAX_ATTEMPTS) { - execSync(`sleep 0.${attempt}`, { timeout: 1000 }); // 0.1s, 0.2s delays + // Use synchronous sleep without shell command (security: avoid shell injection vectors) + const sleepMs = attempt * 100; // 100ms, 200ms, 300ms delays + const sharedBuffer = new SharedArrayBuffer(4); + const sharedArray = new Int32Array(sharedBuffer); + Atomics.wait(sharedArray, 0, 0, sleepMs); } } } diff --git a/src/session.ts b/src/session.ts index 92b5ff99..41d6f7e7 100644 --- a/src/session.ts +++ b/src/session.ts @@ -993,6 +993,11 @@ export class Session extends EventEmitter { .replace(CTRL_L_PATTERN, ''); // Remove Ctrl+L if (!data) return; // Skip if only filtered sequences + // Strip ANSI once for all downstream pattern-matching operations + // This avoids redundant regex operations in parseTokensFromStatusLine, + // parseClaudeCodeInfo, parseTaskDescriptionsFromTerminalData, and working detection + const cleanData = data.replace(ANSI_ESCAPE_PATTERN_FULL, ''); + // BufferAccumulator handles auto-trimming when max size exceeded this._terminalBuffer.append(data); this._lastActivityAt = Date.now(); @@ -1000,21 +1005,21 @@ export class Session extends EventEmitter { this.emit('terminal', data); this.emit('output', data); - // Forward to Ralph tracker to detect Ralph loops and todos + // Forward to Ralph tracker to detect Ralph loops and todos (handles own stripping) this._ralphTracker.processTerminalData(data); - // Forward to Bash tool parser to detect file-viewing commands + // Forward to Bash tool parser to detect file-viewing commands (handles own stripping) this._bashToolParser.processTerminalData(data); // Parse token count from status line (e.g., "123.4k tokens" or "5234 tokens") - this.parseTokensFromStatusLine(data); + this.parseTokensFromStatusLine(cleanData); // Parse Claude Code CLI info (version, model, account type) from startup - this.parseClaudeCodeInfo(data); + this.parseClaudeCodeInfo(cleanData); // Parse task descriptions from terminal output (e.g., "Explore(Check files)") // This enables correlating subagent windows with their short descriptions - this.parseTaskDescriptionsFromTerminalData(data); + this.parseTaskDescriptionsFromTerminalData(cleanData); // Detect if Claude is working or at prompt // The prompt line contains "❯" when waiting for input @@ -1042,14 +1047,13 @@ export class Session extends EventEmitter { } // Detect when Claude starts working (thinking, writing, etc) - // Strip ANSI/OSC sequences to avoid false positives from window titles like "3 File Reading Task" - const cleanDataForWorkingCheck = data.replace(ANSI_ESCAPE_PATTERN_FULL, ''); - if (cleanDataForWorkingCheck.includes('Thinking') || cleanDataForWorkingCheck.includes('Writing') || - cleanDataForWorkingCheck.includes('Reading') || cleanDataForWorkingCheck.includes('Running') || - cleanDataForWorkingCheck.includes('⠋') || cleanDataForWorkingCheck.includes('⠙') || - cleanDataForWorkingCheck.includes('⠹') || cleanDataForWorkingCheck.includes('⠸') || - cleanDataForWorkingCheck.includes('⠼') || cleanDataForWorkingCheck.includes('⠴') || - cleanDataForWorkingCheck.includes('⠦') || cleanDataForWorkingCheck.includes('⠧')) { + // Using pre-cleaned data avoids false positives from window titles like "3 File Reading Task" + if (cleanData.includes('Thinking') || cleanData.includes('Writing') || + cleanData.includes('Reading') || cleanData.includes('Running') || + cleanData.includes('⠋') || cleanData.includes('⠙') || + cleanData.includes('⠹') || cleanData.includes('⠸') || + cleanData.includes('⠼') || cleanData.includes('⠴') || + cleanData.includes('⠦') || cleanData.includes('⠧')) { if (!this._isWorking) { this._isWorking = true; this._status = 'busy'; @@ -1504,31 +1508,43 @@ export class Session extends EventEmitter { } /** - * Parse task descriptions from raw terminal data (may contain multiple lines). - * Called from interactive mode's onData handler. + * Parse task descriptions from terminal data (may contain multiple lines). + * Called from interactive mode's onData handler with ANSI-stripped data. + * @param cleanData - Terminal data with ANSI codes already stripped */ - private parseTaskDescriptionsFromTerminalData(data: string): void { + private parseTaskDescriptionsFromTerminalData(cleanData: string): void { // Quick pre-check: skip if no parentheses present - if (!data.includes('(') || !data.includes(')')) return; + if (!cleanData.includes('(') || !cleanData.includes(')')) return; - // Split by newlines and process each line - const lines = data.split(NEWLINE_SPLIT_PATTERN); + // Split by newlines and process each line (data already ANSI-stripped) + const lines = cleanData.split(NEWLINE_SPLIT_PATTERN); for (const line of lines) { - this.parseTaskDescriptionsFromLine(line); + this.parseTaskDescriptionsDirect(line); } } /** - * Parse task descriptions from terminal output. + * Parse task descriptions from terminal output line. * Claude Code outputs Task tool calls as "ToolName(Description)" in the terminal. * We capture these descriptions to use as window titles for subagents. + * Called from processOutput() with potentially non-cleaned data. */ private parseTaskDescriptionsFromLine(line: string): void { // Quick pre-check: skip expensive regex if no common tool patterns present if (!line.includes('(') || !line.includes(')')) return; - // Strip ANSI codes before matching - terminal output has embedded codes like [1mExplore[0m + // Strip ANSI codes - may still be present from processOutput() path const cleanLine = line.replace(ANSI_ESCAPE_PATTERN_FULL, ''); + this.parseTaskDescriptionsDirect(cleanLine); + } + + /** + * Parse task descriptions from a pre-cleaned line (no ANSI codes). + * Internal method used by both parseTaskDescriptionsFromTerminalData and parseTaskDescriptionsFromLine. + */ + private parseTaskDescriptionsDirect(cleanLine: string): void { + // Quick pre-check: skip expensive regex if no common tool patterns present + if (!cleanLine.includes('(') || !cleanLine.includes(')')) return; // Reset regex lastIndex for global pattern TASK_TOOL_PATTERN.lastIndex = 0; @@ -1611,15 +1627,13 @@ export class Session extends EventEmitter { // - Max tokens per session: 500k (Claude's context is ~200k) // - Max delta per update: 100k (prevents sudden jumps from parsing errors) // - Rejects "M" suffix values > 0.5 (500k) to prevent false matches - private parseTokensFromStatusLine(data: string): void { + private parseTokensFromStatusLine(cleanData: string): void { // Quick pre-check: skip expensive regex if "token" not present (performance optimization) - if (!data.includes('token')) return; - - // Remove ANSI escape codes for cleaner parsing (use pre-compiled pattern) - const cleanData = data.replace(ANSI_ESCAPE_PATTERN_FULL, ''); + if (!cleanData.includes('token')) return; // Match patterns: "123.4k tokens", "5234 tokens", "1.2M tokens" // The status line typically shows total tokens like "1.2k tokens" near the prompt + // Note: ANSI codes are already stripped by caller for performance const tokenMatch = cleanData.match(TOKEN_PATTERN); if (tokenMatch) { @@ -1672,16 +1686,15 @@ export class Session extends EventEmitter { // Parse Claude Code CLI info from terminal startup output // Extracts version, model, and account type for display in Claudeman UI - private parseClaudeCodeInfo(data: string): void { + // Note: Expects cleanData with ANSI codes already stripped by caller + private parseClaudeCodeInfo(cleanData: string): void { // Only parse once per session (during startup) if (this._cliInfoParsed) return; // Quick pre-checks - if (!data.includes('Claude') && !data.includes('current:') && !data.includes('Opus') && !data.includes('Sonnet')) { + if (!cleanData.includes('Claude') && !cleanData.includes('current:') && !cleanData.includes('Opus') && !cleanData.includes('Sonnet')) { return; } - - const cleanData = data.replace(ANSI_ESCAPE_PATTERN_FULL, ''); let changed = false; // Match "Claude Code v2.1.27" or "Claude Code vX.Y.Z" @@ -1921,7 +1934,8 @@ export class Session extends EventEmitter { this._status = 'busy'; this._lastActivityAt = Date.now(); this.runPrompt(input).catch(err => { - this.emit('error', err.message); + const errorMsg = err instanceof Error ? err.message : String(err); + this.emit('error', errorMsg); }); } diff --git a/src/state-store.ts b/src/state-store.ts index c97631ae..a3e8e0c0 100644 --- a/src/state-store.ts +++ b/src/state-store.ts @@ -74,7 +74,9 @@ export class StateStore { private ensureDir(): void { const dir = dirname(this.filePath); if (!existsSync(dir)) { - mkdirSync(dir, { recursive: true }); + // Use restrictive permissions (0o700) - owner only can read/write/traverse + // State files may contain sensitive session data + mkdirSync(dir, { recursive: true, mode: 0o700 }); } } @@ -584,7 +586,7 @@ export class StateStore { if (!this.ralphStateDirty) { return; } - this.ralphStateDirty = false; + // Clear dirty flag only on success to enable retry on failure this.ensureDir(); const data = Object.fromEntries(this.ralphStates); // Atomic write: write to temp file, then rename (atomic on POSIX) @@ -594,13 +596,17 @@ export class StateStore { json = JSON.stringify(data, null, 2); } catch (err) { console.error('[StateStore] Failed to serialize Ralph state (circular reference or invalid data):', err); - throw err; + // Keep dirty flag true for retry - don't throw, let caller continue + return; } try { writeFileSync(tempPath, json, 'utf-8'); renameSync(tempPath, this.ralphStatePath); + // Success - clear dirty flag + this.ralphStateDirty = false; } catch (err) { console.error('[StateStore] Failed to write Ralph state file:', err); + // Keep dirty flag true for retry on next save // Try to clean up temp file on error try { if (existsSync(tempPath)) { @@ -609,7 +615,7 @@ export class StateStore { } catch (cleanupErr) { console.warn('[StateStore] Failed to cleanup temp file during Ralph state save error:', cleanupErr); } - throw err; + // Don't throw - let caller continue, retry on next save } } @@ -655,8 +661,28 @@ export class StateStore { /** Flushes all pending saves (main and inner state). Call before shutdown. */ flushAll(): void { - this.saveNow(); - this.saveRalphStatesNow(); + // Save both states, catching errors to ensure both are attempted + let mainError: unknown = null; + let ralphError: unknown = null; + + try { + this.saveNow(); + } catch (err) { + mainError = err; + console.error('[StateStore] Error flushing main state:', err); + } + + try { + this.saveRalphStatesNow(); + } catch (err) { + ralphError = err; + console.error('[StateStore] Error flushing Ralph state:', err); + } + + // Log summary if any errors occurred + if (mainError || ralphError) { + console.warn('[StateStore] flushAll completed with errors - some state may not be persisted'); + } } } diff --git a/src/subagent-watcher.ts b/src/subagent-watcher.ts index 62aad4ee..7a4870e8 100644 --- a/src/subagent-watcher.ts +++ b/src/subagent-watcher.ts @@ -12,7 +12,7 @@ import { createInterface } from 'node:readline'; import { homedir } from 'node:os'; import { join, basename } from 'node:path'; import { execSync } from 'node:child_process'; -import { PENDING_TOOL_CALL_TTL_MS } from './config/map-limits.js'; +import { PENDING_TOOL_CALL_TTL_MS, MAX_PENDING_TOOL_CALLS } from './config/map-limits.js'; // ========== Types ========== @@ -165,6 +165,11 @@ export class SubagentWatcher extends EventEmitter { // Map of agentId -> Map of toolUseId -> { toolName, timestamp } (for linking tool_result to tool_call) // Includes timestamp for TTL-based cleanup of orphaned entries private pendingToolCalls = new Map>(); + // Guard to prevent concurrent liveness checks (prevents duplicate completed events) + private _isCheckingLiveness = false; + // Store error handlers for FSWatchers to enable proper cleanup (prevent memory leaks) + private dirWatcherErrorHandlers = new Map void>(); + private fileWatcherErrorHandlers = new Map void>(); constructor() { super(); @@ -222,19 +227,29 @@ export class SubagentWatcher extends EventEmitter { if (this.livenessInterval) return; this.livenessInterval = setInterval(async () => { - for (const [agentId, info] of this.agentInfo) { - if (info.status === 'active' || info.status === 'idle') { - const alive = await this.checkSubagentAlive(agentId); - if (!alive) { - info.status = 'completed'; - // Clean up pendingToolCalls for this agent to prevent memory leak - this.pendingToolCalls.delete(agentId); - this.emit('subagent:completed', info); + // Guard: prevent concurrent liveness checks (avoids duplicate completed events) + if (this._isCheckingLiveness) return; + this._isCheckingLiveness = true; + + try { + for (const [agentId, info] of this.agentInfo) { + // Re-check status in case another check completed this agent + if (info.status === 'active' || info.status === 'idle') { + const alive = await this.checkSubagentAlive(agentId); + if (!alive && (info.status === 'active' || info.status === 'idle')) { + // Double-check status after async call to prevent race + info.status = 'completed'; + // Clean up pendingToolCalls for this agent to prevent memory leak + this.pendingToolCalls.delete(agentId); + this.emit('subagent:completed', info); + } } } + // Periodically clean up stale completed agents (older than 24 hours) + this.cleanupStaleAgents(); + } finally { + this._isCheckingLiveness = false; } - // Periodically clean up stale completed agents (older than 24 hours) - this.cleanupStaleAgents(); }, LIVENESS_CHECK_MS); } @@ -288,11 +303,25 @@ export class SubagentWatcher extends EventEmitter { this.livenessInterval = null; } + // Remove error handlers before closing watchers to prevent memory leak + for (const [filePath, handler] of this.fileWatcherErrorHandlers) { + const watcher = this.fileWatchers.get(filePath); + if (watcher) watcher.off('error', handler); + } + this.fileWatcherErrorHandlers.clear(); + for (const watcher of this.fileWatchers.values()) { watcher.close(); } this.fileWatchers.clear(); + // Remove error handlers before closing watchers to prevent memory leak + for (const [dir, handler] of this.dirWatcherErrorHandlers) { + const watcher = this.dirWatchers.get(dir); + if (watcher) watcher.off('error', handler); + } + this.dirWatcherErrorHandlers.clear(); + for (const watcher of this.dirWatchers.values()) { watcher.close(); } @@ -560,7 +589,7 @@ export class SubagentWatcher extends EventEmitter { for (const pidStr of pids) { const pid = parseInt(pidStr, 10); - if (isNaN(pid)) continue; + if (Number.isNaN(pid)) continue; try { // Check /proc/{pid}/environ for session ID @@ -899,11 +928,15 @@ export class SubagentWatcher extends EventEmitter { }); // Handle watcher errors to prevent unhandled exceptions - watcher.on('error', (error) => { + // Store handler reference for proper cleanup + const errorHandler = (error: Error) => { this.emit('subagent:error', error instanceof Error ? error : new Error(String(error))); + this.dirWatcherErrorHandlers.delete(dir); this.dirWatchers.delete(dir); this.knownSubagentDirs.delete(dir); - }); + }; + watcher.on('error', errorHandler); + this.dirWatcherErrorHandlers.set(dir, errorHandler); this.dirWatchers.set(dir, watcher); } catch { @@ -974,6 +1007,9 @@ export class SubagentWatcher extends EventEmitter { // Read existing content this.tailFile(filePath, agentId, sessionId, 0).then((position) => { this.filePositions.set(filePath, position); + }).catch((err) => { + // Log but don't throw - non-critical background operation + console.warn(`[SubagentWatcher] Failed to read initial content for ${agentId}:`, err); }); // Watch for changes @@ -1026,10 +1062,14 @@ export class SubagentWatcher extends EventEmitter { }); // Handle watcher errors to prevent unhandled exceptions - watcher.on('error', (error) => { + // Store handler reference for proper cleanup + const errorHandler = (error: Error) => { this.emit('subagent:error', error instanceof Error ? error : new Error(String(error)), agentId); + this.fileWatcherErrorHandlers.delete(filePath); this.fileWatchers.delete(filePath); - }); + }; + watcher.on('error', errorHandler); + this.fileWatcherErrorHandlers.set(filePath, errorHandler); this.fileWatchers.set(filePath, watcher); this.resetIdleTimer(agentId); @@ -1175,7 +1215,21 @@ export class SubagentWatcher extends EventEmitter { if (!this.pendingToolCalls.has(agentId)) { this.pendingToolCalls.set(agentId, new Map()); } - this.pendingToolCalls.get(agentId)!.set(content.id, { + const agentCalls = this.pendingToolCalls.get(agentId)!; + // Enforce size limit to prevent memory leak from rapid tool calls + if (agentCalls.size >= MAX_PENDING_TOOL_CALLS) { + // Evict oldest entry by timestamp + let oldestId: string | null = null; + let oldestTime = Infinity; + for (const [id, call] of agentCalls) { + if (call.timestamp < oldestTime) { + oldestTime = call.timestamp; + oldestId = id; + } + } + if (oldestId) agentCalls.delete(oldestId); + } + agentCalls.set(content.id, { toolName: content.name, timestamp: Date.now(), }); diff --git a/src/task-tracker.ts b/src/task-tracker.ts index 3f851519..33e1b48c 100644 --- a/src/task-tracker.ts +++ b/src/task-tracker.ts @@ -21,6 +21,7 @@ */ import { EventEmitter } from 'node:events'; +import { assertNever } from './utils/index.js'; // ========== Configuration Constants ========== @@ -68,6 +69,54 @@ const COMPLETE_PATTERNS = [ // ========== Type Definitions ========== +/** + * Content block for tool_use messages from Claude. + * Emitted when Claude invokes a tool. + */ +interface ClaudeToolUseBlock { + type: 'tool_use'; + /** Unique identifier for this tool invocation */ + id: string; + /** Name of the tool being invoked */ + name: string; + /** Parameters passed to the tool */ + input?: { + description?: string; + prompt?: string; + subagent_type?: string; + [key: string]: unknown; + }; +} + +/** + * Content block for tool_result messages from Claude. + * Emitted when a tool completes execution. + */ +interface ClaudeToolResultBlock { + type: 'tool_result'; + /** ID of the tool_use this result corresponds to */ + tool_use_id: string; + /** Whether the tool execution resulted in an error */ + is_error?: boolean; + /** Result content (string or structured data) */ + content?: string | unknown; +} + +/** + * Union type for content blocks we care about. + */ +type ClaudeContentBlock = ClaudeToolUseBlock | ClaudeToolResultBlock | { type: string }; + +/** + * Claude JSON message structure. + * This is the format Claude Code outputs for streaming events. + */ +interface ClaudeMessage { + message?: { + content?: ClaudeContentBlock[]; + }; +} + /** * Represents a background task spawned by Claude Code. * @@ -189,14 +238,14 @@ export class TaskTracker extends EventEmitter { * @fires taskCompleted - When a task finishes successfully * @fires taskFailed - When a task finishes with error */ - processMessage(msg: any): void { + processMessage(msg: ClaudeMessage | null | undefined): void { if (!msg || !msg.message?.content) return; for (const block of msg.message.content) { - if (block.type === 'tool_use' && block.name === 'Task') { - this.handleTaskToolUse(block); + if (block.type === 'tool_use' && (block as ClaudeToolUseBlock).name === 'Task') { + this.handleTaskToolUse(block as ClaudeToolUseBlock); } else if (block.type === 'tool_result') { - this.handleToolResult(block); + this.handleToolResult(block as ClaudeToolResultBlock); } } } @@ -256,7 +305,7 @@ export class TaskTracker extends EventEmitter { * @param block - The tool_use content block * @fires taskCreated */ - private handleTaskToolUse(block: any): void { + private handleTaskToolUse(block: ClaudeToolUseBlock): void { const toolUseId = block.id; const params = block.input || {}; @@ -305,7 +354,7 @@ export class TaskTracker extends EventEmitter { * @fires taskCompleted - If result is success * @fires taskFailed - If result is error */ - private handleToolResult(block: any): void { + private handleToolResult(block: ClaudeToolResultBlock): void { const toolUseId = block.tool_use_id; const task = this.tasks.get(toolUseId); @@ -416,7 +465,7 @@ export class TaskTracker extends EventEmitter { * @fires taskCreated */ private createTaskFromTerminal(agentType: string, _context: string): void { - const taskId = `terminal-${Date.now()}-${Math.random().toString(36).substr(2, 9)}`; + const taskId = `terminal-${Date.now()}-${Math.random().toString(36).slice(2, 11)}`; const parentId = this.taskStack.length > 0 ? this.taskStack[this.taskStack.length - 1] : null; const task: BackgroundTask = { @@ -547,6 +596,8 @@ export class TaskTracker extends EventEmitter { case 'running': running++; break; case 'completed': completed++; break; case 'failed': failed++; break; + default: + assertNever(task.status, `Unhandled BackgroundTask status: ${task.status}`); } } diff --git a/src/transcript-watcher.ts b/src/transcript-watcher.ts index 43c0b6e1..bc94fd9a 100644 --- a/src/transcript-watcher.ts +++ b/src/transcript-watcher.ts @@ -230,6 +230,17 @@ export class TranscriptWatcher extends EventEmitter { } }); + // Add error handler to prevent unhandled errors and fall back to polling + this.fileWatcher.on('error', (err) => { + this.emit('transcript:error', err as Error); + this.fileWatcher?.close(); + this.fileWatcher = null; + // Fall back to polling on error + if (this._isRunning) { + this.startPolling(); + } + }); + // Initial read this.processNewContent(); } catch (err) { diff --git a/src/types.ts b/src/types.ts index 0618161a..f9e115c5 100644 --- a/src/types.ts +++ b/src/types.ts @@ -763,19 +763,12 @@ export interface HookEventRequest { // ========== API Response Types ========== /** - * Standard API response wrapper + * Standard API response wrapper (discriminated union for type safety) * @template T Type of the data payload */ -export interface ApiResponse { - /** Whether the request succeeded */ - success: boolean; - /** Error message if failed */ - error?: string; - /** Error code for programmatic handling */ - errorCode?: ApiErrorCode; - /** Response data payload */ - data?: T; -} +export type ApiResponse = + | { success: true; data?: T } + | { success: false; error: string; errorCode: ApiErrorCode }; /** * Creates a standardized error response @@ -783,7 +776,7 @@ export interface ApiResponse { * @param details Optional detailed error message * @returns Formatted error response */ -export function createErrorResponse(code: ApiErrorCode, details?: string): ApiResponse { +export function createErrorResponse(code: ApiErrorCode, details?: string): ApiResponse { return { success: false, error: details || ErrorMessages[code], diff --git a/src/utils/index.ts b/src/utils/index.ts index c4a31700..4ac32420 100644 --- a/src/utils/index.ts +++ b/src/utils/index.ts @@ -33,3 +33,4 @@ export { fuzzyPhraseMatch, todoContentHash, } from './string-similarity.js'; +export { assertNever } from './type-safety.js'; diff --git a/src/utils/lru-map.ts b/src/utils/lru-map.ts index 3c061d69..e7b44b69 100644 --- a/src/utils/lru-map.ts +++ b/src/utils/lru-map.ts @@ -199,10 +199,15 @@ export class LRUMap extends Map { const now = Date.now(); const cutoff = now - maxAge; let evicted = 0; + let deletedNewest = false; // Iterate from oldest to newest for (const [key, value] of super.entries()) { if (getTimestamp(value) < cutoff) { + // Track if we're deleting the newest key + if (key === this._newestKey) { + deletedNewest = true; + } super.delete(key); this.onEvict?.(key, value); evicted++; @@ -213,6 +218,14 @@ export class LRUMap extends Map { } } + // Update _newestKey if we deleted it (find new newest from remaining entries) + if (deletedNewest) { + this._newestKey = undefined; + for (const k of super.keys()) { + this._newestKey = k; // Last one becomes newest + } + } + return evicted; } diff --git a/src/utils/type-safety.ts b/src/utils/type-safety.ts new file mode 100644 index 00000000..eedeb1d8 --- /dev/null +++ b/src/utils/type-safety.ts @@ -0,0 +1,35 @@ +/** + * @fileoverview Type safety utilities for exhaustive checking. + * + * @module utils/type-safety + */ + +/** + * Assert that a value should never occur at runtime. + * Used in switch/case default blocks to ensure exhaustive handling. + * + * If TypeScript sees a code path where `value` is not `never`, + * it will report a compile-time error, catching missing cases. + * + * @param value - The value that should be of type `never` + * @param message - Optional custom error message + * @throws Error if called at runtime (indicates unhandled case) + * + * @example + * type Status = 'active' | 'inactive' | 'pending'; + * function handleStatus(status: Status): void { + * switch (status) { + * case 'active': // handle active + * break; + * case 'inactive': // handle inactive + * break; + * case 'pending': // handle pending + * break; + * default: + * assertNever(status); // Compile error if case is missing + * } + * } + */ +export function assertNever(value: never, message?: string): never { + throw new Error(message ?? `Unexpected value: ${JSON.stringify(value)}`); +} diff --git a/src/web/public/app.js b/src/web/public/app.js index bf3be20d..8a4c053a 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -1444,8 +1444,12 @@ class ClaudemanApp { }; window.addEventListener('resize', throttledResize); - const resizeObserver = new ResizeObserver(throttledResize); - resizeObserver.observe(container); + // Store resize observer for cleanup (prevents memory leak on terminal re-init) + if (this.terminalResizeObserver) { + this.terminalResizeObserver.disconnect(); + } + this.terminalResizeObserver = new ResizeObserver(throttledResize); + this.terminalResizeObserver.observe(container); // Handle keyboard input with batching for rapid keystrokes this._pendingInput = ''; @@ -1860,10 +1864,11 @@ class ClaudemanApp { } }, true); // Use capture phase to handle before terminal - // Token stats click handler + // Token stats click handler (with guard to prevent duplicate handlers on reconnect) const tokenEl = this.$('headerTokens'); - if (tokenEl) { + if (tokenEl && !tokenEl._statsHandlerAttached) { tokenEl.classList.add('clickable'); + tokenEl._statsHandlerAttached = true; tokenEl.addEventListener('click', () => this.openTokenStats()); } @@ -2913,6 +2918,18 @@ class ClaudemanApp { clearInterval(this.notificationManager.titleFlashInterval); this.notificationManager.titleFlashInterval = null; } + // Clear notification manager grouping timeouts (prevents orphaned timers) + if (this.notificationManager?.groupingMap) { + for (const { timeout } of this.notificationManager.groupingMap.values()) { + clearTimeout(timeout); + } + this.notificationManager.groupingMap.clear(); + } + // Disconnect terminal resize observer (prevents memory leak on reconnect) + if (this.terminalResizeObserver) { + this.terminalResizeObserver.disconnect(); + this.terminalResizeObserver = null; + } // Clear any other orphaned timers if (this.planLoadingTimer) { clearInterval(this.planLoadingTimer); diff --git a/src/web/schemas.ts b/src/web/schemas.ts new file mode 100644 index 00000000..05259d6b --- /dev/null +++ b/src/web/schemas.ts @@ -0,0 +1,141 @@ +/** + * @fileoverview Zod validation schemas for API routes + * + * This module contains Zod schemas for validating API request bodies. + * Schemas are used in src/web/server.ts route handlers. + * + * @module web/schemas + */ + +import { z } from 'zod'; + +// ========== Session Routes ========== + +/** + * Schema for POST /api/sessions + * Creates a new session with optional working directory, mode, and name. + */ +export const CreateSessionSchema = z.object({ + workingDir: z.string().optional(), + mode: z.enum(['claude', 'shell']).optional(), + name: z.string().max(100).optional(), +}); + +/** + * Schema for POST /api/sessions/:id/run + * Runs a prompt in a session. + */ +export const RunPromptSchema = z.object({ + prompt: z.string().min(1).max(100000), +}); + +/** + * Schema for POST /api/sessions/:id/input + * Sends input to an interactive session. + */ +export const SessionInputSchema = z.object({ + input: z.string(), + useScreen: z.boolean().optional(), +}); + +/** + * Schema for POST /api/sessions/:id/resize + * Resizes a session's terminal. + */ +export const ResizeSchema = z.object({ + cols: z.number().int().min(1).max(500), + rows: z.number().int().min(1).max(200), +}); + +// ========== Case Routes ========== + +/** + * Schema for POST /api/cases + * Creates a new case folder. + */ +export const CreateCaseSchema = z.object({ + name: z.string().regex(/^[a-zA-Z0-9_-]+$/, 'Invalid case name format. Use only letters, numbers, hyphens, underscores.'), + description: z.string().max(1000).optional(), +}); + +// ========== Quick Start ========== + +/** + * Schema for POST /api/quick-start + * Creates case (if needed) and starts interactive session. + */ +export const QuickStartSchema = z.object({ + caseName: z.string().regex(/^[a-zA-Z0-9_-]+$/, 'Invalid case name format. Use only letters, numbers, hyphens, underscores.').optional(), + mode: z.enum(['claude', 'shell']).optional(), +}); + +// ========== Hook Events ========== + +/** + * Schema for POST /api/hook-event + * Receives Claude Code hook events. + */ +export const HookEventSchema = z.object({ + event: z.enum(['permission_prompt', 'elicitation_dialog', 'idle_prompt', 'stop']), + sessionId: z.string().min(1), + data: z.unknown().optional(), +}); + +// ========== Configuration ========== + +/** + * Schema for respawn configuration (partial updates allowed) + * Used in PUT /api/config and respawn endpoints. + */ +export const RespawnConfigSchema = z.object({ + idleTimeoutMs: z.number().int().min(1000).max(600000).optional(), + updatePrompt: z.string().max(10000).optional(), + interStepDelayMs: z.number().int().min(100).max(60000).optional(), + enabled: z.boolean().optional(), + sendClear: z.boolean().optional(), + sendUpdate: z.boolean().optional(), + sendInit: z.boolean().optional(), + sendKickstart: z.boolean().optional(), + kickstartPrompt: z.string().max(10000).optional(), + aiIdleCheckEnabled: z.boolean().optional(), + aiIdleCheckTimeoutMs: z.number().int().min(10000).max(300000).optional(), + aiIdleCheckModel: z.string().max(100).optional(), + completionConfirmMs: z.number().int().min(1000).max(60000).optional(), + noOutputTimeoutMs: z.number().int().min(5000).max(600000).optional(), + maxIterations: z.number().int().min(0).max(10000).optional(), + stuckStateWarningMs: z.number().int().min(60000).max(3600000).optional(), + autoAcceptEnabled: z.boolean().optional(), + autoAcceptDelayMs: z.number().int().min(1000).max(60000).optional(), + planModeEnabled: z.boolean().optional(), + planCheckTimeoutMs: z.number().int().min(10000).max(300000).optional(), + planCheckModel: z.string().max(100).optional(), +}).strict(); + +/** + * Schema for PUT /api/config + * Updates application configuration with whitelist of allowed fields. + */ +export const ConfigUpdateSchema = z.object({ + pollIntervalMs: z.number().int().min(100).max(60000).optional(), + defaultTimeoutMs: z.number().int().min(1000).max(3600000).optional(), + maxConcurrentSessions: z.number().int().min(1).max(50).optional(), + respawn: RespawnConfigSchema.optional(), +}).strict(); + +/** + * Schema for PUT /api/settings + * User settings with allowed fields only. + */ +export const SettingsUpdateSchema = z.object({ + defaultClaudeMdPath: z.string().max(500).optional(), + lastUsedCase: z.string().max(200).optional(), + // Add other known settings fields as needed +}).passthrough(); // Allow additional fields but validate known ones + +/** + * Schema for POST /api/sessions/:id/input with length limit + */ +export const SessionInputWithLimitSchema = z.object({ + input: z.string().max(100000), // 100KB max input + useScreen: z.boolean().optional(), +}); diff --git a/src/web/server.ts b/src/web/server.ts index 67e7c7b7..cee651ea 100644 --- a/src/web/server.ts +++ b/src/web/server.ts @@ -12,7 +12,7 @@ import Fastify, { FastifyInstance, FastifyReply } from 'fastify'; import fastifyStatic from '@fastify/static'; -import path, { join, dirname, resolve } from 'node:path'; +import path, { join, dirname, resolve, relative, isAbsolute } from 'node:path'; import { fileURLToPath } from 'node:url'; import { existsSync, mkdirSync, writeFileSync, readdirSync, readFileSync, rmSync } from 'node:fs'; import fs from 'node:fs/promises'; @@ -42,15 +42,8 @@ import { getErrorMessage, ApiErrorCode, createErrorResponse, - type CreateSessionRequest, - type RunPromptRequest, - type SessionInputRequest, - type ResizeRequest, - type CreateCaseRequest, - type QuickStartRequest, type CreateScheduledRunRequest, type QuickRunRequest, - type HookEventRequest, type ApiResponse, type SessionResponse, type QuickStartResponse, @@ -60,6 +53,18 @@ import { type ImageDetectedEvent, DEFAULT_NICE_CONFIG, } from '../types.js'; +import { + CreateSessionSchema, + RunPromptSchema, + SessionInputSchema, + ResizeSchema, + CreateCaseSchema, + QuickStartSchema, + HookEventSchema, + ConfigUpdateSchema, + RespawnConfigSchema, +} from './schemas.js'; +import { StaleExpirationMap } from '../utils/index.js'; const __dirname = dirname(fileURLToPath(import.meta.url)); @@ -317,15 +322,20 @@ export class WebServer extends EventEmitter { // Terminal batching for performance private terminalBatches: Map = new Map(); private terminalBatchTimer: NodeJS.Timeout | null = null; - // Adaptive batching: track rapid events to extend batch window - private lastTerminalEventTime: Map = new Map(); - private adaptiveBatchInterval: number = TERMINAL_BATCH_INTERVAL; + // Adaptive batching: track rapid events to extend batch window (per-session) + // StaleExpirationMap auto-cleans entries for sessions that stop generating output + private lastTerminalEventTime: StaleExpirationMap = new StaleExpirationMap({ + ttlMs: 5 * 60 * 1000, // 5 minutes - auto-expire stale session timing data + refreshOnGet: false, // Don't refresh on reads, only on explicit sets + }); + // Per-session adaptive batch intervals (sessions with rapid output get longer batches) + private adaptiveBatchIntervals: Map = new Map(); // Scheduled runs cleanup timer private scheduledCleanupTimer: NodeJS.Timeout | null = null; // SSE event batching private outputBatches: Map = new Map(); private outputBatchTimer: NodeJS.Timeout | null = null; - private taskUpdateBatches: Map = new Map(); + private taskUpdateBatches: Map = new Map(); private taskUpdateBatchTimer: NodeJS.Timeout | null = null; // State update batching (reduce expensive toDetailedState() serialization) private stateUpdatePending: Set = new Set(); @@ -551,11 +561,12 @@ export class WebServer extends EventEmitter { }); this.app.put('/api/config', async (req) => { - const body = req.body as Record | null; - if (!body || typeof body !== 'object') { - return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Request body must be a JSON object'); + // Validate request body against schema to prevent arbitrary config injection + const parseResult = ConfigUpdateSchema.safeParse(req.body); + if (!parseResult.success) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, `Invalid config: ${parseResult.error.message}`); } - this.store.setConfig(body as Partial>); + this.store.setConfig(parseResult.data as Partial>); return { success: true, config: this.store.getConfig() }; }); @@ -635,10 +646,14 @@ export class WebServer extends EventEmitter { this.app.post('/api/sessions', async (req): Promise => { // Prevent unbounded session creation if (this.sessions.size >= MAX_CONCURRENT_SESSIONS) { - return { success: false, error: `Maximum concurrent sessions (${MAX_CONCURRENT_SESSIONS}) reached. Delete some sessions first.` }; + return createErrorResponse(ApiErrorCode.OPERATION_FAILED, `Maximum concurrent sessions (${MAX_CONCURRENT_SESSIONS}) reached. Delete some sessions first.`); } - const body = req.body as CreateSessionRequest & { mode?: 'claude' | 'shell'; name?: string }; + const result = CreateSessionSchema.safeParse(req.body); + if (!result.success) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, result.error.issues[0]?.message ?? 'Validation failed'); + } + const body = result.data; const workingDir = body.workingDir || process.cwd(); const globalNice = this.getGlobalNiceConfig(); const session = new Session({ @@ -705,7 +720,7 @@ export class WebServer extends EventEmitter { const killScreen = query.killScreen !== 'false'; // Default to true if (!this.sessions.has(id)) { - return { success: false, error: 'Session not found' }; + return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found'); } await this.cleanupSession(id, killScreen); @@ -769,7 +784,7 @@ export class WebServer extends EventEmitter { const session = this.sessions.get(id); if (!session) { - return { success: false, error: 'Session not found' }; + return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found'); } return { @@ -949,9 +964,10 @@ export class WebServer extends EventEmitter { return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Missing path parameter'); } - // Validate path is within working directory (security) + // Validate path is within working directory (security: proper path traversal check) const fullPath = resolve(session.workingDir, filePath); - if (!fullPath.startsWith(session.workingDir)) { + const relativePath = relative(session.workingDir, fullPath); + if (relativePath.startsWith('..') || isAbsolute(relativePath)) { return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Path must be within working directory'); } @@ -978,8 +994,15 @@ export class WebServer extends EventEmitter { }; } - // Read text file with line limit - const maxLines = parseInt(lines || '500', 10); + // Validate file size before reading (DoS protection - prevent memory exhaustion) + const MAX_TEXT_FILE_SIZE = 10 * 1024 * 1024; // 10MB + if (stat.size > MAX_TEXT_FILE_SIZE) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, `File too large (${Math.round(stat.size / 1024 / 1024)}MB > ${MAX_TEXT_FILE_SIZE / 1024 / 1024}MB limit)`); + } + + // Read text file with line limit (bounded to prevent DoS) + const MAX_LINES_LIMIT = 10000; + const maxLines = Math.min(parseInt(lines || '500', 10) || 500, MAX_LINES_LIMIT); const content = await fs.readFile(fullPath, 'utf-8'); const allLines = content.split('\n'); const truncatedContent = allLines.length > maxLines; @@ -1008,23 +1031,32 @@ export class WebServer extends EventEmitter { const session = this.sessions.get(id); if (!session) { - reply.code(404).send({ success: false, error: 'Session not found' }); + reply.code(404).send(createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found')); return; } if (!filePath) { - reply.code(400).send({ success: false, error: 'Missing path parameter' }); + reply.code(400).send(createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Missing path parameter')); return; } - // Validate path is within working directory + // Validate path is within working directory (security: proper path traversal check) const fullPath = resolve(session.workingDir, filePath); - if (!fullPath.startsWith(session.workingDir)) { - reply.code(400).send({ success: false, error: 'Path must be within working directory' }); + const relativePath = relative(session.workingDir, fullPath); + if (relativePath.startsWith('..') || isAbsolute(relativePath)) { + reply.code(400).send(createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Path must be within working directory')); return; } try { + // Validate file size before reading (DoS protection - prevent memory exhaustion) + const MAX_RAW_FILE_SIZE = 50 * 1024 * 1024; // 50MB for raw files + const stat = await fs.stat(fullPath); + if (stat.size > MAX_RAW_FILE_SIZE) { + reply.code(400).send(createErrorResponse(ApiErrorCode.INVALID_INPUT, `File too large (${Math.round(stat.size / 1024 / 1024)}MB > ${MAX_RAW_FILE_SIZE / 1024 / 1024}MB limit)`)); + return; + } + const ext = filePath.split('.').pop()?.toLowerCase() || ''; const mimeTypes: Record = { 'png': 'image/png', 'jpg': 'image/jpeg', 'jpeg': 'image/jpeg', 'gif': 'image/gif', @@ -1038,7 +1070,7 @@ export class WebServer extends EventEmitter { reply.header('Content-Type', mimeTypes[ext] || 'application/octet-stream'); reply.send(content); } catch (err) { - reply.code(500).send({ success: false, error: `Failed to read file: ${getErrorMessage(err)}` }); + reply.code(500).send(createErrorResponse(ApiErrorCode.OPERATION_FAILED, `Failed to read file: ${getErrorMessage(err)}`)); } }); @@ -1049,12 +1081,12 @@ export class WebServer extends EventEmitter { const session = this.sessions.get(id); if (!session) { - reply.code(404).send({ success: false, error: 'Session not found' }); + reply.code(404).send(createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found')); return; } if (!filePath) { - reply.code(400).send({ success: false, error: 'Missing path parameter' }); + reply.code(400).send(createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Missing path parameter')); return; } @@ -1376,7 +1408,11 @@ export class WebServer extends EventEmitter { // Run prompt in session this.app.post('/api/sessions/:id/run', async (req): Promise => { const { id } = req.params as { id: string }; - const { prompt } = req.body as RunPromptRequest; + const result = RunPromptSchema.safeParse(req.body); + if (!result.success) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, result.error.issues[0]?.message ?? 'Validation failed'); + } + const { prompt } = result.data; const session = this.sessions.get(id); if (!session) { @@ -1455,17 +1491,17 @@ export class WebServer extends EventEmitter { // useScreen: true uses writeViaScreen which is more reliable for programmatic input this.app.post('/api/sessions/:id/input', async (req): Promise => { const { id } = req.params as { id: string }; - const { input, useScreen } = req.body as SessionInputRequest & { useScreen?: boolean }; + const result = SessionInputSchema.safeParse(req.body); + if (!result.success) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, result.error.issues[0]?.message ?? 'Validation failed'); + } + const { input, useScreen } = result.data; const session = this.sessions.get(id); if (!session) { return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found'); } - if (input === undefined || input === null) { - return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Input is required'); - } - const inputStr = String(input); if (inputStr.length > MAX_INPUT_LENGTH) { return createErrorResponse(ApiErrorCode.INVALID_INPUT, `Input exceeds maximum length (${MAX_INPUT_LENGTH} bytes)`); @@ -1491,16 +1527,18 @@ export class WebServer extends EventEmitter { // Resize session terminal this.app.post('/api/sessions/:id/resize', async (req): Promise => { const { id } = req.params as { id: string }; - const { cols, rows } = req.body as ResizeRequest; + const result = ResizeSchema.safeParse(req.body); + if (!result.success) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, result.error.issues[0]?.message ?? 'Validation failed'); + } + const { cols, rows } = result.data; const session = this.sessions.get(id); if (!session) { return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found'); } - if (!Number.isInteger(cols) || !Number.isInteger(rows) || cols < 1 || rows < 1) { - return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'cols and rows must be positive integers'); - } + // Note: Zod already validates that cols and rows are positive integers within bounds if (cols > MAX_TERMINAL_COLS || rows > MAX_TERMINAL_ROWS) { return createErrorResponse(ApiErrorCode.INVALID_INPUT, `Terminal dimensions exceed maximum (${MAX_TERMINAL_COLS}x${MAX_TERMINAL_ROWS})`); } @@ -1665,7 +1703,12 @@ export class WebServer extends EventEmitter { // Update respawn configuration (works with or without running controller) this.app.put('/api/sessions/:id/respawn/config', async (req) => { const { id } = req.params as { id: string }; - const config = req.body as Partial; + // Validate respawn config to prevent arbitrary field injection + const parseResult = RespawnConfigSchema.safeParse(req.body); + if (!parseResult.success) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, `Invalid respawn config: ${parseResult.error.message}`); + } + const config = parseResult.data as Partial; const session = this.sessions.get(id); if (!session) { @@ -1993,7 +2036,7 @@ export class WebServer extends EventEmitter { const run = this.scheduledRuns.get(id); if (!run) { - return { success: false, error: 'Scheduled run not found' }; + return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Scheduled run not found'); } await this.stopScheduledRun(id); @@ -2047,31 +2090,33 @@ export class WebServer extends EventEmitter { } } } - } catch { - // Ignore errors reading linked cases + } catch (err) { + // Log but don't fail - linked cases are optional + console.warn('[Server] Failed to read linked cases:', err); } return cases; }); - this.app.post('/api/cases', async (req): Promise<{ success: boolean; case?: { name: string; path: string }; error?: string }> => { - const { name, description } = req.body as CreateCaseRequest; - - if (!name || !/^[a-zA-Z0-9_-]+$/.test(name)) { - return { success: false, error: 'Invalid case name. Use only letters, numbers, hyphens, underscores.' }; + this.app.post('/api/cases', async (req): Promise> => { + const result = CreateCaseSchema.safeParse(req.body); + if (!result.success) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, result.error.issues[0]?.message ?? 'Validation failed'); } + const { name, description } = result.data; const casePath = join(casesDir, name); - // Security: Path traversal protection - ensure resolved path is within casesDir + // Security: Path traversal protection - use relative path check const resolvedPath = resolve(casePath); const resolvedBase = resolve(casesDir); - if (!resolvedPath.startsWith(resolvedBase + '/') && resolvedPath !== resolvedBase) { - return { success: false, error: 'Invalid case path' }; + const relPath = relative(resolvedBase, resolvedPath); + if (relPath.startsWith('..') || isAbsolute(relPath)) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid case path'); } if (existsSync(casePath)) { - return { success: false, error: 'Case already exists' }; + return createErrorResponse(ApiErrorCode.ALREADY_EXISTS, 'Case already exists'); } try { @@ -2088,22 +2133,22 @@ export class WebServer extends EventEmitter { this.broadcast('case:created', { name, path: casePath }); - return { success: true, case: { name, path: casePath } }; + return { success: true, data: { case: { name, path: casePath } } }; } catch (err) { - return { success: false, error: getErrorMessage(err) }; + return createErrorResponse(ApiErrorCode.OPERATION_FAILED, getErrorMessage(err)); } }); // Link an existing folder as a case - this.app.post('/api/cases/link', async (req): Promise<{ success: boolean; case?: { name: string; path: string }; error?: string }> => { + this.app.post('/api/cases/link', async (req): Promise> => { const { name, path: folderPath } = req.body as { name: string; path: string }; if (!name || !/^[a-zA-Z0-9_-]+$/.test(name)) { - return { success: false, error: 'Invalid case name. Use only letters, numbers, hyphens, underscores.' }; + return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid case name. Use only letters, numbers, hyphens, underscores.'); } if (!folderPath) { - return { success: false, error: 'Folder path is required.' }; + return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Folder path is required.'); } // Expand ~ to home directory @@ -2113,13 +2158,13 @@ export class WebServer extends EventEmitter { // Validate the folder exists if (!existsSync(expandedPath)) { - return { success: false, error: `Folder not found: ${expandedPath}` }; + return createErrorResponse(ApiErrorCode.NOT_FOUND, `Folder not found: ${expandedPath}`); } // Check if case name already exists in casesDir const casePath = join(casesDir, name); if (existsSync(casePath)) { - return { success: false, error: 'A case with this name already exists in claudeman-cases.' }; + return createErrorResponse(ApiErrorCode.ALREADY_EXISTS, 'A case with this name already exists in claudeman-cases.'); } // Load existing linked cases @@ -2135,7 +2180,7 @@ export class WebServer extends EventEmitter { // Check if name is already linked if (linkedCases[name]) { - return { success: false, error: `Case "${name}" is already linked to ${linkedCases[name]}` }; + return createErrorResponse(ApiErrorCode.ALREADY_EXISTS, `Case "${name}" is already linked to ${linkedCases[name]}`); } // Save the linked case @@ -2147,9 +2192,9 @@ export class WebServer extends EventEmitter { } writeFileSync(linkedCasesFile, JSON.stringify(linkedCases, null, 2)); this.broadcast('case:linked', { name, path: expandedPath }); - return { success: true, case: { name, path: expandedPath } }; + return { success: true, data: { case: { name, path: expandedPath } } }; } catch (err) { - return { success: false, error: getErrorMessage(err) }; + return createErrorResponse(ApiErrorCode.OPERATION_FAILED, getErrorMessage(err)); } }); @@ -2276,12 +2321,14 @@ export class WebServer extends EventEmitter { } } - const stats = { - total: todos.length, - pending: todos.filter(t => t.status === 'pending').length, - inProgress: todos.filter(t => t.status === 'in_progress').length, - completed: todos.filter(t => t.status === 'completed').length, - }; + // Calculate stats in a single pass for better performance + let pending = 0, inProgress = 0, completed = 0; + for (const t of todos) { + if (t.status === 'pending') pending++; + else if (t.status === 'in_progress') inProgress++; + else if (t.status === 'completed') completed++; + } + const stats = { total: todos.length, pending, inProgress, completed }; return { success: true, @@ -2302,19 +2349,19 @@ export class WebServer extends EventEmitter { return { success: false, error: `Maximum concurrent sessions (${MAX_CONCURRENT_SESSIONS}) reached.` }; } - const { caseName = 'testcase', mode = 'claude' } = req.body as QuickStartRequest; - - // Validate case name - if (!/^[a-zA-Z0-9_-]+$/.test(caseName)) { - return { success: false, error: 'Invalid case name. Use only letters, numbers, hyphens, underscores.' }; + const result = QuickStartSchema.safeParse(req.body); + if (!result.success) { + return { success: false, error: result.error.issues[0]?.message ?? 'Validation failed' }; } + const { caseName = 'testcase', mode = 'claude' } = result.data; const casePath = join(casesDir, caseName); - // Security: Path traversal protection - ensure resolved path is within casesDir + // Security: Path traversal protection - use relative path check const resolvedPath = resolve(casePath); const resolvedBase = resolve(casesDir); - if (!resolvedPath.startsWith(resolvedBase + '/') && resolvedPath !== resolvedBase) { + const relPath = relative(resolvedBase, resolvedPath); + if (relPath.startsWith('..') || isAbsolute(relPath)) { return { success: false, error: 'Invalid case path' }; } @@ -2389,11 +2436,13 @@ export class WebServer extends EventEmitter { mkdirSync(dir, { recursive: true }); } // Use async write to avoid blocking event loop - fs.writeFile(settingsFilePath, JSON.stringify(settings, null, 2)).catch(() => { - // Non-critical, ignore settings save errors + fs.writeFile(settingsFilePath, JSON.stringify(settings, null, 2)).catch((err) => { + // Non-critical but log for debugging + console.warn('[Server] Failed to save settings (lastUsedCase):', err); }); - } catch { - // Non-critical, ignore settings save errors + } catch (err) { + // Non-critical but log for debugging + console.warn('[Server] Failed to prepare settings update:', err); } return { @@ -2626,10 +2675,11 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; if (caseName) { const casesDir = join(homedir(), 'claudeman-cases'); const casePath = join(casesDir, caseName); - // Security: Path traversal protection + // Security: Path traversal protection - use relative path check const resolvedCase = resolve(casePath); const resolvedBase = resolve(casesDir); - if (resolvedCase.startsWith(resolvedBase) && existsSync(casePath)) { + const relPath = relative(resolvedBase, resolvedCase); + if (!relPath.startsWith('..') && !isAbsolute(relPath) && existsSync(casePath)) { outputDir = join(casePath, 'ralph-wizard'); // Clear old ralph-wizard directory to ensure fresh prompts for each generation @@ -2748,10 +2798,11 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; const casesDir = join(homedir(), 'claudeman-cases'); const casePath = join(casesDir, caseName); - // Security: Path traversal protection + // Security: Path traversal protection - use relative path check const resolvedCase = resolve(casePath); const resolvedBase = resolve(casesDir); - if (!resolvedCase.startsWith(resolvedBase)) { + const relPath = relative(resolvedBase, resolvedCase); + if (relPath.startsWith('..') || isAbsolute(relPath)) { return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid case name'); } @@ -2800,10 +2851,11 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; reply.header('Pragma', 'no-cache'); reply.header('Expires', '0'); - // Security: Path traversal protection for case name + // Security: Path traversal protection for case name - use relative path check const resolvedCase = resolve(casePath); const resolvedBase = resolve(casesDir); - if (!resolvedCase.startsWith(resolvedBase)) { + const relPath = relative(resolvedBase, resolvedCase); + if (relPath.startsWith('..') || isAbsolute(relPath)) { return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid case name'); } @@ -2827,13 +2879,23 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; const content = readFileSync(fullPath, 'utf-8'); const isJson = filePath.endsWith('.json'); + // Parse JSON content safely (may contain invalid JSON) + let parsed: unknown = null; + if (isJson) { + try { + parsed = JSON.parse(content); + } catch { + // Invalid JSON - return null for parsed, content still available as raw string + } + } + return { success: true, data: { content, filePath: decodedPath, isJson, - parsed: isJson ? JSON.parse(content) : null, + parsed, }, }; }); @@ -3236,12 +3298,12 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; // ========== Hook Events ========== this.app.post('/api/hook-event', async (req) => { - const { event, sessionId, data } = req.body as HookEventRequest; - const validEvents = ['idle_prompt', 'permission_prompt', 'elicitation_dialog', 'stop'] as const; - if (!event || !validEvents.includes(event as typeof validEvents[number])) { - return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Invalid event type'); + const result = HookEventSchema.safeParse(req.body); + if (!result.success) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, result.error.issues[0]?.message ?? 'Validation failed'); } - if (!sessionId || !this.sessions.has(sessionId)) { + const { event, sessionId, data } = result.data; + if (!this.sessions.has(sessionId)) { return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found'); } @@ -3269,7 +3331,7 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; } // Sanitize forwarded data: only include known safe fields, limit size - const safeData = sanitizeHookData(data); + const safeData = sanitizeHookData(data as Record | undefined); this.broadcast(`hook:${event}`, { sessionId, timestamp: Date.now(), ...safeData }); // Track in run summary @@ -3507,6 +3569,8 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; this.outputBatches.delete(sessionId); this.taskUpdateBatches.delete(sessionId); this.stateUpdatePending.delete(sessionId); + this.lastTerminalEventTime.delete(sessionId); + this.adaptiveBatchIntervals.delete(sessionId); // Reset Ralph tracker on the session before cleanup if (session) { @@ -4032,6 +4096,11 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; console.log(`[Server] Restored respawn controller for session ${session.id} from ${source} (will start in ${Math.ceil(delayMs / 1000)}s)`); const timer = setTimeout(() => { this.pendingRespawnStarts.delete(session.id); + // Verify session still exists (may have been deleted during grace period) + if (!this.sessions.has(session.id)) { + console.log(`[Server] Skipping restored respawn start - session ${session.id} no longer exists`); + return; + } // Double-check controller still exists and is stopped const ctrl = this.respawnControllers.get(session.id); if (ctrl && ctrl.state === 'stopped') { @@ -4112,8 +4181,16 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; this.scheduledRuns.set(id, run); this.broadcast('scheduled:created', run); - // Start the run loop - this.runScheduledLoop(id); + // Start the run loop (fire-and-forget with error handling) + this.runScheduledLoop(id).catch((err) => { + console.error(`[WebServer] Scheduled run ${id} failed:`, err); + const failedRun = this.scheduledRuns.get(id); + if (failedRun && failedRun.status === 'running') { + failedRun.status = 'stopped'; + failedRun.logs.push(`[${new Date().toISOString()}] Error: ${err instanceof Error ? err.message : String(err)}`); + this.broadcast('scheduled:stopped', { id, reason: 'error' }); + } + }); return run; } @@ -4378,21 +4455,23 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; const newBatch = existing + data; this.terminalBatches.set(sessionId, newBatch); - // Adaptive batching: detect rapid events and extend batch window + // Adaptive batching: detect rapid events and extend batch window (per-session) const now = Date.now(); - const lastEvent = this.lastTerminalEventTime.get(sessionId) || 0; + const lastEvent = this.lastTerminalEventTime.get(sessionId) ?? 0; const eventGap = now - lastEvent; this.lastTerminalEventTime.set(sessionId, now); - // Adjust batch interval based on event frequency + // Adjust batch interval based on event frequency (per-session) // Rapid events (<10ms gap) = 50ms batch, moderate (<20ms) = 32ms, else 16ms + let sessionInterval: number; if (eventGap > 0 && eventGap < 10) { - this.adaptiveBatchInterval = 50; + sessionInterval = 50; } else if (eventGap > 0 && eventGap < 20) { - this.adaptiveBatchInterval = 32; + sessionInterval = 32; } else { - this.adaptiveBatchInterval = TERMINAL_BATCH_INTERVAL; + sessionInterval = TERMINAL_BATCH_INTERVAL; } + this.adaptiveBatchIntervals.set(sessionId, sessionInterval); // Flush immediately if batch is large for responsiveness if (newBatch.length > BATCH_FLUSH_THRESHOLD) { @@ -4405,17 +4484,27 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; } // Start batch timer if not already running (uses adaptive interval) + // Pick the minimum interval across all pending sessions for responsiveness if (!this.terminalBatchTimer) { + const batchInterval = Math.min( + ...Array.from(this.adaptiveBatchIntervals.values()), + TERMINAL_BATCH_INTERVAL + ); this.terminalBatchTimer = setTimeout(() => { this.flushTerminalBatches(); this.terminalBatchTimer = null; - // Reset adaptive interval after flush - this.adaptiveBatchInterval = TERMINAL_BATCH_INTERVAL; - }, this.adaptiveBatchInterval); + // Clear per-session intervals after flush (they'll be recalculated on next event) + this.adaptiveBatchIntervals.clear(); + }, batchInterval); } } private flushTerminalBatches(): void { + // Skip if server is stopping (timer may have been queued before stop() was called) + if (this._isStopping) { + this.terminalBatches.clear(); + return; + } for (const [sessionId, data] of this.terminalBatches) { if (data.length > 0) { // Wrap with DEC mode 2026 synchronized output markers @@ -4446,6 +4535,11 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; } private flushOutputBatches(): void { + // Skip if server is stopping (timer may have been queued before stop() was called) + if (this._isStopping) { + this.outputBatches.clear(); + return; + } for (const [sessionId, data] of this.outputBatches) { if (data.length > 0) { this.broadcast('session:output', { id: sessionId, data }); @@ -4454,12 +4548,15 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; this.outputBatches.clear(); } - // Batch task:updated events at 100ms - only send latest update per session + // Batch task:updated events at 100ms - only send latest update per task + // Key is sessionId:taskId to avoid collisions when multiple tasks update concurrently private batchTaskUpdate(sessionId: string, task: BackgroundTask): void { // Skip if server is stopping if (this._isStopping) return; - this.taskUpdateBatches.set(sessionId, task); + // Use composite key to avoid losing updates when multiple tasks update in same batch window + const key = `${sessionId}:${task.id}`; + this.taskUpdateBatches.set(key, { sessionId, task }); if (!this.taskUpdateBatchTimer) { this.taskUpdateBatchTimer = setTimeout(() => { @@ -4470,7 +4567,12 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; } private flushTaskUpdateBatches(): void { - for (const [sessionId, task] of this.taskUpdateBatches) { + // Skip if server is stopping (timer may have been queued before stop() was called) + if (this._isStopping) { + this.taskUpdateBatches.clear(); + return; + } + for (const [, { sessionId, task }] of this.taskUpdateBatches) { this.broadcast('task:updated', { sessionId, task }); } this.taskUpdateBatches.clear(); @@ -4496,6 +4598,11 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; } private flushStateUpdates(): void { + // Skip if server is stopping (timer may have been queued before stop() was called) + if (this._isStopping) { + this.stateUpdatePending.clear(); + return; + } for (const sessionId of this.stateUpdatePending) { const session = this.sessions.get(sessionId); if (session) { @@ -4905,6 +5012,21 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; // Stop image watcher imageWatcher.stop(); + // Destroy file stream manager (clears cleanup timer and kills remaining tail processes) + fileStreamManager.destroy(); + + // Clear remaining Maps that accumulate session references + this.respawnTimers.clear(); + this.runSummaryTrackers.clear(); + this.transcriptWatchers.clear(); + this.sessionListenerRefs.clear(); + this.scheduledRuns.clear(); + // Dispose StaleExpirationMap (stops internal cleanup timer) + this.lastTerminalEventTime.dispose(); + this.adaptiveBatchIntervals.clear(); + this.activePlanOrchestrators.clear(); + this.cleaningUp.clear(); + await this.app.close(); } } diff --git a/test/e2e/e2e.config.ts b/test/e2e/e2e.config.ts index 026a056c..f82d7270 100644 --- a/test/e2e/e2e.config.ts +++ b/test/e2e/e2e.config.ts @@ -17,6 +17,7 @@ export const E2E_PORTS = { MOBILE_SAFARI: 3191, MOBILE_COMPREHENSIVE: 3192, MOBILE_EDGE_CASES: 3193, + PLAN_GENERATION: 3194, } as const; // Mobile device viewports for responsive testing diff --git a/test/e2e/screenshots/baselines/agent-windows-visible.png b/test/e2e/screenshots/baselines/agent-windows-visible.png new file mode 100644 index 00000000..903710ec Binary files /dev/null and b/test/e2e/screenshots/baselines/agent-windows-visible.png differ diff --git a/test/e2e/screenshots/baselines/respawn-active.png b/test/e2e/screenshots/baselines/respawn-active.png new file mode 100644 index 00000000..3d54f844 Binary files /dev/null and b/test/e2e/screenshots/baselines/respawn-active.png differ diff --git a/test/e2e/workflows/mobile-comprehensive.e2e.ts b/test/e2e/workflows/mobile-comprehensive.e2e.ts index 6454f2b2..f304626e 100644 --- a/test/e2e/workflows/mobile-comprehensive.e2e.ts +++ b/test/e2e/workflows/mobile-comprehensive.e2e.ts @@ -19,9 +19,9 @@ import { type ServerFixture, type MobileBrowserFixture, } from '../fixtures/index.js'; -import { E2E_TIMEOUTS, MOBILE_VIEWPORTS, generateCaseName } from '../e2e.config.js'; +import { E2E_TIMEOUTS, MOBILE_VIEWPORTS, generateCaseName, E2E_PORTS } from '../e2e.config.js'; -const PORT = 3192; +const PORT = E2E_PORTS.MOBILE_COMPREHENSIVE; let serverFixture: ServerFixture | null = null; let cleanup: CleanupTracker; diff --git a/test/e2e/workflows/plan-generation.e2e.ts b/test/e2e/workflows/plan-generation.e2e.ts index 7ec3db43..01bdd5c1 100644 --- a/test/e2e/workflows/plan-generation.e2e.ts +++ b/test/e2e/workflows/plan-generation.e2e.ts @@ -22,9 +22,9 @@ import { type ServerFixture, type BrowserFixture, } from '../fixtures/index.js'; -import { E2E_TIMEOUTS } from '../e2e.config.js'; +import { E2E_TIMEOUTS, E2E_PORTS } from '../e2e.config.js'; -const PORT = 3191; // Next available port per CLAUDE.md +const PORT = E2E_PORTS.PLAN_GENERATION; let serverFixture: ServerFixture | null = null; let cleanup: CleanupTracker;