fix: memory leaks, race conditions, and stability improvements

- SSE event listener cleanup: Track all listeners and remove on reconnect
  to prevent listener accumulation (60+ listeners per connect cycle)
- Floating window cleanup: Remove drag listeners, ResizeObservers on close
- Session creation mutex: Prevent race conditions from concurrent requests
- Screen kill verification: Confirm processes are dead after kill signals
- State store circuit breaker: Prevent cascading failures with backup/recovery
- PTY spawn error handling: Wrap all spawn calls with proper error emission
- SSE graceful shutdown: Notify clients before server stops
- Async file writes: Replace blocking writeFileSync in hot paths
- DocumentFragment optimization: Batch DOM updates for plan rendering
- Timer cleanup: Clear notification/countdown timers on SSE reconnect

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
arkon
2026-01-27 20:24:00 +01:00
co-authored by Claude Opus 4.5
parent de45239fb7
commit e27b473e77
7 changed files with 1624 additions and 313 deletions
+372 -14
View File
@@ -20,12 +20,63 @@ import { ScreenManager } from './screen-manager.js';
// Types
// ============================================================================
/** Development phase in TDD cycle */
export type PlanPhase = 'setup' | 'test' | 'impl' | 'verify' | 'review';
/** Task execution status */
export type PlanTaskStatus = 'pending' | 'in_progress' | 'completed' | 'failed' | 'blocked';
/**
* Enhanced plan item with verification, dependencies, and execution tracking.
* Supports TDD workflow, failure tracking, and plan versioning.
*/
export interface PlanItem {
/** Unique identifier (e.g., "P0-001") */
id?: string;
/** Task description */
content: string;
/** Criticality level */
priority: 'P0' | 'P1' | 'P2' | null;
/** Which subagent generated this item */
source?: string;
/** Why this task is needed */
rationale?: string;
/** Legacy numeric phase (1-4) */
phase?: number;
// === NEW: Verification ===
/** How to know it's done (e.g., "npm test passes", "endpoint returns 200") */
verificationCriteria?: string;
/** Command to run for verification (e.g., "npm test -- --grep='auth'") */
testCommand?: string;
// === NEW: Dependencies ===
/** IDs of tasks that must complete first */
dependencies?: string[];
// === NEW: Execution tracking ===
/** Current execution status */
status?: PlanTaskStatus;
/** How many times attempted */
attempts?: number;
/** Most recent failure reason */
lastError?: string;
/** Timestamp of completion */
completedAt?: number;
// === NEW: Metadata ===
/** Estimated complexity */
complexity?: 'low' | 'medium' | 'high';
/** How to undo if needed */
rollbackStrategy?: string;
/** Plan version this belongs to */
version?: number;
/** TDD phase category */
tddPhase?: PlanPhase;
/** ID of paired test/impl task */
pairedWith?: string;
/** Checklist items for review tasks (tddPhase: 'review') */
reviewChecklist?: string[];
}
export interface SubagentResult {
@@ -131,7 +182,7 @@ Return ONLY a JSON array:
Generate 10-20 items. Think about the complete system architecture.`;
const TESTING_SPECIALIST_PROMPT = `You are a TDD Specialist designing a comprehensive test strategy.
const TESTING_SPECIALIST_PROMPT = `You are a TDD Specialist designing a comprehensive test strategy with verification criteria.
## YOUR TASK
Design test coverage for this task:
@@ -145,14 +196,26 @@ Following Test-Driven Development methodology:
2. Plan integration tests for feature interactions
3. Identify edge cases and boundary conditions
4. Consider error scenarios and failure modes
5. Plan verification steps
5. For each test, specify HOW to verify it passes
## OUTPUT FORMAT
Return ONLY a JSON array:
[
{"category": "unit|integration|edge-case|error|verification", "content": "test description", "rationale": "what it validates"}
{
"category": "unit|integration|edge-case|error|verification",
"content": "Write test for user login with valid credentials",
"rationale": "Validates happy path authentication flow",
"verificationCriteria": "Test passes: POST /auth/login returns 200 with JWT token",
"testCommand": "npm test -- --grep='login valid'",
"pairedImpl": "Implement login endpoint handler"
}
]
CRITICAL: Every test item MUST include:
- verificationCriteria: How to know the test passes (observable outcome)
- testCommand: The actual command to run (npm test, pytest, etc.)
- pairedImpl: The implementation step this test validates
Generate 12-25 items. Tests should be written BEFORE implementation.`;
const RISK_ANALYST_PROMPT = `You are a Risk Analyst identifying potential issues and blockers.
@@ -178,6 +241,45 @@ Return ONLY a JSON array:
Generate 8-15 items. Being proactive about risks prevents surprises.`;
// @ts-expect-error Reserved for future use - code review specialist prompt
const CODE_REVIEWER_PROMPT = `You are a Code Review Specialist designing post-implementation review tasks.
## YOUR TASK
Design code review steps for implementations in this task:
## TASK DESCRIPTION
{TASK}
## INSTRUCTIONS
For each implementation identified, create a review task that checks:
1. **Best Practices**: Language-specific conventions and idioms
2. **Security**: OWASP top 10, input validation, authentication
3. **Performance**: Time complexity, memory usage, N+1 queries
4. **Error Handling**: Edge cases covered, meaningful error messages
5. **Code Quality**: DRY, SOLID principles, readability
6. **Type Safety**: Proper typing, no implicit any, null checks
## REVIEW TASK GUIDELINES
- Review tasks run AFTER implementation, BEFORE merge
- Each review should be specific and actionable
- Include what to look for and how to verify
- Reference language-specific linting tools where applicable
## OUTPUT FORMAT
Return ONLY a JSON array:
[
{
"category": "security|performance|quality|error-handling|best-practices|type-safety",
"content": "Review authentication handler for XSS vulnerabilities",
"rationale": "User input flows through auth - must sanitize",
"verificationCriteria": "No unescaped user input, all inputs validated",
"reviewChecklist": ["Check input sanitization", "Verify CSRF tokens", "Review session handling"],
"implToReview": "Implement authentication handler"
}
]
Generate 5-10 review tasks. Code review catches bugs that tests miss.`;
const VERIFICATION_PROMPT = `You are a Plan Verification Expert reviewing an implementation plan for completeness and quality.
## ORIGINAL TASK
@@ -187,30 +289,84 @@ const VERIFICATION_PROMPT = `You are a Plan Verification Expert reviewing an imp
{PLAN}
## YOUR MISSION
Review this plan and:
Review and enhance this plan:
1. Assign priorities (P0=critical/blocking, P1=required, P2=enhancement)
2. Identify any gaps or missing steps
3. Check logical ordering (tests before implementation, setup before coding)
4. Flag potential issues or warnings
5. Calculate an overall quality score (0.0-1.0)
2. Add verification criteria to EVERY task (how to know it's done)
3. Pair test tasks with implementation tasks (TDD cycle)
4. Add dependencies where one task blocks another
5. Identify gaps and calculate quality score
## PRIORITY GUIDELINES
- P0: Foundation tasks, type definitions, project setup, blocking dependencies
- P1: Core implementation, tests, main features, error handling
- P2: Polish, optimization, documentation, nice-to-have features
## TDD + REVIEW CYCLE RULES
The complete cycle is: test → impl → review
- Every implementation task should have a corresponding test task AND review task
- Test task comes BEFORE its paired implementation task
- Review task comes AFTER the implementation it reviews
- Use "pairedWith" to link test ↔ implementation ↔ review
- Verification criteria should reference test results where applicable
## REVIEW TASK REQUIREMENTS
After EVERY implementation task, add a review task that checks:
- Best practices for the language/framework
- Security vulnerabilities (OWASP top 10)
- Performance concerns
- Error handling completeness
- Code quality (DRY, SOLID, readability)
## OUTPUT FORMAT
Return ONLY a JSON object:
{
"validatedPlan": [
{"content": "step description", "priority": "P0|P1|P2", "rationale": "why this priority"}
{
"id": "P0-001",
"content": "Write failing test for user authentication",
"priority": "P0",
"tddPhase": "test",
"verificationCriteria": "Test file exists, test fails with 'not implemented'",
"testCommand": "npm test -- --grep='auth'",
"pairedWith": "P0-002",
"dependencies": [],
"complexity": "low"
},
{
"id": "P0-002",
"content": "Implement user authentication handler",
"priority": "P0",
"tddPhase": "impl",
"verificationCriteria": "npm test -- --grep='auth' passes",
"pairedWith": "P0-001",
"dependencies": ["P0-001"],
"complexity": "medium"
},
{
"id": "P0-003",
"content": "Review auth implementation for security and best practices",
"priority": "P0",
"tddPhase": "review",
"verificationCriteria": "No security issues found, follows TypeScript best practices",
"reviewChecklist": ["Input validation", "XSS prevention", "Session security", "Error handling"],
"pairedWith": "P0-002",
"dependencies": ["P0-002"],
"complexity": "low"
}
],
"gaps": ["missing requirement 1", "missing test coverage for X"],
"warnings": ["consider Y before Z", "potential issue with..."],
"qualityScore": 0.85
}
Be critical but constructive. A thorough review catches issues early.`;
CRITICAL REQUIREMENTS:
1. EVERY task MUST have verificationCriteria (how to verify completion)
2. Implementation tasks MUST have a paired test task AND a review task
3. Review tasks MUST have a reviewChecklist with specific items to check
4. Dependencies must form a valid DAG (no cycles)
5. Use sequential IDs: P0-001, P0-002, P0-003, P1-001, etc.
Be critical but constructive. A thorough review catches issues that tests miss.`;
// ============================================================================
// Main Orchestrator Class
@@ -270,11 +426,19 @@ export class PlanOrchestrator {
totalCost += 0.01; // Verification cost estimate
// Phase 4: Ensure all impl tasks have review tasks
onProgress?.('review-injection', 'Ensuring review tasks for all implementations...');
const planWithReviews = this.ensureReviewTasks(verificationResult.validatedPlan);
const reviewsAdded = planWithReviews.length - verificationResult.validatedPlan.length;
if (reviewsAdded > 0) {
onProgress?.('review-injection', `Added ${reviewsAdded} auto-review task(s)`);
}
const totalDurationMs = Date.now() - startTime;
return {
success: true,
items: verificationResult.validatedPlan,
items: planWithReviews,
costUsd: totalCost,
metadata: {
subagentResults,
@@ -573,19 +737,51 @@ export class PlanOrchestrator {
const parsed = JSON.parse(jsonMatch[0]);
const validatedPlan: PlanItem[] = (parsed.validatedPlan || []).map((item: unknown) => {
const validatedPlan: PlanItem[] = (parsed.validatedPlan || []).map((item: unknown, idx: number) => {
if (typeof item !== 'object' || item === null) {
return { content: String(item), priority: 'P1' as const };
return { content: String(item), priority: 'P1' as const, id: `task-${idx}` };
}
const obj = item as Record<string, unknown>;
let priority: PlanItem['priority'] = null;
if (obj.priority === 'P0' || obj.priority === 'P1' || obj.priority === 'P2') {
priority = obj.priority;
}
// Parse TDD phase (now includes 'review')
let tddPhase: PlanItem['tddPhase'];
if (obj.tddPhase === 'setup' || obj.tddPhase === 'test' || obj.tddPhase === 'impl' || obj.tddPhase === 'verify' || obj.tddPhase === 'review') {
tddPhase = obj.tddPhase;
}
// Parse complexity
let complexity: PlanItem['complexity'];
if (obj.complexity === 'low' || obj.complexity === 'medium' || obj.complexity === 'high') {
complexity = obj.complexity;
}
// Parse reviewChecklist for review tasks
let reviewChecklist: string[] | undefined;
if (Array.isArray(obj.reviewChecklist)) {
reviewChecklist = obj.reviewChecklist.map(String);
}
return {
id: obj.id ? String(obj.id) : `task-${idx}`,
content: String(obj.content || ''),
priority,
rationale: obj.rationale ? String(obj.rationale) : undefined,
// Enhanced fields
verificationCriteria: obj.verificationCriteria ? String(obj.verificationCriteria) : undefined,
testCommand: obj.testCommand ? String(obj.testCommand) : undefined,
tddPhase,
pairedWith: obj.pairedWith ? String(obj.pairedWith) : undefined,
dependencies: Array.isArray(obj.dependencies) ? obj.dependencies.map(String) : undefined,
complexity,
reviewChecklist,
// Execution tracking defaults
status: 'pending' as PlanTaskStatus,
attempts: 0,
version: 1,
};
});
@@ -611,13 +807,21 @@ export class PlanOrchestrator {
/**
* Fallback verification when the verification subagent fails.
* Assigns heuristic priorities and adds basic verification criteria.
*/
private fallbackVerification(items: PlanItem[]): VerificationResult {
return {
validatedPlan: items.map(item => ({
validatedPlan: items.map((item, idx) => ({
...item,
id: item.id || `task-${idx}`,
priority: item.phase === 1 ? 'P0' as const :
item.phase === 4 ? 'P2' as const : 'P1' as const,
// Add default verification criteria based on content
verificationCriteria: item.verificationCriteria ||
this.inferVerificationCriteria(item.content),
status: 'pending' as PlanTaskStatus,
attempts: 0,
version: 1,
})),
gaps: [],
warnings: ['Verification subagent failed - using heuristic priorities'],
@@ -625,6 +829,34 @@ export class PlanOrchestrator {
};
}
/**
* Infer verification criteria from task content.
*/
private inferVerificationCriteria(content: string): string {
const lower = content.toLowerCase();
if (lower.includes('test')) {
return 'Tests pass without errors';
}
if (lower.includes('implement') || lower.includes('create') || lower.includes('add')) {
return 'Code compiles, no type errors';
}
if (lower.includes('fix') || lower.includes('debug')) {
return 'Issue is resolved, tests pass';
}
if (lower.includes('refactor')) {
return 'Code refactored, all tests still pass';
}
if (lower.includes('document') || lower.includes('readme')) {
return 'Documentation exists and is accurate';
}
if (lower.includes('config') || lower.includes('setup')) {
return 'Configuration is valid, app starts';
}
return 'Task completed successfully';
}
/**
* Create a timeout promise.
*/
@@ -633,4 +865,130 @@ export class PlanOrchestrator {
setTimeout(() => reject(new Error(`Timeout after ${ms}ms`)), ms);
});
}
/**
* Inject review tasks for any implementation tasks that don't have them.
* Called as a post-processing step to ensure the test → impl → review cycle is complete.
*
* @param items - The validated plan items
* @returns Updated plan items with review tasks added
*/
ensureReviewTasks(items: PlanItem[]): PlanItem[] {
const result: PlanItem[] = [];
const implTasksNeedingReview: Map<string, PlanItem> = new Map();
// First pass: identify impl tasks and their existing review pairs
const reviewPairs = new Set<string>();
for (const item of items) {
if (item.tddPhase === 'review' && item.pairedWith) {
reviewPairs.add(item.pairedWith);
}
}
// Second pass: collect impl tasks without review pairs
for (const item of items) {
if (item.tddPhase === 'impl' && item.id && !reviewPairs.has(item.id)) {
implTasksNeedingReview.set(item.id, item);
}
}
// Third pass: build result with injected review tasks
for (const item of items) {
result.push(item);
// If this is an impl task needing review, inject one after it
if (item.id && implTasksNeedingReview.has(item.id)) {
const reviewId = this.generateReviewId(item.id);
const reviewTask = this.createReviewTask(item, reviewId);
result.push(reviewTask);
}
}
return result;
}
/**
* Generate a review task ID from an impl task ID.
* P0-002 → P0-002-R, task-5 → task-5-R
*/
private generateReviewId(implId: string): string {
return `${implId}-R`;
}
/**
* Create a review task for an implementation task.
*/
private createReviewTask(implTask: PlanItem, reviewId: string): PlanItem {
const reviewChecklist = this.generateReviewChecklist(implTask.content);
return {
id: reviewId,
content: `Review: ${implTask.content}`,
priority: implTask.priority,
tddPhase: 'review',
pairedWith: implTask.id,
dependencies: implTask.id ? [implTask.id] : [],
verificationCriteria: 'Code review complete, no issues found or all issues addressed',
reviewChecklist,
status: 'pending',
attempts: 0,
version: implTask.version || 1,
complexity: 'low',
};
}
/**
* Generate a review checklist based on the implementation task content.
*/
private generateReviewChecklist(content: string): string[] {
const lower = content.toLowerCase();
const checklist: string[] = [];
// Always include these
checklist.push('Code compiles without errors');
checklist.push('No TypeScript/linting warnings');
// Security checks for certain patterns
if (lower.includes('auth') || lower.includes('login') || lower.includes('password')) {
checklist.push('Input validation implemented');
checklist.push('No sensitive data in logs');
checklist.push('Secure session handling');
}
if (lower.includes('api') || lower.includes('endpoint') || lower.includes('route')) {
checklist.push('Request validation');
checklist.push('Error responses do not leak internals');
checklist.push('Rate limiting considered');
}
if (lower.includes('database') || lower.includes('query') || lower.includes('sql')) {
checklist.push('Parameterized queries used');
checklist.push('No N+1 query issues');
}
if (lower.includes('file') || lower.includes('path') || lower.includes('upload')) {
checklist.push('Path traversal prevented');
checklist.push('File type validation');
}
if (lower.includes('user') || lower.includes('input')) {
checklist.push('XSS prevention');
checklist.push('Input sanitization');
}
// Performance checks
if (lower.includes('loop') || lower.includes('iterate') || lower.includes('array')) {
checklist.push('Algorithm complexity is appropriate');
}
// Error handling
checklist.push('Error cases handled');
checklist.push('Meaningful error messages');
// Code quality
checklist.push('Code is readable and maintainable');
checklist.push('No duplicate logic');
return checklist;
}
}
+75 -21
View File
@@ -365,6 +365,38 @@ export class ScreenManager extends EventEmitter {
return pids;
}
// Check if a process is still alive
private isProcessAlive(pid: number): boolean {
try {
// signal 0 doesn't kill, just checks if process exists
process.kill(pid, 0);
return true;
} catch {
return false;
}
}
// Verify all PIDs are dead, with retry
private async verifyProcessesDead(pids: number[], maxWaitMs: number = 1000): Promise<boolean> {
const startTime = Date.now();
const checkInterval = 50;
while (Date.now() - startTime < maxWaitMs) {
const aliveCount = pids.filter(pid => this.isProcessAlive(pid)).length;
if (aliveCount === 0) {
return true;
}
await new Promise(resolve => setTimeout(resolve, checkInterval));
}
// Log any processes that are still alive
const stillAlive = pids.filter(pid => this.isProcessAlive(pid));
if (stillAlive.length > 0) {
console.warn(`[ScreenManager] ${stillAlive.length} processes still alive after kill: ${stillAlive.join(', ')}`);
}
return stillAlive.length === 0;
}
// Kill a screen session and all its child processes
async killScreen(sessionId: string): Promise<boolean> {
const screen = this.screens.get(sessionId);
@@ -377,40 +409,54 @@ export class ScreenManager extends EventEmitter {
console.log(`[ScreenManager] Killing screen ${screen.screenName} (PID ${currentPid})`);
// Collect all PIDs to track (for verification)
const allPids: number[] = [currentPid];
// Strategy 1: Find and kill all child processes recursively
const childPids = this.getChildPids(currentPid);
// Re-check for children before each kill attempt (they may have changed)
let childPids = this.getChildPids(currentPid);
if (childPids.length > 0) {
console.log(`[ScreenManager] Found ${childPids.length} child processes to kill`);
allPids.push(...childPids);
// Kill children in reverse order (deepest first) with SIGTERM
for (const childPid of childPids.reverse()) {
try {
process.kill(childPid, 'SIGTERM');
} catch {
// Process may already be dead
for (const childPid of [...childPids].reverse()) {
if (this.isProcessAlive(childPid)) {
try {
process.kill(childPid, 'SIGTERM');
} catch {
// Process may already be dead
}
}
}
// Give processes a moment to terminate gracefully
await new Promise(resolve => setTimeout(resolve, SCREEN_KILL_WAIT_MS));
// Force kill any remaining children
// Re-check which children are still alive and force kill them
childPids = this.getChildPids(currentPid);
for (const childPid of childPids) {
try {
process.kill(childPid, 'SIGKILL');
} catch {
// Process already terminated
if (this.isProcessAlive(childPid)) {
try {
process.kill(childPid, 'SIGKILL');
} catch {
// Process already terminated
}
}
}
}
// Strategy 2: Kill the entire process group (catches any orphans we missed)
try {
process.kill(-currentPid, 'SIGTERM');
await new Promise(resolve => setTimeout(resolve, GRACEFUL_SHUTDOWN_WAIT_MS));
process.kill(-currentPid, 'SIGKILL');
} catch {
// Process group may not exist or already terminated
if (this.isProcessAlive(currentPid)) {
try {
process.kill(-currentPid, 'SIGTERM');
await new Promise(resolve => setTimeout(resolve, GRACEFUL_SHUTDOWN_WAIT_MS));
if (this.isProcessAlive(currentPid)) {
process.kill(-currentPid, 'SIGKILL');
}
} catch {
// Process group may not exist or already terminated
}
}
// Strategy 3: Kill screen session by name
@@ -423,10 +469,18 @@ export class ScreenManager extends EventEmitter {
}
// Strategy 4: Direct kill by PID as final fallback
try {
process.kill(currentPid, 'SIGKILL');
} catch {
// Already dead
if (this.isProcessAlive(currentPid)) {
try {
process.kill(currentPid, 'SIGKILL');
} catch {
// Already dead
}
}
// Verify all processes are dead (with timeout)
const allDead = await this.verifyProcessesDead(allPids, 2000);
if (!allDead) {
console.error(`[ScreenManager] Warning: Some processes may still be alive for screen ${screen.screenName}`);
}
this.screens.delete(sessionId);
+56 -34
View File
@@ -61,6 +61,9 @@ export class SessionManager extends EventEmitter {
private sessionHandlers: Map<string, SessionHandlers> = new Map();
private store = getStore();
// Mutex for session creation to prevent race conditions
private _sessionCreationLock: Promise<void> | null = null;
/**
* Creates a new SessionManager and loads previous session state.
*/
@@ -84,54 +87,73 @@ export class SessionManager extends EventEmitter {
/**
* Creates and starts a new Claude session.
* Uses mutex to prevent race conditions when multiple requests arrive simultaneously.
*
* @param workingDir - Working directory for the session
* @returns The newly created session
* @throws Error if max concurrent sessions limit reached
*/
async createSession(workingDir: string): Promise<Session> {
const config = this.store.getConfig();
if (this.sessions.size >= config.maxConcurrentSessions) {
throw new Error(`Maximum concurrent sessions (${config.maxConcurrentSessions}) reached`);
// Wait for any pending session creation to complete (mutex pattern)
while (this._sessionCreationLock) {
await this._sessionCreationLock;
}
const session = new Session({ workingDir });
// Create a new lock promise that others will wait on
let unlock: () => void;
this._sessionCreationLock = new Promise<void>(resolve => {
unlock = resolve;
});
// Set up event forwarding with stored handlers for cleanup
const handlers: SessionHandlers = {
output: (data: string) => {
this.emit('sessionOutput', session.id, data);
this.updateSessionState(session);
},
error: (data: string) => {
this.emit('sessionError', session.id, data);
this.updateSessionState(session);
},
completion: (phrase: string) => {
this.emit('sessionCompletion', session.id, phrase);
},
exit: () => {
this.emit('sessionStopped', session.id);
this.updateSessionState(session);
},
};
try {
const config = this.store.getConfig();
session.on('output', handlers.output);
session.on('error', handlers.error);
session.on('completion', handlers.completion);
session.on('exit', handlers.exit);
// Check limit INSIDE the lock to prevent race conditions
if (this.sessions.size >= config.maxConcurrentSessions) {
throw new Error(`Maximum concurrent sessions (${config.maxConcurrentSessions}) reached`);
}
// Store handlers for later cleanup
this.sessionHandlers.set(session.id, handlers);
const session = new Session({ workingDir });
await session.start();
// Set up event forwarding with stored handlers for cleanup
const handlers: SessionHandlers = {
output: (data: string) => {
this.emit('sessionOutput', session.id, data);
this.updateSessionState(session);
},
error: (data: string) => {
this.emit('sessionError', session.id, data);
this.updateSessionState(session);
},
completion: (phrase: string) => {
this.emit('sessionCompletion', session.id, phrase);
},
exit: () => {
this.emit('sessionStopped', session.id);
this.updateSessionState(session);
},
};
this.sessions.set(session.id, session);
this.store.setSession(session.id, session.toState());
session.on('output', handlers.output);
session.on('error', handlers.error);
session.on('completion', handlers.completion);
session.on('exit', handlers.exit);
this.emit('sessionStarted', session);
return session;
// Store handlers for later cleanup
this.sessionHandlers.set(session.id, handlers);
await session.start();
this.sessions.set(session.id, session);
this.store.setSession(session.id, session.toState());
this.emit('sessionStarted', session);
return session;
} finally {
// Release the lock so other createSession calls can proceed
this._sessionCreationLock = null;
unlock!();
}
}
/**
+95 -63
View File
@@ -888,15 +888,21 @@ export class Session extends EventEmitter {
}
// Attach to the screen session via PTY
this.ptyProcess = pty.spawn('screen', [
'-x', this._screenSession!.screenName
], {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: { ...process.env, TERM: 'xterm-256color' },
});
try {
this.ptyProcess = pty.spawn('screen', [
'-x', this._screenSession!.screenName
], {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: { ...process.env, TERM: 'xterm-256color' },
});
} catch (spawnErr) {
console.error('[Session] Failed to spawn PTY for screen attachment:', spawnErr);
this.emit('error', `Failed to attach to screen: ${spawnErr}`);
throw spawnErr;
}
// For NEW screens: wait for prompt to appear then clean buffer
// For RESTORED screens: don't do anything - client will fetch buffer on tab switch
@@ -941,23 +947,30 @@ export class Session extends EventEmitter {
// Fallback to direct PTY if screen is not used
if (!this.ptyProcess) {
this.ptyProcess = pty.spawn('claude', [
'--dangerously-skip-permissions'
], {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: {
...process.env,
PATH: getAugmentedPath(),
TERM: 'xterm-256color',
// Inform Claude it's running within Claudeman (helps prevent self-termination)
CLAUDEMAN_SCREEN: '1',
CLAUDEMAN_SESSION_ID: this.id,
CLAUDEMAN_API_URL: process.env.CLAUDEMAN_API_URL || 'http://localhost:3000',
},
});
try {
this.ptyProcess = pty.spawn('claude', [
'--dangerously-skip-permissions'
], {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: {
...process.env,
PATH: getAugmentedPath(),
TERM: 'xterm-256color',
// Inform Claude it's running within Claudeman (helps prevent self-termination)
CLAUDEMAN_SCREEN: '1',
CLAUDEMAN_SESSION_ID: this.id,
CLAUDEMAN_API_URL: process.env.CLAUDEMAN_API_URL || 'http://localhost:3000',
},
});
} catch (spawnErr) {
console.error('[Session] Failed to spawn Claude PTY:', spawnErr);
this._status = 'stopped';
this.emit('error', `Failed to start Claude: ${spawnErr}`);
throw new Error(`Failed to spawn Claude process: ${spawnErr}`);
}
}
this._pid = this.ptyProcess.pid;
@@ -1101,15 +1114,21 @@ export class Session extends EventEmitter {
}
// Attach to the screen session via PTY
this.ptyProcess = pty.spawn('screen', [
'-x', this._screenSession!.screenName
], {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: { ...process.env, TERM: 'xterm-256color' },
});
try {
this.ptyProcess = pty.spawn('screen', [
'-x', this._screenSession!.screenName
], {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: { ...process.env, TERM: 'xterm-256color' },
});
} catch (spawnErr) {
console.error('[Session] Failed to spawn PTY for shell screen attachment:', spawnErr);
this.emit('error', `Failed to attach to screen: ${spawnErr}`);
throw spawnErr;
}
// For NEW screens: clear by sending 'clear' command to the shell
// For RESTORED screens: don't clear - we want to see the existing output
@@ -1130,19 +1149,26 @@ export class Session extends EventEmitter {
// Fallback to direct PTY if screen is not used
if (!this.ptyProcess) {
this.ptyProcess = pty.spawn(shell, [], {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: {
...process.env,
TERM: 'xterm-256color',
CLAUDEMAN_SCREEN: '1',
CLAUDEMAN_SESSION_ID: this.id,
CLAUDEMAN_API_URL: process.env.CLAUDEMAN_API_URL || 'http://localhost:3000',
},
});
try {
this.ptyProcess = pty.spawn(shell, [], {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: {
...process.env,
TERM: 'xterm-256color',
CLAUDEMAN_SCREEN: '1',
CLAUDEMAN_SESSION_ID: this.id,
CLAUDEMAN_API_URL: process.env.CLAUDEMAN_API_URL || 'http://localhost:3000',
},
});
} catch (spawnErr) {
console.error('[Session] Failed to spawn shell PTY:', spawnErr);
this._status = 'stopped';
this.emit('error', `Failed to start shell: ${spawnErr}`);
throw new Error(`Failed to spawn shell process: ${spawnErr}`);
}
}
this._pid = this.ptyProcess.pid;
@@ -1251,21 +1277,27 @@ export class Session extends EventEmitter {
}
args.push(prompt);
this.ptyProcess = pty.spawn('claude', args, {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: {
...process.env,
PATH: getAugmentedPath(),
TERM: 'xterm-256color',
// Inform Claude it's running within Claudeman
CLAUDEMAN_SCREEN: '1',
CLAUDEMAN_SESSION_ID: this.id,
CLAUDEMAN_API_URL: process.env.CLAUDEMAN_API_URL || 'http://localhost:3000',
},
});
try {
this.ptyProcess = pty.spawn('claude', args, {
name: 'xterm-256color',
cols: 120,
rows: 40,
cwd: this.workingDir,
env: {
...process.env,
PATH: getAugmentedPath(),
TERM: 'xterm-256color',
// Inform Claude it's running within Claudeman
CLAUDEMAN_SCREEN: '1',
CLAUDEMAN_SESSION_ID: this.id,
CLAUDEMAN_API_URL: process.env.CLAUDEMAN_API_URL || 'http://localhost:3000',
},
});
} catch (spawnErr) {
console.error('[Session] Failed to spawn Claude PTY for runPrompt:', spawnErr);
this.emit('error', `Failed to spawn Claude: ${spawnErr instanceof Error ? spawnErr.message : String(spawnErr)}`);
throw spawnErr;
}
this._pid = this.ptyProcess.pid;
console.log('[Session] PTY spawned with PID:', this._pid);
+101 -3
View File
@@ -43,6 +43,9 @@ const SAVE_DEBOUNCE_MS = 500;
* store.saveNow();
* ```
*/
/** Maximum consecutive save failures before circuit breaker opens */
const MAX_CONSECUTIVE_FAILURES = 3;
export class StateStore {
private state: AppState;
private filePath: string;
@@ -55,6 +58,10 @@ export class StateStore {
private ralphStateSaveTimeout: NodeJS.Timeout | null = null;
private ralphStateDirty: boolean = false;
// Circuit breaker for save failures (prevents hammering disk on persistent errors)
private consecutiveSaveFailures: number = 0;
private circuitBreakerOpen: boolean = false;
constructor(filePath?: string) {
this.filePath = filePath || join(homedir(), '.claudeman', 'state.json');
this.ralphStatePath = this.filePath.replace('.json', '-inner.json');
@@ -109,6 +116,7 @@ export class StateStore {
/**
* Immediately writes state to disk using atomic write pattern.
* Writes to temp file first, then renames to prevent corruption on crash.
* Includes backup mechanism and circuit breaker for reliability.
* Use when guaranteed persistence is required (e.g., before shutdown).
*/
saveNow(): void {
@@ -119,32 +127,122 @@ export class StateStore {
if (!this.dirty) {
return;
}
// Circuit breaker: stop attempting writes after too many failures
if (this.circuitBreakerOpen) {
console.warn('[StateStore] Circuit breaker open - skipping save (too many consecutive failures)');
return;
}
this.dirty = false;
this.ensureDir();
// Atomic write: write to temp file, then rename (atomic on POSIX)
const tempPath = this.filePath + '.tmp';
const backupPath = this.filePath + '.bak';
let json: string;
// Step 1: Serialize state (validates it's JSON-safe)
try {
json = JSON.stringify(this.state, null, 2);
} catch (err) {
console.error('[StateStore] Failed to serialize state (circular reference or invalid data):', err);
throw err;
this.consecutiveSaveFailures++;
if (this.consecutiveSaveFailures >= MAX_CONSECUTIVE_FAILURES) {
console.error('[StateStore] Circuit breaker OPEN - serialization failing repeatedly');
this.circuitBreakerOpen = true;
}
// Don't throw - this prevents crashing the app
// Mark dirty again so we can retry later
this.dirty = true;
return;
}
// Step 2: Create backup of current state file (if exists)
try {
if (existsSync(this.filePath)) {
// Read current file and verify it's valid JSON before backing up
const currentContent = readFileSync(this.filePath, 'utf-8');
JSON.parse(currentContent); // Validate
writeFileSync(backupPath, currentContent, 'utf-8');
}
} catch {
// Backup failed - current file may be corrupt, continue with write
console.warn('[StateStore] Could not create backup (current file may be corrupt)');
}
// Step 3: Atomic write: write to temp file, then rename
try {
writeFileSync(tempPath, json, 'utf-8');
renameSync(tempPath, this.filePath);
// Success! Reset failure counter
this.consecutiveSaveFailures = 0;
if (this.circuitBreakerOpen) {
console.log('[StateStore] Circuit breaker CLOSED - save succeeded');
this.circuitBreakerOpen = false;
}
} catch (err) {
console.error('[StateStore] Failed to write state file:', err);
this.consecutiveSaveFailures++;
// Try to clean up temp file on error
try {
if (existsSync(tempPath)) {
unlinkSync(tempPath);
}
} catch { /* ignore cleanup errors */ }
throw err;
// Check circuit breaker threshold
if (this.consecutiveSaveFailures >= MAX_CONSECUTIVE_FAILURES) {
console.error('[StateStore] Circuit breaker OPEN - writes failing repeatedly');
this.circuitBreakerOpen = true;
}
// Mark dirty so we retry later (don't throw to avoid crashing app)
this.dirty = true;
}
}
/**
* Attempt to recover state from backup file.
* Call this if main state file is corrupt.
*/
recoverFromBackup(): boolean {
const backupPath = this.filePath + '.bak';
try {
if (existsSync(backupPath)) {
const backupContent = readFileSync(backupPath, 'utf-8');
const parsed = JSON.parse(backupContent) as Partial<AppState>;
const initial = createInitialState();
this.state = {
...initial,
...parsed,
sessions: { ...parsed.sessions },
tasks: { ...parsed.tasks },
ralphLoop: { ...initial.ralphLoop, ...parsed.ralphLoop },
config: { ...initial.config, ...parsed.config },
};
console.log('[StateStore] Successfully recovered state from backup');
// Reset circuit breaker after successful recovery
this.circuitBreakerOpen = false;
this.consecutiveSaveFailures = 0;
return true;
}
} catch (err) {
console.error('[StateStore] Failed to recover from backup:', err);
}
return false;
}
/**
* Reset the circuit breaker (for manual intervention).
*/
resetCircuitBreaker(): void {
this.circuitBreakerOpen = false;
this.consecutiveSaveFailures = 0;
console.log('[StateStore] Circuit breaker manually reset');
}
/** Flushes any pending main state save. Call before shutdown. */
flush(): void {
this.saveNow();
+741 -149
View File
File diff suppressed because it is too large Load Diff
+184 -29
View File
@@ -2197,7 +2197,10 @@ export class WebServer extends EventEmitter {
if (!existsSync(dir)) {
mkdirSync(dir, { recursive: true });
}
writeFileSync(settingsFilePath, JSON.stringify(settings, null, 2));
// Use async write to avoid blocking event loop
fs.writeFile(settingsFilePath, JSON.stringify(settings, null, 2)).catch(() => {
// Non-critical, ignore settings save errors
});
} catch {
// Non-critical, ignore settings save errors
}
@@ -2221,10 +2224,8 @@ export class WebServer extends EventEmitter {
detailLevel?: 'brief' | 'standard' | 'detailed';
}
interface PlanItem {
content: string;
priority: 'P0' | 'P1' | 'P2' | null;
}
// Use enhanced PlanItem from orchestrator (has verification, dependencies, tracking)
type PlanItem = import('../plan-orchestrator.js').PlanItem;
this.app.post('/api/generate-plan', async (req): Promise<ApiResponse> => {
const {
@@ -2294,34 +2295,35 @@ For EACH feature:
- Final verification that ALL requirements are met
## OUTPUT FORMAT
Return ONLY a JSON array. Each item:
Return ONLY a JSON array. Each item MUST have:
- id: unique identifier (e.g., "P0-001", "P1-002")
- content: specific action (verb phrase, 15-120 chars, be descriptive!)
- priority: "P0" (critical/blocking), "P1" (required), "P2" (enhancement)
- verificationCriteria: HOW to verify this step is complete (required!)
- tddPhase: "setup" | "test" | "impl" | "verify"
- dependencies: array of task IDs this depends on (empty if none)
## EXAMPLE OUTPUT
[
{"content": "Create project structure with src/, tests/, and config directories", "priority": "P0"},
{"content": "Define TypeScript interfaces for User, Session, and AuthToken types", "priority": "P0"},
{"content": "Write failing unit tests for password hashing (valid password, empty, too short)", "priority": "P0"},
{"content": "Implement password hashing with bcrypt, configurable salt rounds", "priority": "P0"},
{"content": "Run password tests and debug until all pass", "priority": "P0"},
{"content": "Write failing tests for JWT token generation and validation", "priority": "P0"},
{"content": "Implement JWT service with access/refresh token support", "priority": "P0"},
{"content": "Run JWT tests and verify token expiration handling works", "priority": "P0"},
{"content": "Write integration tests for login flow (valid creds, invalid, locked account)", "priority": "P1"},
{"content": "Implement login endpoint with rate limiting and audit logging", "priority": "P1"},
{"content": "Add error handling for network failures and database timeouts", "priority": "P1"},
{"content": "Run full test suite and fix any failures", "priority": "P1"},
{"content": "Verify all original requirements are implemented and tested", "priority": "P1"}
{"id": "P0-001", "content": "Create project structure with src/, tests/, and config directories", "priority": "P0", "verificationCriteria": "Directories exist, package.json initialized", "tddPhase": "setup", "dependencies": []},
{"id": "P0-002", "content": "Define TypeScript interfaces for User, Session, and AuthToken types", "priority": "P0", "verificationCriteria": "Types compile without errors, exported from types.ts", "tddPhase": "setup", "dependencies": ["P0-001"]},
{"id": "P0-003", "content": "Write failing unit tests for password hashing (valid password, empty, too short)", "priority": "P0", "verificationCriteria": "Tests exist, fail with 'not implemented'", "tddPhase": "test", "dependencies": ["P0-002"]},
{"id": "P0-004", "content": "Implement password hashing with bcrypt, configurable salt rounds", "priority": "P0", "verificationCriteria": "npm test -- --grep='password' passes", "tddPhase": "impl", "dependencies": ["P0-003"]},
{"id": "P0-005", "content": "Write failing tests for JWT token generation and validation", "priority": "P0", "verificationCriteria": "Tests exist, fail with 'not implemented'", "tddPhase": "test", "dependencies": ["P0-004"]},
{"id": "P0-006", "content": "Implement JWT service with access/refresh token support", "priority": "P0", "verificationCriteria": "npm test -- --grep='JWT' passes", "tddPhase": "impl", "dependencies": ["P0-005"]},
{"id": "P1-001", "content": "Write integration tests for login flow (valid creds, invalid, locked account)", "priority": "P1", "verificationCriteria": "Integration tests exist, fail until endpoint implemented", "tddPhase": "test", "dependencies": ["P0-006"]},
{"id": "P1-002", "content": "Implement login endpoint with rate limiting and audit logging", "priority": "P1", "verificationCriteria": "All login tests pass, endpoint returns 200/401 correctly", "tddPhase": "impl", "dependencies": ["P1-001"]},
{"id": "P1-003", "content": "Run full test suite and verify all tests pass", "priority": "P1", "verificationCriteria": "npm test exits with code 0, coverage > 80%", "tddPhase": "verify", "dependencies": ["P1-002"]}
]
## CRITICAL RULES
1. EVERY implementation step should have a corresponding test step BEFORE it
2. Include "Run tests and debug/fix" steps after implementation blocks
3. Be SPECIFIC - not "Add tests" but "Write tests for X covering Y and Z"
4. Think about what could fail and add defensive steps
5. End with verification that ALL original requirements are met
6. Use P0 for foundation and core features, P1 for required work, P2 for nice-to-have
1. EVERY task MUST have verificationCriteria - this is non-negotiable!
2. EVERY implementation step should have a corresponding test step BEFORE it
3. Use tddPhase: "test" for writing tests, "impl" for implementation
4. Dependencies must form a valid DAG - no cycles
5. Be SPECIFIC - not "Add tests" but "Write tests for X covering Y and Z"
6. End with verification that ALL original requirements are met
7. Use P0 for foundation and core features, P1 for required work, P2 for nice-to-have
NOW: Generate the implementation plan for the task above. Think step by step.`;
@@ -2352,10 +2354,18 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, 'Invalid response - expected array');
}
// Validate and normalize items
// Validate and normalize items with enhanced fields
items = parsed.map((item: unknown, idx: number) => {
if (typeof item !== 'object' || item === null) {
return { content: `Step ${idx + 1}`, priority: null };
return {
id: `task-${idx}`,
content: `Step ${idx + 1}`,
priority: null,
verificationCriteria: 'Task completed successfully',
status: 'pending' as const,
attempts: 0,
version: 1,
};
}
const obj = item as Record<string, unknown>;
const content = typeof obj.content === 'string' ? obj.content.slice(0, 200) : `Step ${idx + 1}`;
@@ -2363,7 +2373,26 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
if (obj.priority === 'P0' || obj.priority === 'P1' || obj.priority === 'P2') {
priority = obj.priority;
}
return { content, priority };
// Parse tddPhase
let tddPhase: 'setup' | 'test' | 'impl' | 'verify' | undefined;
if (obj.tddPhase === 'setup' || obj.tddPhase === 'test' || obj.tddPhase === 'impl' || obj.tddPhase === 'verify') {
tddPhase = obj.tddPhase;
}
return {
id: obj.id ? String(obj.id) : `task-${idx}`,
content,
priority,
verificationCriteria: typeof obj.verificationCriteria === 'string'
? obj.verificationCriteria
: 'Task completed successfully',
tddPhase,
dependencies: Array.isArray(obj.dependencies) ? obj.dependencies.map(String) : [],
status: 'pending' as const,
attempts: 0,
version: 1,
};
});
// No artificial limit - let Claude generate what's needed
@@ -2436,6 +2465,123 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
}
});
// ============ Plan Management Endpoints ============
// These endpoints support runtime plan adaptation with checkpoints, failure tracking, and versioning
// Update a specific plan task (status, attempts, errors)
this.app.patch('/api/sessions/:id/plan/task/:taskId', async (req) => {
const { id, taskId } = req.params as { id: string; taskId: string };
const session = this.sessions.get(id);
if (!session) {
return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found');
}
const tracker = session.ralphTracker;
if (!tracker) {
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, 'Ralph tracker not available');
}
const update = req.body as {
status?: 'pending' | 'in_progress' | 'completed' | 'failed' | 'blocked';
error?: string;
incrementAttempts?: boolean;
};
const result = tracker.updatePlanTask(taskId, update);
if (!result.success) {
return createErrorResponse(ApiErrorCode.NOT_FOUND, result.error || 'Task not found');
}
this.broadcast('session:planTaskUpdate', { sessionId: id, taskId, update: result.task });
return { success: true, data: result.task };
});
// Trigger a checkpoint review (at iterations 5, 10, 20, etc.)
this.app.post('/api/sessions/:id/plan/checkpoint', async (req) => {
const { id } = req.params as { id: string };
const session = this.sessions.get(id);
if (!session) {
return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found');
}
const tracker = session.ralphTracker;
if (!tracker) {
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, 'Ralph tracker not available');
}
const checkpoint = tracker.generateCheckpointReview();
this.broadcast('session:planCheckpoint', { sessionId: id, checkpoint });
return { success: true, data: checkpoint };
});
// Get plan version history
this.app.get('/api/sessions/:id/plan/history', async (req) => {
const { id } = req.params as { id: string };
const session = this.sessions.get(id);
if (!session) {
return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found');
}
const tracker = session.ralphTracker;
if (!tracker) {
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, 'Ralph tracker not available');
}
return { success: true, data: tracker.getPlanHistory() };
});
// Rollback to a previous plan version
this.app.post('/api/sessions/:id/plan/rollback/:version', async (req) => {
const { id, version } = req.params as { id: string; version: string };
const session = this.sessions.get(id);
if (!session) {
return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found');
}
const tracker = session.ralphTracker;
if (!tracker) {
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, 'Ralph tracker not available');
}
const result = tracker.rollbackToVersion(parseInt(version, 10));
if (!result.success) {
return createErrorResponse(ApiErrorCode.NOT_FOUND, result.error || 'Version not found');
}
this.broadcast('session:planRollback', { sessionId: id, version: parseInt(version, 10) });
return { success: true, data: result.plan };
});
// Add a new task to the plan (for runtime adaptation)
this.app.post('/api/sessions/:id/plan/task', async (req) => {
const { id } = req.params as { id: string };
const session = this.sessions.get(id);
if (!session) {
return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Session not found');
}
const tracker = session.ralphTracker;
if (!tracker) {
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, 'Ralph tracker not available');
}
const task = req.body as {
content: string;
priority?: 'P0' | 'P1' | 'P2';
verificationCriteria?: string;
dependencies?: string[];
insertAfter?: string; // Task ID to insert after
};
if (!task.content) {
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Task content is required');
}
const result = tracker.addPlanTask(task);
this.broadcast('session:planTaskAdded', { sessionId: id, task: result.task });
return { success: true, data: result.task };
});
// ============ App Settings Endpoints ============
const settingsPath = join(homedir(), '.claudeman', 'settings.json');
@@ -4281,7 +4427,16 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
this.sseHealthCheckTimer = null;
}
// Clear all SSE clients
// Gracefully close all SSE connections before clearing
for (const client of this.sseClients) {
try {
// Send a final event to notify clients of shutdown
this.sendSSE(client, 'server:shutdown', { reason: 'Server stopping' });
client.raw.end();
} catch {
// Client may already be disconnected
}
}
this.sseClients.clear();
// Clear batch timers