From add8f8c84b89786deece02486003dfe14c3f91b7 Mon Sep 17 00:00:00 2001 From: arkon Date: Thu, 22 Jan 2026 00:08:32 +0100 Subject: [PATCH] fix: memory leak prevention and resource cleanup improvements SSE Client Management: - Add periodic health check (30s) to detect dead SSE connections - Clean up clients with destroyed/non-writable sockets - Clear all SSE clients on server stop Session Lifecycle: - Add _isStopped flag to prevent new timers after session stops - Check flag in auto-compact/clear callbacks to avoid timer races - Check flag in line buffer flush timer creation Server Shutdown: - Add _isStopping flag to prevent timer creation during shutdown - Skip batch timer creation in terminal/output/task/state handlers - Clear SSE health check timer on stop Ralph Tracker: - Call clearDebounceTimers() in clear() method (was missing) - Add MAX_COMPLETION_PHRASE_ENTRIES (50) to trim phrase map - Add MAX_LINE_BUFFER_SIZE (64KB) to prevent unbounded growth - Trim line buffer when exceeding max size Task Tracker: - Add timestamp to pending tool uses for age-based cleanup - Add PENDING_TOOL_USE_MAX_AGE_MS (1 hour) expiry - Add MAX_PENDING_TOOL_USES (100) limit - Clean up old entries on new tool use Co-Authored-By: Claude Opus 4.5 --- src/ralph-tracker.ts | 36 ++++++++++++++++++++++++ src/session.ts | 53 ++++++++++++++++++++++++------------ src/task-tracker.ts | 54 ++++++++++++++++++++++++++++++++++-- src/web/server.ts | 65 ++++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 187 insertions(+), 21 deletions(-) diff --git a/src/ralph-tracker.ts b/src/ralph-tracker.ts index 93243d62..d0b0f7c8 100644 --- a/src/ralph-tracker.ts +++ b/src/ralph-tracker.ts @@ -49,6 +49,17 @@ const CLEANUP_THROTTLE_MS = 30 * 1000; */ const EVENT_DEBOUNCE_MS = 50; +/** + * Maximum number of completion phrase entries to track. + * Prevents unbounded growth if many unique phrases are seen. + */ +const MAX_COMPLETION_PHRASE_ENTRIES = 50; + +/** + * Maximum line buffer size to prevent unbounded growth from long lines. + */ +const MAX_LINE_BUFFER_SIZE = 64 * 1024; + // ========== Pre-compiled Regex Patterns ========== // Pre-compiled for performance (avoid re-compilation on each call) @@ -546,6 +557,12 @@ export class RalphTracker extends EventEmitter { // Buffer data for line-based processing this._lineBuffer += cleanData; + // Prevent unbounded line buffer growth from very long lines + if (this._lineBuffer.length > MAX_LINE_BUFFER_SIZE) { + // Truncate to last portion to preserve recent data + this._lineBuffer = this._lineBuffer.slice(-MAX_LINE_BUFFER_SIZE / 2); + } + // Process complete lines const lines = this._lineBuffer.split('\n'); this._lineBuffer = lines.pop() || ''; // Keep incomplete line in buffer @@ -870,6 +887,23 @@ export class RalphTracker extends EventEmitter { const count = (this._completionPhraseCount.get(phrase) || 0) + 1; this._completionPhraseCount.set(phrase, count); + // Trim completion phrase map if it exceeds the limit + if (this._completionPhraseCount.size > MAX_COMPLETION_PHRASE_ENTRIES) { + // Keep only the most important entries (current expected phrase and highest counts) + const entries = Array.from(this._completionPhraseCount.entries()); + entries.sort((a, b) => b[1] - a[1]); // Sort by count descending + this._completionPhraseCount.clear(); + // Keep top half of entries + const keepCount = Math.floor(MAX_COMPLETION_PHRASE_ENTRIES / 2); + for (let i = 0; i < Math.min(keepCount, entries.length); i++) { + this._completionPhraseCount.set(entries[i][0], entries[i][1]); + } + // Always keep the expected phrase if set + if (this._loopState.completionPhrase && !this._completionPhraseCount.has(this._loopState.completionPhrase)) { + this._completionPhraseCount.set(this._loopState.completionPhrase, 1); + } + } + // Store phrase on first occurrence if (!this._loopState.completionPhrase) { this._loopState.completionPhrase = phrase; @@ -1340,6 +1374,8 @@ export class RalphTracker extends EventEmitter { * @fires todoUpdate - With empty array */ clear(): void { + // Clear debounce timers to prevent stale emissions after clear + this.clearDebounceTimers(); this._loopState = createInitialRalphTrackerState(); // This sets enabled: false this._todos.clear(); this._lineBuffer = ''; diff --git a/src/session.ts b/src/session.ts index 705ed161..510b3d0e 100644 --- a/src/session.ts +++ b/src/session.ts @@ -319,6 +319,8 @@ export class Session extends EventEmitter { private _screenManager: ScreenManager | null = null; private _screenSession: ScreenSession | null = null; private _useScreen: boolean = false; + // Flag to prevent new timers after session is stopped + private _isStopped: boolean = false; // Ralph tracking (Ralph Wiggum loops and todo lists inside Claude Code) private _ralphTracker: RalphTracker; @@ -1071,10 +1073,10 @@ export class Session extends EventEmitter { } // Start flush timer if not running (handles partial lines after 100ms) - if (!this._lineBufferFlushTimer && this._lineBuffer.length > 0) { + if (!this._lineBufferFlushTimer && this._lineBuffer.length > 0 && !this._isStopped) { this._lineBufferFlushTimer = setTimeout(() => { this._lineBufferFlushTimer = null; - if (this._lineBuffer.length > 0) { + if (this._lineBuffer.length > 0 && !this._isStopped) { // Flush partial line as text output this._textOutput.append(this._lineBuffer); this._lineBuffer = ''; @@ -1188,7 +1190,7 @@ export class Session extends EventEmitter { // Check if we should auto-compact based on token threshold private checkAutoCompact(): void { - if (!this._autoCompactEnabled || this._isCompacting || this._isClearing) return; + if (!this._autoCompactEnabled || this._isCompacting || this._isClearing || this._isStopped) return; const totalTokens = this._totalInputTokens + this._totalOutputTokens; if (totalTokens >= this._autoCompactThreshold) { @@ -1198,7 +1200,7 @@ export class Session extends EventEmitter { // Wait for Claude to be idle before compacting const checkAndCompact = () => { // Check if session is still valid (not stopped) - if (!this._isCompacting) return; + if (!this._isCompacting || this._isStopped) return; if (!this._isWorking) { // Send /compact command with optional prompt @@ -1213,24 +1215,30 @@ export class Session extends EventEmitter { }); // Wait a moment then re-enable (longer than clear since compact takes time) - this._autoCompactTimer = setTimeout(() => { - this._autoCompactTimer = null; - this._isCompacting = false; - }, 10000); + if (!this._isStopped) { + this._autoCompactTimer = setTimeout(() => { + this._autoCompactTimer = null; + this._isCompacting = false; + }, 10000); + } } else { // Check again in 2 seconds - this._autoCompactTimer = setTimeout(checkAndCompact, 2000); + if (!this._isStopped) { + this._autoCompactTimer = setTimeout(checkAndCompact, 2000); + } } }; // Start checking after a short delay - this._autoCompactTimer = setTimeout(checkAndCompact, 1000); + if (!this._isStopped) { + this._autoCompactTimer = setTimeout(checkAndCompact, 1000); + } } } // Check if we should auto-clear based on token threshold private checkAutoClear(): void { - if (!this._autoClearEnabled || this._isClearing || this._isCompacting) return; + if (!this._autoClearEnabled || this._isClearing || this._isCompacting || this._isStopped) return; const totalTokens = this._totalInputTokens + this._totalOutputTokens; if (totalTokens >= this._autoClearThreshold) { @@ -1240,7 +1248,7 @@ export class Session extends EventEmitter { // Wait for Claude to be idle before clearing const checkAndClear = () => { // Check if session is still valid (not stopped) - if (!this._isClearing) return; + if (!this._isClearing || this._isStopped) return; if (!this._isWorking) { // Send /clear command @@ -1251,18 +1259,24 @@ export class Session extends EventEmitter { this.emit('autoClear', { tokens: totalTokens, threshold: this._autoClearThreshold }); // Wait a moment then re-enable - this._autoClearTimer = setTimeout(() => { - this._autoClearTimer = null; - this._isClearing = false; - }, 5000); + if (!this._isStopped) { + this._autoClearTimer = setTimeout(() => { + this._autoClearTimer = null; + this._isClearing = false; + }, 5000); + } } else { // Check again in 2 seconds - this._autoClearTimer = setTimeout(checkAndClear, 2000); + if (!this._isStopped) { + this._autoClearTimer = setTimeout(checkAndClear, 2000); + } } }; // Start checking after a short delay - this._autoClearTimer = setTimeout(checkAndClear, 1000); + if (!this._isStopped) { + this._autoClearTimer = setTimeout(checkAndClear, 1000); + } } } @@ -1384,6 +1398,9 @@ export class Session extends EventEmitter { * ``` */ async stop(killScreen: boolean = true): Promise { + // Set stopped flag first to prevent new timers from being created + this._isStopped = true; + // Clear activity timeout to prevent memory leak if (this.activityTimeout) { clearTimeout(this.activityTimeout); diff --git a/src/task-tracker.ts b/src/task-tracker.ts index cecc9b86..ec356b8c 100644 --- a/src/task-tracker.ts +++ b/src/task-tracker.ts @@ -30,6 +30,19 @@ import { EventEmitter } from 'node:events'; */ const MAX_COMPLETED_TASKS = 100; +/** + * Maximum age for pending tool uses (in milliseconds). + * Entries older than this are cleaned up to prevent unbounded growth. + * Default: 1 hour + */ +const PENDING_TOOL_USE_MAX_AGE_MS = 60 * 60 * 1000; + +/** + * Maximum number of pending tool uses to allow. + * Prevents unbounded growth if tool_results never arrive. + */ +const MAX_PENDING_TOOL_USES = 100; + // ========== Pre-compiled Regex Patterns ========== /** @@ -155,8 +168,8 @@ export class TaskTracker extends EventEmitter { /** Stack of active task IDs for tracking nesting depth */ private taskStack: string[] = []; - /** Pending tool_use blocks waiting for results */ - private pendingToolUses: Map = new Map(); + /** Pending tool_use blocks waiting for results (with timestamp for cleanup) */ + private pendingToolUses: Map = new Map(); /** * Creates a new TaskTracker instance. @@ -254,7 +267,10 @@ export class TaskTracker extends EventEmitter { const parentId = this.taskStack.length > 0 ? this.taskStack[this.taskStack.length - 1] : null; // Store pending tool use - task starts when we see activity - this.pendingToolUses.set(toolUseId, { description, subagentType, parentId }); + this.pendingToolUses.set(toolUseId, { description, subagentType, parentId, createdAt: Date.now() }); + + // Clean up old pending entries to prevent unbounded growth + this.cleanupOldPendingToolUses(); // Create the task immediately const task: BackgroundTask = { @@ -320,6 +336,38 @@ export class TaskTracker extends EventEmitter { this.pendingToolUses.delete(toolUseId); } + /** + * Remove old pending tool uses that never received results. + * Prevents unbounded growth if tool_results never arrive. + */ + private cleanupOldPendingToolUses(): void { + const now = Date.now(); + const toDelete: string[] = []; + + // Remove entries older than PENDING_TOOL_USE_MAX_AGE_MS + for (const [id, entry] of this.pendingToolUses) { + if (now - entry.createdAt > PENDING_TOOL_USE_MAX_AGE_MS) { + toDelete.push(id); + } + } + + // If still over limit after age-based cleanup, remove oldest entries + if (this.pendingToolUses.size - toDelete.length > MAX_PENDING_TOOL_USES) { + const entries = Array.from(this.pendingToolUses.entries()) + .filter(([id]) => !toDelete.includes(id)) + .sort((a, b) => a[1].createdAt - b[1].createdAt); + + const removeCount = entries.length - MAX_PENDING_TOOL_USES; + for (let i = 0; i < removeCount; i++) { + toDelete.push(entries[i][0]); + } + } + + for (const id of toDelete) { + this.pendingToolUses.delete(id); + } + } + /** * Remove old completed/failed tasks when exceeding the limit * Keeps running tasks and the most recent completed tasks diff --git a/src/web/server.ts b/src/web/server.ts index dde671ca..f48c48a8 100644 --- a/src/web/server.ts +++ b/src/web/server.ts @@ -73,6 +73,8 @@ const SCHEDULED_CLEANUP_INTERVAL = 5 * 60 * 1000; const SCHEDULED_RUN_MAX_AGE = 60 * 60 * 1000; // Maximum concurrent sessions to prevent resource exhaustion const MAX_CONCURRENT_SESSIONS = 50; +// SSE client health check interval (every 30 seconds) +const SSE_HEALTH_CHECK_INTERVAL = 30 * 1000; /** * Auto-configure Ralph tracker for a session. @@ -148,6 +150,10 @@ export class WebServer extends EventEmitter { // State update batching (reduce expensive toDetailedState() serialization) private stateUpdatePending: Set = new Set(); private stateUpdateTimer: NodeJS.Timeout | null = null; + // SSE client health check timer + private sseHealthCheckTimer: NodeJS.Timeout | null = null; + // Flag to prevent new timers during shutdown + private _isStopping: boolean = false; constructor(port: number = 3000) { super(); @@ -1683,6 +1689,9 @@ export class WebServer extends EventEmitter { // Batch terminal data for better performance (60fps) // Flushes immediately if batch > 1KB for snappier response to large outputs private batchTerminalData(sessionId: string, data: string): void { + // Skip if server is stopping + if (this._isStopping) return; + const existing = this.terminalBatches.get(sessionId) || ''; const newBatch = existing + data; this.terminalBatches.set(sessionId, newBatch); @@ -1717,6 +1726,9 @@ export class WebServer extends EventEmitter { // Batch session:output events at 50ms for better performance private batchOutputData(sessionId: string, data: string): void { + // Skip if server is stopping + if (this._isStopping) return; + const existing = this.outputBatches.get(sessionId) || ''; this.outputBatches.set(sessionId, existing + data); @@ -1739,6 +1751,9 @@ export class WebServer extends EventEmitter { // Batch task:updated events at 100ms - only send latest update per session private batchTaskUpdate(sessionId: string, task: BackgroundTask): void { + // Skip if server is stopping + if (this._isStopping) return; + this.taskUpdateBatches.set(sessionId, task); if (!this.taskUpdateBatchTimer) { @@ -1762,6 +1777,9 @@ export class WebServer extends EventEmitter { * and only serialize once per STATE_UPDATE_DEBOUNCE_INTERVAL. */ private broadcastSessionStateDebounced(sessionId: string): void { + // Skip if server is stopping + if (this._isStopping) return; + this.stateUpdatePending.add(sessionId); if (!this.stateUpdateTimer) { @@ -1783,6 +1801,36 @@ export class WebServer extends EventEmitter { this.stateUpdatePending.clear(); } + /** + * Clean up dead SSE clients that may not have properly disconnected. + * This prevents memory leaks from abruptly terminated connections. + */ + private cleanupDeadSSEClients(): void { + const deadClients: FastifyReply[] = []; + + for (const client of this.sseClients) { + try { + // Check if the underlying socket is still writable + const socket = client.raw.socket || (client.raw as any).connection; + if (!socket || socket.destroyed || !socket.writable) { + deadClients.push(client); + } + } catch { + // Error accessing socket means client is dead + deadClients.push(client); + } + } + + // Remove dead clients + for (const client of deadClients) { + this.sseClients.delete(client); + } + + if (deadClients.length > 0) { + console.log(`[Server] Cleaned up ${deadClients.length} dead SSE client(s)`); + } + } + async start(): Promise { await this.setupRoutes(); await this.app.listen({ port: this.port, host: '0.0.0.0' }); @@ -1793,6 +1841,11 @@ export class WebServer extends EventEmitter { this.cleanupScheduledRuns(); }, SCHEDULED_CLEANUP_INTERVAL); + // Start SSE client health check timer (prevents memory leaks from dead connections) + this.sseHealthCheckTimer = setInterval(() => { + this.cleanupDeadSSEClients(); + }, SSE_HEALTH_CHECK_INTERVAL); + // Restore screen sessions from previous run await this.restoreScreenSessions(); } @@ -1895,6 +1948,18 @@ export class WebServer extends EventEmitter { } async stop(): Promise { + // Set stopping flag to prevent new timer creation during shutdown + this._isStopping = true; + + // Clear SSE health check timer + if (this.sseHealthCheckTimer) { + clearInterval(this.sseHealthCheckTimer); + this.sseHealthCheckTimer = null; + } + + // Clear all SSE clients + this.sseClients.clear(); + // Clear batch timers if (this.terminalBatchTimer) { clearTimeout(this.terminalBatchTimer);