mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-10-08 08:29:42 +02:00
fix: memory leaks and performance optimizations
- WebServer: Store session listener refs for explicit cleanup on delete - RespawnController: Fix clearInterval bug (was using clearTimeout) - RespawnController: Add try/catch to interval callbacks - PlanOrchestrator: Guarantee progressInterval cleanup with try/finally - SubagentWatcher: Guard idle timer against deleted agents - Session: Add threshold validation for auto-clear/compact (1k-500k) - Session: Add token sum overflow check in restoreTokens - Session: Use LRUMap for _recentTaskDescriptions auto-eviction Commit from flight QR118 in an Airbus A350-1000, stable connection thanks to Starlink! Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
+24
-18
@@ -394,17 +394,17 @@ export class PlanOrchestrator {
|
|||||||
|
|
||||||
this.runningSessions.add(session);
|
this.runningSessions.add(session);
|
||||||
|
|
||||||
|
const prompt = RESEARCH_AGENT_PROMPT.replace('{TASK}', taskDescription);
|
||||||
|
|
||||||
|
// Start progress interval before try block to ensure cleanup in finally
|
||||||
|
const progressInterval = setInterval(() => {
|
||||||
|
const elapsed = Math.floor((Date.now() - startTime) / 1000);
|
||||||
|
onSubagent?.({ type: 'progress', agentId, agentType: 'research', model: MODEL, status: 'running', detail: `${elapsed}s elapsed` });
|
||||||
|
}, 30000);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
const prompt = RESEARCH_AGENT_PROMPT.replace('{TASK}', taskDescription);
|
|
||||||
|
|
||||||
const progressInterval = setInterval(() => {
|
|
||||||
const elapsed = Math.floor((Date.now() - startTime) / 1000);
|
|
||||||
onSubagent?.({ type: 'progress', agentId, agentType: 'research', model: MODEL, status: 'running', detail: `${elapsed}s elapsed` });
|
|
||||||
}, 30000);
|
|
||||||
|
|
||||||
const { result: response } = await session.runPrompt(prompt, { model: MODEL });
|
const { result: response } = await session.runPrompt(prompt, { model: MODEL });
|
||||||
|
|
||||||
clearInterval(progressInterval);
|
|
||||||
this.runningSessions.delete(session);
|
this.runningSessions.delete(session);
|
||||||
|
|
||||||
const durationMs = Date.now() - startTime;
|
const durationMs = Date.now() - startTime;
|
||||||
@@ -446,6 +446,9 @@ export class PlanOrchestrator {
|
|||||||
const error = err instanceof Error ? err.message : String(err);
|
const error = err instanceof Error ? err.message : String(err);
|
||||||
onSubagent?.({ type: 'failed', agentId, agentType: 'research', model: MODEL, status: 'failed', error, durationMs });
|
onSubagent?.({ type: 'failed', agentId, agentType: 'research', model: MODEL, status: 'failed', error, durationMs });
|
||||||
return { success: false, findings: { externalResources: [], codebasePatterns: [], technicalRecommendations: [], potentialChallenges: [], recommendedTools: [] }, enrichedTaskDescription: taskDescription, error, durationMs };
|
return { success: false, findings: { externalResources: [], codebasePatterns: [], technicalRecommendations: [], potentialChallenges: [], recommendedTools: [] }, enrichedTaskDescription: taskDescription, error, durationMs };
|
||||||
|
} finally {
|
||||||
|
// Always clear the progress interval to prevent memory leaks
|
||||||
|
clearInterval(progressInterval);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -473,19 +476,19 @@ export class PlanOrchestrator {
|
|||||||
|
|
||||||
this.runningSessions.add(session);
|
this.runningSessions.add(session);
|
||||||
|
|
||||||
|
const prompt = PLANNER_PROMPT
|
||||||
|
.replace('{TASK}', taskDescription)
|
||||||
|
.replace('{RESEARCH_CONTEXT}', researchContext || '');
|
||||||
|
|
||||||
|
// Start progress interval before try block to ensure cleanup in finally
|
||||||
|
const progressInterval = setInterval(() => {
|
||||||
|
const elapsed = Math.floor((Date.now() - startTime) / 1000);
|
||||||
|
onSubagent?.({ type: 'progress', agentId, agentType: 'planner', model: MODEL, status: 'running', detail: `${elapsed}s elapsed` });
|
||||||
|
}, 30000);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
const prompt = PLANNER_PROMPT
|
|
||||||
.replace('{TASK}', taskDescription)
|
|
||||||
.replace('{RESEARCH_CONTEXT}', researchContext || '');
|
|
||||||
|
|
||||||
const progressInterval = setInterval(() => {
|
|
||||||
const elapsed = Math.floor((Date.now() - startTime) / 1000);
|
|
||||||
onSubagent?.({ type: 'progress', agentId, agentType: 'planner', model: MODEL, status: 'running', detail: `${elapsed}s elapsed` });
|
|
||||||
}, 30000);
|
|
||||||
|
|
||||||
const { result: response } = await session.runPrompt(prompt, { model: MODEL });
|
const { result: response } = await session.runPrompt(prompt, { model: MODEL });
|
||||||
|
|
||||||
clearInterval(progressInterval);
|
|
||||||
this.runningSessions.delete(session);
|
this.runningSessions.delete(session);
|
||||||
|
|
||||||
const durationMs = Date.now() - startTime;
|
const durationMs = Date.now() - startTime;
|
||||||
@@ -520,6 +523,9 @@ export class PlanOrchestrator {
|
|||||||
const error = err instanceof Error ? err.message : String(err);
|
const error = err instanceof Error ? err.message : String(err);
|
||||||
onSubagent?.({ type: 'failed', agentId, agentType: 'planner', model: MODEL, status: 'failed', error, durationMs });
|
onSubagent?.({ type: 'failed', agentId, agentType: 'planner', model: MODEL, status: 'failed', error, durationMs });
|
||||||
return { success: false, error };
|
return { success: false, error };
|
||||||
|
} finally {
|
||||||
|
// Always clear the progress interval to prevent memory leaks
|
||||||
|
clearInterval(progressInterval);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1160,8 +1160,12 @@ export class RespawnController extends EventEmitter {
|
|||||||
private startDetectionUpdates(): void {
|
private startDetectionUpdates(): void {
|
||||||
this.stopDetectionUpdates();
|
this.stopDetectionUpdates();
|
||||||
this.detectionUpdateTimer = setInterval(() => {
|
this.detectionUpdateTimer = setInterval(() => {
|
||||||
if (this._state !== 'stopped') {
|
try {
|
||||||
this.emit('detectionUpdate', this.getDetectionStatus());
|
if (this._state !== 'stopped') {
|
||||||
|
this.emit('detectionUpdate', this.getDetectionStatus());
|
||||||
|
}
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`[RespawnController] Error in detectionUpdateTimer:`, err);
|
||||||
}
|
}
|
||||||
}, 500);
|
}, 500);
|
||||||
}
|
}
|
||||||
@@ -1724,7 +1728,7 @@ export class RespawnController extends EventEmitter {
|
|||||||
this.hookConfirmTimer = null;
|
this.hookConfirmTimer = null;
|
||||||
}
|
}
|
||||||
if (this.stuckStateTimer) {
|
if (this.stuckStateTimer) {
|
||||||
clearTimeout(this.stuckStateTimer);
|
clearInterval(this.stuckStateTimer);
|
||||||
this.stuckStateTimer = null;
|
this.stuckStateTimer = null;
|
||||||
}
|
}
|
||||||
// Clear all tracked timers
|
// Clear all tracked timers
|
||||||
@@ -1743,7 +1747,7 @@ export class RespawnController extends EventEmitter {
|
|||||||
|
|
||||||
// Clear existing timer
|
// Clear existing timer
|
||||||
if (this.stuckStateTimer) {
|
if (this.stuckStateTimer) {
|
||||||
clearTimeout(this.stuckStateTimer);
|
clearInterval(this.stuckStateTimer);
|
||||||
this.stuckStateTimer = null;
|
this.stuckStateTimer = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1751,7 +1755,11 @@ export class RespawnController extends EventEmitter {
|
|||||||
const checkIntervalMs = Math.min(this.config.stuckStateWarningMs, 60000); // Check every minute max
|
const checkIntervalMs = Math.min(this.config.stuckStateWarningMs, 60000); // Check every minute max
|
||||||
|
|
||||||
this.stuckStateTimer = setInterval(() => {
|
this.stuckStateTimer = setInterval(() => {
|
||||||
this.checkStuckState();
|
try {
|
||||||
|
this.checkStuckState();
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`[RespawnController] Error in stuckStateTimer:`, err);
|
||||||
|
}
|
||||||
}, checkIntervalMs);
|
}, checkIntervalMs);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+45
-25
@@ -27,6 +27,7 @@ import { RalphTracker } from './ralph-tracker.js';
|
|||||||
import { BashToolParser } from './bash-tool-parser.js';
|
import { BashToolParser } from './bash-tool-parser.js';
|
||||||
import { ScreenManager } from './screen-manager.js';
|
import { ScreenManager } from './screen-manager.js';
|
||||||
import { BufferAccumulator } from './utils/buffer-accumulator.js';
|
import { BufferAccumulator } from './utils/buffer-accumulator.js';
|
||||||
|
import { LRUMap } from './utils/lru-map.js';
|
||||||
import {
|
import {
|
||||||
ANSI_ESCAPE_PATTERN_FULL,
|
ANSI_ESCAPE_PATTERN_FULL,
|
||||||
TOKEN_PATTERN,
|
TOKEN_PATTERN,
|
||||||
@@ -299,6 +300,10 @@ export class Session extends EventEmitter {
|
|||||||
readonly createdAt: number;
|
readonly createdAt: number;
|
||||||
readonly mode: SessionMode;
|
readonly mode: SessionMode;
|
||||||
|
|
||||||
|
/** Maximum number of task descriptions to keep (LRUMap handles size limit automatically) */
|
||||||
|
private static readonly MAX_TASK_DESCRIPTIONS = 100;
|
||||||
|
private static readonly TASK_DESCRIPTION_MAX_AGE_MS = 30000; // Keep descriptions for 30 seconds
|
||||||
|
|
||||||
private _name: string;
|
private _name: string;
|
||||||
private ptyProcess: pty.IPty | null = null;
|
private ptyProcess: pty.IPty | null = null;
|
||||||
private _pid: number | null = null;
|
private _pid: number | null = null;
|
||||||
@@ -403,8 +408,8 @@ export class Session extends EventEmitter {
|
|||||||
|
|
||||||
// Task descriptions parsed from terminal output (e.g., "Explore(Description)")
|
// Task descriptions parsed from terminal output (e.g., "Explore(Description)")
|
||||||
// Used to correlate with SubagentWatcher discoveries for better window titles
|
// Used to correlate with SubagentWatcher discoveries for better window titles
|
||||||
private _recentTaskDescriptions: Map<number, string> = new Map(); // timestamp -> description
|
// Uses LRUMap for automatic eviction at MAX_TASK_DESCRIPTIONS limit
|
||||||
private static readonly TASK_DESCRIPTION_MAX_AGE_MS = 30000; // Keep descriptions for 30 seconds
|
private _recentTaskDescriptions: LRUMap<number, string> = new LRUMap({ maxSize: Session.MAX_TASK_DESCRIPTIONS });
|
||||||
|
|
||||||
constructor(config: Partial<SessionConfig> & {
|
constructor(config: Partial<SessionConfig> & {
|
||||||
workingDir: string;
|
workingDir: string;
|
||||||
@@ -636,11 +641,16 @@ export class Session extends EventEmitter {
|
|||||||
* Called when recovering sessions after server restart.
|
* Called when recovering sessions after server restart.
|
||||||
*/
|
*/
|
||||||
restoreTokens(inputTokens: number, outputTokens: number, totalCost: number): void {
|
restoreTokens(inputTokens: number, outputTokens: number, totalCost: number): void {
|
||||||
// Sanity check: reject absurdly large values
|
// Sanity check: reject absurdly large individual values
|
||||||
if (inputTokens > MAX_SESSION_TOKENS || outputTokens > MAX_SESSION_TOKENS) {
|
if (inputTokens > MAX_SESSION_TOKENS || outputTokens > MAX_SESSION_TOKENS) {
|
||||||
console.warn(`[Session ${this.id}] Rejected absurd restored tokens: input=${inputTokens}, output=${outputTokens}`);
|
console.warn(`[Session ${this.id}] Rejected absurd restored tokens: input=${inputTokens}, output=${outputTokens}`);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
// Check token sum doesn't overflow MAX_SESSION_TOKENS
|
||||||
|
if (inputTokens + outputTokens > MAX_SESSION_TOKENS) {
|
||||||
|
console.warn(`[Session ${this.id}] Rejected token sum overflow: input=${inputTokens} + output=${outputTokens} = ${inputTokens + outputTokens} > ${MAX_SESSION_TOKENS}`);
|
||||||
|
return;
|
||||||
|
}
|
||||||
// Reject negative values
|
// Reject negative values
|
||||||
if (inputTokens < 0 || outputTokens < 0 || totalCost < 0) {
|
if (inputTokens < 0 || outputTokens < 0 || totalCost < 0) {
|
||||||
console.warn(`[Session ${this.id}] Rejected negative restored tokens: input=${inputTokens}, output=${outputTokens}, cost=${totalCost}`);
|
console.warn(`[Session ${this.id}] Rejected negative restored tokens: input=${inputTokens}, output=${outputTokens}, cost=${totalCost}`);
|
||||||
@@ -668,10 +678,25 @@ export class Session extends EventEmitter {
|
|||||||
this._name = value;
|
this._name = value;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Minimum valid threshold for auto-clear/compact (1000 tokens) */
|
||||||
|
private static readonly MIN_AUTO_THRESHOLD = 1000;
|
||||||
|
/** Maximum valid threshold for auto-clear/compact (500k tokens) */
|
||||||
|
private static readonly MAX_AUTO_THRESHOLD = 500_000;
|
||||||
|
/** Default auto-clear threshold when invalid value provided */
|
||||||
|
private static readonly DEFAULT_AUTO_CLEAR_THRESHOLD = 140_000;
|
||||||
|
/** Default auto-compact threshold when invalid value provided */
|
||||||
|
private static readonly DEFAULT_AUTO_COMPACT_THRESHOLD = 110_000;
|
||||||
|
|
||||||
setAutoClear(enabled: boolean, threshold?: number): void {
|
setAutoClear(enabled: boolean, threshold?: number): void {
|
||||||
this._autoClearEnabled = enabled;
|
this._autoClearEnabled = enabled;
|
||||||
if (threshold !== undefined) {
|
if (threshold !== undefined) {
|
||||||
this._autoClearThreshold = threshold;
|
// Validate threshold bounds
|
||||||
|
if (threshold < Session.MIN_AUTO_THRESHOLD || threshold > Session.MAX_AUTO_THRESHOLD) {
|
||||||
|
console.warn(`[Session ${this.id}] Invalid autoClear threshold ${threshold}, must be between ${Session.MIN_AUTO_THRESHOLD} and ${Session.MAX_AUTO_THRESHOLD}. Using default ${Session.DEFAULT_AUTO_CLEAR_THRESHOLD}.`);
|
||||||
|
this._autoClearThreshold = Session.DEFAULT_AUTO_CLEAR_THRESHOLD;
|
||||||
|
} else {
|
||||||
|
this._autoClearThreshold = threshold;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -690,7 +715,13 @@ export class Session extends EventEmitter {
|
|||||||
setAutoCompact(enabled: boolean, threshold?: number, prompt?: string): void {
|
setAutoCompact(enabled: boolean, threshold?: number, prompt?: string): void {
|
||||||
this._autoCompactEnabled = enabled;
|
this._autoCompactEnabled = enabled;
|
||||||
if (threshold !== undefined) {
|
if (threshold !== undefined) {
|
||||||
this._autoCompactThreshold = threshold;
|
// Validate threshold bounds
|
||||||
|
if (threshold < Session.MIN_AUTO_THRESHOLD || threshold > Session.MAX_AUTO_THRESHOLD) {
|
||||||
|
console.warn(`[Session ${this.id}] Invalid autoCompact threshold ${threshold}, must be between ${Session.MIN_AUTO_THRESHOLD} and ${Session.MAX_AUTO_THRESHOLD}. Using default ${Session.DEFAULT_AUTO_COMPACT_THRESHOLD}.`);
|
||||||
|
this._autoCompactThreshold = Session.DEFAULT_AUTO_COMPACT_THRESHOLD;
|
||||||
|
} else {
|
||||||
|
this._autoCompactThreshold = threshold;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if (prompt !== undefined) {
|
if (prompt !== undefined) {
|
||||||
this._autoCompactPrompt = prompt;
|
this._autoCompactPrompt = prompt;
|
||||||
@@ -1515,32 +1546,21 @@ export class Session extends EventEmitter {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Maximum number of task descriptions to keep */
|
|
||||||
private static readonly MAX_TASK_DESCRIPTIONS = 100;
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Remove task descriptions older than TASK_DESCRIPTION_MAX_AGE_MS.
|
* Remove task descriptions older than TASK_DESCRIPTION_MAX_AGE_MS.
|
||||||
* Also enforces MAX_TASK_DESCRIPTIONS size limit.
|
* Size limit is handled automatically by LRUMap eviction on set().
|
||||||
*/
|
*/
|
||||||
private cleanupOldTaskDescriptions(): void {
|
private cleanupOldTaskDescriptions(): void {
|
||||||
const cutoff = Date.now() - Session.TASK_DESCRIPTION_MAX_AGE_MS;
|
const cutoff = Date.now() - Session.TASK_DESCRIPTION_MAX_AGE_MS;
|
||||||
// Collect keys to delete first, then delete (avoids modifying Map during iteration)
|
// Keys are timestamps - iterate and delete expired entries
|
||||||
const keysToDelete: number[] = [];
|
// LRUMap maintains insertion order, so we can break early once we find a non-expired entry
|
||||||
for (const [timestamp] of this._recentTaskDescriptions) {
|
for (const timestamp of this._recentTaskDescriptions.keysInOrder()) {
|
||||||
if (timestamp < cutoff) {
|
if (timestamp < cutoff) {
|
||||||
keysToDelete.push(timestamp);
|
this._recentTaskDescriptions.delete(timestamp);
|
||||||
}
|
} else {
|
||||||
}
|
// Keys are ordered by insertion time (which is the timestamp)
|
||||||
for (const key of keysToDelete) {
|
// Once we find a non-expired one, all subsequent are also non-expired
|
||||||
this._recentTaskDescriptions.delete(key);
|
break;
|
||||||
}
|
|
||||||
|
|
||||||
// Enforce size limit by removing oldest entries
|
|
||||||
if (this._recentTaskDescriptions.size > Session.MAX_TASK_DESCRIPTIONS) {
|
|
||||||
const sortedKeys = Array.from(this._recentTaskDescriptions.keys()).sort((a, b) => a - b);
|
|
||||||
const keysToRemove = sortedKeys.slice(0, this._recentTaskDescriptions.size - Session.MAX_TASK_DESCRIPTIONS);
|
|
||||||
for (const key of keysToRemove) {
|
|
||||||
this._recentTaskDescriptions.delete(key);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1327,8 +1327,14 @@ export class SubagentWatcher extends EventEmitter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
const timer = setTimeout(() => {
|
const timer = setTimeout(() => {
|
||||||
|
// Guard against race condition: agent may have been deleted before timer fires
|
||||||
const info = this.agentInfo.get(agentId);
|
const info = this.agentInfo.get(agentId);
|
||||||
if (info && info.status === 'active') {
|
if (!info) {
|
||||||
|
// Agent was deleted - clean up timer reference
|
||||||
|
this.idleTimers.delete(agentId);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (info.status === 'active') {
|
||||||
info.status = 'idle';
|
info.status = 'idle';
|
||||||
}
|
}
|
||||||
}, IDLE_TIMEOUT_MS);
|
}, IDLE_TIMEOUT_MS);
|
||||||
|
|||||||
+271
-176
@@ -268,6 +268,35 @@ function getOrCreateSelfSignedCert(): { key: string; cert: string } {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Stored listener references for session cleanup (prevents memory leaks) */
|
||||||
|
interface SessionListenerRefs {
|
||||||
|
output: (data: string) => void;
|
||||||
|
terminal: (data: string) => void;
|
||||||
|
clearTerminal: () => void;
|
||||||
|
message: (msg: ClaudeMessage) => void;
|
||||||
|
error: (error: string) => void;
|
||||||
|
completion: (result: string, cost: number) => void;
|
||||||
|
exit: (code: number | null) => void;
|
||||||
|
working: () => void;
|
||||||
|
idle: () => void;
|
||||||
|
taskCreated: (task: BackgroundTask) => void;
|
||||||
|
taskUpdated: (task: BackgroundTask) => void;
|
||||||
|
taskCompleted: (task: BackgroundTask) => void;
|
||||||
|
taskFailed: (task: BackgroundTask, error: string) => void;
|
||||||
|
autoClear: (data: { tokens: number; threshold: number }) => void;
|
||||||
|
autoCompact: (data: { tokens: number; threshold: number; prompt?: string }) => void;
|
||||||
|
cliInfoUpdated: (data: { version?: string; model?: string; accountType?: string; latestVersion?: string }) => void;
|
||||||
|
ralphLoopUpdate: (state: RalphTrackerState) => void;
|
||||||
|
ralphTodoUpdate: (todos: RalphTodoItem[]) => void;
|
||||||
|
ralphCompletionDetected: (phrase: string) => void;
|
||||||
|
ralphStatusBlockDetected: (block: import('../types.js').RalphStatusBlock) => void;
|
||||||
|
ralphCircuitBreakerUpdate: (status: import('../types.js').CircuitBreakerStatus) => void;
|
||||||
|
ralphExitGateMet: (data: { completionIndicators: number; exitSignal: boolean }) => void;
|
||||||
|
bashToolStart: (tool: ActiveBashTool) => void;
|
||||||
|
bashToolEnd: (tool: ActiveBashTool) => void;
|
||||||
|
bashToolsUpdate: (tools: ActiveBashTool[]) => void;
|
||||||
|
}
|
||||||
|
|
||||||
export class WebServer extends EventEmitter {
|
export class WebServer extends EventEmitter {
|
||||||
private app: FastifyInstance;
|
private app: FastifyInstance;
|
||||||
private sessions: Map<string, Session> = new Map();
|
private sessions: Map<string, Session> = new Map();
|
||||||
@@ -275,6 +304,8 @@ export class WebServer extends EventEmitter {
|
|||||||
private respawnTimers: Map<string, { timer: NodeJS.Timeout; endAt: number; startedAt: number }> = new Map();
|
private respawnTimers: Map<string, { timer: NodeJS.Timeout; endAt: number; startedAt: number }> = new Map();
|
||||||
private runSummaryTrackers: Map<string, RunSummaryTracker> = new Map();
|
private runSummaryTrackers: Map<string, RunSummaryTracker> = new Map();
|
||||||
private transcriptWatchers: Map<string, TranscriptWatcher> = new Map();
|
private transcriptWatchers: Map<string, TranscriptWatcher> = new Map();
|
||||||
|
// Store session listener references for explicit cleanup (prevents memory leaks)
|
||||||
|
private sessionListenerRefs: Map<string, SessionListenerRefs> = new Map();
|
||||||
private scheduledRuns: Map<string, ScheduledRun> = new Map();
|
private scheduledRuns: Map<string, ScheduledRun> = new Map();
|
||||||
private sseClients: Set<FastifyReply> = new Set();
|
private sseClients: Set<FastifyReply> = new Set();
|
||||||
private store = getStore();
|
private store = getStore();
|
||||||
@@ -3509,6 +3540,37 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
|
|||||||
console.log(`[Server] Added to global stats: ${session.inputTokens + session.outputTokens} tokens, $${session.totalCost.toFixed(4)} from session ${sessionId}`);
|
console.log(`[Server] Added to global stats: ${session.inputTokens + session.outputTokens} tokens, $${session.totalCost.toFixed(4)} from session ${sessionId}`);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Explicitly remove stored listeners to break closure references (prevents memory leak)
|
||||||
|
const listeners = this.sessionListenerRefs.get(sessionId);
|
||||||
|
if (listeners) {
|
||||||
|
session.off('output', listeners.output);
|
||||||
|
session.off('terminal', listeners.terminal);
|
||||||
|
session.off('clearTerminal', listeners.clearTerminal);
|
||||||
|
session.off('message', listeners.message);
|
||||||
|
session.off('error', listeners.error);
|
||||||
|
session.off('completion', listeners.completion);
|
||||||
|
session.off('exit', listeners.exit);
|
||||||
|
session.off('working', listeners.working);
|
||||||
|
session.off('idle', listeners.idle);
|
||||||
|
session.off('taskCreated', listeners.taskCreated);
|
||||||
|
session.off('taskUpdated', listeners.taskUpdated);
|
||||||
|
session.off('taskCompleted', listeners.taskCompleted);
|
||||||
|
session.off('taskFailed', listeners.taskFailed);
|
||||||
|
session.off('autoClear', listeners.autoClear);
|
||||||
|
session.off('autoCompact', listeners.autoCompact);
|
||||||
|
session.off('cliInfoUpdated', listeners.cliInfoUpdated);
|
||||||
|
session.off('ralphLoopUpdate', listeners.ralphLoopUpdate);
|
||||||
|
session.off('ralphTodoUpdate', listeners.ralphTodoUpdate);
|
||||||
|
session.off('ralphCompletionDetected', listeners.ralphCompletionDetected);
|
||||||
|
session.off('ralphStatusBlockDetected', listeners.ralphStatusBlockDetected);
|
||||||
|
session.off('ralphCircuitBreakerUpdate', listeners.ralphCircuitBreakerUpdate);
|
||||||
|
session.off('ralphExitGateMet', listeners.ralphExitGateMet);
|
||||||
|
session.off('bashToolStart', listeners.bashToolStart);
|
||||||
|
session.off('bashToolEnd', listeners.bashToolEnd);
|
||||||
|
session.off('bashToolsUpdate', listeners.bashToolsUpdate);
|
||||||
|
this.sessionListenerRefs.delete(sessionId);
|
||||||
|
}
|
||||||
|
|
||||||
session.removeAllListeners();
|
session.removeAllListeners();
|
||||||
// Close any active file streams for this session
|
// Close any active file streams for this session
|
||||||
fileStreamManager.closeSessionStreams(sessionId);
|
fileStreamManager.closeSessionStreams(sessionId);
|
||||||
@@ -3540,211 +3602,244 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
|
|||||||
imageWatcher.watchSession(session.id, session.workingDir);
|
imageWatcher.watchSession(session.id, session.workingDir);
|
||||||
}
|
}
|
||||||
|
|
||||||
session.on('output', (data) => {
|
// Store all listener references for explicit cleanup on session delete
|
||||||
// Use batching for better performance at high throughput
|
// This prevents memory leaks from closure references keeping objects alive
|
||||||
this.batchOutputData(session.id, data);
|
const listeners: SessionListenerRefs = {
|
||||||
});
|
output: (data) => {
|
||||||
|
// Use batching for better performance at high throughput
|
||||||
|
this.batchOutputData(session.id, data);
|
||||||
|
},
|
||||||
|
|
||||||
session.on('terminal', (data) => {
|
terminal: (data) => {
|
||||||
// Use batching for better performance at high throughput
|
// Use batching for better performance at high throughput
|
||||||
this.batchTerminalData(session.id, data);
|
this.batchTerminalData(session.id, data);
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('clearTerminal', () => {
|
clearTerminal: () => {
|
||||||
// Tell clients to clear their terminal (after screen attach)
|
// Tell clients to clear their terminal (after screen attach)
|
||||||
this.broadcast('session:clearTerminal', { id: session.id });
|
this.broadcast('session:clearTerminal', { id: session.id });
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('message', (msg: ClaudeMessage) => {
|
message: (msg: ClaudeMessage) => {
|
||||||
this.broadcast('session:message', { id: session.id, message: msg });
|
this.broadcast('session:message', { id: session.id, message: msg });
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('error', (error) => {
|
error: (error) => {
|
||||||
this.broadcast('session:error', { id: session.id, error });
|
this.broadcast('session:error', { id: session.id, error });
|
||||||
// Track in run summary
|
// Track in run summary
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
if (tracker) tracker.recordError('Session error', String(error));
|
if (tracker) tracker.recordError('Session error', String(error));
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('completion', (result, cost) => {
|
completion: (result, cost) => {
|
||||||
this.broadcast('session:completion', { id: session.id, result, cost });
|
this.broadcast('session:completion', { id: session.id, result, cost });
|
||||||
this.broadcast('session:updated', this.getSessionStateWithRespawn(session));
|
|
||||||
this.persistSessionState(session);
|
|
||||||
// Track tokens in run summary (completion event has updated token values)
|
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
|
||||||
if (tracker) tracker.recordTokens(session.inputTokens, session.outputTokens);
|
|
||||||
});
|
|
||||||
|
|
||||||
session.on('exit', (code) => {
|
|
||||||
// Wrap in try/catch to ensure cleanup always happens
|
|
||||||
try {
|
|
||||||
this.broadcast('session:exit', { id: session.id, code });
|
|
||||||
this.broadcast('session:updated', this.getSessionStateWithRespawn(session));
|
this.broadcast('session:updated', this.getSessionStateWithRespawn(session));
|
||||||
this.persistSessionState(session);
|
this.persistSessionState(session);
|
||||||
} catch (err) {
|
// Track tokens in run summary (completion event has updated token values)
|
||||||
console.error(`[Server] Error broadcasting session exit for ${session.id}:`, err);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
}
|
if (tracker) tracker.recordTokens(session.inputTokens, session.outputTokens);
|
||||||
|
},
|
||||||
|
|
||||||
// Always clean up respawn controller, even if broadcast failed
|
exit: (code) => {
|
||||||
try {
|
// Wrap in try/catch to ensure cleanup always happens
|
||||||
const controller = this.respawnControllers.get(session.id);
|
try {
|
||||||
if (controller) {
|
this.broadcast('session:exit', { id: session.id, code });
|
||||||
controller.stop();
|
this.broadcast('session:updated', this.getSessionStateWithRespawn(session));
|
||||||
controller.removeAllListeners();
|
this.persistSessionState(session);
|
||||||
this.respawnControllers.delete(session.id);
|
} catch (err) {
|
||||||
|
console.error(`[Server] Error broadcasting session exit for ${session.id}:`, err);
|
||||||
}
|
}
|
||||||
// Also clean up the respawn timer to prevent orphaned timers
|
|
||||||
const timerInfo = this.respawnTimers.get(session.id);
|
// Always clean up respawn controller, even if broadcast failed
|
||||||
if (timerInfo) {
|
try {
|
||||||
clearTimeout(timerInfo.timer);
|
const controller = this.respawnControllers.get(session.id);
|
||||||
this.respawnTimers.delete(session.id);
|
if (controller) {
|
||||||
|
controller.stop();
|
||||||
|
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);
|
||||||
}
|
}
|
||||||
} catch (err) {
|
},
|
||||||
console.error(`[Server] Error cleaning up respawn controller for ${session.id}:`, err);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
session.on('working', () => {
|
working: () => {
|
||||||
this.broadcast('session:working', { id: session.id });
|
this.broadcast('session:working', { id: session.id });
|
||||||
// Track in run summary
|
// Track in run summary
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
if (tracker) {
|
if (tracker) {
|
||||||
tracker.recordWorking();
|
tracker.recordWorking();
|
||||||
tracker.recordTokens(session.inputTokens, session.outputTokens);
|
tracker.recordTokens(session.inputTokens, session.outputTokens);
|
||||||
}
|
}
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('idle', () => {
|
idle: () => {
|
||||||
this.broadcast('session:idle', { id: session.id });
|
this.broadcast('session:idle', { id: session.id });
|
||||||
// Use debounced state update (idle can fire frequently)
|
// Use debounced state update (idle can fire frequently)
|
||||||
this.broadcastSessionStateDebounced(session.id);
|
this.broadcastSessionStateDebounced(session.id);
|
||||||
// Track in run summary
|
// Track in run summary
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
if (tracker) {
|
if (tracker) {
|
||||||
tracker.recordIdle();
|
tracker.recordIdle();
|
||||||
tracker.recordTokens(session.inputTokens, session.outputTokens);
|
tracker.recordTokens(session.inputTokens, session.outputTokens);
|
||||||
}
|
}
|
||||||
});
|
},
|
||||||
|
|
||||||
// Background task events - use debounced state updates to reduce serialization overhead
|
// Background task events - use debounced state updates to reduce serialization overhead
|
||||||
session.on('taskCreated', (task: BackgroundTask) => {
|
taskCreated: (task: BackgroundTask) => {
|
||||||
this.broadcast('task:created', { sessionId: session.id, task });
|
this.broadcast('task:created', { sessionId: session.id, task });
|
||||||
this.broadcastSessionStateDebounced(session.id);
|
this.broadcastSessionStateDebounced(session.id);
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('taskUpdated', (task: BackgroundTask) => {
|
taskUpdated: (task: BackgroundTask) => {
|
||||||
// Use batching for better performance at high update rates
|
// Use batching for better performance at high update rates
|
||||||
this.batchTaskUpdate(session.id, task);
|
this.batchTaskUpdate(session.id, task);
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('taskCompleted', (task: BackgroundTask) => {
|
taskCompleted: (task: BackgroundTask) => {
|
||||||
this.broadcast('task:completed', { sessionId: session.id, task });
|
this.broadcast('task:completed', { sessionId: session.id, task });
|
||||||
this.broadcastSessionStateDebounced(session.id);
|
this.broadcastSessionStateDebounced(session.id);
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('taskFailed', (task: BackgroundTask, error: string) => {
|
taskFailed: (task: BackgroundTask, error: string) => {
|
||||||
this.broadcast('task:failed', { sessionId: session.id, task, error });
|
this.broadcast('task:failed', { sessionId: session.id, task, error });
|
||||||
this.broadcastSessionStateDebounced(session.id);
|
this.broadcastSessionStateDebounced(session.id);
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('autoClear', (data: { tokens: number; threshold: number }) => {
|
autoClear: (data: { tokens: number; threshold: number }) => {
|
||||||
this.broadcast('session:autoClear', { sessionId: session.id, ...data });
|
this.broadcast('session:autoClear', { sessionId: session.id, ...data });
|
||||||
this.broadcastSessionStateDebounced(session.id);
|
this.broadcastSessionStateDebounced(session.id);
|
||||||
// Track in run summary
|
// Track in run summary
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
if (tracker) tracker.recordAutoClear(data.tokens, data.threshold);
|
if (tracker) tracker.recordAutoClear(data.tokens, data.threshold);
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('autoCompact', (data: { tokens: number; threshold: number; prompt?: string }) => {
|
autoCompact: (data: { tokens: number; threshold: number; prompt?: string }) => {
|
||||||
this.broadcast('session:autoCompact', { sessionId: session.id, ...data });
|
this.broadcast('session:autoCompact', { sessionId: session.id, ...data });
|
||||||
this.broadcastSessionStateDebounced(session.id);
|
this.broadcastSessionStateDebounced(session.id);
|
||||||
// Track in run summary
|
// Track in run summary
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
if (tracker) tracker.recordAutoCompact(data.tokens, data.threshold);
|
if (tracker) tracker.recordAutoCompact(data.tokens, data.threshold);
|
||||||
});
|
},
|
||||||
|
|
||||||
// Claude Code CLI info parsed from terminal (version, model, account)
|
// Claude Code CLI info parsed from terminal (version, model, account)
|
||||||
session.on('cliInfoUpdated', (data: { version?: string; model?: string; accountType?: string; latestVersion?: string }) => {
|
cliInfoUpdated: (data: { version?: string; model?: string; accountType?: string; latestVersion?: string }) => {
|
||||||
this.broadcast('session:cliInfo', { sessionId: session.id, ...data });
|
this.broadcast('session:cliInfo', { sessionId: session.id, ...data });
|
||||||
this.broadcastSessionStateDebounced(session.id);
|
this.broadcastSessionStateDebounced(session.id);
|
||||||
});
|
},
|
||||||
|
|
||||||
// Ralph tracking events
|
// Ralph tracking events
|
||||||
session.on('ralphLoopUpdate', (state: RalphTrackerState) => {
|
ralphLoopUpdate: (state: RalphTrackerState) => {
|
||||||
this.broadcast('session:ralphLoopUpdate', { sessionId: session.id, state });
|
this.broadcast('session:ralphLoopUpdate', { sessionId: session.id, state });
|
||||||
// Persist Ralph state
|
// Persist Ralph state
|
||||||
this.store.updateRalphState(session.id, { loop: state });
|
this.store.updateRalphState(session.id, { loop: state });
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('ralphTodoUpdate', (todos: RalphTodoItem[]) => {
|
ralphTodoUpdate: (todos: RalphTodoItem[]) => {
|
||||||
this.broadcast('session:ralphTodoUpdate', { sessionId: session.id, todos });
|
this.broadcast('session:ralphTodoUpdate', { sessionId: session.id, todos });
|
||||||
// Persist Ralph state
|
// Persist Ralph state
|
||||||
this.store.updateRalphState(session.id, { todos });
|
this.store.updateRalphState(session.id, { todos });
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('ralphCompletionDetected', (phrase: string) => {
|
ralphCompletionDetected: (phrase: string) => {
|
||||||
this.broadcast('session:ralphCompletionDetected', { sessionId: session.id, phrase });
|
this.broadcast('session:ralphCompletionDetected', { sessionId: session.id, phrase });
|
||||||
// Track in run summary
|
// Track in run summary
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
if (tracker) tracker.recordRalphCompletion(phrase);
|
if (tracker) tracker.recordRalphCompletion(phrase);
|
||||||
});
|
},
|
||||||
|
|
||||||
// RALPH_STATUS block events
|
// RALPH_STATUS block events
|
||||||
session.on('ralphStatusBlockDetected', (block: import('../types.js').RalphStatusBlock) => {
|
ralphStatusBlockDetected: (block: import('../types.js').RalphStatusBlock) => {
|
||||||
this.broadcast('session:ralphStatusUpdate', { sessionId: session.id, block });
|
this.broadcast('session:ralphStatusUpdate', { sessionId: session.id, block });
|
||||||
// Track in run summary
|
// Track in run summary
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
if (tracker) {
|
if (tracker) {
|
||||||
tracker.addEvent(
|
tracker.addEvent(
|
||||||
block.status === 'BLOCKED' ? 'warning' : 'idle_detected',
|
block.status === 'BLOCKED' ? 'warning' : 'idle_detected',
|
||||||
block.status === 'BLOCKED' ? 'warning' : 'info',
|
block.status === 'BLOCKED' ? 'warning' : 'info',
|
||||||
`Ralph Status: ${block.status}`,
|
`Ralph Status: ${block.status}`,
|
||||||
`Tasks: ${block.tasksCompletedThisLoop}, Files: ${block.filesModified}, Tests: ${block.testsStatus}`
|
`Tasks: ${block.tasksCompletedThisLoop}, Files: ${block.filesModified}, Tests: ${block.testsStatus}`
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('ralphCircuitBreakerUpdate', (status: import('../types.js').CircuitBreakerStatus) => {
|
ralphCircuitBreakerUpdate: (status: import('../types.js').CircuitBreakerStatus) => {
|
||||||
this.broadcast('session:circuitBreakerUpdate', { sessionId: session.id, status });
|
this.broadcast('session:circuitBreakerUpdate', { sessionId: session.id, status });
|
||||||
// Track state changes in run summary
|
// Track state changes in run summary
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
if (tracker && status.state === 'OPEN') {
|
if (tracker && status.state === 'OPEN') {
|
||||||
tracker.addEvent(
|
tracker.addEvent(
|
||||||
'warning',
|
'warning',
|
||||||
'warning',
|
'warning',
|
||||||
'Circuit Breaker Opened',
|
'Circuit Breaker Opened',
|
||||||
status.reason
|
status.reason
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('ralphExitGateMet', (data: { completionIndicators: number; exitSignal: boolean }) => {
|
ralphExitGateMet: (data: { completionIndicators: number; exitSignal: boolean }) => {
|
||||||
this.broadcast('session:exitGateMet', { sessionId: session.id, ...data });
|
this.broadcast('session:exitGateMet', { sessionId: session.id, ...data });
|
||||||
// Track in run summary
|
// Track in run summary
|
||||||
const tracker = this.runSummaryTrackers.get(session.id);
|
const tracker = this.runSummaryTrackers.get(session.id);
|
||||||
if (tracker) {
|
if (tracker) {
|
||||||
tracker.addEvent(
|
tracker.addEvent(
|
||||||
'ralph_completion',
|
'ralph_completion',
|
||||||
'success',
|
'success',
|
||||||
'Exit Gate Met',
|
'Exit Gate Met',
|
||||||
`Indicators: ${data.completionIndicators}, EXIT_SIGNAL: ${data.exitSignal}`
|
`Indicators: ${data.completionIndicators}, EXIT_SIGNAL: ${data.exitSignal}`
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
});
|
},
|
||||||
|
|
||||||
// Bash tool tracking events (for clickable file paths)
|
// Bash tool tracking events (for clickable file paths)
|
||||||
session.on('bashToolStart', (tool: ActiveBashTool) => {
|
bashToolStart: (tool: ActiveBashTool) => {
|
||||||
this.broadcast('session:bashToolStart', { sessionId: session.id, tool });
|
this.broadcast('session:bashToolStart', { sessionId: session.id, tool });
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('bashToolEnd', (tool: ActiveBashTool) => {
|
bashToolEnd: (tool: ActiveBashTool) => {
|
||||||
this.broadcast('session:bashToolEnd', { sessionId: session.id, tool });
|
this.broadcast('session:bashToolEnd', { sessionId: session.id, tool });
|
||||||
});
|
},
|
||||||
|
|
||||||
session.on('bashToolsUpdate', (tools: ActiveBashTool[]) => {
|
bashToolsUpdate: (tools: ActiveBashTool[]) => {
|
||||||
this.broadcast('session:bashToolsUpdate', { sessionId: session.id, tools });
|
this.broadcast('session:bashToolsUpdate', { sessionId: session.id, tools });
|
||||||
});
|
},
|
||||||
|
};
|
||||||
|
|
||||||
|
// Store listener refs for cleanup
|
||||||
|
this.sessionListenerRefs.set(session.id, listeners);
|
||||||
|
|
||||||
|
// Attach all listeners to the session
|
||||||
|
session.on('output', listeners.output);
|
||||||
|
session.on('terminal', listeners.terminal);
|
||||||
|
session.on('clearTerminal', listeners.clearTerminal);
|
||||||
|
session.on('message', listeners.message);
|
||||||
|
session.on('error', listeners.error);
|
||||||
|
session.on('completion', listeners.completion);
|
||||||
|
session.on('exit', listeners.exit);
|
||||||
|
session.on('working', listeners.working);
|
||||||
|
session.on('idle', listeners.idle);
|
||||||
|
session.on('taskCreated', listeners.taskCreated);
|
||||||
|
session.on('taskUpdated', listeners.taskUpdated);
|
||||||
|
session.on('taskCompleted', listeners.taskCompleted);
|
||||||
|
session.on('taskFailed', listeners.taskFailed);
|
||||||
|
session.on('autoClear', listeners.autoClear);
|
||||||
|
session.on('autoCompact', listeners.autoCompact);
|
||||||
|
session.on('cliInfoUpdated', listeners.cliInfoUpdated);
|
||||||
|
session.on('ralphLoopUpdate', listeners.ralphLoopUpdate);
|
||||||
|
session.on('ralphTodoUpdate', listeners.ralphTodoUpdate);
|
||||||
|
session.on('ralphCompletionDetected', listeners.ralphCompletionDetected);
|
||||||
|
session.on('ralphStatusBlockDetected', listeners.ralphStatusBlockDetected);
|
||||||
|
session.on('ralphCircuitBreakerUpdate', listeners.ralphCircuitBreakerUpdate);
|
||||||
|
session.on('ralphExitGateMet', listeners.ralphExitGateMet);
|
||||||
|
session.on('bashToolStart', listeners.bashToolStart);
|
||||||
|
session.on('bashToolEnd', listeners.bashToolEnd);
|
||||||
|
session.on('bashToolsUpdate', listeners.bashToolsUpdate);
|
||||||
}
|
}
|
||||||
|
|
||||||
private setupRespawnListeners(sessionId: string, controller: RespawnController): void {
|
private setupRespawnListeners(sessionId: string, controller: RespawnController): void {
|
||||||
|
|||||||
Reference in New Issue
Block a user