diff --git a/src/bash-tool-parser.ts b/src/bash-tool-parser.ts index 4ec5cdfa..76586c4e 100644 --- a/src/bash-tool-parser.ts +++ b/src/bash-tool-parser.ts @@ -470,120 +470,91 @@ export class BashToolParser extends EventEmitter { * Process a single pre-stripped line of terminal output. */ private processCleanLine(cleanLine: string): void { - // Check for tool start + if (this._handleToolStart(cleanLine)) return; + if (this._handleToolCompletion(cleanLine)) return; + if (this._handleTextCommand(cleanLine)) return; + this._handleLogFileMention(cleanLine); + } + + private _handleToolStart(cleanLine: string): boolean { const startMatch = cleanLine.match(BASH_TOOL_START_PATTERN); - if (startMatch) { - const command = startMatch[1]; - const timeout = startMatch[2]?.trim(); + if (!startMatch) return false; - // Check if this is a file-viewing command - if (this.isFileViewerCommand(command)) { - const filePaths = this.extractFilePaths(command); + const command = startMatch[1]; + const timeout = startMatch[2]?.trim(); - // Skip if any file path is already tracked (cross-pattern dedup) - if (filePaths.some((fp) => this.isFilePathTracked(fp))) { - return; - } + if (!this.isFileViewerCommand(command)) return true; - if (filePaths.length > 0) { - const tool: ActiveBashTool = { - id: uuidv4(), - command, - filePaths, - timeout, - startedAt: Date.now(), - status: 'running', - sessionId: this._sessionId, - }; + const filePaths = this.extractFilePaths(command); - // Enforce max tools limit - if (this._activeTools.size >= MAX_ACTIVE_TOOLS) { - // Remove oldest tool (O(n) min-scan instead of O(n log n) sort) - let oldestKey: string | undefined; - let oldestTime = Infinity; - for (const [key, entry] of this._activeTools) { - if (entry.startedAt < oldestTime) { - oldestTime = entry.startedAt; - oldestKey = key; - } - } - if (oldestKey) { - this._activeTools.delete(oldestKey); - } + // Skip if any file path is already tracked (cross-pattern dedup) + if (filePaths.some((fp) => this.isFilePathTracked(fp))) return true; + + if (filePaths.length > 0) { + const tool = this._createActiveTool(command, filePaths, 'running', timeout); + + // Enforce max tools limit + if (this._activeTools.size >= MAX_ACTIVE_TOOLS) { + // Remove oldest tool (O(n) min-scan instead of O(n log n) sort) + let oldestKey: string | undefined; + let oldestTime = Infinity; + for (const [key, entry] of this._activeTools) { + if (entry.startedAt < oldestTime) { + oldestTime = entry.startedAt; + oldestKey = key; } - - this._activeTools.set(tool.id, tool); - this._lastToolId = tool.id; - - this.emit('toolStart', tool); - this.scheduleUpdate(); } - } - return; - } - - // Check for tool completion - if (TOOL_COMPLETION_PATTERN.test(cleanLine) && this._lastToolId) { - const tool = this._activeTools.get(this._lastToolId); - if (tool && tool.status === 'running') { - tool.status = 'completed'; - this.emit('toolEnd', tool); - this.scheduleUpdate(); - - // Remove completed tool after a short delay to allow UI to show completion - this.cleanup.setTimeout( - () => { - if (this._destroyed) return; - this._activeTools.delete(tool.id); - this.scheduleUpdate(); - }, - 2000, - { description: 'auto-remove completed tool' } - ); - } - this._lastToolId = null; - return; - } - - // Fallback: Check for command suggestions in plain text (e.g., "tail -f /tmp/file.log") - const textCmdMatch = cleanLine.match(TEXT_COMMAND_PATTERN); - if (textCmdMatch) { - const filePath = textCmdMatch[2]; - - // Create a suggestion tool (marked as 'suggestion' status) - const tool: ActiveBashTool = { - id: uuidv4(), - command: cleanLine.trim(), - filePaths: [filePath], - timeout: undefined, - startedAt: Date.now(), - status: 'running', // Shows as clickable - sessionId: this._sessionId, - }; - - // Don't add if file path already tracked (cross-pattern dedup) - if (this.isFilePathTracked(filePath)) { - return; + if (oldestKey) { + this._activeTools.delete(oldestKey); + } } this._activeTools.set(tool.id, tool); + this._lastToolId = tool.id; + this.emit('toolStart', tool); this.scheduleUpdate(); - - // Auto-remove suggestions after 30 seconds - this.cleanup.setTimeout( - () => { - if (this._destroyed) return; - this._activeTools.delete(tool.id); - this.scheduleUpdate(); - }, - 30000, - { description: 'auto-remove suggestion tool' } - ); - return; } - // Last fallback: Check for log file paths mentioned anywhere in the line + return true; + } + + private _handleToolCompletion(cleanLine: string): boolean { + if (!TOOL_COMPLETION_PATTERN.test(cleanLine) || !this._lastToolId) return false; + + const tool = this._activeTools.get(this._lastToolId); + if (tool && tool.status === 'running') { + tool.status = 'completed'; + this.emit('toolEnd', tool); + this.scheduleUpdate(); + + this._scheduleAutoRemove(tool.id, 2000, 'auto-remove completed tool'); + } + this._lastToolId = null; + return true; + } + + private _handleTextCommand(cleanLine: string): boolean { + const textCmdMatch = cleanLine.match(TEXT_COMMAND_PATTERN); + if (!textCmdMatch) return false; + + const filePath = textCmdMatch[2]; + + // Don't add if file path already tracked (cross-pattern dedup) + if (this.isFilePathTracked(filePath)) return true; + + const tool = this._createActiveTool(cleanLine.trim(), [filePath], 'running'); + + this._activeTools.set(tool.id, tool); + this.emit('toolStart', tool); + this.scheduleUpdate(); + + // Auto-remove suggestions after 30 seconds + this._scheduleAutoRemove(tool.id, 30000, 'auto-remove suggestion tool'); + return true; + } + + private _handleLogFileMention(cleanLine: string): void { LOG_FILE_MENTION_PATTERN.lastIndex = 0; let logMatch; while ((logMatch = LOG_FILE_MENTION_PATTERN.exec(cleanLine)) !== null) { @@ -595,33 +566,46 @@ export class BashToolParser extends EventEmitter { // Skip if file path already tracked (cross-pattern dedup) if (this.isFilePathTracked(filePath)) continue; - const tool: ActiveBashTool = { - id: uuidv4(), - command: `View: ${filePath}`, - filePaths: [filePath], - timeout: undefined, - startedAt: Date.now(), - status: 'running', - sessionId: this._sessionId, - }; + const tool = this._createActiveTool(`View: ${filePath}`, [filePath], 'running'); this._activeTools.set(tool.id, tool); this.emit('toolStart', tool); this.scheduleUpdate(); // Auto-remove after 60 seconds - this.cleanup.setTimeout( - () => { - if (this._destroyed) return; - this._activeTools.delete(tool.id); - this.scheduleUpdate(); - }, - 60000, - { description: 'auto-remove log file tool' } - ); + this._scheduleAutoRemove(tool.id, 60000, 'auto-remove log file tool'); } } + private _createActiveTool( + command: string, + filePaths: string[], + status: ActiveBashTool['status'], + timeout?: string + ): ActiveBashTool { + return { + id: uuidv4(), + command, + filePaths, + timeout, + startedAt: Date.now(), + status, + sessionId: this._sessionId, + }; + } + + private _scheduleAutoRemove(toolId: string, delayMs: number, description: string): void { + this.cleanup.setTimeout( + () => { + if (this._destroyed) return; + this._activeTools.delete(toolId); + this.scheduleUpdate(); + }, + delayMs, + { description } + ); + } + /** * Check if a command is a file-viewing command worth tracking. */ diff --git a/src/orchestrator-loop.ts b/src/orchestrator-loop.ts index 19def538..2d7ccdc2 100644 --- a/src/orchestrator-loop.ts +++ b/src/orchestrator-loop.ts @@ -499,27 +499,34 @@ export class OrchestratorLoop extends EventEmitter { } } - private handleTaskCompleted(queueTaskId: string): void { + private _finalizeTask(queueTaskId: string, status: 'completed' | 'failed', error?: string): OrchestratorTask | null { const orchTask = this.findOrchestratorTaskByQueueId(queueTaskId); + if (!orchTask) return null; + + orchTask.status = status; + if (status === 'completed') { + orchTask.completedAt = Date.now(); + this.stats.totalTasksCompleted++; + } else { + orchTask.error = error ?? null; + this.stats.totalTasksFailed++; + } + this.persist(); + return orchTask; + } + + private handleTaskCompleted(queueTaskId: string): void { + const orchTask = this._finalizeTask(queueTaskId, 'completed'); if (!orchTask) return; - orchTask.status = 'completed'; - orchTask.completedAt = Date.now(); - this.stats.totalTasksCompleted++; - this.persist(); this.emit('taskCompleted', orchTask); - this.checkPhaseCompletion(); } private handleTaskFailed(queueTaskId: string, error: string): void { - const orchTask = this.findOrchestratorTaskByQueueId(queueTaskId); + const orchTask = this._finalizeTask(queueTaskId, 'failed', error); if (!orchTask) return; - orchTask.status = 'failed'; - orchTask.error = error; - this.stats.totalTasksFailed++; - this.persist(); this.emit('taskFailed', orchTask, error); // Check if we should retry the task or fail the phase @@ -553,19 +560,20 @@ export class OrchestratorLoop extends EventEmitter { }, this.config.phaseTimeoutMs); } + private _clearTimer( + timerKey: 'phasePollTimer' | 'phaseTimeoutTimer' | 'postPhaseTimer', + clearFn: typeof clearInterval | typeof clearTimeout + ): void { + if (this[timerKey]) { + clearFn(this[timerKey]); + this[timerKey] = null; + } + } + private clearPhasePoll(): void { - if (this.phasePollTimer) { - clearInterval(this.phasePollTimer); - this.phasePollTimer = null; - } - if (this.phaseTimeoutTimer) { - clearTimeout(this.phaseTimeoutTimer); - this.phaseTimeoutTimer = null; - } - if (this.postPhaseTimer) { - clearTimeout(this.postPhaseTimer); - this.postPhaseTimer = null; - } + this._clearTimer('phasePollTimer', clearInterval); + this._clearTimer('phaseTimeoutTimer', clearTimeout); + this._clearTimer('postPhaseTimer', clearTimeout); } private pollPhaseStatus(phase: OrchestratorPhase): void { diff --git a/src/plan-orchestrator.ts b/src/plan-orchestrator.ts index 7c206498..4c649821 100644 --- a/src/plan-orchestrator.ts +++ b/src/plan-orchestrator.ts @@ -231,6 +231,49 @@ export class PlanOrchestrator { return md; } + private _extractJsonFromResponse(response: string): string | null { + let jsonMatch = response.match(/```(?:json)?\s*(\{[\s\S]*?\})\s*```/); + if (jsonMatch) { + jsonMatch = [jsonMatch[1]]; // Use captured group (inside code block) + } else { + jsonMatch = response.match(/\{[\s\S]*\}/); + } + return jsonMatch ? jsonMatch[0] : null; + } + + private _emitAgentFailure( + onSubagent: SubagentCallback | undefined, + agentId: string, + agentType: 'research' | 'planner', + model: string, + error: string, + durationMs: number + ): void { + onSubagent?.({ + type: 'failed', + agentId, + agentType, + model, + status: 'failed', + error, + durationMs, + }); + } + + private _formatResearchSection( + parts: string[], + title: string, + items: unknown[], + formatter: (item: unknown) => string[] + ): void { + if (items.length === 0) return; + parts.push(title); + for (const item of items.slice(0, 5)) { + parts.push(...formatter(item)); + } + parts.push(''); + } + async cancel(): Promise { this.cancelled = true; // Stop all running sessions and await cleanup to prevent PTY process leaks @@ -322,32 +365,23 @@ export class PlanOrchestrator { const parts: string[] = ['## Research Context\n']; - if (research.findings.externalResources.length > 0) { - parts.push('### External Resources'); - for (const r of research.findings.externalResources.slice(0, 5)) { - parts.push(`- ${r.title}${r.url ? ` (${r.url})` : ''}`); - if (r.keyInsights.length > 0) { - parts.push(` Key insights: ${r.keyInsights.slice(0, 3).join(', ')}`); - } + this._formatResearchSection(parts, '### External Resources', research.findings.externalResources, (item) => { + const r = item as ResearchResult['findings']['externalResources'][number]; + const lines = [`- ${r.title}${r.url ? ` (${r.url})` : ''}`]; + if (r.keyInsights.length > 0) { + lines.push(` Key insights: ${r.keyInsights.slice(0, 3).join(', ')}`); } - parts.push(''); - } + return lines; + }); - if (research.findings.codebasePatterns.length > 0) { - parts.push('### Existing Codebase Patterns'); - for (const p of research.findings.codebasePatterns.slice(0, 5)) { - parts.push(`- ${p.pattern} at ${p.location}`); - } - parts.push(''); - } + this._formatResearchSection(parts, '### Existing Codebase Patterns', research.findings.codebasePatterns, (item) => { + const p = item as ResearchResult['findings']['codebasePatterns'][number]; + return [`- ${p.pattern} at ${p.location}`]; + }); - if (research.findings.technicalRecommendations.length > 0) { - parts.push('### Recommendations'); - for (const r of research.findings.technicalRecommendations.slice(0, 5)) { - parts.push(`- ${r}`); - } - parts.push(''); - } + this._formatResearchSection(parts, '### Recommendations', research.findings.technicalRecommendations, (item) => [ + `- ${item as string}`, + ]); return parts.join('\n'); } @@ -420,27 +454,14 @@ export class PlanOrchestrator { ); // Extract JSON from response — try multiple strategies - let jsonMatch = response.match(/```(?:json)?\s*(\{[\s\S]*?\})\s*```/); - if (jsonMatch) { - jsonMatch = [jsonMatch[1]]; // Use captured group (inside code block) - } else { - jsonMatch = response.match(/\{[\s\S]*\}/); - } + const jsonStr = this._extractJsonFromResponse(response); - if (!jsonMatch) { + if (!jsonStr) { console.error( `[PlanOrchestrator] No JSON found in research response. Full response:`, response.substring(0, 2000) ); - onSubagent?.({ - type: 'failed', - agentId, - agentType: 'research', - model: this.researchModel, - status: 'failed', - error: 'No JSON found', - durationMs, - }); + this._emitAgentFailure(onSubagent, agentId, 'research', this.researchModel, 'No JSON found', durationMs); return { success: false, findings: { @@ -456,17 +477,9 @@ export class PlanOrchestrator { }; } - const parsed = tryParseJSON(jsonMatch[0]); + const parsed = tryParseJSON(jsonStr); if (!parsed.success) { - onSubagent?.({ - type: 'failed', - agentId, - agentType: 'research', - model: this.researchModel, - status: 'failed', - error: parsed.error, - durationMs, - }); + this._emitAgentFailure(onSubagent, agentId, 'research', this.researchModel, parsed.error!, durationMs); return { success: false, findings: { @@ -511,15 +524,7 @@ export class PlanOrchestrator { } catch (err) { const durationMs = Date.now() - startTime; const error = err instanceof Error ? err.message : String(err); - onSubagent?.({ - type: 'failed', - agentId, - agentType: 'research', - model: this.researchModel, - status: 'failed', - error, - durationMs, - }); + this._emitAgentFailure(onSubagent, agentId, 'research', this.researchModel, error, durationMs); return { success: false, findings: { @@ -608,41 +613,20 @@ export class PlanOrchestrator { ); // Extract JSON from response — try multiple strategies - let jsonMatch = response.match(/```(?:json)?\s*(\{[\s\S]*?\})\s*```/); - if (jsonMatch) { - jsonMatch = [jsonMatch[1]]; // Use captured group (inside code block) - } else { - jsonMatch = response.match(/\{[\s\S]*\}/); - } + const jsonStr = this._extractJsonFromResponse(response); - if (!jsonMatch) { + if (!jsonStr) { console.error( `[PlanOrchestrator] No JSON found in planner response. Full response:`, response.substring(0, 2000) ); - onSubagent?.({ - type: 'failed', - agentId, - agentType: 'planner', - model: this.plannerModel, - status: 'failed', - error: 'No JSON found', - durationMs, - }); + this._emitAgentFailure(onSubagent, agentId, 'planner', this.plannerModel, 'No JSON found', durationMs); return { success: false, error: 'No JSON in response' }; } - const parsed = tryParseJSON(jsonMatch[0]); + const parsed = tryParseJSON(jsonStr); if (!parsed.success) { - onSubagent?.({ - type: 'failed', - agentId, - agentType: 'planner', - model: this.plannerModel, - status: 'failed', - error: parsed.error, - durationMs, - }); + this._emitAgentFailure(onSubagent, agentId, 'planner', this.plannerModel, parsed.error!, durationMs); return { success: false, error: parsed.error }; } @@ -668,15 +652,7 @@ export class PlanOrchestrator { } catch (err) { const durationMs = Date.now() - startTime; const error = err instanceof Error ? err.message : String(err); - onSubagent?.({ - type: 'failed', - agentId, - agentType: 'planner', - model: this.plannerModel, - status: 'failed', - error, - durationMs, - }); + this._emitAgentFailure(onSubagent, agentId, 'planner', this.plannerModel, error, durationMs); return { success: false, error }; } finally { // Always clean up session and progress interval — centralizing here diff --git a/src/ralph-status-parser.ts b/src/ralph-status-parser.ts index d4993fd9..61ce795b 100644 --- a/src/ralph-status-parser.ts +++ b/src/ralph-status-parser.ts @@ -91,6 +91,67 @@ const COMPLETION_INDICATOR_PATTERNS = [ /project\s+(?:is\s+)?(?:completed?|done|finished)/i, ]; +interface FieldParser { + pattern: RegExp; + field: keyof RalphStatusBlock; + validate: (value: string) => boolean; + transform: (value: string) => T; + errorMsg: (value: string) => string; +} + +const FIELD_PARSERS: FieldParser[] = [ + { + pattern: RALPH_STATUS_FIELD_PATTERN, + field: 'status', + validate: (v) => ['IN_PROGRESS', 'COMPLETE', 'BLOCKED'].includes(v.toUpperCase()), + transform: (v) => v.toUpperCase() as RalphStatusValue, + errorMsg: (v) => `Invalid STATUS value: "${v}". Expected: IN_PROGRESS, COMPLETE, or BLOCKED`, + }, + { + pattern: RALPH_TASKS_COMPLETED_PATTERN, + field: 'tasksCompletedThisLoop', + validate: (v) => !Number.isNaN(parseInt(v, 10)) && parseInt(v, 10) >= 0, + transform: (v) => parseInt(v, 10), + errorMsg: (v) => `Invalid TASKS_COMPLETED_THIS_LOOP value: "${v}". Expected: non-negative integer`, + }, + { + pattern: RALPH_FILES_MODIFIED_PATTERN, + field: 'filesModified', + validate: (v) => !Number.isNaN(parseInt(v, 10)) && parseInt(v, 10) >= 0, + transform: (v) => parseInt(v, 10), + errorMsg: (v) => `Invalid FILES_MODIFIED value: "${v}". Expected: non-negative integer`, + }, + { + pattern: RALPH_TESTS_STATUS_PATTERN, + field: 'testsStatus', + validate: (v) => ['PASSING', 'FAILING', 'NOT_RUN'].includes(v.toUpperCase()), + transform: (v) => v.toUpperCase() as RalphTestsStatus, + errorMsg: (v) => `Invalid TESTS_STATUS value: "${v}". Expected: PASSING, FAILING, or NOT_RUN`, + }, + { + pattern: RALPH_WORK_TYPE_PATTERN, + field: 'workType', + validate: (v) => ['IMPLEMENTATION', 'TESTING', 'DOCUMENTATION', 'REFACTORING'].includes(v.toUpperCase()), + transform: (v) => v.toUpperCase() as RalphWorkType, + errorMsg: (v) => + `Invalid WORK_TYPE value: "${v}". Expected: IMPLEMENTATION, TESTING, DOCUMENTATION, or REFACTORING`, + }, + { + pattern: RALPH_EXIT_SIGNAL_PATTERN, + field: 'exitSignal', + validate: () => true, + transform: (v) => v.toLowerCase() === 'true', + errorMsg: () => '', + }, + { + pattern: RALPH_RECOMMENDATION_PATTERN, + field: 'recommendation', + validate: () => true, + transform: (v) => v.trim(), + errorMsg: () => '', + }, +]; + /** * RalphStatusParser - Parses RALPH_STATUS blocks and manages circuit breaker. * @@ -303,85 +364,21 @@ export class RalphStatusParser extends EventEmitter { const trimmedLine = line.trim(); if (!trimmedLine) continue; - // Track whether this line matched any known field let matched = false; - // STATUS field (required) - const statusMatch = trimmedLine.match(RALPH_STATUS_FIELD_PATTERN); - if (statusMatch) { - const value = statusMatch[1].toUpperCase(); - if (['IN_PROGRESS', 'COMPLETE', 'BLOCKED'].includes(value)) { - block.status = value as RalphStatusValue; - } else { - parseErrors.push(`Invalid STATUS value: "${value}". Expected: IN_PROGRESS, COMPLETE, or BLOCKED`); + for (const parser of FIELD_PARSERS) { + const match = trimmedLine.match(parser.pattern); + if (match) { + const rawValue = match[1]; + if (parser.validate(rawValue)) { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + (block as any)[parser.field] = parser.transform(rawValue); + } else { + parseErrors.push(parser.errorMsg(rawValue)); + } + matched = true; + break; } - matched = true; - } - - // TASKS_COMPLETED_THIS_LOOP field - const tasksMatch = trimmedLine.match(RALPH_TASKS_COMPLETED_PATTERN); - if (tasksMatch) { - const value = parseInt(tasksMatch[1], 10); - if (!Number.isNaN(value) && value >= 0) { - block.tasksCompletedThisLoop = value; - } else { - parseErrors.push( - `Invalid TASKS_COMPLETED_THIS_LOOP value: "${tasksMatch[1]}". Expected: non-negative integer` - ); - } - matched = true; - } - - // FILES_MODIFIED field - const filesMatch = trimmedLine.match(RALPH_FILES_MODIFIED_PATTERN); - if (filesMatch) { - const value = parseInt(filesMatch[1], 10); - if (!Number.isNaN(value) && value >= 0) { - block.filesModified = value; - } else { - parseErrors.push(`Invalid FILES_MODIFIED value: "${filesMatch[1]}". Expected: non-negative integer`); - } - matched = true; - } - - // TESTS_STATUS field - const testsMatch = trimmedLine.match(RALPH_TESTS_STATUS_PATTERN); - if (testsMatch) { - const value = testsMatch[1].toUpperCase(); - if (['PASSING', 'FAILING', 'NOT_RUN'].includes(value)) { - block.testsStatus = value as RalphTestsStatus; - } else { - parseErrors.push(`Invalid TESTS_STATUS value: "${value}". Expected: PASSING, FAILING, or NOT_RUN`); - } - matched = true; - } - - // WORK_TYPE field - const workMatch = trimmedLine.match(RALPH_WORK_TYPE_PATTERN); - if (workMatch) { - const value = workMatch[1].toUpperCase(); - if (['IMPLEMENTATION', 'TESTING', 'DOCUMENTATION', 'REFACTORING'].includes(value)) { - block.workType = value as RalphWorkType; - } else { - parseErrors.push( - `Invalid WORK_TYPE value: "${value}". Expected: IMPLEMENTATION, TESTING, DOCUMENTATION, or REFACTORING` - ); - } - matched = true; - } - - // EXIT_SIGNAL field - const exitMatch = trimmedLine.match(RALPH_EXIT_SIGNAL_PATTERN); - if (exitMatch) { - block.exitSignal = exitMatch[1].toLowerCase() === 'true'; - matched = true; - } - - // RECOMMENDATION field - const recMatch = trimmedLine.match(RALPH_RECOMMENDATION_PATTERN); - if (recMatch) { - block.recommendation = recMatch[1].trim(); - matched = true; } // Track unknown fields for debugging (only if looks like a field) @@ -475,38 +472,9 @@ export class RalphStatusParser extends EventEmitter { const prevState = this._circuitBreaker.state; if (hasProgress) { - // Progress detected - reset counters, possibly close circuit - this._circuitBreaker.consecutiveNoProgress = 0; - this._circuitBreaker.consecutiveSameError = 0; - this._circuitBreaker.lastProgressIteration = this._cycleCount; - - if (this._circuitBreaker.state === 'HALF_OPEN') { - this._circuitBreaker.state = 'CLOSED'; - this._circuitBreaker.reason = 'Progress detected, circuit closed'; - this._circuitBreaker.reasonCode = 'progress_detected'; - } + this._handleProgressDetected(); } else { - // No progress - this._circuitBreaker.consecutiveNoProgress++; - - // State transitions based on consecutive no-progress - if (this._circuitBreaker.state === 'CLOSED') { - if (this._circuitBreaker.consecutiveNoProgress >= 3) { - this._circuitBreaker.state = 'OPEN'; - this._circuitBreaker.reason = `No progress for ${this._circuitBreaker.consecutiveNoProgress} iterations`; - this._circuitBreaker.reasonCode = 'no_progress_open'; - } else if (this._circuitBreaker.consecutiveNoProgress >= 2) { - this._circuitBreaker.state = 'HALF_OPEN'; - this._circuitBreaker.reason = 'Warning: no progress detected'; - this._circuitBreaker.reasonCode = 'no_progress_warning'; - } - } else if (this._circuitBreaker.state === 'HALF_OPEN') { - if (this._circuitBreaker.consecutiveNoProgress >= 3) { - this._circuitBreaker.state = 'OPEN'; - this._circuitBreaker.reason = `No progress for ${this._circuitBreaker.consecutiveNoProgress} iterations`; - this._circuitBreaker.reasonCode = 'no_progress_open'; - } - } + this._handleNoProgress(); } // Track tests failure @@ -535,6 +503,40 @@ export class RalphStatusParser extends EventEmitter { } } + private _handleProgressDetected(): void { + this._circuitBreaker.consecutiveNoProgress = 0; + this._circuitBreaker.consecutiveSameError = 0; + this._circuitBreaker.lastProgressIteration = this._cycleCount; + + if (this._circuitBreaker.state === 'HALF_OPEN') { + this._circuitBreaker.state = 'CLOSED'; + this._circuitBreaker.reason = 'Progress detected, circuit closed'; + this._circuitBreaker.reasonCode = 'progress_detected'; + } + } + + private _handleNoProgress(): void { + this._circuitBreaker.consecutiveNoProgress++; + + if (this._circuitBreaker.state === 'CLOSED') { + if (this._circuitBreaker.consecutiveNoProgress >= 3) { + this._circuitBreaker.state = 'OPEN'; + this._circuitBreaker.reason = `No progress for ${this._circuitBreaker.consecutiveNoProgress} iterations`; + this._circuitBreaker.reasonCode = 'no_progress_open'; + } else if (this._circuitBreaker.consecutiveNoProgress >= 2) { + this._circuitBreaker.state = 'HALF_OPEN'; + this._circuitBreaker.reason = 'Warning: no progress detected'; + this._circuitBreaker.reasonCode = 'no_progress_warning'; + } + } else if (this._circuitBreaker.state === 'HALF_OPEN') { + if (this._circuitBreaker.consecutiveNoProgress >= 3) { + this._circuitBreaker.state = 'OPEN'; + this._circuitBreaker.reason = `No progress for ${this._circuitBreaker.consecutiveNoProgress} iterations`; + this._circuitBreaker.reasonCode = 'no_progress_open'; + } + } + } + /** * Check line for completion indicators (natural language patterns). * Used for dual-condition exit gate. diff --git a/src/respawn-controller.ts b/src/respawn-controller.ts index 7aa17336..3e9c6fc1 100644 --- a/src/respawn-controller.ts +++ b/src/respawn-controller.ts @@ -860,27 +860,29 @@ export class RespawnController extends EventEmitter { } }; - // Ensure timeouts are positive - validatePositiveTimeout('idleTimeoutMs'); - validatePositiveTimeout('completionConfirmMs'); - validatePositiveTimeout('noOutputTimeoutMs'); - validatePositiveTimeout('autoAcceptDelayMs', true); - validatePositiveTimeout('interStepDelayMs'); + const REQUIRED_TIMEOUT_FIELDS = [ + 'idleTimeoutMs', + 'completionConfirmMs', + 'noOutputTimeoutMs', + 'interStepDelayMs', + 'aiIdleCheckTimeoutMs', + 'aiIdleCheckMaxContext', + 'aiPlanCheckTimeoutMs', + 'aiPlanCheckMaxContext', + ] as const; + for (const field of REQUIRED_TIMEOUT_FIELDS) { + validatePositiveTimeout(field); + } + + const ALLOW_ZERO_FIELDS = ['autoAcceptDelayMs', 'aiIdleCheckCooldownMs', 'aiPlanCheckCooldownMs'] as const; + for (const field of ALLOW_ZERO_FIELDS) { + validatePositiveTimeout(field, true); + } // Ensure completion confirm doesn't exceed no-output timeout if (c.completionConfirmMs > c.noOutputTimeoutMs) { c.completionConfirmMs = c.noOutputTimeoutMs; } - - // Ensure AI check timeouts are positive - validatePositiveTimeout('aiIdleCheckTimeoutMs'); - validatePositiveTimeout('aiIdleCheckCooldownMs', true); - validatePositiveTimeout('aiIdleCheckMaxContext'); - - // Ensure plan check timeouts are positive - validatePositiveTimeout('aiPlanCheckTimeoutMs'); - validatePositiveTimeout('aiPlanCheckCooldownMs', true); - validatePositiveTimeout('aiPlanCheckMaxContext'); } /** Wire up AI checker events to controller events (removes existing listeners first to prevent duplicates) */ @@ -1385,98 +1387,13 @@ export class RespawnController extends EventEmitter { this.lastTokenChangeTime = now; } - // Detect completion message FIRST (Layer 1) - PRIMARY DETECTION - // Check this before working patterns because completion message indicates - // the work is done, even if working patterns are still in the rolling window - if (isCompletionMessage(data)) { - // Clear the rolling window - completion marks a transition point - this.clearWorkingPatternWindow(); - this.workingDetected = false; - this.completionMessageTime = now; - this.cancelAutoAcceptTimer(); // Normal idle flow handles this - this.log(`Completion message detected: "${data.trim().substring(0, 50)}..."`); + // Layer 1: Completion message (PRIMARY) — checked before working patterns + if (this._detectCompletionMessage(data, now)) return; - // In watching state, start completion confirmation timer - if (this._state === 'watching') { - this.startCompletionConfirmTimer(); - return; - } + // Layer 4: Working patterns + if (this._detectWorkingPattern(data, now)) return; - // 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'); - break; - case 'waiting_clear': - this.checkClearComplete(); // /clear is quick, no need to wait - break; - case 'waiting_init': - this.startStepConfirmTimer('init'); - break; - 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; - } - - // Detect working patterns (Layer 4) - const isWorking = this.checkWorkingPattern(data); - if (isWorking) { - this.workingDetected = true; - this.promptDetected = false; - this.elicitationDetected = false; // Clear on new work cycle - this.resetHookState(); // Clear hook signals on new work - this.lastWorkingPatternTime = now; - - // Cancel hook confirmation timer if running - this.cancelTrackedTimer('hook-confirm', 'working patterns detected'); - - // Cancel any pending completion confirmation - this.cancelCompletionConfirm(); - - // Cancel any pending step confirmation (Claude is still working) - this.cancelStepConfirm(); - - // If AI check is running, cancel it (Claude is working) - if (this._state === 'ai_checking') { - this.log('Working patterns detected during AI check, cancelling'); - this.aiChecker.cancel(); - this.setState('watching'); - } - - // Cancel plan check if running (Claude started working) - if (this.planChecker.status === 'checking') { - this.log('Working patterns detected during plan check, cancelling'); - this.planChecker.cancel(); - } - - // If we're monitoring init and work started, go to watching (no kickstart needed) - if (this._state === 'monitoring_init') { - this.log('/init triggered work, skipping kickstart'); - this.emit('stepCompleted', 'init'); - this.completeCycle(); - } - return; - } - - // In confirming_idle or ai_checking state, substantial output cancels the flow. - // This prevents false triggers when Claude pauses briefly mid-work. + // Substantial output during confirming_idle/ai_checking cancels the flow if (this._state === 'confirming_idle' || this._state === 'ai_checking') { // Strip ANSI escape codes to check if there's real content ANSI_ESCAPE_PATTERN_SIMPLE.lastIndex = 0; @@ -1497,43 +1414,137 @@ export class RespawnController extends EventEmitter { } } - // Legacy fallback: detect prompt characters (still useful for waiting_* states) - const hasPrompt = PROMPT_PATTERNS.some((pattern) => data.includes(pattern)); - if (hasPrompt) { - this.promptDetected = true; - this.workingDetected = false; + // Legacy fallback: prompt detection + this._detectPrompt(data); + } - // Handle legacy detection in waiting states - also use confirmation timers - switch (this._state) { - case 'waiting_update': - this.startStepConfirmTimer('update'); - break; - case 'waiting_clear': - this.checkClearComplete(); // /clear is quick, no need to wait - break; - case 'waiting_init': - this.startStepConfirmTimer('init'); - break; - case 'monitoring_init': - this.checkMonitoringInitIdle(); - break; - 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}`); - } + private _detectCompletionMessage(data: string, now: number): boolean { + if (!isCompletionMessage(data)) return false; + + // Clear the rolling window - completion marks a transition point + this.clearWorkingPatternWindow(); + this.workingDetected = false; + this.completionMessageTime = now; + this.cancelAutoAcceptTimer(); // Normal idle flow handles this + this.log(`Completion message detected: "${data.trim().substring(0, 50)}..."`); + + // In watching state, start completion confirmation timer + if (this._state === 'watching') { + this.startCompletionConfirmTimer(); + return true; + } + + // 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'); + break; + case 'waiting_clear': + this.checkClearComplete(); // /clear is quick, no need to wait + break; + case 'waiting_init': + this.startStepConfirmTimer('init'); + break; + 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 true; + } + + private _detectWorkingPattern(data: string, now: number): boolean { + const isWorking = this.checkWorkingPattern(data); + if (!isWorking) return false; + + this.workingDetected = true; + this.promptDetected = false; + this.elicitationDetected = false; // Clear on new work cycle + this.resetHookState(); // Clear hook signals on new work + this.lastWorkingPatternTime = now; + + // Cancel hook confirmation timer if running + this.cancelTrackedTimer('hook-confirm', 'working patterns detected'); + + // Cancel any pending completion confirmation + this.cancelCompletionConfirm(); + + // Cancel any pending step confirmation (Claude is still working) + this.cancelStepConfirm(); + + // If AI check is running, cancel it (Claude is working) + if (this._state === 'ai_checking') { + this.log('Working patterns detected during AI check, cancelling'); + this.aiChecker.cancel(); + this.setState('watching'); + } + + // Cancel plan check if running (Claude started working) + if (this.planChecker.status === 'checking') { + this.log('Working patterns detected during plan check, cancelling'); + this.planChecker.cancel(); + } + + // If we're monitoring init and work started, go to watching (no kickstart needed) + if (this._state === 'monitoring_init') { + this.log('/init triggered work, skipping kickstart'); + this.emit('stepCompleted', 'init'); + this.completeCycle(); + } + return true; + } + + private _detectPrompt(data: string): void { + const hasPrompt = PROMPT_PATTERNS.some((pattern) => data.includes(pattern)); + if (!hasPrompt) return; + + this.promptDetected = true; + this.workingDetected = false; + + // Handle legacy detection in waiting states - also use confirmation timers + switch (this._state) { + case 'waiting_update': + this.startStepConfirmTimer('update'); + break; + case 'waiting_clear': + this.checkClearComplete(); // /clear is quick, no need to wait + break; + case 'waiting_init': + this.startStepConfirmTimer('init'); + break; + case 'monitoring_init': + this.checkMonitoringInitIdle(); + break; + 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}`); } } diff --git a/src/session.ts b/src/session.ts index 2c499d1d..b85d24c3 100644 --- a/src/session.ts +++ b/src/session.ts @@ -871,6 +871,71 @@ export class Session extends EventEmitter { * session.write('help me with this code\r'); * ``` */ + private async _setupOrAttachMuxSession(options: { + respawnPaneOptions: import('./mux-interface.js').RespawnPaneOptions; + createSessionOptions: import('./mux-interface.js').CreateSessionOptions; + spawnErrLabel: string; + }): Promise<{ isRestored: boolean }> { + const mux = this._mux!; + + // Verify stale mux session — tmux may have been destroyed (e.g., killed externally) + if (this._muxSession && !mux.muxSessionExists(this._muxSession.muxName)) { + console.log('[Session] Stale mux session detected (tmux gone):', this._muxSession.muxName); + this._muxSession = null; + } + + // Check if session exists but pane is dead (remain-on-exit keeps it alive) + // Respawn the pane instead of creating a whole new session — preserves tmux scrollback + let needsNewSession = false; + if (this._muxSession && mux.isPaneDead(this._muxSession.muxName)) { + console.log('[Session] Dead pane detected, respawning:', this._muxSession.muxName); + const newPid = await mux.respawnPane(options.respawnPaneOptions); + if (!newPid) { + console.error('[Session] Failed to respawn pane, will create new session'); + needsNewSession = true; + } else { + // Wait a moment for the respawned process to fully start + await new Promise((resolve) => setTimeout(resolve, MUX_STARTUP_DELAY_MS)); + } + } + + // Check if we already have a mux session (restored session) + const isRestored = this._muxSession !== null && !needsNewSession; + if (isRestored) { + console.log('[Session] Attaching to existing mux session:', this._muxSession!.muxName); + } else { + // Create a new mux session + this._muxSession = await mux.createSession(options.createSessionOptions); + console.log('[Session] Created mux session:', this._muxSession.muxName); + // No extra sleep — createSession() already waits for tmux readiness + } + + // Attach to the mux session via PTY + try { + this.ptyProcess = pty.spawn(mux.getAttachCommand(), mux.getAttachArgs(this._muxSession!.muxName), { + name: 'xterm-256color', + cols: 120, + rows: 40, + cwd: this.workingDir, + env: buildMuxAttachEnv(), + }); + } catch (spawnErr) { + console.error(`[Session] Failed to spawn PTY for ${options.spawnErrLabel}:`, spawnErr); + this.emit('error', `Failed to attach to mux session: ${spawnErr}`); + throw spawnErr; + } + + return { isRestored }; + } + + private _handleTerminalOutput(data: string): void { + // BufferAccumulator handles auto-trimming when max size exceeded + this._terminalBuffer.append(data); + this._lastActivityAt = Date.now(); + this.emit('terminal', data); + this.emit('output', data); + } + async startInteractive(): Promise { if (this.ptyProcess) { throw new Error('Session already has a running process'); @@ -886,18 +951,8 @@ export class Session extends EventEmitter { // If mux wrapping is enabled, create or attach to a mux session if (this._useMux && this._mux) { try { - // Verify stale mux session — tmux may have been destroyed (e.g., killed externally) - if (this._muxSession && !this._mux.muxSessionExists(this._muxSession.muxName)) { - console.log('[Session] Stale mux session detected (tmux gone):', this._muxSession.muxName); - this._muxSession = null; - } - - // Check if session exists but pane is dead (remain-on-exit keeps it alive) - // Respawn the pane instead of creating a whole new session — preserves tmux scrollback - let needsNewSession = false; - if (this._muxSession && this._mux.isPaneDead(this._muxSession.muxName)) { - console.log('[Session] Dead pane detected, respawning:', this._muxSession.muxName); - const newPid = await this._mux.respawnPane({ + const { isRestored } = await this._setupOrAttachMuxSession({ + respawnPaneOptions: { sessionId: this.id, workingDir: this.workingDir, mode: this.mode, @@ -907,23 +962,8 @@ export class Session extends EventEmitter { allowedTools: this._allowedTools, openCodeConfig: this._openCodeConfig, resumeSessionId: this._resumeSessionId, - }); - if (!newPid) { - console.error('[Session] Failed to respawn pane, will create new session'); - needsNewSession = true; - } else { - // Wait a moment for the respawned process to fully start - await new Promise((resolve) => setTimeout(resolve, MUX_STARTUP_DELAY_MS)); - } - } - - // Check if we already have a mux session (restored session) - const isRestoredSession = this._muxSession !== null && !needsNewSession; - if (isRestoredSession) { - console.log('[Session] Attaching to existing mux session:', this._muxSession!.muxName); - } else { - // Create a new mux session - this._muxSession = await this._mux.createSession({ + }, + createSessionOptions: { sessionId: this.id, workingDir: this.workingDir, mode: this.mode, @@ -934,36 +974,16 @@ export class Session extends EventEmitter { allowedTools: this._allowedTools, openCodeConfig: this._openCodeConfig, resumeSessionId: this._resumeSessionId, - }); - console.log('[Session] Created mux session:', this._muxSession.muxName); - // No extra sleep — createSession() already waits for tmux readiness - } + }, + spawnErrLabel: 'mux attachment', + }); - // Attach to the mux session via PTY - try { - this.ptyProcess = pty.spawn( - this._mux.getAttachCommand(), - this._mux.getAttachArgs(this._muxSession!.muxName), - { - name: 'xterm-256color', - cols: 120, - rows: 40, - cwd: this.workingDir, - env: buildMuxAttachEnv(), - } - ); - - // Set claudeSessionId — when resuming, the Claude conversation ID is the resumed one. - this._claudeSessionId = this._resumeSessionId || this.id; - } catch (spawnErr) { - console.error('[Session] Failed to spawn PTY for mux attachment:', spawnErr); - this.emit('error', `Failed to attach to mux session: ${spawnErr}`); - throw spawnErr; - } + // Set claudeSessionId — when resuming, the Claude conversation ID is the resumed one. + this._claudeSessionId = this._resumeSessionId || this.id; // For NEW mux sessions: wait for readiness then clean buffer // For RESTORED mux sessions: don't do anything - client will fetch buffer on tab switch - if (!isRestoredSession) { + if (!isRestored) { if (this.mode === 'opencode') { // OpenCode uses Bubble Tea TUI — no ❯ prompt to detect. // Wait for TUI to stabilize (output stops changing), then mark ready. @@ -1050,12 +1070,7 @@ export class Session extends EventEmitter { const data = rawData.replace(FOCUS_ESCAPE_FILTER, '').replace(CTRL_L_PATTERN, ''); // Remove Ctrl+L if (!data) return; // Skip if only filtered sequences - // BufferAccumulator handles auto-trimming when max size exceeded - this._terminalBuffer.append(data); - this._lastActivityAt = Date.now(); - - this.emit('terminal', data); - this.emit('output', data); + this._handleTerminalOutput(data); // === Idle/working detection runs on every chunk (latency-sensitive) === // Detect if Claude is working or at prompt @@ -1263,69 +1278,26 @@ export class Session extends EventEmitter { // If mux wrapping is enabled, create or attach to a mux session if (this._useMux && this._mux) { try { - // Verify stale mux session — tmux may have been destroyed externally - if (this._muxSession && !this._mux.muxSessionExists(this._muxSession.muxName)) { - console.log('[Session] Stale mux session detected (tmux gone):', this._muxSession.muxName); - this._muxSession = null; - } - - // Check if session exists but pane is dead (remain-on-exit keeps it alive) - let needsNewSession = false; - if (this._muxSession && this._mux.isPaneDead(this._muxSession.muxName)) { - console.log('[Session] Dead pane detected, respawning:', this._muxSession.muxName); - const newPid = await this._mux.respawnPane({ + const { isRestored } = await this._setupOrAttachMuxSession({ + respawnPaneOptions: { sessionId: this.id, workingDir: this.workingDir, mode: 'shell', niceConfig: this._niceConfig, - }); - if (!newPid) { - console.error('[Session] Failed to respawn pane, will create new session'); - needsNewSession = true; - } else { - await new Promise((resolve) => setTimeout(resolve, MUX_STARTUP_DELAY_MS)); - } - } - - // Check if we already have a mux session (restored session) - const isRestoredSession = this._muxSession !== null && !needsNewSession; - if (isRestoredSession) { - console.log('[Session] Attaching to existing mux session:', this._muxSession!.muxName); - } else { - // Create a new mux session - this._muxSession = await this._mux.createSession({ + }, + createSessionOptions: { sessionId: this.id, workingDir: this.workingDir, mode: 'shell', name: this._name, niceConfig: this._niceConfig, - }); - console.log('[Session] Created mux session:', this._muxSession.muxName); - // No extra sleep — createSession() already waits for tmux readiness - } - - // Attach to the mux session via PTY - try { - this.ptyProcess = pty.spawn( - this._mux.getAttachCommand(), - this._mux.getAttachArgs(this._muxSession!.muxName), - { - name: 'xterm-256color', - cols: 120, - rows: 40, - cwd: this.workingDir, - env: buildMuxAttachEnv(), - } - ); - } catch (spawnErr) { - console.error('[Session] Failed to spawn PTY for shell mux attachment:', spawnErr); - this.emit('error', `Failed to attach to mux session: ${spawnErr}`); - throw spawnErr; - } + }, + spawnErrLabel: 'shell mux attachment', + }); // For NEW sessions: clear by sending 'clear' command to the shell // For RESTORED sessions: don't clear - we want to see the existing output - if (!isRestoredSession) { + if (!isRestored) { setTimeout(() => { if (this.ptyProcess) { this._terminalBuffer.clear(); @@ -1366,12 +1338,7 @@ export class Session extends EventEmitter { const data = rawData.replace(FOCUS_ESCAPE_FILTER, ''); if (!data) return; // Skip if only focus sequences - // BufferAccumulator handles auto-trimming when max size exceeded - this._terminalBuffer.append(data); - this._lastActivityAt = Date.now(); - - this.emit('terminal', data); - this.emit('output', data); + this._handleTerminalOutput(data); }); this.ptyProcess.onExit(({ exitCode }) => { @@ -1479,12 +1446,7 @@ export class Session extends EventEmitter { const data = rawData.replace(FOCUS_ESCAPE_FILTER, ''); if (!data) return; // Skip if only focus sequences - // BufferAccumulator handles auto-trimming when max size exceeded - this._terminalBuffer.append(data); - this._lastActivityAt = Date.now(); - - this.emit('terminal', data); - this.emit('output', data); + this._handleTerminalOutput(data); // Also try to parse JSON lines for structured data this.processOutput(data); diff --git a/src/state-store.ts b/src/state-store.ts index f53d784b..ec94d48b 100644 --- a/src/state-store.ts +++ b/src/state-store.ts @@ -116,6 +116,26 @@ export class StateStore { this.loadRalphStates(); } + private _mergeWithInitialState(parsed: Partial): AppState { + const initial = createInitialState(); + return { + ...initial, + ...parsed, + sessions: { ...parsed.sessions }, + tasks: { ...parsed.tasks }, + ralphLoop: { ...initial.ralphLoop, ...parsed.ralphLoop }, + config: { ...initial.config, ...parsed.config }, + }; + } + + private _resetCircuitBreaker(): void { + this.consecutiveSaveFailures = 0; + if (this.circuitBreakerOpen) { + console.log('[StateStore] Circuit breaker CLOSED - save succeeded'); + this.circuitBreakerOpen = false; + } + } + private ensureDir(): void { const dir = dirname(this.filePath); if (!existsSync(dir)) { @@ -132,15 +152,7 @@ export class StateStore { if (existsSync(path)) { const data = readFileSync(path, 'utf-8'); const parsed = JSON.parse(data) as Partial; - const initial = createInitialState(); - const result = { - ...initial, - ...parsed, - sessions: { ...parsed.sessions }, - tasks: { ...parsed.tasks }, - ralphLoop: { ...initial.ralphLoop, ...parsed.ralphLoop }, - config: { ...initial.config, ...parsed.config }, - }; + const result = this._mergeWithInitialState(parsed); if (path !== this.filePath) { console.warn(`[StateStore] Recovered state from backup: ${path}`); } @@ -321,11 +333,7 @@ export class StateStore { await writeFile(tempPath, json, 'utf-8'); await rename(tempPath, this.filePath); - this.consecutiveSaveFailures = 0; - if (this.circuitBreakerOpen) { - console.log('[StateStore] Circuit breaker CLOSED - save succeeded'); - this.circuitBreakerOpen = false; - } + this._resetCircuitBreaker(); } catch (err) { console.error('[StateStore] Failed to write state file:', err); // Re-mark dirty so the data is retried on the next save cycle @@ -385,11 +393,7 @@ export class StateStore { renameSync(tempPath, this.filePath); // Clear dirty flag only AFTER successful write this.dirty = false; - this.consecutiveSaveFailures = 0; - if (this.circuitBreakerOpen) { - console.log('[StateStore] Circuit breaker CLOSED - save succeeded'); - this.circuitBreakerOpen = false; - } + this._resetCircuitBreaker(); } catch (err) { console.error('[StateStore] Failed to write state file:', err); this.consecutiveSaveFailures++; @@ -415,15 +419,7 @@ export class StateStore { if (existsSync(backupPath)) { const backupContent = readFileSync(backupPath, 'utf-8'); const parsed = JSON.parse(backupContent) as Partial; - const initial = createInitialState(); - this.state = { - ...initial, - ...parsed, - sessions: { ...parsed.sessions }, - tasks: { ...parsed.tasks }, - ralphLoop: { ...initial.ralphLoop, ...parsed.ralphLoop }, - config: { ...initial.config, ...parsed.config }, - }; + this.state = this._mergeWithInitialState(parsed); console.log('[StateStore] Successfully recovered state from backup'); // Reset circuit breaker after successful recovery this.circuitBreakerOpen = false; diff --git a/src/subagent-watcher.ts b/src/subagent-watcher.ts index 715186a9..8d5e0674 100644 --- a/src/subagent-watcher.ts +++ b/src/subagent-watcher.ts @@ -988,6 +988,24 @@ export class SubagentWatcher extends EventEmitter { return truncated.replace(/[.!?,:\s]+$/, ''); } + private async _resolveDescription( + projectHash: string, + sessionId: string, + agentId: string, + filePath: string, + fallbackText?: string + ): Promise { + // First try parent transcript (most reliable) + const fromParent = await this.extractDescriptionFromParentTranscript(projectHash, sessionId, agentId); + if (fromParent) return fromParent; + + // Fallback: inline text (from processEntry) or file extraction + if (fallbackText) { + return this.extractSmartTitle(fallbackText); + } + return this.extractDescriptionFromFile(filePath); + } + /** * Extract the short description from the parent session's transcript. * This is the most reliable method because it reads the actual Task tool result @@ -1245,18 +1263,13 @@ export class SubagentWatcher extends EventEmitter { // Retry description extraction if missing (race condition fix) if (!existingInfo.description) { - // First try parent transcript (most reliable) - let extractedDescription = await this.extractDescriptionFromParentTranscript( + const extractedDescription = await this._resolveDescription( existingInfo.projectHash, existingInfo.sessionId, - agentId + agentId, + filePath ); - // Fallback to subagent file - if (!extractedDescription) { - extractedDescription = await this.extractDescriptionFromFile(filePath); - } if (extractedDescription) { - // Check if this is an internal agent - if so, remove it if (this.isInternalAgent(extractedDescription)) { this.removeAgent(agentId); return; @@ -1307,13 +1320,7 @@ export class SubagentWatcher extends EventEmitter { } // Extract description - prefer reading from parent transcript (most reliable) - // The parent transcript has the exact Task tool call with description parameter - let description = await this.extractDescriptionFromParentTranscript(projectHash, sessionId, agentId); - - // Fallback: extract a smart title from the subagent's prompt if parent lookup failed - if (!description) { - description = await this.extractDescriptionFromFile(filePath); - } + const description = await this._resolveDescription(projectHash, sessionId, agentId, filePath); // Skip internal Claude Code agents (e.g., suggestion mode) - not real subagents if (this.isInternalAgent(description)) { @@ -1406,43 +1413,11 @@ export class SubagentWatcher extends EventEmitter { private async processEntry(entry: SubagentTranscriptEntry, agentId: string, sessionId: string): Promise { const info = this.agentInfo.get(agentId); - // Extract model from assistant messages (first one sets the model) - if (info && entry.type === 'assistant' && entry.message?.model && !info.model) { - info.model = entry.message.model; - info.modelShort = this.extractModelShort(entry.message.model); - this.emit('subagent:updated', info); - } + if (info) { + this._processModelInfo(entry, info); + this._processTokenInfo(entry, info); - // Aggregate token usage from messages - if (info && entry.message?.usage) { - if (entry.message.usage.input_tokens) { - info.totalInputTokens = (info.totalInputTokens || 0) + entry.message.usage.input_tokens; - } - if (entry.message.usage.output_tokens) { - info.totalOutputTokens = (info.totalOutputTokens || 0) + entry.message.usage.output_tokens; - } - } - - // Check if this is first user message and description is missing - if (info && !info.description && entry.type === 'user' && entry.message?.content) { - // First try parent transcript (most reliable) - let description = await this.extractDescriptionFromParentTranscript(info.projectHash, info.sessionId, agentId); - // Fallback: extract smart title from the prompt content - if (!description) { - const text = this.extractFirstTextContent(entry.message.content); - if (text) { - description = this.extractSmartTitle(text); - } - } - if (description) { - // Check if this is an internal agent - if so, remove it - if (this.isInternalAgent(description)) { - this.removeAgent(agentId); - return; - } - info.description = description; - this.emit('subagent:updated', info); - } + if (await this._processDescription(entry, agentId, info)) return; } if (entry.type === 'progress' && entry.data) { @@ -1453,7 +1428,6 @@ export class SubagentWatcher extends EventEmitter { progressType: entry.data.type, query: entry.data.query, resultCount: entry.data.resultCount, - // Extract hook event info if present hookEvent: entry.data.hookEvent, hookName: entry.data.hookName || @@ -1463,115 +1437,166 @@ export class SubagentWatcher extends EventEmitter { }; this.emit('subagent:progress', progress); } else if (entry.type === 'assistant' && 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 > 0) { - const message: SubagentMessage = { + this._processAssistantContent(entry, agentId, sessionId); + } else if (entry.type === 'user' && entry.message?.content) { + this._processUserContent(entry, agentId, sessionId); + } + } + + private _processModelInfo(entry: SubagentTranscriptEntry, agent: SubagentInfo): void { + if (entry.type === 'assistant' && entry.message?.model && !agent.model) { + agent.model = entry.message.model; + agent.modelShort = this.extractModelShort(entry.message.model); + this.emit('subagent:updated', agent); + } + } + + private _processTokenInfo(entry: SubagentTranscriptEntry, agent: SubagentInfo): void { + if (!entry.message?.usage) return; + if (entry.message.usage.input_tokens) { + agent.totalInputTokens = (agent.totalInputTokens || 0) + entry.message.usage.input_tokens; + } + if (entry.message.usage.output_tokens) { + agent.totalOutputTokens = (agent.totalOutputTokens || 0) + entry.message.usage.output_tokens; + } + } + + private async _processDescription( + entry: SubagentTranscriptEntry, + agentId: string, + agent: SubagentInfo + ): Promise { + if (agent.description || entry.type !== 'user' || !entry.message?.content) return false; + + const fallbackText = this.extractFirstTextContent(entry.message.content); + const description = await this._resolveDescription( + agent.projectHash, + agent.sessionId, + agentId, + agent.filePath, + fallbackText + ); + if (description) { + if (this.isInternalAgent(description)) { + this.removeAgent(agentId); + return true; + } + agent.description = description; + this.emit('subagent:updated', agent); + } + return false; + } + + private _processAssistantContent(entry: SubagentTranscriptEntry, agentId: string, sessionId: string): void { + const messageContent = entry.message!.content; + if (typeof messageContent === 'string') { + const text = messageContent.trim(); + if (text.length > 0) { + const message: SubagentMessage = { + agentId, + sessionId, + timestamp: entry.timestamp, + role: 'assistant', + text: text.substring(0, MESSAGE_TEXT_LIMIT), + }; + this.emit('subagent:message', message); + } + } else { + for (const content of messageContent) { + if (content.type === 'tool_use' && content.name) { + // Store toolUseId for linking to results, with timestamp for TTL cleanup + if (content.id) { + if (!this.pendingToolCalls.has(agentId)) { + this.pendingToolCalls.set(agentId, new Map()); + } + const agentCalls = this.pendingToolCalls.get(agentId)!; + // Enforce size limit to prevent memory leak from rapid tool calls + if (agentCalls.size >= MAX_PENDING_TOOL_CALLS) { + // FIFO eviction: delete first (oldest) entry using Map insertion order + const firstKey = agentCalls.keys().next().value; + if (firstKey !== undefined) agentCalls.delete(firstKey); + } + agentCalls.set(content.id, { + toolName: content.name, + timestamp: Date.now(), + }); + } + + const toolCall: SubagentToolCall = { agentId, sessionId, timestamp: entry.timestamp, - role: 'assistant', - text: text.substring(0, MESSAGE_TEXT_LIMIT), + tool: content.name, + input: this.getTruncatedInput(content.name, content.input || {}), + toolUseId: content.id, + fullInput: content.input || {}, }; - this.emit('subagent:message', message); - } - } else { - for (const content of entry.message.content) { - if (content.type === 'tool_use' && content.name) { - // Store toolUseId for linking to results, with timestamp for TTL cleanup - if (content.id) { - if (!this.pendingToolCalls.has(agentId)) { - this.pendingToolCalls.set(agentId, new Map()); - } - const agentCalls = this.pendingToolCalls.get(agentId)!; - // Enforce size limit to prevent memory leak from rapid tool calls - if (agentCalls.size >= MAX_PENDING_TOOL_CALLS) { - // FIFO eviction: delete first (oldest) entry using Map insertion order - const firstKey = agentCalls.keys().next().value; - if (firstKey !== undefined) agentCalls.delete(firstKey); - } - agentCalls.set(content.id, { - toolName: content.name, - timestamp: Date.now(), - }); - } + this.emit('subagent:tool_call', toolCall); - const toolCall: SubagentToolCall = { + // Update tool call count + const agentInfo = this.agentInfo.get(agentId); + if (agentInfo) { + agentInfo.toolCallCount++; + } + } else if (content.type === 'tool_result' && content.tool_use_id) { + this.emitToolResult( + { tool_use_id: content.tool_use_id, content: content.content, is_error: content.is_error }, + agentId, + sessionId, + entry.timestamp + ); + } else if (content.type === 'text' && content.text) { + const text = content.text.trim(); + if (text.length > 0) { + const message: SubagentMessage = { agentId, sessionId, timestamp: entry.timestamp, - tool: content.name, - input: this.getTruncatedInput(content.name, content.input || {}), - toolUseId: content.id, - fullInput: content.input || {}, + role: 'assistant', + text: text.substring(0, MESSAGE_TEXT_LIMIT), }; - this.emit('subagent:tool_call', toolCall); - - // Update tool call count - const agentInfo = this.agentInfo.get(agentId); - if (agentInfo) { - agentInfo.toolCallCount++; - } - } else if (content.type === 'tool_result' && content.tool_use_id) { - // Extract tool result - this.emitToolResult( - { tool_use_id: content.tool_use_id, content: content.content, is_error: content.is_error }, - agentId, - sessionId, - entry.timestamp - ); - } else if (content.type === 'text' && content.text) { - const text = content.text.trim(); - if (text.length > 0) { - const message: SubagentMessage = { - agentId, - sessionId, - timestamp: entry.timestamp, - role: 'assistant', - text: text.substring(0, MESSAGE_TEXT_LIMIT), // Limit text length - }; - this.emit('subagent:message', message); - } + this.emit('subagent:message', message); } } } - } else if (entry.type === 'user' && entry.message?.content) { - // Handle both string and array content formats - also check for tool_result in user messages - if (typeof entry.message.content === 'string') { - const userText = entry.message.content.trim(); - if (userText.length > 0 && userText.length < 500) { - const message: SubagentMessage = { + } + } + + private _processUserContent(entry: SubagentTranscriptEntry, agentId: string, sessionId: string): void { + const messageContent = entry.message!.content; + if (typeof messageContent === 'string') { + const userText = messageContent.trim(); + if (userText.length > 0 && userText.length < 500) { + const message: SubagentMessage = { + agentId, + sessionId, + timestamp: entry.timestamp, + role: 'user', + text: userText, + }; + this.emit('subagent:message', message); + } + } else { + // Check for tool_result blocks in user messages (common pattern) + for (const content of messageContent) { + if (content.type === 'tool_result' && content.tool_use_id) { + this.emitToolResult( + { tool_use_id: content.tool_use_id, content: content.content, is_error: content.is_error }, agentId, sessionId, - timestamp: entry.timestamp, - role: 'user', - text: userText, - }; - this.emit('subagent:message', message); - } - } else { - // 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) { - this.emitToolResult( - { tool_use_id: content.tool_use_id, content: content.content, is_error: content.is_error }, + entry.timestamp + ); + } else if (content.type === 'text' && content.text) { + const userText = content.text.trim(); + if (userText.length > 0 && userText.length < 500) { + const message: SubagentMessage = { agentId, sessionId, - entry.timestamp - ); - } else if (content.type === 'text' && content.text) { - const userText = content.text.trim(); - if (userText.length > 0 && userText.length < 500) { - const message: SubagentMessage = { - agentId, - sessionId, - timestamp: entry.timestamp, - role: 'user', - text: userText, - }; - this.emit('subagent:message', message); - } + timestamp: entry.timestamp, + role: 'user', + text: userText, + }; + this.emit('subagent:message', message); } } } diff --git a/src/web/public/app.js b/src/web/public/app.js index 27013e81..35b10fb5 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -939,15 +939,7 @@ class CodemanApp { if (data.id === this.activeSessionId) { this.terminal.writeln(`\x1b[1;31m Error: ${data.error}\x1b[0m`); } - const session = this.sessions.get(data.id); - this.notificationManager?.notify({ - urgency: 'critical', - category: 'session-error', - sessionId: data.id, - sessionName: session?.name || this.getShortId(data.id), - title: 'Session Error', - message: data.error || 'Unknown error', - }); + this._notifySession(data.id, 'critical', 'session-error', 'Session Error', data.error || 'Unknown error'); } _onSessionExit(data) { @@ -960,14 +952,7 @@ class CodemanApp { } // Notify on unexpected exit (non-zero code) if (data.code && data.code !== 0) { - this.notificationManager?.notify({ - urgency: 'critical', - category: 'session-crash', - sessionId: data.id, - sessionName: session?.name || this.getShortId(data.id), - title: 'Session Crashed', - message: `Exited with code ${data.code}`, - }); + this._notifySession(data.id, 'critical', 'session-crash', 'Session Crashed', `Exited with code ${data.code}`); } } @@ -984,15 +969,7 @@ class CodemanApp { const threshold = this.notificationManager?.preferences?.stuckThresholdMs || 600000; clearTimeout(this.idleTimers.get(data.id)); this.idleTimers.set(data.id, setTimeout(() => { - const s = this.sessions.get(data.id); - this.notificationManager?.notify({ - urgency: 'warning', - category: 'session-stuck', - sessionId: data.id, - sessionName: s?.name || this.getShortId(data.id), - title: 'Session Idle', - message: `Idle for ${Math.round(threshold / 60000)}+ minutes`, - }); + this._notifySession(data.id, 'warning', 'session-stuck', 'Session Idle', `Idle for ${Math.round(threshold / 60000)}+ minutes`); this.idleTimers.delete(data.id); }, threshold)); } @@ -1023,15 +1000,7 @@ class CodemanApp { this.showToast(`Auto-cleared at ${data.tokens.toLocaleString()} tokens`, 'info'); this.updateRespawnTokens(0); } - const session = this.sessions.get(data.sessionId); - this.notificationManager?.notify({ - urgency: 'info', - category: 'auto-clear', - sessionId: data.sessionId, - sessionName: session?.name || this.getShortId(data.sessionId), - title: 'Auto-Cleared', - message: `Context reset at ${(data.tokens || 0).toLocaleString()} tokens`, - }); + this._notifySession(data.sessionId, 'info', 'auto-clear', 'Auto-Cleared', `Context reset at ${(data.tokens || 0).toLocaleString()} tokens`); } _onSessionCliInfo(data) { @@ -1983,6 +1952,18 @@ class CodemanApp { return this.getShortId(session.id); } + _notifySession(sessionId, urgency, category, title, message) { + const session = this.sessions.get(sessionId); + this.notificationManager?.notify({ + urgency, + category, + sessionId, + sessionName: session?.name || this.getShortId(sessionId), + title, + message, + }); + } + /** * Clean up state from the previous session before switching tabs. * Handles: WebSocket teardown, CJK clear, flicker filter, tab completion, diff --git a/src/web/public/panels-ui.js b/src/web/public/panels-ui.js index 4de9b489..1226acf3 100644 --- a/src/web/public/panels-ui.js +++ b/src/web/public/panels-ui.js @@ -14,6 +14,13 @@ */ Object.assign(CodemanApp.prototype, { + _addActivityEntry(agentId, entry, maxSize = 50) { + const activity = this.subagentActivity.get(agentId) || []; + activity.push(entry); + if (activity.length > maxSize) activity.shift(); + this.subagentActivity.set(agentId, activity); + }, + // Tasks _onTaskCreated(data) { this.renderSessionTabs(); @@ -106,15 +113,7 @@ Object.assign(CodemanApp.prototype, { // Notify about new subagent discovery const parentId = this.subagentParentMap.get(data.agentId); - const parentSession = parentId ? this.sessions.get(parentId) : null; - this.notificationManager?.notify({ - urgency: 'info', - category: 'subagent-spawn', - sessionId: parentId || data.sessionId, - sessionName: parentSession?.name || parentId || data.sessionId, - title: 'Subagent Spawned', - message: data.description || 'New background agent started', - }); + this._notifySession(parentId || data.sessionId, 'info', 'subagent-spawn', 'Subagent Spawned', data.description || 'New background agent started'); }, _onSubagentUpdated(data) { @@ -135,10 +134,7 @@ Object.assign(CodemanApp.prototype, { }, _onSubagentToolCall(data) { - const activity = this.subagentActivity.get(data.agentId) || []; - activity.push({ type: 'tool', ...data }); - if (activity.length > 50) activity.shift(); // Keep last 50 entries - this.subagentActivity.set(data.agentId, activity); + this._addActivityEntry(data.agentId, { type: 'tool', ...data }); if (this.activeSubagentId === data.agentId) { this.renderSubagentDetail(); } @@ -150,10 +146,7 @@ Object.assign(CodemanApp.prototype, { }, _onSubagentProgress(data) { - const activity = this.subagentActivity.get(data.agentId) || []; - activity.push({ type: 'progress', ...data }); - if (activity.length > 50) activity.shift(); - this.subagentActivity.set(data.agentId, activity); + this._addActivityEntry(data.agentId, { type: 'progress', ...data }); if (this.activeSubagentId === data.agentId) { this.renderSubagentDetail(); } @@ -164,10 +157,7 @@ Object.assign(CodemanApp.prototype, { }, _onSubagentMessage(data) { - const activity = this.subagentActivity.get(data.agentId) || []; - activity.push({ type: 'message', ...data }); - if (activity.length > 50) activity.shift(); - this.subagentActivity.set(data.agentId, activity); + this._addActivityEntry(data.agentId, { type: 'message', ...data }); if (this.activeSubagentId === data.agentId) { this.renderSubagentDetail(); } @@ -190,10 +180,7 @@ Object.assign(CodemanApp.prototype, { } // Add to activity stream - const activity = this.subagentActivity.get(data.agentId) || []; - activity.push({ type: 'tool_result', ...data }); - if (activity.length > 50) activity.shift(); - this.subagentActivity.set(data.agentId, activity); + this._addActivityEntry(data.agentId, { type: 'tool_result', ...data }); if (this.activeSubagentId === data.agentId) { this.renderSubagentDetail(); @@ -224,15 +211,7 @@ Object.assign(CodemanApp.prototype, { // Notify about subagent completion const parentId = this.subagentParentMap.get(data.agentId); - const parentSession = parentId ? this.sessions.get(parentId) : null; - this.notificationManager?.notify({ - urgency: 'info', - category: 'subagent-complete', - sessionId: parentId || existing?.sessionId || data.sessionId, - sessionName: parentSession?.name || parentId || data.sessionId, - title: 'Subagent Completed', - message: existing?.description || data.description || 'Background agent finished', - }); + this._notifySession(parentId || existing?.sessionId || data.sessionId, 'info', 'subagent-complete', 'Subagent Completed', existing?.description || data.description || 'Background agent finished'); // Clean up activity/tool data for completed agents after 5 minutes // This prevents memory leaks from long-running sessions with many subagents diff --git a/src/web/public/ralph-panel.js b/src/web/public/ralph-panel.js index 1bf39b02..dbd74afa 100644 --- a/src/web/public/ralph-panel.js +++ b/src/web/public/ralph-panel.js @@ -45,15 +45,7 @@ Object.assign(CodemanApp.prototype, { this.updateRalphState(data.sessionId, existing); } - const session = this.sessions.get(data.sessionId); - this.notificationManager?.notify({ - urgency: 'warning', - category: 'ralph-complete', - sessionId: data.sessionId, - sessionName: session?.name || this.getShortId(data.sessionId), - title: 'Loop Complete', - message: `Completion: ${data.phrase || 'unknown'}`, - }); + this._notifySession(data.sessionId, 'warning', 'ralph-complete', 'Loop Complete', `Completion: ${data.phrase || 'unknown'}`); }, _onRalphStatusUpdate(data) { @@ -68,28 +60,12 @@ Object.assign(CodemanApp.prototype, { this.updateRalphState(data.sessionId, { circuitBreaker: data.status }); // Notify if circuit breaker opens if (data.status.state === 'OPEN') { - const session = this.sessions.get(data.sessionId); - this.notificationManager?.notify({ - urgency: 'critical', - category: 'circuit-breaker', - sessionId: data.sessionId, - sessionName: session?.name || this.getShortId(data.sessionId), - title: 'Circuit Breaker Open', - message: data.status.reason || 'Loop stuck - no progress detected', - }); + this._notifySession(data.sessionId, 'critical', 'circuit-breaker', 'Circuit Breaker Open', data.status.reason || 'Loop stuck - no progress detected'); } }, _onExitGateMet(data) { - const session = this.sessions.get(data.sessionId); - this.notificationManager?.notify({ - urgency: 'warning', - category: 'exit-gate', - sessionId: data.sessionId, - sessionName: session?.name || this.getShortId(data.sessionId), - title: 'Exit Gate Met', - message: `Loop ready to exit (indicators: ${data.completionIndicators})`, - }); + this._notifySession(data.sessionId, 'warning', 'exit-gate', 'Exit Gate Met', `Loop ready to exit (indicators: ${data.completionIndicators})`); }, // Bash tools diff --git a/src/web/public/respawn-ui.js b/src/web/public/respawn-ui.js index 930ca824..05179b54 100644 --- a/src/web/public/respawn-ui.js +++ b/src/web/public/respawn-ui.js @@ -44,21 +44,13 @@ Object.assign(CodemanApp.prototype, { }, _onRespawnBlocked(data) { - const session = this.sessions.get(data.sessionId); const reasonMap = { circuit_breaker_open: 'Circuit Breaker Open', exit_signal: 'Exit Signal Detected', status_blocked: 'Claude Reported BLOCKED', }; const title = reasonMap[data.reason] || 'Respawn Blocked'; - this.notificationManager?.notify({ - urgency: 'critical', - category: 'respawn-blocked', - sessionId: data.sessionId, - sessionName: session?.name || this.getShortId(data.sessionId), - title, - message: data.details, - }); + this._notifySession(data.sessionId, 'critical', 'respawn-blocked', title, data.details); // Update respawn panel to show blocked state if (data.sessionId === this.activeSessionId) { const stateEl = document.getElementById('respawnStateLabel'); @@ -71,14 +63,7 @@ Object.assign(CodemanApp.prototype, { _onRespawnAutoAcceptSent(data) { const session = this.sessions.get(data.sessionId); - this.notificationManager?.notify({ - urgency: 'info', - category: 'auto-accept', - sessionId: data.sessionId, - sessionName: session?.name || this.getShortId(data.sessionId), - title: 'Plan Accepted', - message: `Accepted plan mode for ${session?.name || 'session'}`, - }); + this._notifySession(data.sessionId, 'info', 'auto-accept', 'Plan Accepted', `Accepted plan mode for ${session?.name || 'session'}`); }, _onRespawnDetectionUpdate(data) { @@ -144,15 +129,7 @@ Object.assign(CodemanApp.prototype, { }, _onRespawnError(data) { - const session = this.sessions.get(data.sessionId); - this.notificationManager?.notify({ - urgency: 'critical', - category: 'session-error', - sessionId: data.sessionId, - sessionName: session?.name || data.sessionId, - title: 'Respawn Error', - message: data.error || data.message || 'Respawn encountered an error', - }); + this._notifySession(data.sessionId, 'critical', 'session-error', 'Respawn Error', data.error || data.message || 'Respawn encountered an error'); }, _onRespawnActionLog(data) { diff --git a/src/web/public/settings-ui.js b/src/web/public/settings-ui.js index 523935d2..d3110875 100644 --- a/src/web/public/settings-ui.js +++ b/src/web/public/settings-ui.js @@ -14,92 +14,46 @@ Object.assign(CodemanApp.prototype, { // Hooks (Claude Code hook events) _onHookIdlePrompt(data) { - const session = this.sessions.get(data.sessionId); // Always track pending hook - alert will show when switching away from session if (data.sessionId) { this.setPendingHook(data.sessionId, 'idle_prompt'); } - this.notificationManager?.notify({ - urgency: 'warning', - category: 'hook-idle', - sessionId: data.sessionId, - sessionName: session?.name || data.sessionId, - title: 'Waiting for Input', - message: data.message || 'Claude is idle and waiting for a prompt', - }); + this._notifySession(data.sessionId, 'warning', 'hook-idle', 'Waiting for Input', data.message || 'Claude is idle and waiting for a prompt'); }, _onHookPermissionPrompt(data) { - const session = this.sessions.get(data.sessionId); // Always track pending hook - action alerts need user interaction to clear if (data.sessionId) { this.setPendingHook(data.sessionId, 'permission_prompt'); } const toolInfo = data.tool ? `${data.tool}${data.command ? ': ' + data.command : data.file ? ': ' + data.file : ''}` : ''; - this.notificationManager?.notify({ - urgency: 'critical', - category: 'hook-permission', - sessionId: data.sessionId, - sessionName: session?.name || data.sessionId, - title: 'Permission Required', - message: toolInfo || 'Claude needs tool approval to continue', - }); + this._notifySession(data.sessionId, 'critical', 'hook-permission', 'Permission Required', toolInfo || 'Claude needs tool approval to continue'); }, _onHookElicitationDialog(data) { - const session = this.sessions.get(data.sessionId); // Always track pending hook - action alerts need user interaction to clear if (data.sessionId) { this.setPendingHook(data.sessionId, 'elicitation_dialog'); } - this.notificationManager?.notify({ - urgency: 'critical', - category: 'hook-elicitation', - sessionId: data.sessionId, - sessionName: session?.name || data.sessionId, - title: 'Question Asked', - message: data.question || 'Claude is asking a question and waiting for your answer', - }); + this._notifySession(data.sessionId, 'critical', 'hook-elicitation', 'Question Asked', data.question || 'Claude is asking a question and waiting for your answer'); }, _onHookStop(data) { - const session = this.sessions.get(data.sessionId); // Clear all pending hooks when Claude finishes responding if (data.sessionId) { this.clearPendingHooks(data.sessionId); } - this.notificationManager?.notify({ - urgency: 'info', - category: 'hook-stop', - sessionId: data.sessionId, - sessionName: session?.name || data.sessionId, - title: 'Response Complete', - message: data.reason || 'Claude has finished responding', - }); + this._notifySession(data.sessionId, 'info', 'hook-stop', 'Response Complete', data.reason || 'Claude has finished responding'); }, _onHookTeammateIdle(data) { const session = this.sessions.get(data.sessionId); - this.notificationManager?.notify({ - urgency: 'warning', - category: 'hook-teammate-idle', - sessionId: data.sessionId, - sessionName: session?.name || data.sessionId, - title: 'Teammate Idle', - message: `A teammate is idle in ${session?.name || data.sessionId}`, - }); + this._notifySession(data.sessionId, 'warning', 'hook-teammate-idle', 'Teammate Idle', `A teammate is idle in ${session?.name || data.sessionId}`); }, _onHookTaskCompleted(data) { const session = this.sessions.get(data.sessionId); - this.notificationManager?.notify({ - urgency: 'info', - category: 'hook-task-completed', - sessionId: data.sessionId, - sessionName: session?.name || data.sessionId, - title: 'Task Completed', - message: `A team task completed in ${session?.name || data.sessionId}`, - }); + this._notifySession(data.sessionId, 'info', 'hook-task-completed', 'Task Completed', `A team task completed in ${session?.name || data.sessionId}`); }, @@ -517,16 +471,17 @@ Object.assign(CodemanApp.prototype, { } }, - _updateTunnelUrlDisplay(url) { - const row = document.getElementById('tunnelUrlRow'); - const display = document.getElementById('tunnelUrlDisplay'); + _updateTunnelUrlRow(rowId, displayId, url, suffix = '') { + const row = document.getElementById(rowId); + const display = document.getElementById(displayId); if (!row || !display) return; if (url) { + const fullUrl = url + suffix; row.style.display = ''; - display.textContent = url; + display.textContent = fullUrl; display.onclick = () => { - navigator.clipboard.writeText(url).then(() => { - this.showToast('Tunnel URL copied', 'success'); + navigator.clipboard.writeText(fullUrl).then(() => { + this.showToast(`${suffix ? 'Upload' : 'Tunnel'} URL copied`, 'success'); }); }; } else { @@ -534,24 +489,11 @@ Object.assign(CodemanApp.prototype, { display.textContent = ''; display.onclick = null; } - // Upload URL row - const uploadRow = document.getElementById('tunnelUploadUrlRow'); - const uploadDisplay = document.getElementById('tunnelUploadUrlDisplay'); - if (!uploadRow || !uploadDisplay) return; - if (url) { - const uploadUrl = url + '/upload.html'; - uploadRow.style.display = ''; - uploadDisplay.textContent = uploadUrl; - uploadDisplay.onclick = () => { - navigator.clipboard.writeText(uploadUrl).then(() => { - this.showToast('Upload URL copied', 'success'); - }); - }; - } else { - uploadRow.style.display = 'none'; - uploadDisplay.textContent = ''; - uploadDisplay.onclick = null; - } + }, + + _updateTunnelUrlDisplay(url) { + this._updateTunnelUrlRow('tunnelUrlRow', 'tunnelUrlDisplay', url); + this._updateTunnelUrlRow('tunnelUploadUrlRow', 'tunnelUploadUrlDisplay', url, '/upload.html'); }, showTunnelQR() { diff --git a/src/web/route-helpers.ts b/src/web/route-helpers.ts index ad79b803..a4972b97 100644 --- a/src/web/route-helpers.ts +++ b/src/web/route-helpers.ts @@ -185,6 +185,26 @@ export function sanitizeHookData(data: Record | null | undefine return safeFields; } +/** + * Toggles a service (watcher/manager) on or off based on an enabled flag. + * Logs start/stop to console with the given label. Runs an optional callback after starting. + */ +export function toggleService( + enabled: boolean, + service: { isRunning(): boolean; start(): void; stop(): void }, + label: string, + onStart?: () => void +): void { + if (enabled && !service.isRunning()) { + service.start(); + onStart?.(); + console.log(`${label} started via settings change`); + } else if (!enabled && service.isRunning()) { + service.stop(); + console.log(`${label} stopped via settings change`); + } +} + /** * Auto-configure Ralph tracker for a session. * diff --git a/src/web/routes/orchestrator-routes.ts b/src/web/routes/orchestrator-routes.ts index 5f330b0a..fa0d4f6f 100644 --- a/src/web/routes/orchestrator-routes.ts +++ b/src/web/routes/orchestrator-routes.ts @@ -39,43 +39,37 @@ export function registerOrchestratorRoutes(app: FastifyInstance, ctx: Orchestrat return loop; } + const EVENT_MAP: [string, (typeof SseEvent)[keyof typeof SseEvent], string[]][] = [ + ['stateChanged', SseEvent.OrchestratorStateChanged, ['state', 'prevState']], + ['planProgress', SseEvent.OrchestratorPlanProgress, ['phase', 'detail']], + ['planReady', SseEvent.OrchestratorPlanReady, ['plan']], + ['phaseStarted', SseEvent.OrchestratorPhaseStarted, ['phase']], + ['phaseCompleted', SseEvent.OrchestratorPhaseCompleted, ['phase']], + ['phaseFailed', SseEvent.OrchestratorPhaseFailed, ['phase', 'reason']], + ['taskAssigned', SseEvent.OrchestratorTaskAssigned, ['task', 'sessionId']], + ['taskCompleted', SseEvent.OrchestratorTaskCompleted, ['task']], + ['taskFailed', SseEvent.OrchestratorTaskFailed, ['task', 'error']], + ['completed', SseEvent.OrchestratorCompleted, ['stats']], + ]; + let forwardingLoop: import('../../orchestrator-loop.js').OrchestratorLoop | null = null; function setupEventForwarding(loop: import('../../orchestrator-loop.js').OrchestratorLoop) { if (forwardingLoop === loop) return; // Already attached to this loop instance forwardingLoop = loop; - loop.on('stateChanged', (state, prevState) => { - ctx.broadcast(SseEvent.OrchestratorStateChanged, { state, prevState }); - }); - loop.on('planProgress', (phase, detail) => { - ctx.broadcast(SseEvent.OrchestratorPlanProgress, { phase, detail }); - }); - loop.on('planReady', (plan) => { - ctx.broadcast(SseEvent.OrchestratorPlanReady, { plan }); - }); - loop.on('phaseStarted', (phase) => { - ctx.broadcast(SseEvent.OrchestratorPhaseStarted, { phase }); - }); - loop.on('phaseCompleted', (phase) => { - ctx.broadcast(SseEvent.OrchestratorPhaseCompleted, { phase }); - }); - loop.on('phaseFailed', (phase, reason) => { - ctx.broadcast(SseEvent.OrchestratorPhaseFailed, { phase, reason }); - }); + for (const [event, sseEvent, argNames] of EVENT_MAP) { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + loop.on(event, (...args: any[]) => { + const payload: Record = {}; + argNames.forEach((name, i) => { + payload[name] = args[i]; + }); + ctx.broadcast(sseEvent, payload); + }); + } + // Special cases with non-trivial payload transforms loop.on('verificationResult', (phase, result) => { ctx.broadcast(SseEvent.OrchestratorVerification, { phaseId: phase.id, result }); }); - loop.on('taskAssigned', (task, sessionId) => { - ctx.broadcast(SseEvent.OrchestratorTaskAssigned, { task, sessionId }); - }); - loop.on('taskCompleted', (task) => { - ctx.broadcast(SseEvent.OrchestratorTaskCompleted, { task }); - }); - loop.on('taskFailed', (task, error) => { - ctx.broadcast(SseEvent.OrchestratorTaskFailed, { task, error }); - }); - loop.on('completed', (stats) => { - ctx.broadcast(SseEvent.OrchestratorCompleted, { stats }); - }); loop.on('error', (error) => { ctx.broadcast(SseEvent.OrchestratorError, { error: error.message }); }); diff --git a/src/web/routes/system-routes.ts b/src/web/routes/system-routes.ts index d6e99d70..a041f799 100644 --- a/src/web/routes/system-routes.ts +++ b/src/web/routes/system-routes.ts @@ -24,7 +24,14 @@ import { import { subagentWatcher } from '../../subagent-watcher.js'; import { imageWatcher } from '../../image-watcher.js'; import { getLifecycleLog } from '../../session-lifecycle-log.js'; -import { findSessionOrFail, formatUptime, parseBody, readJsonConfig, SETTINGS_PATH } from '../route-helpers.js'; +import { + findSessionOrFail, + formatUptime, + parseBody, + readJsonConfig, + toggleService, + SETTINGS_PATH, +} from '../route-helpers.js'; import { SseEvent } from '../sse-events.js'; import type { SessionPort, EventPort, ConfigPort, InfraPort, AuthPort } from '../ports/index.js'; import { AUTH_COOKIE_NAME } from '../middleware/auth.js'; @@ -281,15 +288,20 @@ export function registerSystemRoutes( // ========== Stats ========== - app.get('/api/stats', async () => { - const activeSessionTokens: Record = {}; + function collectActiveTokens(): Record { + const tokens: Record = {}; for (const [sessionId, session] of ctx.sessions) { - activeSessionTokens[sessionId] = { + tokens[sessionId] = { inputTokens: session.inputTokens, outputTokens: session.outputTokens, totalCost: session.totalCost, }; } + return tokens; + } + + app.get('/api/stats', async () => { + const activeSessionTokens = collectActiveTokens(); return { success: true, stats: ctx.store.getAggregateStats(activeSessionTokens), @@ -298,14 +310,7 @@ export function registerSystemRoutes( }); app.get('/api/token-stats', async () => { - const activeSessionTokens: Record = {}; - for (const [sessionId, session] of ctx.sessions) { - activeSessionTokens[sessionId] = { - inputTokens: session.inputTokens, - outputTokens: session.outputTokens, - totalCost: session.totalCost, - }; - } + const activeSessionTokens = collectActiveTokens(); return { success: true, daily: ctx.store.getDailyStats(30), @@ -414,30 +419,17 @@ export function registerSystemRoutes( await fs.writeFile(SETTINGS_PATH, JSON.stringify(merged, null, 2)); // Handle subagent tracking toggle dynamically - const subagentEnabled = settings.subagentTrackingEnabled ?? true; - if (subagentEnabled && !subagentWatcher.isRunning()) { - subagentWatcher.start(); - console.log('Subagent watcher started via settings change'); - } else if (!subagentEnabled && subagentWatcher.isRunning()) { - subagentWatcher.stop(); - console.log('Subagent watcher stopped via settings change'); - } + toggleService((settings.subagentTrackingEnabled as boolean) ?? true, subagentWatcher, 'Subagent watcher'); // Handle image watcher toggle dynamically - const imageWatcherEnabled = settings.imageWatcherEnabled ?? false; - if (imageWatcherEnabled && !imageWatcher.isRunning()) { - imageWatcher.start(); + toggleService((settings.imageWatcherEnabled as boolean) ?? false, imageWatcher, 'Image watcher', () => { // Re-watch all active sessions that have image watcher enabled for (const session of ctx.sessions.values()) { if (session.imageWatcherEnabled) { imageWatcher.watchSession(session.id, session.workingDir); } } - console.log('Image watcher started via settings change'); - } else if (!imageWatcherEnabled && imageWatcher.isRunning()) { - imageWatcher.stop(); - console.log('Image watcher stopped via settings change'); - } + }); // Handle tunnel toggle dynamically if ('tunnelEnabled' in settings) {