From 3145eac6d936a824b335854f28f9058f29d55a0c Mon Sep 17 00:00:00 2001 From: arkon Date: Wed, 25 Mar 2026 23:21:38 +0100 Subject: [PATCH] refactor: extract helper methods to reduce duplication and improve readability DRY up repeated patterns across 7 core files: - state-store: extract serializeState() and split assembleStateJson() into 3 focused methods - session: extract _resetBuffers(), _clearAllTimers(), _handleJsonMessage() - ralph-tracker: extract completeAllTodos() (was 4x duplicated), emitValidationWarning(), similarity constants - subagent-watcher: extract markSubagentAsCompleted(), extractFirstTextContent(), emitToolResult(), findOldestInactiveAgent() - respawn-controller: extract recoveryResetToWatching(), canAutoAccept(), formatRemainingSeconds(), validatePositiveTimeout() - tmux-manager: replace 15 path.includes() checks with single UNSAFE_PATH_CHARS regex - session-auto-ops: extract executeWhenIdle() shared retry helper for checkAutoCompact/checkAutoClear Co-Authored-By: Claude Opus 4.6 --- src/ralph-tracker.ts | 124 ++++++++++---------- src/respawn-controller.ts | 140 ++++++++++++++--------- src/session-auto-ops.ts | 149 +++++++++++++++--------- src/session.ts | 232 +++++++++++++++++++------------------- src/state-store.ts | 98 ++++++++-------- src/subagent-watcher.ts | 177 +++++++++++++++-------------- src/tmux-manager.ts | 21 +--- 7 files changed, 507 insertions(+), 434 deletions(-) diff --git a/src/ralph-tracker.ts b/src/ralph-tracker.ts index ff408d87..8cf5a4c0 100644 --- a/src/ralph-tracker.ts +++ b/src/ralph-tracker.ts @@ -100,6 +100,18 @@ const TODO_CLEANUP_INTERVAL_MS = INACTIVITY_TIMEOUT_MS; */ const TODO_SIMILARITY_THRESHOLD = 0.85; +/** + * Similarity threshold for short todo content (<30 chars). + * Higher threshold reduces false positive deduplication of short strings. + */ +const SIMILARITY_THRESHOLD_SHORT = 0.95; + +/** + * Similarity threshold for medium-length todo content (30-60 chars). + * Slightly relaxed compared to short strings. + */ +const SIMILARITY_THRESHOLD_MEDIUM = 0.9; + /** * Debounce interval for event emissions (milliseconds). * Prevents UI jitter from rapid consecutive updates. @@ -1299,6 +1311,24 @@ export class RalphTracker extends EventEmitter { this.detectTodoItems(trimmed); } + /** + * Mark all tracked todos as completed and emit todoUpdate if any changed. + * @returns true if any todo was updated + */ + private completeAllTodos(): boolean { + let updated = false; + for (const todo of this._todos.values()) { + if (todo.status !== 'completed') { + todo.status = 'completed'; + updated = true; + } + } + if (updated) { + this.emit('todoUpdate', this.todos); + } + return updated; + } + /** * Detect "all tasks complete" messages. */ @@ -1318,16 +1348,7 @@ export class RalphTracker extends EventEmitter { return; } - let updated = false; - for (const todo of this._todos.values()) { - if (todo.status !== 'completed') { - todo.status = 'completed'; - updated = true; - } - } - if (updated) { - this.emit('todoUpdate', this.todos); - } + this.completeAllTodos(); if (this._loopState.completionPhrase) { this._loopState.active = false; @@ -1425,16 +1446,7 @@ export class RalphTracker extends EventEmitter { if (bareCount > 1) return; - let updated = false; - for (const todo of this._todos.values()) { - if (todo.status !== 'completed') { - todo.status = 'completed'; - updated = true; - } - } - if (updated) { - this.emit('todoUpdate', this.todos); - } + this.completeAllTodos(); this._loopState.active = false; this._loopState.lastActivity = Date.now(); @@ -1480,16 +1492,7 @@ export class RalphTracker extends EventEmitter { if (canonicalCount >= 2 || this._loopState.active) { this._loopState.active = false; this._loopState.lastActivity = Date.now(); - let updated = false; - for (const todo of this._todos.values()) { - if (todo.status !== 'completed') { - todo.status = 'completed'; - updated = true; - } - } - if (updated) { - this.emit('todoUpdate', this.todos); - } + this.completeAllTodos(); this.emit('completionDetected', matchedPhrase); this.emit('loopUpdate', this.loopState); return; @@ -1497,16 +1500,7 @@ export class RalphTracker extends EventEmitter { } if (this._loopState.active || count >= 2) { - let updated = false; - for (const todo of this._todos.values()) { - if (todo.status !== 'completed') { - todo.status = 'completed'; - updated = true; - } - } - if (updated) { - this.emit('todoUpdate', this.todos); - } + this.completeAllTodos(); this._loopState.active = false; this._loopState.lastActivity = Date.now(); @@ -1532,41 +1526,39 @@ export class RalphTracker extends EventEmitter { const suggestedPhrase = `${phrase}_${uniqueSuffix}`; if (COMMON_COMPLETION_PHRASES.has(normalized)) { - console.warn( - `[RalphTracker] Warning: Completion phrase "${phrase}" is very common and may cause false positives. Consider using: "${suggestedPhrase}"` - ); - this.emit('phraseValidationWarning', { - phrase, - reason: 'common', - suggestedPhrase, - }); + this.emitValidationWarning(phrase, 'common', suggestedPhrase); return; } if (normalized.length < MIN_RECOMMENDED_PHRASE_LENGTH) { - console.warn( - `[RalphTracker] Warning: Completion phrase "${phrase}" is too short (${normalized.length} chars). Consider using: "${suggestedPhrase}"` - ); - this.emit('phraseValidationWarning', { - phrase, - reason: 'short', - suggestedPhrase, - }); + this.emitValidationWarning(phrase, 'short', suggestedPhrase); return; } if (/^\d+$/.test(normalized)) { - console.warn( - `[RalphTracker] Warning: Completion phrase "${phrase}" is numeric-only and may cause false positives. Consider using: "${suggestedPhrase}"` - ); - this.emit('phraseValidationWarning', { - phrase, - reason: 'numeric', - suggestedPhrase, - }); + this.emitValidationWarning(phrase, 'numeric', suggestedPhrase); } } + /** + * Emit a phrase validation warning with a console message and event. + */ + private emitValidationWarning(phrase: string, reason: 'common' | 'short' | 'numeric', suggestedPhrase: string): void { + const descriptions: Record<'common' | 'short' | 'numeric', string> = { + common: 'is very common and may cause false positives', + short: `is too short (${phrase.toUpperCase().replace(/[\s_\-.]+/g, '').length} chars)`, + numeric: 'is numeric-only and may cause false positives', + }; + console.warn( + `[RalphTracker] Warning: Completion phrase "${phrase}" ${descriptions[reason]}. Consider using: "${suggestedPhrase}"` + ); + this.emit('phraseValidationWarning', { + phrase, + reason, + suggestedPhrase, + }); + } + /** * Activate the loop if not already active. */ @@ -1977,9 +1969,9 @@ export class RalphTracker extends EventEmitter { let threshold: number; if (normalized.length < 30) { - threshold = 0.95; + threshold = SIMILARITY_THRESHOLD_SHORT; } else if (normalized.length < 60) { - threshold = 0.9; + threshold = SIMILARITY_THRESHOLD_MEDIUM; } else { threshold = TODO_SIMILARITY_THRESHOLD; } diff --git a/src/respawn-controller.ts b/src/respawn-controller.ts index b3b33e03..7aa17336 100644 --- a/src/respawn-controller.ts +++ b/src/respawn-controller.ts @@ -551,6 +551,14 @@ export interface RespawnEvents { respawnBlocked: (data: { reason: string; details: string }) => void; } +/** + * Convert milliseconds to a non-negative whole number of seconds for countdown display. + * Rounds up so that e.g. 1200 ms shows as 2 s (never under-reports remaining time). + */ +function formatRemainingSeconds(ms: number): number { + return Math.max(0, Math.ceil(ms / 1000)); +} + /** Default configuration values */ const DEFAULT_CONFIG: RespawnConfig = { idleTimeoutMs: 10000, // 10 seconds of no activity after prompt (legacy, still used as fallback) @@ -839,12 +847,25 @@ export class RespawnController extends EventEmitter { private validateConfig(): void { const c = this.config; + /** + * Validate that a timeout value is positive (or non-negative when allowZero is true). + * Falls back to the DEFAULT_CONFIG value if invalid. + */ + const validatePositiveTimeout = (field: keyof RespawnConfig, allowZero = false): void => { + const value = c[field] as number; + const invalid = allowZero ? value < 0 : value <= 0; + if (invalid) { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + (c as any)[field] = DEFAULT_CONFIG[field]; + } + }; + // Ensure timeouts are positive - if (c.idleTimeoutMs <= 0) c.idleTimeoutMs = DEFAULT_CONFIG.idleTimeoutMs; - if (c.completionConfirmMs <= 0) c.completionConfirmMs = DEFAULT_CONFIG.completionConfirmMs; - if (c.noOutputTimeoutMs <= 0) c.noOutputTimeoutMs = DEFAULT_CONFIG.noOutputTimeoutMs; - if (c.autoAcceptDelayMs < 0) c.autoAcceptDelayMs = DEFAULT_CONFIG.autoAcceptDelayMs; - if (c.interStepDelayMs <= 0) c.interStepDelayMs = DEFAULT_CONFIG.interStepDelayMs; + validatePositiveTimeout('idleTimeoutMs'); + validatePositiveTimeout('completionConfirmMs'); + validatePositiveTimeout('noOutputTimeoutMs'); + validatePositiveTimeout('autoAcceptDelayMs', true); + validatePositiveTimeout('interStepDelayMs'); // Ensure completion confirm doesn't exceed no-output timeout if (c.completionConfirmMs > c.noOutputTimeoutMs) { @@ -852,14 +873,14 @@ export class RespawnController extends EventEmitter { } // Ensure AI check timeouts are positive - if (c.aiIdleCheckTimeoutMs <= 0) c.aiIdleCheckTimeoutMs = DEFAULT_CONFIG.aiIdleCheckTimeoutMs; - if (c.aiIdleCheckCooldownMs < 0) c.aiIdleCheckCooldownMs = DEFAULT_CONFIG.aiIdleCheckCooldownMs; - if (c.aiIdleCheckMaxContext <= 0) c.aiIdleCheckMaxContext = DEFAULT_CONFIG.aiIdleCheckMaxContext; + validatePositiveTimeout('aiIdleCheckTimeoutMs'); + validatePositiveTimeout('aiIdleCheckCooldownMs', true); + validatePositiveTimeout('aiIdleCheckMaxContext'); // Ensure plan check timeouts are positive - if (c.aiPlanCheckTimeoutMs <= 0) c.aiPlanCheckTimeoutMs = DEFAULT_CONFIG.aiPlanCheckTimeoutMs; - if (c.aiPlanCheckCooldownMs < 0) c.aiPlanCheckCooldownMs = DEFAULT_CONFIG.aiPlanCheckCooldownMs; - if (c.aiPlanCheckMaxContext <= 0) c.aiPlanCheckMaxContext = DEFAULT_CONFIG.aiPlanCheckMaxContext; + validatePositiveTimeout('aiPlanCheckTimeoutMs'); + validatePositiveTimeout('aiPlanCheckCooldownMs', true); + validatePositiveTimeout('aiPlanCheckMaxContext'); } /** Wire up AI checker events to controller events (removes existing listeners first to prevent duplicates) */ @@ -987,11 +1008,11 @@ export class RespawnController extends EventEmitter { waitingFor = 'AI verdict (IDLE or WORKING)'; } else if (this._state === 'confirming_idle') { statusText = `Confirming idle (${confidence}% confidence)`; - waitingFor = `${Math.max(0, Math.ceil((this.config.completionConfirmMs - msSinceLastOutput) / 1000))}s more silence`; + waitingFor = `${formatRemainingSeconds(this.config.completionConfirmMs - msSinceLastOutput)}s more silence`; } else if (this._state === 'watching') { const aiState = this.aiChecker.getState(); if (aiState.status === 'cooldown') { - const remaining = Math.ceil(this.aiChecker.getCooldownRemainingMs() / 1000); + const remaining = formatRemainingSeconds(this.aiChecker.getCooldownRemainingMs()); statusText = `AI Check: WORKING (cooldown ${remaining}s)`; waitingFor = 'Cooldown to expire'; } else if (completionMessageDetected) { @@ -1798,24 +1819,28 @@ export class RespawnController extends EventEmitter { case 'sending_init': case 'sending_kickstart': // For sending states, retry the send - this.log('Recovery: returning to watching state'); - this.setState('watching'); - this.startNoOutputTimer(); - this.startPreFilterTimer(); - if (this.config.autoAcceptPrompts) { - this.startAutoAcceptTimer(); - } + this.recoveryResetToWatching('returning to watching state'); break; default: // Fallback: reset to watching - this.log('Recovery: fallback to watching state'); - this.setState('watching'); - this.startNoOutputTimer(); - this.startPreFilterTimer(); - if (this.config.autoAcceptPrompts) { - this.startAutoAcceptTimer(); - } + this.recoveryResetToWatching('fallback to watching state'); + } + } + + /** + * Reset the controller to watching state during stuck-state recovery. + * Sets state to watching and restarts all detection timers. + * + * @param reason - Human-readable reason for the reset (logged) + */ + private recoveryResetToWatching(reason: string): void { + this.log(`Recovery: ${reason}`); + this.setState('watching'); + this.startNoOutputTimer(); + this.startPreFilterTimer(); + if (this.config.autoAcceptPrompts) { + this.startAutoAcceptTimer(); } } @@ -2051,7 +2076,7 @@ export class RespawnController extends EventEmitter { // If on cooldown, don't start check - wait for cooldown to expire if (this.aiChecker.isOnCooldown()) { this.log( - `AI check on cooldown (${Math.ceil(this.aiChecker.getCooldownRemainingMs() / 1000)}s remaining), waiting...` + `AI check on cooldown (${formatRemainingSeconds(this.aiChecker.getCooldownRemainingMs())}s remaining), waiting...` ); return; } @@ -2198,36 +2223,15 @@ export class RespawnController extends EventEmitter { * @fires planCheckStarted */ private tryAutoAccept(): void { - // Only auto-accept in watching state (not during a respawn cycle) - if (this._state !== 'watching') return; + if (!this.canAutoAccept()) return; - // Don't auto-accept if a completion message was detected (normal idle handles it) - if (this.completionMessageTime !== null) return; - - // Don't auto-accept if disabled - if (!this.config.autoAcceptPrompts) return; - - // Don't auto-accept if we haven't received any output yet (prevents spurious Enter on fresh start) - if (!this.hasReceivedOutput) return; - - // Don't auto-accept if an elicitation dialog (AskUserQuestion) was detected - if (this.elicitationDetected) { - this.log('Skipping auto-accept: elicitation dialog detected (AskUserQuestion)'); - return; - } - - // Stage 1: Pre-filter — check if buffer looks like plan mode const buffer = this.terminalBuffer.value; - if (!this.isPlanModePreFilterMatch(buffer)) { - this.log('Skipping auto-accept: pre-filter did not match plan mode patterns'); - return; - } // Stage 2: AI confirmation (if enabled and available) if (this.config.aiPlanCheckEnabled && this.planChecker.status !== 'disabled') { if (this.planChecker.isOnCooldown()) { this.log( - `Skipping auto-accept: plan checker on cooldown (${Math.ceil(this.planChecker.getCooldownRemainingMs() / 1000)}s remaining)` + `Skipping auto-accept: plan checker on cooldown (${formatRemainingSeconds(this.planChecker.getCooldownRemainingMs())}s remaining)` ); return; } @@ -2244,6 +2248,40 @@ export class RespawnController extends EventEmitter { this.sendAutoAcceptEnter(); } + /** + * Check whether all preconditions for auto-accept are met. + * Validates state, config, and pre-filter conditions before attempting auto-accept. + * + * @returns True if auto-accept should proceed to the AI confirmation stage + */ + private canAutoAccept(): boolean { + // Only auto-accept in watching state (not during a respawn cycle) + if (this._state !== 'watching') return false; + + // Don't auto-accept if a completion message was detected (normal idle handles it) + if (this.completionMessageTime !== null) return false; + + // Don't auto-accept if disabled + if (!this.config.autoAcceptPrompts) return false; + + // Don't auto-accept if we haven't received any output yet (prevents spurious Enter on fresh start) + if (!this.hasReceivedOutput) return false; + + // Don't auto-accept if an elicitation dialog (AskUserQuestion) was detected + if (this.elicitationDetected) { + this.log('Skipping auto-accept: elicitation dialog detected (AskUserQuestion)'); + return false; + } + + // Stage 1: Pre-filter — check if buffer looks like plan mode + if (!this.isPlanModePreFilterMatch(this.terminalBuffer.value)) { + this.log('Skipping auto-accept: pre-filter did not match plan mode patterns'); + return false; + } + + return true; + } + /** * Check if the terminal buffer matches plan mode pre-filter patterns. * Only checks the last 2000 chars (plan mode UI appears at the bottom). diff --git a/src/session-auto-ops.ts b/src/session-auto-ops.ts index e1330ef8..48d3f61f 100644 --- a/src/session-auto-ops.ts +++ b/src/session-auto-ops.ts @@ -27,6 +27,57 @@ const COMPACT_COOLDOWN_MS = 10000; /** Cooldown after clear completes before re-enabling (5 seconds) */ const CLEAR_COOLDOWN_MS = 5000; +/** + * Executes an action when the session becomes idle, retrying if currently working. + * + * @param action - The async action to execute once idle + * @param isActive - Returns whether this operation is still active (not cancelled) + * @param isWorking - Returns whether the session is currently working + * @param isStopped - Returns whether the session has been stopped + * @param retryMs - Delay between retry attempts when working + * @param cooldownMs - Delay after action completes before calling onCooldownDone + * @param setTimer - Stores the timer reference for cleanup + * @param onCooldownDone - Called after cooldown to reset state + */ +async function executeWhenIdle( + action: () => Promise, + isActive: () => boolean, + isWorking: () => boolean, + isStopped: () => boolean, + retryMs: number, + cooldownMs: number, + setTimer: (timer: NodeJS.Timeout | null) => void, + onCooldownDone: () => void +): Promise { + if (isStopped()) return; + if (!isActive()) return; + + if (!isWorking()) { + if (isStopped()) return; + + await action(); + + if (!isStopped()) { + setTimer( + setTimeout(() => { + if (isStopped()) return; + setTimer(null); + onCooldownDone(); + }, cooldownMs) + ); + } + } else { + if (!isStopped()) { + setTimer( + setTimeout( + () => executeWhenIdle(action, isActive, isWorking, isStopped, retryMs, cooldownMs, setTimer, onCooldownDone), + retryMs + ) + ); + } + } +} + /** Minimum valid threshold for auto-clear/compact (1000 tokens) */ const MIN_AUTO_THRESHOLD = 1000; @@ -181,37 +232,35 @@ export class SessionAutoOps extends EventEmitter { `[SessionAutoOps] Auto-compact triggered: ${totalTokens} tokens >= ${this._autoCompactThreshold} threshold` ); - const checkAndCompact = async () => { - if (this.callbacks.isStopped()) return; - if (!this._isCompacting) return; - - if (!this.callbacks.isWorking()) { - if (this.callbacks.isStopped()) return; - - const compactCmd = this._autoCompactPrompt ? `/compact ${this._autoCompactPrompt}\r` : '/compact\r'; - await this.callbacks.writeCommand(compactCmd); - this.emit('autoCompact', { - tokens: totalTokens, - threshold: this._autoCompactThreshold, - prompt: this._autoCompactPrompt || undefined, - }); - - if (!this.callbacks.isStopped()) { - this._autoCompactTimer = setTimeout(() => { - if (this.callbacks.isStopped()) return; - this._autoCompactTimer = null; - this._isCompacting = false; - }, COMPACT_COOLDOWN_MS); - } - } else { - if (!this.callbacks.isStopped()) { - this._autoCompactTimer = setTimeout(checkAndCompact, AUTO_RETRY_DELAY_MS); - } - } + const action = async () => { + const compactCmd = this._autoCompactPrompt ? `/compact ${this._autoCompactPrompt}\r` : '/compact\r'; + await this.callbacks.writeCommand(compactCmd); + this.emit('autoCompact', { + tokens: totalTokens, + threshold: this._autoCompactThreshold, + prompt: this._autoCompactPrompt || undefined, + }); }; if (!this.callbacks.isStopped()) { - this._autoCompactTimer = setTimeout(checkAndCompact, AUTO_INITIAL_DELAY_MS); + this._autoCompactTimer = setTimeout( + () => + executeWhenIdle( + action, + () => this._isCompacting, + () => this.callbacks.isWorking(), + () => this.callbacks.isStopped(), + AUTO_RETRY_DELAY_MS, + COMPACT_COOLDOWN_MS, + (timer) => { + this._autoCompactTimer = timer; + }, + () => { + this._isCompacting = false; + } + ), + AUTO_INITIAL_DELAY_MS + ); } } } @@ -231,32 +280,30 @@ export class SessionAutoOps extends EventEmitter { `[SessionAutoOps] Auto-clear triggered: ${totalTokens} tokens >= ${this._autoClearThreshold} threshold` ); - const checkAndClear = async () => { - if (this.callbacks.isStopped()) return; - if (!this._isClearing) return; - - if (!this.callbacks.isWorking()) { - if (this.callbacks.isStopped()) return; - - await this.callbacks.writeCommand('/clear\r'); - this.emit('autoClear', { tokens: totalTokens, threshold: this._autoClearThreshold }); - - if (!this.callbacks.isStopped()) { - this._autoClearTimer = setTimeout(() => { - if (this.callbacks.isStopped()) return; - this._autoClearTimer = null; - this._isClearing = false; - }, CLEAR_COOLDOWN_MS); - } - } else { - if (!this.callbacks.isStopped()) { - this._autoClearTimer = setTimeout(checkAndClear, AUTO_RETRY_DELAY_MS); - } - } + const action = async () => { + await this.callbacks.writeCommand('/clear\r'); + this.emit('autoClear', { tokens: totalTokens, threshold: this._autoClearThreshold }); }; if (!this.callbacks.isStopped()) { - this._autoClearTimer = setTimeout(checkAndClear, AUTO_INITIAL_DELAY_MS); + this._autoClearTimer = setTimeout( + () => + executeWhenIdle( + action, + () => this._isClearing, + () => this.callbacks.isWorking(), + () => this.callbacks.isStopped(), + AUTO_RETRY_DELAY_MS, + CLEAR_COOLDOWN_MS, + (timer) => { + this._autoClearTimer = timer; + }, + () => { + this._isClearing = false; + } + ), + AUTO_INITIAL_DELAY_MS + ); } } } diff --git a/src/session.ts b/src/session.ts index 1296fe5f..242bc03a 100644 --- a/src/session.ts +++ b/src/session.ts @@ -876,13 +876,7 @@ export class Session extends EventEmitter { throw new Error('Session already has a running process'); } - this._status = 'busy'; - this._terminalBuffer.clear(); - this._textOutput.clear(); - this._errorBuffer = ''; - this._messages = []; - this._lineBuffer = ''; - this._lastActivityAt = Date.now(); + this._resetBuffers(); const modeLabel = this.mode === 'opencode' ? 'OpenCode' : 'Claude'; console.log( @@ -1257,13 +1251,7 @@ export class Session extends EventEmitter { throw new Error('Session already has a running process'); } - this._status = 'busy'; - this._terminalBuffer.clear(); - this._textOutput.clear(); - this._errorBuffer = ''; - this._messages = []; - this._lineBuffer = ''; - this._lastActivityAt = Date.now(); + this._resetBuffers(); // Use user's default shell or bash const shell = process.env.SHELL || '/bin/bash'; @@ -1448,13 +1436,7 @@ export class Session extends EventEmitter { return; } - this._status = 'busy'; - this._terminalBuffer.clear(); - this._textOutput.clear(); - this._errorBuffer = ''; - this._messages = []; - this._lineBuffer = ''; - this._lastActivityAt = Date.now(); + this._resetBuffers(); this._promptResolved = false; // Reset race condition guard this.resolvePromise = resolve; @@ -1565,6 +1547,117 @@ export class Session extends EventEmitter { }); } + private _resetBuffers(): void { + this._status = 'busy'; + this._terminalBuffer.clear(); + this._textOutput.clear(); + this._errorBuffer = ''; + this._messages = []; + this._lineBuffer = ''; + this._lastActivityAt = Date.now(); + } + + private _clearAllTimers(): void { + // Clear activity timeout to prevent memory leak + if (this.activityTimeout) { + clearTimeout(this.activityTimeout); + this.activityTimeout = null; + } + + // Clear line buffer flush timer + if (this._lineBufferFlushTimer) { + clearTimeout(this._lineBufferFlushTimer); + this._lineBufferFlushTimer = null; + } + + // Destroy auto-compact/auto-clear automation (clears its timers) + this._autoOps.destroy(); + + // Clear prompt check timers + if (this._promptCheckInterval) { + clearInterval(this._promptCheckInterval); + this._promptCheckInterval = null; + } + if (this._promptCheckTimeout) { + clearTimeout(this._promptCheckTimeout); + this._promptCheckTimeout = null; + } + + // Clear shell idle timer + if (this._shellIdleTimer) { + clearTimeout(this._shellIdleTimer); + this._shellIdleTimer = null; + } + + // Clear expensive processing timer + if (this._expensiveProcessTimer) { + clearTimeout(this._expensiveProcessTimer); + this._expensiveProcessTimer = null; + } + this._pendingCleanData = ''; + } + + private _handleJsonMessage(cleanLine: string, rawLine: string): void { + try { + const msg = JSON.parse(cleanLine) as ClaudeMessage; + this._messages.push(msg); + this.emit('message', msg); + + // Trim messages array for long-running sessions + if (this._messages.length > MAX_MESSAGES) { + this._messages = this._messages.slice(-Math.floor(MAX_MESSAGES * 0.8)); + } + + // Extract Claude session ID from messages (can be in any message type) + // Support both sessionId (camelCase) and session_id (snake_case) + const msgSessionId = + ((msg as unknown as Record).sessionId as string | undefined) ?? msg.session_id; + if (msgSessionId && !this._claudeSessionId) { + this._claudeSessionId = msgSessionId; + } + + // Process message for task tracking + this._taskTracker.processMessage(msg); + + if (msg.type === 'assistant' && msg.message?.content) { + for (const block of msg.message.content) { + if (block.type === 'text' && block.text) { + this._textOutput.append(block.text); + } + } + // Track tokens from usage (with validation) + if (msg.message.usage) { + const inputDelta = msg.message.usage.input_tokens || 0; + const outputDelta = msg.message.usage.output_tokens || 0; + + // Sanity check: max 100k tokens per message (generous limit) + const MAX_TOKENS_PER_MESSAGE = 100_000; + if (inputDelta > 0 && inputDelta <= MAX_TOKENS_PER_MESSAGE) { + this._totalInputTokens += inputDelta; + } + if (outputDelta > 0 && outputDelta <= MAX_TOKENS_PER_MESSAGE) { + this._totalOutputTokens += outputDelta; + } + + // Check if we should auto-compact or auto-clear + this._autoOps.checkAutoCompact(); + this._autoOps.checkAutoClear(); + } + } + + if (msg.type === 'result' && msg.total_cost_usd) { + this._totalCost = msg.total_cost_usd; + } + } catch (parseErr) { + // Not JSON, just regular output - this is expected for non-JSON lines + console.debug( + '[Session] Line not JSON (expected for text output):', + parseErr instanceof Error ? parseErr.message : parseErr + ); + this._textOutput.append(rawLine + '\n'); + } + } + private processOutput(data: string): void { // Early return if session is stopped to prevent any processing or timer creation if (this._isStopped) return; @@ -1606,64 +1699,7 @@ export class Session extends EventEmitter { const cleanLine = trimmed.replace(ANSI_ESCAPE_PATTERN_FULL, ''); if (cleanLine.startsWith('{') && cleanLine.endsWith('}')) { - try { - const msg = JSON.parse(cleanLine) as ClaudeMessage; - this._messages.push(msg); - this.emit('message', msg); - - // Trim messages array for long-running sessions - if (this._messages.length > MAX_MESSAGES) { - this._messages = this._messages.slice(-Math.floor(MAX_MESSAGES * 0.8)); - } - - // Extract Claude session ID from messages (can be in any message type) - // Support both sessionId (camelCase) and session_id (snake_case) - const msgSessionId = - ((msg as unknown as Record).sessionId as string | undefined) ?? msg.session_id; - if (msgSessionId && !this._claudeSessionId) { - this._claudeSessionId = msgSessionId; - } - - // Process message for task tracking - this._taskTracker.processMessage(msg); - - if (msg.type === 'assistant' && msg.message?.content) { - for (const block of msg.message.content) { - if (block.type === 'text' && block.text) { - this._textOutput.append(block.text); - } - } - // Track tokens from usage (with validation) - if (msg.message.usage) { - const inputDelta = msg.message.usage.input_tokens || 0; - const outputDelta = msg.message.usage.output_tokens || 0; - - // Sanity check: max 100k tokens per message (generous limit) - const MAX_TOKENS_PER_MESSAGE = 100_000; - if (inputDelta > 0 && inputDelta <= MAX_TOKENS_PER_MESSAGE) { - this._totalInputTokens += inputDelta; - } - if (outputDelta > 0 && outputDelta <= MAX_TOKENS_PER_MESSAGE) { - this._totalOutputTokens += outputDelta; - } - - // Check if we should auto-compact or auto-clear - this._autoOps.checkAutoCompact(); - this._autoOps.checkAutoClear(); - } - } - - if (msg.type === 'result' && msg.total_cost_usd) { - this._totalCost = msg.total_cost_usd; - } - } catch (parseErr) { - // Not JSON, just regular output - this is expected for non-JSON lines - console.debug( - '[Session] Line not JSON (expected for text output):', - parseErr instanceof Error ? parseErr.message : parseErr - ); - this._textOutput.append(line + '\n'); - } + this._handleJsonMessage(cleanLine, line); } else if (trimmed) { this._textOutput.append(line + '\n'); } @@ -2030,43 +2066,7 @@ export class Session extends EventEmitter { // Set stopped flag first to prevent new timers from being created this._isStopped = true; - // Clear activity timeout to prevent memory leak - if (this.activityTimeout) { - clearTimeout(this.activityTimeout); - this.activityTimeout = null; - } - - // Clear line buffer flush timer - if (this._lineBufferFlushTimer) { - clearTimeout(this._lineBufferFlushTimer); - this._lineBufferFlushTimer = null; - } - - // Destroy auto-compact/auto-clear automation (clears its timers) - this._autoOps.destroy(); - - // Clear prompt check timers - if (this._promptCheckInterval) { - clearInterval(this._promptCheckInterval); - this._promptCheckInterval = null; - } - if (this._promptCheckTimeout) { - clearTimeout(this._promptCheckTimeout); - this._promptCheckTimeout = null; - } - - // Clear shell idle timer - if (this._shellIdleTimer) { - clearTimeout(this._shellIdleTimer); - this._shellIdleTimer = null; - } - - // Clear expensive processing timer - if (this._expensiveProcessTimer) { - clearTimeout(this._expensiveProcessTimer); - this._expensiveProcessTimer = null; - } - this._pendingCleanData = ''; + this._clearAllTimers(); // Immediately cleanup Promise callbacks to prevent orphaned references // during the rest of stop() processing (e.g., if mux kill times out) diff --git a/src/state-store.ts b/src/state-store.ts index f238776d..f53d784b 100644 --- a/src/state-store.ts +++ b/src/state-store.ts @@ -195,16 +195,7 @@ export class StateStore { * Only dirty sessions are re-serialized; clean sessions use cached JSON fragments. */ private assembleStateJson(): string { - // Re-serialize dirty sessions and update cache - for (const id of this.dirtySessions) { - const session = this.state.sessions[id]; - if (session) { - this.cachedSessionJsons.set(id, JSON.stringify(session)); - } else { - this.cachedSessionJsons.delete(id); - } - } - this.dirtySessions.clear(); + this.updateDirtySessionCache(); // Build sessions object from cached fragments const sessionParts: string[] = []; @@ -218,6 +209,25 @@ export class StateStore { sessionParts.push(`${JSON.stringify(id)}:${json}`); } + this.pruneStaleCacheEntries(); + + return this.buildPartialJson(sessionParts); + } + + private updateDirtySessionCache(): void { + // Re-serialize dirty sessions and update cache + for (const id of this.dirtySessions) { + const session = this.state.sessions[id]; + if (session) { + this.cachedSessionJsons.set(id, JSON.stringify(session)); + } else { + this.cachedSessionJsons.delete(id); + } + } + this.dirtySessions.clear(); + } + + private pruneStaleCacheEntries(): void { // Prune stale cache entries (sessions removed via direct state mutation) if (this.cachedSessionJsons.size > Object.keys(this.state.sessions).length) { for (const cachedId of this.cachedSessionJsons.keys()) { @@ -226,7 +236,9 @@ export class StateStore { } } } + } + private buildPartialJson(sessionParts: string[]): string { // Build final JSON: sessions from cache, everything else re-serialized (tiny) const sessionsJson = `{${sessionParts.join(',')}}`; @@ -249,6 +261,28 @@ export class StateStore { return `{${parts.join(',')}}`; } + private serializeState(): string | null { + try { + return this.assembleStateJson(); + } catch (assembleErr) { + // Fallback to full serialization if incremental assembly fails + console.warn('[StateStore] assembleStateJson failed, falling back to full serialize:', assembleErr); + this.cachedSessionJsons.clear(); + this.dirtySessions.clear(); + try { + return JSON.stringify(this.state); + } catch (err) { + console.error('[StateStore] Failed to serialize state (circular reference or invalid data):', err); + this.consecutiveSaveFailures++; + if (this.consecutiveSaveFailures >= MAX_CONSECUTIVE_FAILURES) { + console.error('[StateStore] Circuit breaker OPEN - serialization failing repeatedly'); + this.circuitBreakerOpen = true; + } + return null; + } + } + } + private async _doSaveAsync(): Promise { this.saveDeb.cancel(); if (!this.dirty) { @@ -265,28 +299,10 @@ export class StateStore { const tempPath = this.filePath + '.tmp'; const backupPath = this.filePath + '.bak'; - let json: string; // Step 1: Serialize state (validates it's JSON-safe) - try { - json = this.assembleStateJson(); - } catch (assembleErr) { - // Fallback to full serialization if incremental assembly fails - console.warn('[StateStore] assembleStateJson failed, falling back to full serialize:', assembleErr); - this.cachedSessionJsons.clear(); - this.dirtySessions.clear(); - try { - json = JSON.stringify(this.state); - } catch (err) { - console.error('[StateStore] Failed to serialize state (circular reference or invalid data):', err); - this.consecutiveSaveFailures++; - if (this.consecutiveSaveFailures >= MAX_CONSECUTIVE_FAILURES) { - console.error('[StateStore] Circuit breaker OPEN - serialization failing repeatedly'); - this.circuitBreakerOpen = true; - } - return; - } - } + const json = this.serializeState(); + if (json === null) return; // Clear dirty flag BEFORE async I/O so mutations during write re-set it. // The state snapshot is already captured in `json` above. @@ -351,27 +367,9 @@ export class StateStore { const tempPath = this.filePath + '.tmp'; const backupPath = this.filePath + '.bak'; - let json: string; - try { - json = this.assembleStateJson(); - } catch (assembleErr) { - // Fallback to full serialization if incremental assembly fails - console.warn('[StateStore] assembleStateJson failed, falling back to full serialize:', assembleErr); - this.cachedSessionJsons.clear(); - this.dirtySessions.clear(); - try { - json = JSON.stringify(this.state); - } catch (err) { - console.error('[StateStore] Failed to serialize state (circular reference or invalid data):', err); - this.consecutiveSaveFailures++; - if (this.consecutiveSaveFailures >= MAX_CONSECUTIVE_FAILURES) { - console.error('[StateStore] Circuit breaker OPEN - serialization failing repeatedly'); - this.circuitBreakerOpen = true; - } - return; - } - } + const json = this.serializeState(); + if (json === null) return; // Backup via atomic copy (avoids reading entire file into memory) try { diff --git a/src/subagent-watcher.ts b/src/subagent-watcher.ts index 6079c23c..715186a9 100644 --- a/src/subagent-watcher.ts +++ b/src/subagent-watcher.ts @@ -222,6 +222,83 @@ export class SubagentWatcher extends EventEmitter { return INTERNAL_AGENT_PATTERNS.some((pattern) => pattern.test(description)); } + /** + * Mark a subagent as completed: clear PID, set status, clean up pending tool calls, emit event. + */ + private markSubagentAsCompleted(info: SubagentInfo): void { + info.pid = undefined; + info.status = 'completed'; + this.pendingToolCalls.delete(info.agentId); + this.emit('subagent:completed', info); + } + + /** + * Extract text from message content, handling both string and array formats. + * For array content, returns the text from the first 'text' block. + */ + private extractFirstTextContent( + content: string | Array<{ type: string; text?: string }> | undefined + ): string | undefined { + if (!content) return undefined; + if (typeof content === 'string') { + const trimmed = content.trim(); + return trimmed.length > 0 ? trimmed : undefined; + } + if (Array.isArray(content)) { + const firstContent = content[0]; + if (firstContent?.type === 'text' && firstContent.text) { + const trimmed = firstContent.text.trim(); + return trimmed.length > 0 ? trimmed : undefined; + } + } + return undefined; + } + + /** + * Process a tool_result content block: look up pending tool call, emit tool_result event. + */ + private emitToolResult( + content: { tool_use_id: string; content?: string | Array<{ type: string; text?: string }>; is_error?: boolean }, + agentId: string, + sessionId: string, + timestamp: string + ): void { + const resultContent = this.extractToolResultContent(content.content); + const agentPendingCalls = this.pendingToolCalls.get(agentId); + const pendingCall = agentPendingCalls?.get(content.tool_use_id); + const toolName = pendingCall?.toolName; + // Delete after lookup to prevent memory leak + agentPendingCalls?.delete(content.tool_use_id); + + const toolResult: SubagentToolResult = { + agentId, + sessionId, + timestamp, + toolUseId: content.tool_use_id, + tool: toolName, + preview: resultContent.substring(0, MESSAGE_TEXT_LIMIT), + contentLength: resultContent.length, + isError: content.is_error || false, + }; + this.emit('subagent:tool_result', toolResult); + } + + /** + * Find the oldest inactive (non-active) agent for LRU eviction. + * Returns the agent ID of the oldest inactive agent, or null if all are active. + */ + private findOldestInactiveAgent(): string | null { + let oldestId: string | null = null; + let oldestTime = Infinity; + for (const [id, existing] of this.agentInfo) { + if (existing.status !== 'active' && existing.lastActivityAt < oldestTime) { + oldestTime = existing.lastActivityAt; + oldestId = id; + } + } + return oldestId; + } + /** * Extract short model identifier from full model name */ @@ -307,10 +384,7 @@ export class SubagentWatcher extends EventEmitter { const alive = this.checkSubagentAliveFromPidMap(info, pidMap); if (!alive) { - info.pid = undefined; - info.status = 'completed'; - this.pendingToolCalls.delete(info.agentId); - this.emit('subagent:completed', info); + this.markSubagentAsCompleted(info); } } } @@ -677,10 +751,7 @@ export class SubagentWatcher extends EventEmitter { const pid = await this.findSubagentProcess(info.sessionId); if (pid) { process.kill(pid, 'SIGTERM'); - info.pid = undefined; - info.status = 'completed'; - this.pendingToolCalls.delete(info.agentId); - this.emit('subagent:completed', info); + this.markSubagentAsCompleted(info); return true; } } catch { @@ -688,10 +759,7 @@ export class SubagentWatcher extends EventEmitter { } // Mark as completed even if we couldn't find the process - info.pid = undefined; - info.status = 'completed'; - this.pendingToolCalls.delete(info.agentId); - this.emit('subagent:completed', info); + this.markSubagentAsCompleted(info); return true; } @@ -843,19 +911,9 @@ export class SubagentWatcher extends EventEmitter { } } else if (entry.type === 'user' && entry.message?.content) { // Handle both string and array content formats - if (typeof entry.message.content === 'string') { - const text = entry.message.content.trim(); - if (text.length < 100 && !text.includes('{')) { - lines.push(`${this.formatTime(entry.timestamp)} 📥 User: ${text.substring(0, USER_TEXT_PREVIEW_LENGTH)}`); - } - } else { - const firstContent = entry.message.content[0]; - if (firstContent?.type === 'text' && firstContent.text) { - const text = firstContent.text.trim(); - if (text.length < 100 && !text.includes('{')) { - lines.push(`${this.formatTime(entry.timestamp)} 📥 User: ${text.substring(0, USER_TEXT_PREVIEW_LENGTH)}`); - } - } + const text = this.extractFirstTextContent(entry.message.content); + if (text && text.length < 100 && !text.includes('{')) { + lines.push(`${this.formatTime(entry.timestamp)} 📥 User: ${text.substring(0, USER_TEXT_PREVIEW_LENGTH)}`); } } } @@ -1027,15 +1085,7 @@ export class SubagentWatcher extends EventEmitter { try { const entry = JSON.parse(line); if (entry.type === 'user' && entry.message?.content) { - let text: string | undefined; - if (typeof entry.message.content === 'string') { - text = entry.message.content.trim(); - } else if (Array.isArray(entry.message.content)) { - const firstContent = entry.message.content[0]; - if (firstContent?.type === 'text' && firstContent.text) { - text = firstContent.text.trim(); - } - } + const text = this.extractFirstTextContent(entry.message.content); if (text) { resolved = true; rl.close(); @@ -1286,14 +1336,7 @@ export class SubagentWatcher extends EventEmitter { // Enforce MAX_TRACKED_AGENTS during insertion — evict oldest inactive agent if (this.agentInfo.size >= MAX_TRACKED_AGENTS) { - let oldestId: string | null = null; - let oldestTime = Infinity; - for (const [id, existing] of this.agentInfo) { - if (existing.status !== 'active' && existing.lastActivityAt < oldestTime) { - oldestTime = existing.lastActivityAt; - oldestId = id; - } - } + const oldestId = this.findOldestInactiveAgent(); if (oldestId) { this.removeAgent(oldestId); } @@ -1386,15 +1429,7 @@ export class SubagentWatcher extends EventEmitter { let description = await this.extractDescriptionFromParentTranscript(info.projectHash, info.sessionId, agentId); // Fallback: extract smart title from the prompt content if (!description) { - let text: string | undefined; - if (typeof entry.message.content === 'string') { - text = entry.message.content.trim(); - } else if (Array.isArray(entry.message.content)) { - const firstContent = entry.message.content[0]; - if (firstContent?.type === 'text' && firstContent.text) { - text = firstContent.text.trim(); - } - } + const text = this.extractFirstTextContent(entry.message.content); if (text) { description = this.extractSmartTitle(text); } @@ -1480,24 +1515,12 @@ export class SubagentWatcher extends EventEmitter { } } else if (content.type === 'tool_result' && content.tool_use_id) { // Extract tool result - const resultContent = this.extractToolResultContent(content.content); - const agentPendingCalls = this.pendingToolCalls.get(agentId); - const pendingCall = agentPendingCalls?.get(content.tool_use_id); - const toolName = pendingCall?.toolName; - // Delete after lookup to prevent memory leak - agentPendingCalls?.delete(content.tool_use_id); - - const toolResult: SubagentToolResult = { + this.emitToolResult( + { tool_use_id: content.tool_use_id, content: content.content, is_error: content.is_error }, agentId, sessionId, - timestamp: entry.timestamp, - toolUseId: content.tool_use_id, - tool: toolName, - preview: resultContent.substring(0, MESSAGE_TEXT_LIMIT), - contentLength: resultContent.length, - isError: content.is_error || false, - }; - this.emit('subagent:tool_result', toolResult); + entry.timestamp + ); } else if (content.type === 'text' && content.text) { const text = content.text.trim(); if (text.length > 0) { @@ -1531,24 +1554,12 @@ export class SubagentWatcher extends EventEmitter { // Check for tool_result blocks in user messages (common pattern) for (const content of entry.message.content) { if (content.type === 'tool_result' && content.tool_use_id) { - const resultContent = this.extractToolResultContent(content.content); - const agentPendingCalls = this.pendingToolCalls.get(agentId); - const pendingCall = agentPendingCalls?.get(content.tool_use_id); - const toolName = pendingCall?.toolName; - // Delete after lookup to prevent memory leak - agentPendingCalls?.delete(content.tool_use_id); - - const toolResult: SubagentToolResult = { + this.emitToolResult( + { tool_use_id: content.tool_use_id, content: content.content, is_error: content.is_error }, agentId, sessionId, - timestamp: entry.timestamp, - toolUseId: content.tool_use_id, - tool: toolName, - preview: resultContent.substring(0, MESSAGE_TEXT_LIMIT), - contentLength: resultContent.length, - isError: content.is_error || false, - }; - this.emit('subagent:tool_result', toolResult); + entry.timestamp + ); } else if (content.type === 'text' && content.text) { const userText = content.text.trim(); if (userText.length > 0 && userText.length < 500) { diff --git a/src/tmux-manager.ts b/src/tmux-manager.ts index a8681381..640595b3 100644 --- a/src/tmux-manager.ts +++ b/src/tmux-manager.ts @@ -98,6 +98,9 @@ const LEGACY_MUX_NAME_PATTERN = /^claudeman-[a-f0-9-]+$/; /** Regex to validate tmux pane targets (e.g., "%0", "%1", "0", "1") */ const SAFE_PANE_TARGET_PATTERN = /^(%\d+|\d+)$/; +/** Characters unsafe in paths — shell metacharacters, quotes, and control chars */ +const UNSAFE_PATH_CHARS = /[;&|$`(){}<>'"\n\r]/; + /** * Validates that a session name contains only safe characters. * Prevents command injection via malformed session IDs. @@ -111,23 +114,7 @@ function isValidMuxName(name: string): boolean { * Prevents command injection via malformed paths. */ function isValidPath(path: string): boolean { - if ( - path.includes(';') || - path.includes('&') || - path.includes('|') || - path.includes('$') || - path.includes('`') || - path.includes('(') || - path.includes(')') || - path.includes('{') || - path.includes('}') || - path.includes('<') || - path.includes('>') || - path.includes("'") || - path.includes('"') || - path.includes('\n') || - path.includes('\r') - ) { + if (UNSAFE_PATH_CHARS.test(path)) { return false; } if (path.includes('..')) {