diff --git a/CLAUDE.md b/CLAUDE.md index 1ff59435..99a4f9de 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -35,7 +35,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.1578 (must match `package.json` for npm publish) +**Version**: 0.1579 (must match `package.json` for npm publish) ## Project Overview diff --git a/package.json b/package.json index 7a67a4ea..0b4f00f2 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "claudeman", - "version": "0.1578", + "version": "0.1579", "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/session.ts b/src/session.ts index 24025236..c42219d0 100644 --- a/src/session.ts +++ b/src/session.ts @@ -353,6 +353,13 @@ export class Session extends EventEmitter { // Uses LRUMap for automatic eviction at MAX_TASK_DESCRIPTIONS limit private _recentTaskDescriptions: LRUMap = new LRUMap({ maxSize: Session.MAX_TASK_DESCRIPTIONS }); + // Throttle expensive PTY processing (Ralph, bash parser, task descriptions) + // Accumulates clean data between processing windows to avoid running regex on every chunk + private _lastExpensiveProcessTime: number = 0; + private _pendingCleanData: string = ''; + private _expensiveProcessTimer: NodeJS.Timeout | null = null; + private static readonly EXPENSIVE_PROCESS_INTERVAL_MS = 150; // Process at most every 150ms + constructor(config: Partial & { workingDir: string; mode?: SessionMode; @@ -990,17 +997,6 @@ export class Session extends EventEmitter { .replace(CTRL_L_PATTERN, ''); // Remove Ctrl+L if (!data) return; // Skip if only filtered sequences - // Lazy ANSI strip: only compute cleanData when a consumer actually needs it. - // During active streaming, many consumers early-exit (ralph disabled, cli info parsed, - // no 'token' in data, etc.), so we skip the expensive O(n) regex on most chunks. - let _cleanData: string | null = null; - const getCleanData = (): string => { - if (_cleanData === null) { - _cleanData = data.replace(ANSI_ESCAPE_PATTERN_FULL, ''); - } - return _cleanData; - }; - // BufferAccumulator handles auto-trimming when max size exceeded this._terminalBuffer.append(data); this._lastActivityAt = Date.now(); @@ -1008,36 +1004,7 @@ export class Session extends EventEmitter { this.emit('terminal', data); this.emit('output', data); - // Forward to Ralph tracker to detect Ralph loops and todos - // Ralph tracker early-exits when disabled + autoEnableDisabled, skipping ANSI strip - if (this._ralphTracker.enabled || !this._ralphTracker.autoEnableDisabled) { - this._ralphTracker.processCleanData(getCleanData()); - } - - // Forward to Bash tool parser to detect file-viewing commands - // Parser early-exits when disabled, skipping ANSI strip - if (this._bashToolParser.enabled) { - this._bashToolParser.processCleanData(getCleanData()); - } - - // Parse token count from status line (e.g., "123.4k tokens" or "5234 tokens") - // Pre-check on raw data: 'token' won't appear in ANSI sequences - if (data.includes('token')) { - this.parseTokensFromStatusLine(getCleanData()); - } - - // Parse Claude Code CLI info (version, model, account type) from startup - // Gated by _cliInfoParsed — only runs during first few chunks - if (!this._cliInfoParsed) { - this.parseClaudeCodeInfo(getCleanData()); - } - - // Parse task descriptions from terminal output (e.g., "Explore(Check files)") - // Pre-check on raw data: parentheses are safe to check without ANSI strip - if (data.includes('(') && data.includes(')')) { - this.parseTaskDescriptionsFromTerminalData(getCleanData()); - } - + // === Idle/working detection runs on every chunk (latency-sensitive) === // Detect if Claude is working or at prompt // The prompt line contains "❯" when waiting for input if (data.includes('❯') || data.includes('\u276f')) { @@ -1069,20 +1036,52 @@ export class Session extends EventEmitter { data.includes('⠹') || data.includes('⠸') || data.includes('⠼') || data.includes('⠴') || data.includes('⠦') || data.includes('⠧'); - // Slow path: check text keywords on clean data (avoids false positives from window titles) - const hasWorkKeyword = hasSpinner || - getCleanData().includes('Thinking') || getCleanData().includes('Writing') || - getCleanData().includes('Reading') || getCleanData().includes('Running'); - if (hasWorkKeyword) { + if (hasSpinner) { if (!this._isWorking) { this._isWorking = true; this._status = 'busy'; this.emit('working'); } - // Reset timeout and idle confirmation flag since Claude is active this._awaitingIdleConfirmation = false; if (this.activityTimeout) clearTimeout(this.activityTimeout); } + + // === Expensive processing (ANSI strip, Ralph, bash parser) is throttled === + // Instead of running regex-heavy parsers on every PTY chunk, we accumulate + // raw data and process at most every EXPENSIVE_PROCESS_INTERVAL_MS. + // This dramatically reduces CPU load with multiple busy sessions. + const now = Date.now(); + const elapsed = now - this._lastExpensiveProcessTime; + if (elapsed >= Session.EXPENSIVE_PROCESS_INTERVAL_MS) { + // Process immediately — include any previously accumulated data + this._lastExpensiveProcessTime = now; + const accumulated = this._pendingCleanData ? this._pendingCleanData + data : data; + this._pendingCleanData = ''; + if (this._expensiveProcessTimer) { + clearTimeout(this._expensiveProcessTimer); + this._expensiveProcessTimer = null; + } + this._processExpensiveParsers(accumulated); + } else { + // Accumulate for deferred processing + this._pendingCleanData += data; + // Cap accumulated size to prevent unbounded growth + if (this._pendingCleanData.length > 64 * 1024) { + this._pendingCleanData = this._pendingCleanData.slice(-32 * 1024); + } + // Schedule deferred processing if not already scheduled + if (!this._expensiveProcessTimer) { + this._expensiveProcessTimer = setTimeout(() => { + this._expensiveProcessTimer = null; + this._lastExpensiveProcessTime = Date.now(); + const pending = this._pendingCleanData; + this._pendingCleanData = ''; + if (pending) { + this._processExpensiveParsers(pending); + } + }, Session.EXPENSIVE_PROCESS_INTERVAL_MS - elapsed); + } + } }); this.ptyProcess.onExit(({ exitCode }) => { @@ -1104,6 +1103,12 @@ export class Session extends EventEmitter { clearTimeout(this._promptCheckTimeout); this._promptCheckTimeout = null; } + // Clear expensive processing timer and flush any pending data + if (this._expensiveProcessTimer) { + clearTimeout(this._expensiveProcessTimer); + this._expensiveProcessTimer = null; + } + this._pendingCleanData = ''; // If using mux, mark the session as detached but don't kill it if (this._muxSession && this._mux) { this._mux.setAttached(this.id, false); @@ -1112,6 +1117,61 @@ export class Session extends EventEmitter { }); } + /** + * Process expensive parsers (ANSI strip, Ralph, bash tool, token, CLI info, task descriptions). + * Called on a throttled schedule (every EXPENSIVE_PROCESS_INTERVAL_MS) instead of on every + * PTY data chunk. Receives accumulated raw data to process in one batch. + */ + private _processExpensiveParsers(rawData: string): void { + // Lazy ANSI strip: only compute cleanData when a consumer actually needs it. + let _cleanData: string | null = null; + const getCleanData = (): string => { + if (_cleanData === null) { + _cleanData = rawData.replace(ANSI_ESCAPE_PATTERN_FULL, ''); + } + return _cleanData; + }; + + // Forward to Ralph tracker to detect Ralph loops and todos + if (this._ralphTracker.enabled || !this._ralphTracker.autoEnableDisabled) { + this._ralphTracker.processCleanData(getCleanData()); + } + + // Forward to Bash tool parser to detect file-viewing commands + if (this._bashToolParser.enabled) { + this._bashToolParser.processCleanData(getCleanData()); + } + + // Parse token count from status line (e.g., "123.4k tokens" or "5234 tokens") + if (rawData.includes('token')) { + this.parseTokensFromStatusLine(getCleanData()); + } + + // Parse Claude Code CLI info (version, model, account type) from startup + if (!this._cliInfoParsed) { + this.parseClaudeCodeInfo(getCleanData()); + } + + // Parse task descriptions from terminal output (e.g., "Explore(Check files)") + if (rawData.includes('(') && rawData.includes(')')) { + this.parseTaskDescriptionsFromTerminalData(getCleanData()); + } + + // Work keyword detection (text-based, needs clean data) + // Only check if spinner didn't already trigger working state + if (!this._isWorking) { + const cleanData = getCleanData(); + if (cleanData.includes('Thinking') || cleanData.includes('Writing') || + cleanData.includes('Reading') || cleanData.includes('Running')) { + this._isWorking = true; + this._status = 'busy'; + this.emit('working'); + this._awaitingIdleConfirmation = false; + if (this.activityTimeout) clearTimeout(this.activityTimeout); + } + } + } + /** * Starts a plain shell session (bash/zsh) without Claude CLI. * @@ -2067,6 +2127,13 @@ export class Session extends EventEmitter { this._shellIdleTimer = null; } + // Clear expensive processing timer + if (this._expensiveProcessTimer) { + clearTimeout(this._expensiveProcessTimer); + this._expensiveProcessTimer = null; + } + this._pendingCleanData = ''; + // Immediately cleanup Promise callbacks to prevent orphaned references // during the rest of stop() processing (e.g., if mux kill times out) if (this.rejectPromise && !this._promptResolved) { diff --git a/src/web/server.ts b/src/web/server.ts index 60b94d43..04e687ff 100644 --- a/src/web/server.ts +++ b/src/web/server.ts @@ -354,17 +354,13 @@ export class WebServer extends EventEmitter { // Terminal batching for performance private terminalBatches: Map = new Map(); private terminalBatchSizes: Map = new Map(); // Running total avoids O(n) reduce per push - private terminalBatchTimer: NodeJS.Timeout | null = null; + private terminalBatchTimers: Map = new Map(); // Per-session timers (staggered flushes) // 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(); - // Tracked minimum across adaptiveBatchIntervals (avoids spreading into Math.min on every batch) - private _minBatchInterval: number = TERMINAL_BATCH_INTERVAL; // Scheduled runs cleanup timer private scheduledCleanupTimer: NodeJS.Timeout | null = null; // SSE event batching @@ -4059,13 +4055,17 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; this.runSummaryTrackers.delete(sessionId); } - // Clear batches and pending state updates + // Clear batches, per-session timers, and pending state updates this.terminalBatches.delete(sessionId); this.terminalBatchSizes.delete(sessionId); + const batchTimer = this.terminalBatchTimers.get(sessionId); + if (batchTimer) { + clearTimeout(batchTimer); + this.terminalBatchTimers.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) { @@ -4967,7 +4967,8 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; } // Batch terminal data for better performance (60fps) - // Uses adaptive batching: extends batch window when events are rapid-fire + // Uses per-session timers with adaptive intervals to prevent thundering herd: + // each session flushes independently rather than all sessions flushing in one burst. private batchTerminalData(sessionId: string, data: string): void { // Skip if server is stopping if (this._isStopping) return; @@ -4998,56 +4999,55 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; } else { sessionInterval = TERMINAL_BATCH_INTERVAL; } - this.adaptiveBatchIntervals.set(sessionId, sessionInterval); - // Track minimum to avoid O(n) spread on every batch event - if (sessionInterval < this._minBatchInterval) { - this._minBatchInterval = sessionInterval; - } // Flush immediately if batch is large for responsiveness if (totalLength > BATCH_FLUSH_THRESHOLD) { - if (this.terminalBatchTimer) { - clearTimeout(this.terminalBatchTimer); - this.terminalBatchTimer = null; + const existingTimer = this.terminalBatchTimers.get(sessionId); + if (existingTimer) { + clearTimeout(existingTimer); + this.terminalBatchTimers.delete(sessionId); } - this.flushTerminalBatches(); + this.flushSessionTerminalBatch(sessionId); return; } - // Start batch timer if not already running (uses adaptive interval) - // Pick the minimum interval across all pending sessions for responsiveness - if (!this.terminalBatchTimer) { - this.terminalBatchTimer = setTimeout(() => { - this.flushTerminalBatches(); - this.terminalBatchTimer = null; - // Clear per-session intervals after flush (they'll be recalculated on next event) - this.adaptiveBatchIntervals.clear(); - this._minBatchInterval = TERMINAL_BATCH_INTERVAL; - }, this._minBatchInterval); + // Start per-session batch timer if not already running + // Each session flushes independently — prevents one busy session from + // forcing all sessions to flush at its rate (thundering herd) + if (!this.terminalBatchTimers.has(sessionId)) { + this.terminalBatchTimers.set(sessionId, setTimeout(() => { + this.terminalBatchTimers.delete(sessionId); + this.flushSessionTerminalBatch(sessionId); + }, sessionInterval)); } } - private flushTerminalBatches(): void { - // Skip if server is stopping (timer may have been queued before stop() was called) + /** Flush a single session's batched terminal data */ + private flushSessionTerminalBatch(sessionId: string): void { if (this._isStopping) { - this.terminalBatches.clear(); - this.terminalBatchSizes.clear(); + this.terminalBatches.delete(sessionId); + this.terminalBatchSizes.delete(sessionId); return; } - for (const [sessionId, chunks] of this.terminalBatches) { - if (chunks.length > 0) { - // Join chunks only at flush time (avoids O(n^2) string concatenation in batchTerminalData) - const data = chunks.join(''); - // Wrap with DEC mode 2026 synchronized output markers - // Terminal buffers all output between markers and renders atomically, - // eliminating partial-frame flicker from Ink's full-screen redraws. - // Unsupported terminals ignore these sequences harmlessly. - const syncData = DEC_SYNC_START + data + DEC_SYNC_END; - this.broadcast('session:terminal', { id: sessionId, data: syncData }); + const chunks = this.terminalBatches.get(sessionId); + if (chunks && chunks.length > 0) { + // Join chunks only at flush time (avoids O(n^2) string concatenation in batchTerminalData) + const data = chunks.join(''); + // Wrap with DEC mode 2026 synchronized output markers + // Terminal buffers all output between markers and renders atomically, + // eliminating partial-frame flicker from Ink's full-screen redraws. + // Unsupported terminals ignore these sequences harmlessly. + const syncData = DEC_SYNC_START + data + DEC_SYNC_END; + // Fast path: build SSE message directly without JSON.stringify on wrapper object. + // Only the terminal data string needs escaping; sessionId is a UUID (safe to template). + const escapedData = JSON.stringify(syncData); + const message = `event: session:terminal\ndata: {"id":"${sessionId}","data":${escapedData}}\n\n`; + for (const client of this.sseClients) { + this.sendSSEPreformatted(client, message); } } - this.terminalBatches.clear(); - this.terminalBatchSizes.clear(); + this.terminalBatches.delete(sessionId); + this.terminalBatchSizes.delete(sessionId); } // Batch task:updated events at 100ms - only send latest update per task @@ -5455,11 +5455,11 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; this.sseClients.clear(); this.backpressuredClients.clear(); - // Clear batch timers - if (this.terminalBatchTimer) { - clearTimeout(this.terminalBatchTimer); - this.terminalBatchTimer = null; + // Clear per-session batch timers + for (const timer of this.terminalBatchTimers.values()) { + clearTimeout(timer); } + this.terminalBatchTimers.clear(); this.terminalBatches.clear(); this.terminalBatchSizes.clear(); @@ -5606,8 +5606,6 @@ NOW: Generate the implementation plan for the task above. Think step by step.`; this.scheduledRuns.clear(); // Dispose StaleExpirationMap (stops internal cleanup timer) this.lastTerminalEventTime.dispose(); - this.adaptiveBatchIntervals.clear(); - this._minBatchInterval = TERMINAL_BATCH_INTERVAL; this.activePlanOrchestrators.clear(); this.cleaningUp.clear();