diff --git a/src/session-manager.ts b/src/session-manager.ts index 48ef0cfd..1a44953f 100644 --- a/src/session-manager.ts +++ b/src/session-manager.ts @@ -48,8 +48,17 @@ export interface SessionManagerEvents { * @fires SessionManagerEvents.sessionOutput * @fires SessionManagerEvents.sessionCompletion */ +/** Stored event handlers for a session, used for cleanup */ +interface SessionHandlers { + output: (data: string) => void; + error: (data: string) => void; + completion: (phrase: string) => void; + exit: () => void; +} + export class SessionManager extends EventEmitter { private sessions: Map = new Map(); + private sessionHandlers: Map = new Map(); private store = getStore(); /** @@ -89,25 +98,32 @@ export class SessionManager extends EventEmitter { const session = new Session({ workingDir }); - // Set up event forwarding - session.on('output', (data) => { - this.emit('sessionOutput', session.id, data); - this.updateSessionState(session); - }); + // Set up event forwarding with stored handlers for cleanup + const handlers: SessionHandlers = { + output: (data: string) => { + this.emit('sessionOutput', session.id, data); + this.updateSessionState(session); + }, + error: (data: string) => { + this.emit('sessionError', session.id, data); + this.updateSessionState(session); + }, + completion: (phrase: string) => { + this.emit('sessionCompletion', session.id, phrase); + }, + exit: () => { + this.emit('sessionStopped', session.id); + this.updateSessionState(session); + }, + }; - session.on('error', (data) => { - this.emit('sessionError', session.id, data); - this.updateSessionState(session); - }); + session.on('output', handlers.output); + session.on('error', handlers.error); + session.on('completion', handlers.completion); + session.on('exit', handlers.exit); - session.on('completion', (phrase) => { - this.emit('sessionCompletion', session.id, phrase); - }); - - session.on('exit', () => { - this.emit('sessionStopped', session.id); - this.updateSessionState(session); - }); + // Store handlers for later cleanup + this.sessionHandlers.set(session.id, handlers); await session.start(); @@ -136,6 +152,16 @@ export class SessionManager extends EventEmitter { return; } + // Remove event listeners to prevent memory leaks + const handlers = this.sessionHandlers.get(id); + if (handlers) { + session.off('output', handlers.output); + session.off('error', handlers.error); + session.off('completion', handlers.completion); + session.off('exit', handlers.exit); + this.sessionHandlers.delete(id); + } + await session.stop(); this.sessions.delete(id); this.updateSessionState(session); @@ -235,3 +261,14 @@ export function getSessionManager(): SessionManager { } return managerInstance; } + +/** + * Resets the singleton SessionManager instance. + * Primarily used for testing to ensure test isolation. + */ +export async function resetSessionManager(): Promise { + if (managerInstance) { + await managerInstance.stopAllSessions(); + managerInstance = null; + } +} diff --git a/src/state-store.ts b/src/state-store.ts index 222f61fc..5f30a9c2 100644 --- a/src/state-store.ts +++ b/src/state-store.ts @@ -14,7 +14,7 @@ * @module state-store */ -import { readFileSync, writeFileSync, existsSync, mkdirSync } from 'node:fs'; +import { readFileSync, writeFileSync, existsSync, mkdirSync, renameSync } from 'node:fs'; import { homedir } from 'node:os'; import { dirname, join } from 'node:path'; import { AppState, createInitialState, RalphSessionState, createInitialRalphSessionState, GlobalStats, createInitialGlobalStats, TokenStats, TokenUsageEntry } from './types.js'; @@ -107,7 +107,8 @@ export class StateStore { } /** - * Immediately writes state to disk. + * Immediately writes state to disk using atomic write pattern. + * Writes to temp file first, then renames to prevent corruption on crash. * Use when guaranteed persistence is required (e.g., before shutdown). */ saveNow(): void { @@ -120,7 +121,22 @@ export class StateStore { } this.dirty = false; this.ensureDir(); - writeFileSync(this.filePath, JSON.stringify(this.state, null, 2), 'utf-8'); + // Atomic write: write to temp file, then rename (atomic on POSIX) + const tempPath = this.filePath + '.tmp'; + try { + writeFileSync(tempPath, JSON.stringify(this.state, null, 2), 'utf-8'); + renameSync(tempPath, this.filePath); + } catch (err) { + console.error('[StateStore] Failed to save state:', err); + // Try to clean up temp file on error + try { + if (existsSync(tempPath)) { + const { unlinkSync } = require('node:fs'); + unlinkSync(tempPath); + } + } catch { /* ignore cleanup errors */ } + throw err; + } } /** Flushes any pending main state save. Call before shutdown. */ @@ -293,14 +309,13 @@ export class StateStore { * Get or initialize token stats from state. */ getTokenStats(): TokenStats { - const state = this.state as AppState & { tokenStats?: TokenStats }; - if (!state.tokenStats) { - state.tokenStats = { + if (!this.state.tokenStats) { + this.state.tokenStats = { daily: [], lastUpdated: Date.now(), }; } - return state.tokenStats; + return this.state.tokenStats; } /** @@ -425,7 +440,10 @@ export class StateStore { }, SAVE_DEBOUNCE_MS); } - // Immediate save for inner states + /** + * Immediate save for inner states using atomic write pattern. + * Writes to temp file first, then renames to prevent corruption on crash. + */ private saveRalphStatesNow(): void { if (this.ralphStateSaveTimeout) { clearTimeout(this.ralphStateSaveTimeout); @@ -436,11 +454,23 @@ export class StateStore { } this.ralphStateDirty = false; this.ensureDir(); - const data: Record = {}; - for (const [sessionId, state] of this.ralphStates) { - data[sessionId] = state; + const data = Object.fromEntries(this.ralphStates); + // Atomic write: write to temp file, then rename (atomic on POSIX) + const tempPath = this.ralphStatePath + '.tmp'; + try { + writeFileSync(tempPath, JSON.stringify(data, null, 2), 'utf-8'); + renameSync(tempPath, this.ralphStatePath); + } catch (err) { + console.error('[StateStore] Failed to save Ralph state:', err); + // Try to clean up temp file on error + try { + if (existsSync(tempPath)) { + const { unlinkSync } = require('node:fs'); + unlinkSync(tempPath); + } + } catch { /* ignore cleanup errors */ } + throw err; } - writeFileSync(this.ralphStatePath, JSON.stringify(data, null, 2), 'utf-8'); } /** Returns inner state for a session, or null if not found. */ diff --git a/src/subagent-watcher.ts b/src/subagent-watcher.ts index 866e6de2..332189e4 100644 --- a/src/subagent-watcher.ts +++ b/src/subagent-watcher.ts @@ -123,6 +123,7 @@ const CLAUDE_PROJECTS_DIR = join(homedir(), '.claude/projects'); const IDLE_TIMEOUT_MS = 30000; // Consider agent idle after 30s of no activity const POLL_INTERVAL_MS = 1000; // Check for new files every second const LIVENESS_CHECK_MS = 10000; // Check if subagent processes are still alive every 10s +const STALE_AGENT_MAX_AGE_MS = 24 * 60 * 60 * 1000; // Remove completed agents older than 24 hours // ========== SubagentWatcher Class ========== @@ -188,10 +189,14 @@ export class SubagentWatcher extends EventEmitter { const alive = await this.checkSubagentAlive(agentId); if (!alive) { info.status = 'completed'; + // Clean up pendingToolCalls for this agent to prevent memory leak + this.pendingToolCalls.delete(agentId); this.emit('subagent:completed', info); } } } + // Periodically clean up stale completed agents (older than 24 hours) + this.cleanupStaleAgents(); }, LIVENESS_CHECK_MS); } @@ -230,7 +235,7 @@ export class SubagentWatcher extends EventEmitter { } /** - * Stop watching + * Stop watching and clean up all state */ stop(): void { this._isRunning = false; @@ -259,6 +264,12 @@ export class SubagentWatcher extends EventEmitter { clearTimeout(timer); } this.idleTimers.clear(); + + // Clear all state for clean restart + this.filePositions.clear(); + this.agentInfo.clear(); + this.knownSubagentDirs.clear(); + this.pendingToolCalls.clear(); } /** diff --git a/src/types.ts b/src/types.ts index 1f9d0436..27be1f8d 100644 --- a/src/types.ts +++ b/src/types.ts @@ -229,6 +229,8 @@ export interface AppState { config: AppConfig; /** Global statistics (cumulative across all sessions) */ globalStats?: GlobalStats; + /** Daily token usage statistics */ + tokenStats?: TokenStats; } // ========== Respawn Controller Types ========== diff --git a/src/web/public/app.js b/src/web/public/app.js index da24f4ae..08256340 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -965,11 +965,22 @@ class ClaudemanApp { // ========== SSE Connection ========== connectSSE() { + // Close existing EventSource before creating new one to prevent duplicate connections + if (this.eventSource) { + this.eventSource.close(); + this.eventSource = null; + } + this.eventSource = new EventSource('/api/events'); this.eventSource.onopen = () => this.setConnectionStatus('connected'); this.eventSource.onerror = () => { this.setConnectionStatus('disconnected'); + // Close the failed connection before scheduling reconnect + if (this.eventSource) { + this.eventSource.close(); + this.eventSource = null; + } setTimeout(() => this.connectSSE(), 3000); }; diff --git a/src/web/server.ts b/src/web/server.ts index d4c09cd3..377ff847 100644 --- a/src/web/server.ts +++ b/src/web/server.ts @@ -1969,6 +1969,7 @@ export class WebServer extends EventEmitter { private stopTranscriptWatcher(sessionId: string): void { const watcher = this.transcriptWatchers.get(sessionId); if (watcher) { + watcher.removeAllListeners(); // Prevent memory leaks from attached listeners watcher.stop(); this.transcriptWatchers.delete(sessionId); } @@ -2247,6 +2248,12 @@ export class WebServer extends EventEmitter { controller.removeAllListeners(); this.respawnControllers.delete(session.id); } + // Also clean up the respawn timer to prevent orphaned timers + const timerInfo = this.respawnTimers.get(session.id); + if (timerInfo) { + clearTimeout(timerInfo.timer); + this.respawnTimers.delete(session.id); + } } catch (err) { console.error(`[Server] Error cleaning up respawn controller for ${session.id}:`, err); } @@ -2485,13 +2492,9 @@ export class WebServer extends EventEmitter { return session ? session.totalCost : 0; }, stopSession: async (sessionId: string) => { - const session = this.sessions.get(sessionId); - if (session) { - await session.stop(); - this.sessions.delete(sessionId); - this.broadcast('session:deleted', { id: sessionId }); - this.persistSessionState(session); - } + // Use cleanupSession to properly clean up all resources (respawn controllers, + // run summary trackers, file streams, Ralph state, etc.) + await this.cleanupSession(sessionId); }, onSessionCompletion: (sessionId: string, handler: (phrase: string) => void) => { const session = this.sessions.get(sessionId);