mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 12:39:42 +02:00
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 <noreply@anthropic.com>
This commit is contained in:
@@ -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 = '';
|
||||
|
||||
+35
-18
@@ -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<void> {
|
||||
// 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);
|
||||
|
||||
+51
-3
@@ -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<string, { description: string; subagentType: string; parentId: string | null }> = new Map();
|
||||
/** Pending tool_use blocks waiting for results (with timestamp for cleanup) */
|
||||
private pendingToolUses: Map<string, { description: string; subagentType: string; parentId: string | null; createdAt: number }> = 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
|
||||
|
||||
@@ -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<string> = 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<void> {
|
||||
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<void> {
|
||||
// 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);
|
||||
|
||||
Reference in New Issue
Block a user