mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-10-07 07:59:42 +02:00
fix: improve stability with memory leak fixes and atomic writes
Session Manager: - Store event handlers for proper cleanup on session stop - Add resetSessionManager() for test isolation State Store: - Use atomic write pattern (temp file + rename) to prevent corruption - Add tokenStats to AppState interface Subagent Watcher: - Clean up pendingToolCalls on agent completion - Add cleanupStaleAgents() for 24h+ old agents - Clear all state maps on stop() for clean restart Web Server: - Remove listeners from transcript watcher before stopping - Clean up respawn timers when session is deleted - Use cleanupSession() in spawn orchestrator for proper resource cleanup Frontend: - Close existing EventSource before reconnecting to prevent duplicates - Close failed connection before scheduling reconnect Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
+54
-17
@@ -48,8 +48,17 @@ export interface SessionManagerEvents {
|
|||||||
* @fires SessionManagerEvents.sessionOutput
|
* @fires SessionManagerEvents.sessionOutput
|
||||||
* @fires SessionManagerEvents.sessionCompletion
|
* @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 {
|
export class SessionManager extends EventEmitter {
|
||||||
private sessions: Map<string, Session> = new Map();
|
private sessions: Map<string, Session> = new Map();
|
||||||
|
private sessionHandlers: Map<string, SessionHandlers> = new Map();
|
||||||
private store = getStore();
|
private store = getStore();
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -89,25 +98,32 @@ export class SessionManager extends EventEmitter {
|
|||||||
|
|
||||||
const session = new Session({ workingDir });
|
const session = new Session({ workingDir });
|
||||||
|
|
||||||
// Set up event forwarding
|
// Set up event forwarding with stored handlers for cleanup
|
||||||
session.on('output', (data) => {
|
const handlers: SessionHandlers = {
|
||||||
this.emit('sessionOutput', session.id, data);
|
output: (data: string) => {
|
||||||
this.updateSessionState(session);
|
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) => {
|
session.on('output', handlers.output);
|
||||||
this.emit('sessionError', session.id, data);
|
session.on('error', handlers.error);
|
||||||
this.updateSessionState(session);
|
session.on('completion', handlers.completion);
|
||||||
});
|
session.on('exit', handlers.exit);
|
||||||
|
|
||||||
session.on('completion', (phrase) => {
|
// Store handlers for later cleanup
|
||||||
this.emit('sessionCompletion', session.id, phrase);
|
this.sessionHandlers.set(session.id, handlers);
|
||||||
});
|
|
||||||
|
|
||||||
session.on('exit', () => {
|
|
||||||
this.emit('sessionStopped', session.id);
|
|
||||||
this.updateSessionState(session);
|
|
||||||
});
|
|
||||||
|
|
||||||
await session.start();
|
await session.start();
|
||||||
|
|
||||||
@@ -136,6 +152,16 @@ export class SessionManager extends EventEmitter {
|
|||||||
return;
|
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();
|
await session.stop();
|
||||||
this.sessions.delete(id);
|
this.sessions.delete(id);
|
||||||
this.updateSessionState(session);
|
this.updateSessionState(session);
|
||||||
@@ -235,3 +261,14 @@ export function getSessionManager(): SessionManager {
|
|||||||
}
|
}
|
||||||
return managerInstance;
|
return managerInstance;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Resets the singleton SessionManager instance.
|
||||||
|
* Primarily used for testing to ensure test isolation.
|
||||||
|
*/
|
||||||
|
export async function resetSessionManager(): Promise<void> {
|
||||||
|
if (managerInstance) {
|
||||||
|
await managerInstance.stopAllSessions();
|
||||||
|
managerInstance = null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+42
-12
@@ -14,7 +14,7 @@
|
|||||||
* @module state-store
|
* @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 { homedir } from 'node:os';
|
||||||
import { dirname, join } from 'node:path';
|
import { dirname, join } from 'node:path';
|
||||||
import { AppState, createInitialState, RalphSessionState, createInitialRalphSessionState, GlobalStats, createInitialGlobalStats, TokenStats, TokenUsageEntry } from './types.js';
|
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).
|
* Use when guaranteed persistence is required (e.g., before shutdown).
|
||||||
*/
|
*/
|
||||||
saveNow(): void {
|
saveNow(): void {
|
||||||
@@ -120,7 +121,22 @@ export class StateStore {
|
|||||||
}
|
}
|
||||||
this.dirty = false;
|
this.dirty = false;
|
||||||
this.ensureDir();
|
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. */
|
/** Flushes any pending main state save. Call before shutdown. */
|
||||||
@@ -293,14 +309,13 @@ export class StateStore {
|
|||||||
* Get or initialize token stats from state.
|
* Get or initialize token stats from state.
|
||||||
*/
|
*/
|
||||||
getTokenStats(): TokenStats {
|
getTokenStats(): TokenStats {
|
||||||
const state = this.state as AppState & { tokenStats?: TokenStats };
|
if (!this.state.tokenStats) {
|
||||||
if (!state.tokenStats) {
|
this.state.tokenStats = {
|
||||||
state.tokenStats = {
|
|
||||||
daily: [],
|
daily: [],
|
||||||
lastUpdated: Date.now(),
|
lastUpdated: Date.now(),
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
return state.tokenStats;
|
return this.state.tokenStats;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -425,7 +440,10 @@ export class StateStore {
|
|||||||
}, SAVE_DEBOUNCE_MS);
|
}, 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 {
|
private saveRalphStatesNow(): void {
|
||||||
if (this.ralphStateSaveTimeout) {
|
if (this.ralphStateSaveTimeout) {
|
||||||
clearTimeout(this.ralphStateSaveTimeout);
|
clearTimeout(this.ralphStateSaveTimeout);
|
||||||
@@ -436,11 +454,23 @@ export class StateStore {
|
|||||||
}
|
}
|
||||||
this.ralphStateDirty = false;
|
this.ralphStateDirty = false;
|
||||||
this.ensureDir();
|
this.ensureDir();
|
||||||
const data: Record<string, RalphSessionState> = {};
|
const data = Object.fromEntries(this.ralphStates);
|
||||||
for (const [sessionId, state] of this.ralphStates) {
|
// Atomic write: write to temp file, then rename (atomic on POSIX)
|
||||||
data[sessionId] = state;
|
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. */
|
/** Returns inner state for a session, or null if not found. */
|
||||||
|
|||||||
+12
-1
@@ -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 IDLE_TIMEOUT_MS = 30000; // Consider agent idle after 30s of no activity
|
||||||
const POLL_INTERVAL_MS = 1000; // Check for new files every second
|
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 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 ==========
|
// ========== SubagentWatcher Class ==========
|
||||||
|
|
||||||
@@ -188,10 +189,14 @@ export class SubagentWatcher extends EventEmitter {
|
|||||||
const alive = await this.checkSubagentAlive(agentId);
|
const alive = await this.checkSubagentAlive(agentId);
|
||||||
if (!alive) {
|
if (!alive) {
|
||||||
info.status = 'completed';
|
info.status = 'completed';
|
||||||
|
// Clean up pendingToolCalls for this agent to prevent memory leak
|
||||||
|
this.pendingToolCalls.delete(agentId);
|
||||||
this.emit('subagent:completed', info);
|
this.emit('subagent:completed', info);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
// Periodically clean up stale completed agents (older than 24 hours)
|
||||||
|
this.cleanupStaleAgents();
|
||||||
}, LIVENESS_CHECK_MS);
|
}, LIVENESS_CHECK_MS);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -230,7 +235,7 @@ export class SubagentWatcher extends EventEmitter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Stop watching
|
* Stop watching and clean up all state
|
||||||
*/
|
*/
|
||||||
stop(): void {
|
stop(): void {
|
||||||
this._isRunning = false;
|
this._isRunning = false;
|
||||||
@@ -259,6 +264,12 @@ export class SubagentWatcher extends EventEmitter {
|
|||||||
clearTimeout(timer);
|
clearTimeout(timer);
|
||||||
}
|
}
|
||||||
this.idleTimers.clear();
|
this.idleTimers.clear();
|
||||||
|
|
||||||
|
// Clear all state for clean restart
|
||||||
|
this.filePositions.clear();
|
||||||
|
this.agentInfo.clear();
|
||||||
|
this.knownSubagentDirs.clear();
|
||||||
|
this.pendingToolCalls.clear();
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -229,6 +229,8 @@ export interface AppState {
|
|||||||
config: AppConfig;
|
config: AppConfig;
|
||||||
/** Global statistics (cumulative across all sessions) */
|
/** Global statistics (cumulative across all sessions) */
|
||||||
globalStats?: GlobalStats;
|
globalStats?: GlobalStats;
|
||||||
|
/** Daily token usage statistics */
|
||||||
|
tokenStats?: TokenStats;
|
||||||
}
|
}
|
||||||
|
|
||||||
// ========== Respawn Controller Types ==========
|
// ========== Respawn Controller Types ==========
|
||||||
|
|||||||
@@ -965,11 +965,22 @@ class ClaudemanApp {
|
|||||||
// ========== SSE Connection ==========
|
// ========== SSE Connection ==========
|
||||||
|
|
||||||
connectSSE() {
|
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 = new EventSource('/api/events');
|
||||||
|
|
||||||
this.eventSource.onopen = () => this.setConnectionStatus('connected');
|
this.eventSource.onopen = () => this.setConnectionStatus('connected');
|
||||||
this.eventSource.onerror = () => {
|
this.eventSource.onerror = () => {
|
||||||
this.setConnectionStatus('disconnected');
|
this.setConnectionStatus('disconnected');
|
||||||
|
// Close the failed connection before scheduling reconnect
|
||||||
|
if (this.eventSource) {
|
||||||
|
this.eventSource.close();
|
||||||
|
this.eventSource = null;
|
||||||
|
}
|
||||||
setTimeout(() => this.connectSSE(), 3000);
|
setTimeout(() => this.connectSSE(), 3000);
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
+10
-7
@@ -1969,6 +1969,7 @@ export class WebServer extends EventEmitter {
|
|||||||
private stopTranscriptWatcher(sessionId: string): void {
|
private stopTranscriptWatcher(sessionId: string): void {
|
||||||
const watcher = this.transcriptWatchers.get(sessionId);
|
const watcher = this.transcriptWatchers.get(sessionId);
|
||||||
if (watcher) {
|
if (watcher) {
|
||||||
|
watcher.removeAllListeners(); // Prevent memory leaks from attached listeners
|
||||||
watcher.stop();
|
watcher.stop();
|
||||||
this.transcriptWatchers.delete(sessionId);
|
this.transcriptWatchers.delete(sessionId);
|
||||||
}
|
}
|
||||||
@@ -2247,6 +2248,12 @@ export class WebServer extends EventEmitter {
|
|||||||
controller.removeAllListeners();
|
controller.removeAllListeners();
|
||||||
this.respawnControllers.delete(session.id);
|
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) {
|
} catch (err) {
|
||||||
console.error(`[Server] Error cleaning up respawn controller for ${session.id}:`, 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;
|
return session ? session.totalCost : 0;
|
||||||
},
|
},
|
||||||
stopSession: async (sessionId: string) => {
|
stopSession: async (sessionId: string) => {
|
||||||
const session = this.sessions.get(sessionId);
|
// Use cleanupSession to properly clean up all resources (respawn controllers,
|
||||||
if (session) {
|
// run summary trackers, file streams, Ralph state, etc.)
|
||||||
await session.stop();
|
await this.cleanupSession(sessionId);
|
||||||
this.sessions.delete(sessionId);
|
|
||||||
this.broadcast('session:deleted', { id: sessionId });
|
|
||||||
this.persistSessionState(session);
|
|
||||||
}
|
|
||||||
},
|
},
|
||||||
onSessionCompletion: (sessionId: string, handler: (phrase: string) => void) => {
|
onSessionCompletion: (sessionId: string, handler: (phrase: string) => void) => {
|
||||||
const session = this.sessions.get(sessionId);
|
const session = this.sessions.get(sessionId);
|
||||||
|
|||||||
Reference in New Issue
Block a user