mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 12:39:42 +02:00
feat: comprehensive codebase improvements (cleanup, logic, performance)
Code Cleanups: - Consolidated PlanTaskStatus and TddPhase types to types.ts - Standardized node: prefix for Node.js builtin imports - Fixed timer type to NodeJS.Timeout in respawn-controller Logic Improvements: - Fixed AI check race condition with UUID tracking in respawn-controller - Added comprehensive _isStopped guards in session timer callbacks - Made RalphTracker reset events non-reentrant via process.nextTick() - Improved SessionManager mutex pattern reliability Performance Optimizations: - Fixed pendingToolCalls memory leak with TTL cleanup in subagent-watcher - Improved LRUMap.newest() from O(n) to O(1) with _newestKey tracking Documentation: - Enhanced README antiflicker section with 6-layer technical details Version: 0.1441 Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
@@ -16,7 +16,7 @@ When user says "COM":
|
||||
1. Increment version in BOTH `package.json` AND `CLAUDE.md`
|
||||
2. Run: `git add -A && git commit -m "chore: bump version to X.XXXX" && git push && npm run build && systemctl --user restart claudeman-web`
|
||||
|
||||
**Version**: 0.1440 (must match `package.json`)
|
||||
**Version**: 0.1441 (must match `package.json`)
|
||||
|
||||
## Project Overview
|
||||
|
||||
@@ -256,10 +256,6 @@ Use `LRUMap` for bounded caches with eviction, `StaleExpirationMap` for TTL-base
|
||||
| `scripts/monitor-respawn.sh` | Monitor respawn state machine in real-time |
|
||||
| `scripts/postinstall.js` | npm postinstall hook for setup |
|
||||
|
||||
## Deprecated Code
|
||||
|
||||
The TUI (Terminal UI) has been removed in favor of the web interface. Files in `src/tui/` are excluded from compilation via `tsconfig.json`.
|
||||
|
||||
## Memory Leak Prevention
|
||||
|
||||
Frontend runs long (24+ hour sessions); all Maps/timers must be cleaned up.
|
||||
|
||||
@@ -216,6 +216,78 @@ Run **20 parallel sessions** with full visibility:
|
||||
|
||||
---
|
||||
|
||||
### ⚡ Zero-Flicker Terminal Rendering
|
||||
|
||||
**The problem:** Claude Code uses [Ink](https://github.com/vadimdemedes/ink) (React for terminals), which redraws the entire screen on every state change. Without special handling, you'd see constant flickering — unusable for monitoring multiple sessions.
|
||||
|
||||
**The solution:** Claudeman implements a 6-layer antiflicker system that delivers butter-smooth 60fps terminal output:
|
||||
|
||||
```
|
||||
PTY Output → 16ms Batch → DEC 2026 Wrap → SSE → rAF Batch → xterm.js → 60fps Canvas
|
||||
```
|
||||
|
||||
#### How Each Layer Works
|
||||
|
||||
| Layer | Location | Technique | Purpose |
|
||||
|-------|----------|-----------|---------|
|
||||
| **1. Server Batching** | server.ts | 16ms collection window | Combines rapid PTY writes into single packets |
|
||||
| **2. DEC Mode 2026** | server.ts | `\x1b[?2026h`...`\x1b[?2026l` | Marks atomic update boundaries (terminal standard) |
|
||||
| **3. Client rAF** | app.js | `requestAnimationFrame` | Syncs writes to 60Hz display refresh |
|
||||
| **4. Sync Block Parser** | app.js | DEC 2026 extraction | Parses atomic segments for xterm.js |
|
||||
| **5. Flicker Filter** | app.js | Ink pattern detection | Buffers screen-clear sequences (optional) |
|
||||
| **6. Chunked Loading** | app.js | 64KB/frame writes | Large buffers don't freeze UI |
|
||||
|
||||
#### Technical Implementation
|
||||
|
||||
**Server-side (16ms batching + DEC 2026):**
|
||||
```typescript
|
||||
// Accumulate PTY output per-session
|
||||
const newBatch = existing + data;
|
||||
terminalBatches.set(sessionId, newBatch);
|
||||
|
||||
// Flush every 16ms (60fps) or immediately if >1KB
|
||||
if (!terminalBatchTimer) {
|
||||
terminalBatchTimer = setTimeout(() => {
|
||||
for (const [id, data] of terminalBatches) {
|
||||
// Wrap with synchronized output markers
|
||||
const syncData = '\x1b[?2026h' + data + '\x1b[?2026l';
|
||||
broadcast('session:terminal', { id, data: syncData });
|
||||
}
|
||||
terminalBatches.clear();
|
||||
}, 16);
|
||||
}
|
||||
```
|
||||
|
||||
**Client-side (rAF batching + sync block handling):**
|
||||
```javascript
|
||||
batchTerminalWrite(data) {
|
||||
pendingWrites += data;
|
||||
|
||||
if (!writeFrameScheduled) {
|
||||
writeFrameScheduled = true;
|
||||
requestAnimationFrame(() => {
|
||||
// Wait up to 50ms for incomplete sync blocks
|
||||
if (hasStartMarker && !hasEndMarker) {
|
||||
setTimeout(flushPendingWrites, 50);
|
||||
return;
|
||||
}
|
||||
|
||||
// Extract atomic segments, strip markers, write to xterm
|
||||
const segments = extractSyncSegments(pendingWrites);
|
||||
for (const segment of segments) {
|
||||
terminal.write(segment);
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
**Optional flicker filter** detects Ink's screen-clear patterns (`ESC[2J`, `ESC[H ESC[J`) and buffers 50ms of subsequent output for extra smoothness on problematic terminals.
|
||||
|
||||
**Result:** Watch 20 Claude sessions simultaneously without any visual artifacts, even during heavy tool use.
|
||||
|
||||
---
|
||||
|
||||
### 📈 Run Summary ("What Happened While You Were Away")
|
||||
|
||||
Click the chart icon on any session tab to see a complete timeline of what happened:
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "claudeman",
|
||||
"version": "0.1440",
|
||||
"version": "0.1441",
|
||||
"description": "The missing control plane for Claude Code - run 20 autonomous agents with real-time monitoring and session persistence",
|
||||
"type": "module",
|
||||
"main": "dist/index.js",
|
||||
|
||||
@@ -8,10 +8,10 @@
|
||||
* to ensure files are fully written before emitting detection events.
|
||||
*/
|
||||
|
||||
import { EventEmitter } from 'events';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { watch, type FSWatcher } from 'chokidar';
|
||||
import { basename, extname } from 'path';
|
||||
import { statSync } from 'fs';
|
||||
import { basename, extname } from 'node:path';
|
||||
import { statSync } from 'node:fs';
|
||||
import type { ImageDetectedEvent } from './types.js';
|
||||
|
||||
// ========== Types ==========
|
||||
|
||||
@@ -23,16 +23,14 @@ import {
|
||||
RESEARCH_AGENT_PROMPT,
|
||||
PLANNER_PROMPT,
|
||||
} from './prompts/index.js';
|
||||
import { PlanTaskStatus, TddPhase } from './types.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';
|
||||
/** Development phase in TDD cycle (alias for TddPhase) */
|
||||
export type PlanPhase = TddPhase;
|
||||
|
||||
/**
|
||||
* Plan item with TDD structure.
|
||||
@@ -260,7 +258,9 @@ export class PlanOrchestrator {
|
||||
for (const session of this.runningSessions) {
|
||||
try {
|
||||
session.stop();
|
||||
} catch {}
|
||||
} catch (err) {
|
||||
console.error('[PlanOrchestrator] Failed to stop session during cancel:', err);
|
||||
}
|
||||
}
|
||||
this.runningSessions.clear();
|
||||
}
|
||||
|
||||
+23
-25
@@ -31,6 +31,8 @@ import {
|
||||
CompletionConfidence,
|
||||
createInitialRalphTrackerState,
|
||||
createInitialCircuitBreakerStatus,
|
||||
PlanTaskStatus,
|
||||
TddPhase,
|
||||
} from './types.js';
|
||||
import {
|
||||
ANSI_ESCAPE_PATTERN_SIMPLE,
|
||||
@@ -38,14 +40,12 @@ import {
|
||||
todoContentHash,
|
||||
stringSimilarity,
|
||||
} from './utils/index.js';
|
||||
import { MAX_LINE_BUFFER_SIZE } from './config/buffer-limits.js';
|
||||
import { MAX_TODOS_PER_SESSION } from './config/map-limits.js';
|
||||
|
||||
// ========== Enhanced Plan Task Interface ==========
|
||||
|
||||
/** Task execution status for plan tracking */
|
||||
export type PlanTaskStatus = 'pending' | 'in_progress' | 'completed' | 'failed' | 'blocked';
|
||||
|
||||
/** TDD phase categories */
|
||||
export type TddPhase = 'setup' | 'test' | 'impl' | 'verify' | 'review';
|
||||
// Note: PlanTaskStatus and TddPhase are imported from types.ts
|
||||
|
||||
/**
|
||||
* Enhanced plan task with verification criteria, dependencies, and execution tracking.
|
||||
@@ -106,12 +106,7 @@ export interface CheckpointReview {
|
||||
}
|
||||
|
||||
// ========== Configuration Constants ==========
|
||||
|
||||
/**
|
||||
* Maximum number of todo items to track per session.
|
||||
* Older items are removed when this limit is reached.
|
||||
*/
|
||||
const MAX_TODO_ITEMS = 50;
|
||||
// Note: MAX_TODOS_PER_SESSION and MAX_LINE_BUFFER_SIZE are imported from config modules
|
||||
|
||||
/**
|
||||
* Todo items older than this duration (in milliseconds) will be auto-expired.
|
||||
@@ -165,11 +160,6 @@ const COMMON_COMPLETION_PHRASES = new Set([
|
||||
*/
|
||||
const MIN_RECOMMENDED_PHRASE_LENGTH = 6;
|
||||
|
||||
/**
|
||||
* 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)
|
||||
|
||||
@@ -904,9 +894,13 @@ export class RalphTracker extends EventEmitter {
|
||||
this._totalFilesModified = 0;
|
||||
this._totalTasksCompleted = 0;
|
||||
// Keep circuit breaker state on soft reset (it tracks across iterations)
|
||||
// Emit immediately on reset (no debounce)
|
||||
this.emit('loopUpdate', this.loopState);
|
||||
this.emit('todoUpdate', this.todos);
|
||||
// Emit on next tick to prevent listeners from modifying state during reset (non-reentrant)
|
||||
const loopState = this.loopState;
|
||||
const todos = this.todos;
|
||||
process.nextTick(() => {
|
||||
this.emit('loopUpdate', loopState);
|
||||
this.emit('todoUpdate', todos);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -934,9 +928,13 @@ export class RalphTracker extends EventEmitter {
|
||||
this._totalFilesModified = 0;
|
||||
this._totalTasksCompleted = 0;
|
||||
this._circuitBreaker = createInitialCircuitBreakerStatus();
|
||||
// Emit immediately on reset (no debounce)
|
||||
this.emit('loopUpdate', this.loopState);
|
||||
this.emit('todoUpdate', this.todos);
|
||||
// Emit on next tick to prevent listeners from modifying state during reset (non-reentrant)
|
||||
const loopState = this.loopState;
|
||||
const todos = this.todos;
|
||||
process.nextTick(() => {
|
||||
this.emit('loopUpdate', loopState);
|
||||
this.emit('todoUpdate', todos);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -2200,7 +2198,7 @@ export class RalphTracker extends EventEmitter {
|
||||
* - ID is generated from normalized content (stable hash)
|
||||
* - Priority is parsed from content (P0/P1/P2, Critical, High Priority, etc.)
|
||||
* - Existing item: Updates status and timestamp
|
||||
* - New item: Adds to map, evicts oldest if at MAX_TODO_ITEMS
|
||||
* - New item: Adds to map, evicts oldest if at MAX_TODOS_PER_SESSION
|
||||
*
|
||||
* @param content - Raw todo content text
|
||||
* @param status - Status to set
|
||||
@@ -2284,7 +2282,7 @@ export class RalphTracker extends EventEmitter {
|
||||
}
|
||||
|
||||
// Add new todo
|
||||
if (this._todos.size >= MAX_TODO_ITEMS) {
|
||||
if (this._todos.size >= MAX_TODOS_PER_SESSION) {
|
||||
// Remove oldest todo to make room
|
||||
const oldest = this.findOldestTodo();
|
||||
if (oldest) {
|
||||
@@ -2683,7 +2681,7 @@ export class RalphTracker extends EventEmitter {
|
||||
|
||||
/**
|
||||
* Find the todo item with the oldest detectedAt timestamp.
|
||||
* Used for LRU eviction when at MAX_TODO_ITEMS limit.
|
||||
* Used for LRU eviction when at MAX_TODOS_PER_SESSION limit.
|
||||
* @returns Oldest todo item, or undefined if map is empty
|
||||
*/
|
||||
private findOldestTodo(): RalphTodoItem | undefined {
|
||||
|
||||
@@ -35,6 +35,7 @@
|
||||
*/
|
||||
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { randomUUID } from 'node:crypto';
|
||||
import { Session } from './session.js';
|
||||
import { AiIdleChecker, type AiCheckResult, type AiCheckState } from './ai-idle-checker.js';
|
||||
import { AiPlanChecker, type AiPlanCheckResult } from './ai-plan-checker.js';
|
||||
@@ -716,6 +717,9 @@ export class RespawnController extends EventEmitter {
|
||||
/** Timestamp when plan check was started (to detect stale results) */
|
||||
private planCheckStartTime: number = 0;
|
||||
|
||||
/** Unique ID for current AI check request (to detect stale results) */
|
||||
private _currentAiCheckId: string | null = null;
|
||||
|
||||
/** Timer for /clear step fallback (sends /init if no prompt detected) */
|
||||
private clearFallbackTimer: NodeJS.Timeout | null = null;
|
||||
|
||||
@@ -2131,6 +2135,10 @@ export class RespawnController extends EventEmitter {
|
||||
this.logAction('ai-check', `Spawning AI idle checker (${reason})`);
|
||||
this.emit('aiCheckStarted');
|
||||
|
||||
// Generate unique ID for this check to detect stale results
|
||||
const checkId = randomUUID();
|
||||
this._currentAiCheckId = checkId;
|
||||
|
||||
// Get the terminal buffer for analysis
|
||||
const buffer = this.terminalBuffer.value;
|
||||
|
||||
@@ -2141,6 +2149,12 @@ export class RespawnController extends EventEmitter {
|
||||
return;
|
||||
}
|
||||
|
||||
// Validate this is the result for the current check (not a stale one)
|
||||
if (this._currentAiCheckId !== checkId) {
|
||||
this.log(`AI check result ignored (stale check ID: ${checkId.substring(0, 8)})`);
|
||||
return;
|
||||
}
|
||||
|
||||
if (result.verdict === 'IDLE') {
|
||||
// Cancel any pending confirmation timers - AI has spoken
|
||||
this.cancelTrackedTimer('completion-confirm', this.completionConfirmTimer, 'AI verdict: IDLE');
|
||||
@@ -2173,6 +2187,10 @@ export class RespawnController extends EventEmitter {
|
||||
this.startPreFilterTimer();
|
||||
}
|
||||
}).catch((err) => {
|
||||
// Validate this is the error for the current check
|
||||
if (this._currentAiCheckId !== checkId) {
|
||||
return; // Stale check, ignore error
|
||||
}
|
||||
if (this._state === 'ai_checking') {
|
||||
const errorMsg = err instanceof Error ? err.message : String(err);
|
||||
this.logAction('ai-check', `Failed: ${errorMsg.substring(0, 50)}`);
|
||||
|
||||
@@ -100,10 +100,12 @@ export class SessionManager extends EventEmitter {
|
||||
}
|
||||
|
||||
// Create a new lock promise that others will wait on
|
||||
let unlock: () => void;
|
||||
this._sessionCreationLock = new Promise<void>(resolve => {
|
||||
// Define unlock first to ensure it's always in scope before promise assignment
|
||||
let unlock!: () => void;
|
||||
const lockPromise = new Promise<void>(resolve => {
|
||||
unlock = resolve;
|
||||
});
|
||||
this._sessionCreationLock = lockPromise;
|
||||
|
||||
try {
|
||||
const config = this.store.getConfig();
|
||||
@@ -152,7 +154,7 @@ export class SessionManager extends EventEmitter {
|
||||
} finally {
|
||||
// Release the lock so other createSession calls can proceed
|
||||
this._sessionCreationLock = null;
|
||||
unlock!();
|
||||
unlock();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+30
-16
@@ -123,8 +123,9 @@ export function getAugmentedPath(): string {
|
||||
if (result && existsSync(result)) {
|
||||
claudeDir = dirname(result);
|
||||
}
|
||||
} catch {
|
||||
// not in PATH, check common locations
|
||||
} catch (err) {
|
||||
// Claude not in PATH, will check common locations
|
||||
console.warn('[Session] Claude not found via which command, checking common locations:', err instanceof Error ? err.message : err);
|
||||
}
|
||||
|
||||
// Fallback: check common installation directories
|
||||
@@ -1441,8 +1442,9 @@ export class Session extends EventEmitter {
|
||||
if (msg.type === 'result' && msg.total_cost_usd) {
|
||||
this._totalCost = msg.total_cost_usd;
|
||||
}
|
||||
} catch {
|
||||
// Not JSON, just regular output
|
||||
} catch (parseErr) {
|
||||
// Not JSON, just regular output - this is expected for non-JSON lines
|
||||
console.debug('[Session] Line not JSON (expected for text output):', parseErr instanceof Error ? parseErr.message : parseErr);
|
||||
this._textOutput.append(line + '\n');
|
||||
}
|
||||
} else if (trimmed) {
|
||||
@@ -1636,7 +1638,8 @@ 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 || this._isStopped) return;
|
||||
if (this._isStopped) return; // Early exit check
|
||||
if (!this._autoCompactEnabled || this._isCompacting || this._isClearing) return;
|
||||
|
||||
const totalTokens = this._totalInputTokens + this._totalOutputTokens;
|
||||
if (totalTokens >= this._autoCompactThreshold) {
|
||||
@@ -1645,10 +1648,14 @@ 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 || this._isStopped) return;
|
||||
// Check if session is still valid (not stopped) - must be first check
|
||||
if (this._isStopped) return;
|
||||
if (!this._isCompacting) return;
|
||||
|
||||
if (!this._isWorking) {
|
||||
// Re-check stopped state after async operation might have completed
|
||||
if (this._isStopped) return;
|
||||
|
||||
// Send /compact command with optional prompt
|
||||
const compactCmd = this._autoCompactPrompt
|
||||
? `/compact ${this._autoCompactPrompt}\r`
|
||||
@@ -1663,6 +1670,7 @@ export class Session extends EventEmitter {
|
||||
// Wait a moment then re-enable (longer than clear since compact takes time)
|
||||
if (!this._isStopped) {
|
||||
this._autoCompactTimer = setTimeout(() => {
|
||||
if (this._isStopped) return; // Check at callback start
|
||||
this._autoCompactTimer = null;
|
||||
this._isCompacting = false;
|
||||
}, 10000);
|
||||
@@ -1684,7 +1692,8 @@ export class Session extends EventEmitter {
|
||||
|
||||
// Check if we should auto-clear based on token threshold
|
||||
private checkAutoClear(): void {
|
||||
if (!this._autoClearEnabled || this._isClearing || this._isCompacting || this._isStopped) return;
|
||||
if (this._isStopped) return; // Early exit check
|
||||
if (!this._autoClearEnabled || this._isClearing || this._isCompacting) return;
|
||||
|
||||
const totalTokens = this._totalInputTokens + this._totalOutputTokens;
|
||||
if (totalTokens >= this._autoClearThreshold) {
|
||||
@@ -1693,10 +1702,14 @@ 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 || this._isStopped) return;
|
||||
// Check if session is still valid (not stopped) - must be first check
|
||||
if (this._isStopped) return;
|
||||
if (!this._isClearing) return;
|
||||
|
||||
if (!this._isWorking) {
|
||||
// Re-check stopped state after async operation might have completed
|
||||
if (this._isStopped) return;
|
||||
|
||||
// Send /clear command
|
||||
this.writeViaScreen('/clear\r');
|
||||
// Reset token counts
|
||||
@@ -1707,6 +1720,7 @@ export class Session extends EventEmitter {
|
||||
// Wait a moment then re-enable
|
||||
if (!this._isStopped) {
|
||||
this._autoClearTimer = setTimeout(() => {
|
||||
if (this._isStopped) return; // Check at callback start
|
||||
this._autoClearTimer = null;
|
||||
this._isClearing = false;
|
||||
}, 5000);
|
||||
@@ -1921,8 +1935,8 @@ export class Session extends EventEmitter {
|
||||
// First try graceful SIGTERM
|
||||
try {
|
||||
this.ptyProcess.kill();
|
||||
} catch {
|
||||
// Process may already be dead
|
||||
} catch (err) {
|
||||
console.warn('[Session] Failed to send SIGTERM to PTY process (may already be dead):', err);
|
||||
}
|
||||
|
||||
// Give it a moment to terminate gracefully
|
||||
@@ -1933,8 +1947,8 @@ export class Session extends EventEmitter {
|
||||
if (pid) {
|
||||
process.kill(pid, 'SIGKILL');
|
||||
}
|
||||
} catch {
|
||||
// Process already terminated
|
||||
} catch (err) {
|
||||
console.warn('[Session] Failed to send SIGKILL to process (already terminated):', err);
|
||||
}
|
||||
|
||||
// Also try to kill any child processes in the process group
|
||||
@@ -1942,8 +1956,8 @@ export class Session extends EventEmitter {
|
||||
if (pid) {
|
||||
process.kill(-pid, 'SIGKILL');
|
||||
}
|
||||
} catch {
|
||||
// Process group may not exist or already terminated
|
||||
} catch (err) {
|
||||
console.warn('[Session] Failed to send SIGKILL to process group (may not exist):', err);
|
||||
}
|
||||
|
||||
this.ptyProcess = null;
|
||||
|
||||
@@ -1,700 +0,0 @@
|
||||
/**
|
||||
* @fileoverview Type definitions for the spawn1337 Autonomous Agent Protocol.
|
||||
*
|
||||
* Defines all types for the agent spawning system including:
|
||||
* - Task specifications (what the parent writes)
|
||||
* - Agent progress reporting
|
||||
* - Result delivery format
|
||||
* - Bidirectional messaging
|
||||
* - Orchestrator state tracking
|
||||
*
|
||||
* Also includes a simple YAML frontmatter parser and factory functions.
|
||||
*
|
||||
* @module spawn-types
|
||||
*/
|
||||
|
||||
// ========== Agent Task Specification ==========
|
||||
|
||||
/** Priority levels for spawn tasks */
|
||||
export type SpawnPriority = 'low' | 'normal' | 'high' | 'critical';
|
||||
|
||||
/** How results should be delivered */
|
||||
export type SpawnResultDelivery = 'file' | 'notify' | 'both';
|
||||
|
||||
/** Agent execution status */
|
||||
export type SpawnStatus = 'queued' | 'initializing' | 'running' | 'completing' | 'completed' | 'failed' | 'timeout' | 'cancelled';
|
||||
|
||||
/**
|
||||
* Task specification parsed from the .md file's YAML frontmatter.
|
||||
* This is the contract for what the parent LLM writes.
|
||||
*/
|
||||
export interface SpawnTaskSpec {
|
||||
// === Identity ===
|
||||
/** Unique agent identifier (auto-generated if not provided) */
|
||||
agentId: string;
|
||||
/** Human-readable name for this agent */
|
||||
name: string;
|
||||
/** Task type/category */
|
||||
type: 'explore' | 'implement' | 'test' | 'review' | 'refactor' | 'research' | 'generate' | 'fix' | 'general';
|
||||
|
||||
// === Scheduling ===
|
||||
/** Priority for queue ordering */
|
||||
priority: SpawnPriority;
|
||||
/** Dependencies - other agentIds that must complete first */
|
||||
dependsOn?: string[];
|
||||
|
||||
// === Environment ===
|
||||
/** Working directory (relative to parent, or absolute) */
|
||||
workingDir?: string;
|
||||
/** Files to copy/symlink into agent workspace as context */
|
||||
contextFiles?: string[];
|
||||
/** Whether the agent can modify files in the parent's directory */
|
||||
canModifyParentFiles: boolean;
|
||||
/** Additional environment variables for the agent */
|
||||
env?: Record<string, string>;
|
||||
|
||||
// === Resource Governance ===
|
||||
/** Maximum token budget (input + output combined) */
|
||||
maxTokens?: number;
|
||||
/** Maximum cost in USD */
|
||||
maxCost?: number;
|
||||
/** Timeout in minutes */
|
||||
timeoutMinutes: number;
|
||||
|
||||
// === Communication ===
|
||||
/** How to deliver results */
|
||||
resultDelivery: SpawnResultDelivery;
|
||||
/** Completion phrase for RalphTracker (auto-generated if not set) */
|
||||
completionPhrase: string;
|
||||
/** How often the agent should report progress (seconds, 0 = no progress) */
|
||||
progressIntervalSeconds: number;
|
||||
|
||||
// === Output ===
|
||||
/** Expected output format */
|
||||
outputFormat: 'markdown' | 'json' | 'code' | 'structured' | 'freeform';
|
||||
/** Success criteria (included in agent's CLAUDE.md) */
|
||||
successCriteria: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* The full parsed task (spec + instructions body).
|
||||
*/
|
||||
export interface SpawnTask {
|
||||
spec: SpawnTaskSpec;
|
||||
/** The markdown body - actual instructions for the agent */
|
||||
instructions: string;
|
||||
/** Source file path */
|
||||
sourceFile: string;
|
||||
/** Parent session ID that requested this spawn */
|
||||
parentSessionId: string;
|
||||
/** Spawn depth (0 = direct child of user) */
|
||||
depth: number;
|
||||
}
|
||||
|
||||
// ========== Agent Communication ==========
|
||||
|
||||
/**
|
||||
* Progress report written by agent to spawn-comms/progress.json
|
||||
*/
|
||||
export interface AgentProgress {
|
||||
/** Current phase/step description */
|
||||
phase: string;
|
||||
/** Completion percentage (0-100) */
|
||||
percentComplete: number;
|
||||
/** What the agent is currently doing */
|
||||
currentAction: string;
|
||||
/** Todos/subtasks the agent is tracking */
|
||||
subtasks?: Array<{
|
||||
description: string;
|
||||
status: 'pending' | 'in_progress' | 'completed';
|
||||
}>;
|
||||
/** Timestamp of last update */
|
||||
updatedAt: number;
|
||||
/** Files modified so far */
|
||||
filesModified: string[];
|
||||
/** Tokens used so far */
|
||||
tokensUsed: number;
|
||||
/** Cost so far */
|
||||
costSoFar: number;
|
||||
}
|
||||
|
||||
/**
|
||||
* Result delivered by agent on completion.
|
||||
* Written to spawn-comms/result.md as YAML frontmatter + body.
|
||||
*/
|
||||
export interface SpawnResult {
|
||||
// === Status ===
|
||||
/** Final execution status */
|
||||
status: 'completed' | 'failed' | 'timeout' | 'cancelled';
|
||||
/** Error message if failed */
|
||||
error?: string;
|
||||
|
||||
// === Metrics ===
|
||||
/** Total execution duration in ms */
|
||||
durationMs: number;
|
||||
/** Token usage breakdown */
|
||||
tokens: {
|
||||
input: number;
|
||||
output: number;
|
||||
total: number;
|
||||
};
|
||||
/** Total cost in USD */
|
||||
cost: number;
|
||||
|
||||
// === Output ===
|
||||
/** Executive summary (1-3 sentences) */
|
||||
summary: string;
|
||||
/** Full structured output */
|
||||
output: string;
|
||||
/** Files modified or created */
|
||||
filesChanged: Array<{
|
||||
path: string;
|
||||
action: 'created' | 'modified' | 'deleted';
|
||||
summary?: string;
|
||||
}>;
|
||||
/** Any artifacts produced (data files, diagrams, etc.) */
|
||||
artifacts?: Array<{
|
||||
name: string;
|
||||
path: string;
|
||||
type: string;
|
||||
description: string;
|
||||
}>;
|
||||
|
||||
// === Metadata ===
|
||||
/** Agent ID */
|
||||
agentId: string;
|
||||
/** Completion timestamp */
|
||||
completedAt: number;
|
||||
/** Number of respawn cycles if agent used Ralph loop */
|
||||
cycleCount?: number;
|
||||
}
|
||||
|
||||
/**
|
||||
* Message in the bidirectional communication channel.
|
||||
* Written to spawn-comms/messages/NNN-{sender}.md
|
||||
*/
|
||||
export interface SpawnMessage {
|
||||
/** Sequential message number */
|
||||
sequence: number;
|
||||
/** Who sent it */
|
||||
sender: 'parent' | 'agent';
|
||||
/** Message content (markdown) */
|
||||
content: string;
|
||||
/** Timestamp */
|
||||
sentAt: number;
|
||||
/** Whether it's been read by the recipient */
|
||||
read: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* Status report for UI display and API responses.
|
||||
*/
|
||||
export interface AgentStatusReport {
|
||||
agentId: string;
|
||||
name: string;
|
||||
type: string;
|
||||
status: SpawnStatus;
|
||||
priority: SpawnPriority;
|
||||
parentSessionId: string;
|
||||
childSessionId: string | null;
|
||||
depth: number;
|
||||
startedAt: number | null;
|
||||
elapsedMs: number;
|
||||
progress: AgentProgress | null;
|
||||
tokensUsed: number;
|
||||
costSoFar: number;
|
||||
tokenBudget: number | null;
|
||||
costBudget: number | null;
|
||||
timeoutMinutes: number;
|
||||
timeRemainingMs: number;
|
||||
completionPhrase: string;
|
||||
dependsOn: string[];
|
||||
dependencyStatus: 'waiting' | 'ready' | 'n/a';
|
||||
}
|
||||
|
||||
// ========== Tracker State (for SpawnOrchestrator) ==========
|
||||
|
||||
export interface SpawnTrackerState {
|
||||
enabled: boolean;
|
||||
activeCount: number;
|
||||
queuedCount: number;
|
||||
totalSpawned: number;
|
||||
totalCompleted: number;
|
||||
totalFailed: number;
|
||||
maxDepthReached: number;
|
||||
agents: AgentStatusReport[];
|
||||
}
|
||||
|
||||
// ========== Orchestrator Configuration ==========
|
||||
|
||||
export interface SpawnOrchestratorConfig {
|
||||
/** Max concurrent agent sessions (default: 5) */
|
||||
maxConcurrentAgents: number;
|
||||
/** Base directory for agent cases (default: ~/claudeman-cases/) */
|
||||
casesDir: string;
|
||||
/** Default timeout in minutes (default: 30) */
|
||||
defaultTimeoutMinutes: number;
|
||||
/** Max timeout allowed in minutes (default: 120) */
|
||||
maxTimeoutMinutes: number;
|
||||
/** Max agent tree depth (prevent infinite recursion) (default: 3) */
|
||||
maxSpawnDepth: number;
|
||||
/** Progress poll interval in ms (default: 5000) */
|
||||
progressPollIntervalMs: number;
|
||||
}
|
||||
|
||||
// ========== Agent Context (internal orchestrator state) ==========
|
||||
|
||||
export interface AgentContext {
|
||||
/** The parsed task specification */
|
||||
task: SpawnTask;
|
||||
/** The spawned session ID (set after session creation) */
|
||||
sessionId: string | null;
|
||||
/** Resolved working directory for this agent */
|
||||
workingDir: string;
|
||||
/** Communication directory path */
|
||||
commsDir: string;
|
||||
/** Parent session ID */
|
||||
parentSessionId: string;
|
||||
/** Depth in the spawn tree (0 = direct child of user session) */
|
||||
depth: number;
|
||||
/** Timeout timer handle */
|
||||
timeoutTimer: NodeJS.Timeout | null;
|
||||
/** Warning timer handle (fires at 90% of timeout) */
|
||||
warningTimer: NodeJS.Timeout | null;
|
||||
/** Progress poll timer handle */
|
||||
progressTimer: NodeJS.Timeout | null;
|
||||
/** Current status */
|
||||
status: SpawnStatus;
|
||||
/** When the agent started working */
|
||||
startedAt: number | null;
|
||||
/** Token budget remaining (null = unlimited) */
|
||||
tokenBudget: number | null;
|
||||
/** Cost budget remaining (null = unlimited) */
|
||||
costBudget: number | null;
|
||||
}
|
||||
|
||||
// ========== Persisted State ==========
|
||||
|
||||
export interface SpawnPersistedState {
|
||||
config: SpawnOrchestratorConfig;
|
||||
agents: Record<string, {
|
||||
agentId: string;
|
||||
status: SpawnStatus;
|
||||
parentSessionId: string;
|
||||
childSessionId: string | null;
|
||||
depth: number;
|
||||
startedAt: number | null;
|
||||
commsDir: string;
|
||||
workingDir: string;
|
||||
completionPhrase: string;
|
||||
timeoutMinutes: number;
|
||||
}>;
|
||||
}
|
||||
|
||||
// ========== Constants ==========
|
||||
|
||||
/** Maximum concurrent agents */
|
||||
export const MAX_CONCURRENT_AGENTS = 5;
|
||||
/** Default timeout in minutes */
|
||||
export const DEFAULT_TIMEOUT_MINUTES = 30;
|
||||
/** Maximum timeout in minutes */
|
||||
export const MAX_TIMEOUT_MINUTES = 120;
|
||||
/** Maximum spawn depth */
|
||||
export const MAX_SPAWN_DEPTH = 3;
|
||||
/** Progress poll interval in ms */
|
||||
export const PROGRESS_POLL_INTERVAL_MS = 5000;
|
||||
/** Maximum task file size (2MB) */
|
||||
export const MAX_TASK_FILE_SIZE = 2 * 1024 * 1024;
|
||||
/** Maximum context file size (100KB each) */
|
||||
export const MAX_CONTEXT_FILE_SIZE = 100 * 1024;
|
||||
/** Maximum number of context files */
|
||||
export const MAX_CONTEXT_FILES = 20;
|
||||
/** Maximum queue length */
|
||||
export const MAX_QUEUE_LENGTH = 50;
|
||||
/** Budget warning threshold (80%) */
|
||||
export const BUDGET_WARNING_THRESHOLD = 0.8;
|
||||
/** Budget grace period in seconds */
|
||||
export const BUDGET_GRACE_PERIOD_S = 60;
|
||||
/** Agent name max length */
|
||||
export const AGENT_NAME_MAX_LENGTH = 64;
|
||||
/** Message max size (50KB) */
|
||||
export const MESSAGE_MAX_SIZE = 50 * 1024;
|
||||
/** Max messages per channel */
|
||||
export const MAX_MESSAGES_PER_CHANNEL = 100;
|
||||
/** Max tracked agents (LRU) */
|
||||
export const MAX_TRACKED_AGENTS = 200;
|
||||
|
||||
// ========== Factory Functions ==========
|
||||
|
||||
/**
|
||||
* Creates a default SpawnTaskSpec with sensible defaults.
|
||||
*/
|
||||
export function createDefaultSpawnTaskSpec(agentId: string): SpawnTaskSpec {
|
||||
return {
|
||||
agentId,
|
||||
name: agentId,
|
||||
type: 'general',
|
||||
priority: 'normal',
|
||||
canModifyParentFiles: false,
|
||||
timeoutMinutes: DEFAULT_TIMEOUT_MINUTES,
|
||||
resultDelivery: 'both',
|
||||
completionPhrase: `AGENT_${agentId.toUpperCase().replace(/[^A-Z0-9]/g, '_')}_DONE`,
|
||||
progressIntervalSeconds: 30,
|
||||
outputFormat: 'markdown',
|
||||
successCriteria: '',
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates an empty AgentProgress object.
|
||||
*/
|
||||
export function createEmptyAgentProgress(): AgentProgress {
|
||||
return {
|
||||
phase: 'initializing',
|
||||
percentComplete: 0,
|
||||
currentAction: '',
|
||||
subtasks: [],
|
||||
updatedAt: Date.now(),
|
||||
filesModified: [],
|
||||
tokensUsed: 0,
|
||||
costSoFar: 0,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates initial SpawnTrackerState.
|
||||
*/
|
||||
export function createInitialSpawnTrackerState(): SpawnTrackerState {
|
||||
return {
|
||||
enabled: false,
|
||||
activeCount: 0,
|
||||
queuedCount: 0,
|
||||
totalSpawned: 0,
|
||||
totalCompleted: 0,
|
||||
totalFailed: 0,
|
||||
maxDepthReached: 0,
|
||||
agents: [],
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates default orchestrator config.
|
||||
*/
|
||||
export function createDefaultOrchestratorConfig(): SpawnOrchestratorConfig {
|
||||
const homeDir = process.env.HOME || process.env.USERPROFILE || '/tmp';
|
||||
return {
|
||||
maxConcurrentAgents: MAX_CONCURRENT_AGENTS,
|
||||
casesDir: `${homeDir}/claudeman-cases`,
|
||||
defaultTimeoutMinutes: DEFAULT_TIMEOUT_MINUTES,
|
||||
maxTimeoutMinutes: MAX_TIMEOUT_MINUTES,
|
||||
maxSpawnDepth: MAX_SPAWN_DEPTH,
|
||||
progressPollIntervalMs: PROGRESS_POLL_INTERVAL_MS,
|
||||
};
|
||||
}
|
||||
|
||||
// ========== YAML Frontmatter Parser ==========
|
||||
|
||||
/**
|
||||
* Simple YAML frontmatter parser for task spec files.
|
||||
* Handles: strings, numbers, booleans, arrays (block and inline), one-level nested objects.
|
||||
* Does NOT handle: multi-line strings, anchors, aliases, complex nesting.
|
||||
*/
|
||||
export function parseYamlFrontmatter(content: string): { frontmatter: Record<string, unknown>; body: string } | null {
|
||||
const lines = content.split('\n');
|
||||
|
||||
// Must start with ---
|
||||
if (lines[0].trim() !== '---') return null;
|
||||
|
||||
let endIndex = -1;
|
||||
for (let i = 1; i < lines.length; i++) {
|
||||
if (lines[i].trim() === '---') {
|
||||
endIndex = i;
|
||||
break;
|
||||
}
|
||||
}
|
||||
if (endIndex === -1) return null;
|
||||
|
||||
const yamlLines = lines.slice(1, endIndex);
|
||||
const body = lines.slice(endIndex + 1).join('\n').trim();
|
||||
const frontmatter: Record<string, unknown> = {};
|
||||
|
||||
let currentKey: string | null = null;
|
||||
let currentArray: unknown[] | null = null;
|
||||
let currentObject: Record<string, unknown> | null = null;
|
||||
|
||||
for (const line of yamlLines) {
|
||||
// Skip empty lines and comments
|
||||
if (!line.trim() || line.trim().startsWith('#')) continue;
|
||||
|
||||
const indent = line.length - line.trimStart().length;
|
||||
|
||||
// Array item (indented with -)
|
||||
if (indent >= 2 && line.trim().startsWith('- ')) {
|
||||
const value = line.trim().slice(2).trim();
|
||||
if (currentArray && currentKey) {
|
||||
// Check if it's a key: value pair within an array item
|
||||
const kvMatch = value.match(/^(\w+):\s*(.+)$/);
|
||||
if (kvMatch && currentArray.length > 0 && typeof currentArray[currentArray.length - 1] === 'object') {
|
||||
// Add to existing object in array
|
||||
(currentArray[currentArray.length - 1] as Record<string, unknown>)[kvMatch[1]] = parseYamlValue(kvMatch[2]);
|
||||
} else if (kvMatch && value.includes(':')) {
|
||||
// New object in array
|
||||
const obj: Record<string, unknown> = {};
|
||||
obj[kvMatch[1]] = parseYamlValue(kvMatch[2]);
|
||||
currentArray.push(obj);
|
||||
} else {
|
||||
currentArray.push(parseYamlValue(value));
|
||||
}
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
// Indented key: value (nested object or additional array object fields)
|
||||
if (indent >= 2 && currentKey && !line.trim().startsWith('- ')) {
|
||||
const kvMatch = line.trim().match(/^(\w+):\s*(.*)$/);
|
||||
if (kvMatch) {
|
||||
// Switch from array mode to object mode if array is empty
|
||||
if (currentArray && currentArray.length === 0 && !currentObject) {
|
||||
currentArray = null;
|
||||
currentObject = {};
|
||||
}
|
||||
if (!currentObject) {
|
||||
currentObject = {};
|
||||
}
|
||||
currentObject[kvMatch[1]] = parseYamlValue(kvMatch[2]);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
// Top-level key: value
|
||||
const topMatch = line.match(/^(\w+):\s*(.*)$/);
|
||||
if (topMatch) {
|
||||
// Save previous array/object
|
||||
if (currentKey && currentArray) {
|
||||
frontmatter[currentKey] = currentArray;
|
||||
} else if (currentKey && currentObject) {
|
||||
frontmatter[currentKey] = currentObject;
|
||||
}
|
||||
|
||||
currentKey = topMatch[1];
|
||||
const value = topMatch[2].trim();
|
||||
|
||||
if (value === '' || value === '[]') {
|
||||
// Could be start of array or object
|
||||
currentArray = [];
|
||||
currentObject = null;
|
||||
if (value === '[]') {
|
||||
frontmatter[currentKey] = [];
|
||||
currentKey = null;
|
||||
currentArray = null;
|
||||
}
|
||||
} else if (value.startsWith('[') && value.endsWith(']')) {
|
||||
// Inline array
|
||||
const items = value.slice(1, -1).split(',').map(s => parseYamlValue(s.trim()));
|
||||
frontmatter[currentKey] = items;
|
||||
currentKey = null;
|
||||
currentArray = null;
|
||||
currentObject = null;
|
||||
} else {
|
||||
frontmatter[currentKey] = parseYamlValue(value);
|
||||
currentKey = null;
|
||||
currentArray = null;
|
||||
currentObject = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Save last pending array/object
|
||||
if (currentKey && currentArray) {
|
||||
frontmatter[currentKey] = currentArray;
|
||||
} else if (currentKey && currentObject) {
|
||||
frontmatter[currentKey] = currentObject;
|
||||
}
|
||||
|
||||
return { frontmatter, body };
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a single YAML value string into the appropriate JS type.
|
||||
*/
|
||||
function parseYamlValue(value: string): unknown {
|
||||
if (!value || value === '~' || value === 'null') return null;
|
||||
|
||||
// Remove surrounding quotes
|
||||
if ((value.startsWith('"') && value.endsWith('"')) ||
|
||||
(value.startsWith("'") && value.endsWith("'"))) {
|
||||
return value.slice(1, -1);
|
||||
}
|
||||
|
||||
// Booleans
|
||||
if (value === 'true' || value === 'yes') return true;
|
||||
if (value === 'false' || value === 'no') return false;
|
||||
|
||||
// Numbers
|
||||
const num = Number(value);
|
||||
if (!isNaN(num) && value !== '') return num;
|
||||
|
||||
return value;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a task spec file content into a SpawnTaskSpec.
|
||||
* Returns null if parsing fails.
|
||||
*/
|
||||
export function parseTaskSpecFile(content: string, fallbackAgentId: string): { spec: SpawnTaskSpec; instructions: string } | null {
|
||||
const parsed = parseYamlFrontmatter(content);
|
||||
if (!parsed) return null;
|
||||
|
||||
const { frontmatter, body } = parsed;
|
||||
const defaults = createDefaultSpawnTaskSpec(fallbackAgentId);
|
||||
|
||||
const spec: SpawnTaskSpec = {
|
||||
agentId: String(frontmatter.agentId ?? defaults.agentId),
|
||||
name: String(frontmatter.name ?? defaults.name),
|
||||
type: validateType(frontmatter.type) ?? defaults.type,
|
||||
priority: validatePriority(frontmatter.priority) ?? defaults.priority,
|
||||
dependsOn: Array.isArray(frontmatter.dependsOn) ? frontmatter.dependsOn.map(String) : undefined,
|
||||
workingDir: frontmatter.workingDir != null ? String(frontmatter.workingDir) : undefined,
|
||||
contextFiles: Array.isArray(frontmatter.contextFiles) ? frontmatter.contextFiles.map(String) : undefined,
|
||||
canModifyParentFiles: Boolean(frontmatter.canModifyParentFiles ?? defaults.canModifyParentFiles),
|
||||
env: isStringRecord(frontmatter.env) ? frontmatter.env : undefined,
|
||||
maxTokens: typeof frontmatter.maxTokens === 'number' ? frontmatter.maxTokens : undefined,
|
||||
maxCost: typeof frontmatter.maxCost === 'number' ? frontmatter.maxCost : undefined,
|
||||
timeoutMinutes: typeof frontmatter.timeoutMinutes === 'number' ? frontmatter.timeoutMinutes : defaults.timeoutMinutes,
|
||||
resultDelivery: validateResultDelivery(frontmatter.resultDelivery) ?? defaults.resultDelivery,
|
||||
completionPhrase: String(frontmatter.completionPhrase ?? defaults.completionPhrase),
|
||||
progressIntervalSeconds: typeof frontmatter.progressIntervalSeconds === 'number' ? frontmatter.progressIntervalSeconds : defaults.progressIntervalSeconds,
|
||||
outputFormat: validateOutputFormat(frontmatter.outputFormat) ?? defaults.outputFormat,
|
||||
successCriteria: String(frontmatter.successCriteria ?? defaults.successCriteria),
|
||||
};
|
||||
|
||||
// Validate agent name length
|
||||
if (spec.name.length > AGENT_NAME_MAX_LENGTH) {
|
||||
spec.name = spec.name.slice(0, AGENT_NAME_MAX_LENGTH);
|
||||
}
|
||||
|
||||
return { spec, instructions: body };
|
||||
}
|
||||
|
||||
// ========== Validation Helpers ==========
|
||||
|
||||
const VALID_TYPES = ['explore', 'implement', 'test', 'review', 'refactor', 'research', 'generate', 'fix', 'general'] as const;
|
||||
const VALID_PRIORITIES = ['low', 'normal', 'high', 'critical'] as const;
|
||||
const VALID_RESULT_DELIVERIES = ['file', 'notify', 'both'] as const;
|
||||
const VALID_OUTPUT_FORMATS = ['markdown', 'json', 'code', 'structured', 'freeform'] as const;
|
||||
|
||||
function validateType(value: unknown): SpawnTaskSpec['type'] | null {
|
||||
return VALID_TYPES.includes(value as typeof VALID_TYPES[number]) ? value as SpawnTaskSpec['type'] : null;
|
||||
}
|
||||
|
||||
function validatePriority(value: unknown): SpawnPriority | null {
|
||||
return VALID_PRIORITIES.includes(value as typeof VALID_PRIORITIES[number]) ? value as SpawnPriority : null;
|
||||
}
|
||||
|
||||
function validateResultDelivery(value: unknown): SpawnResultDelivery | null {
|
||||
return VALID_RESULT_DELIVERIES.includes(value as typeof VALID_RESULT_DELIVERIES[number]) ? value as SpawnResultDelivery : null;
|
||||
}
|
||||
|
||||
function validateOutputFormat(value: unknown): SpawnTaskSpec['outputFormat'] | null {
|
||||
return VALID_OUTPUT_FORMATS.includes(value as typeof VALID_OUTPUT_FORMATS[number]) ? value as SpawnTaskSpec['outputFormat'] : null;
|
||||
}
|
||||
|
||||
function isStringRecord(value: unknown): value is Record<string, string> {
|
||||
if (!value || typeof value !== 'object') return false;
|
||||
return Object.values(value).every(v => typeof v === 'string');
|
||||
}
|
||||
|
||||
/**
|
||||
* Serialize a SpawnResult to YAML frontmatter + markdown body.
|
||||
*/
|
||||
export function serializeSpawnResult(result: SpawnResult): string {
|
||||
const lines: string[] = ['---'];
|
||||
lines.push(`status: ${result.status}`);
|
||||
if (result.error) lines.push(`error: "${result.error.replace(/"/g, '\\"')}"`);
|
||||
lines.push(`summary: "${result.summary.replace(/"/g, '\\"')}"`);
|
||||
lines.push(`durationMs: ${result.durationMs}`);
|
||||
lines.push(`cost: ${result.cost}`);
|
||||
lines.push(`agentId: ${result.agentId}`);
|
||||
lines.push(`completedAt: ${result.completedAt}`);
|
||||
if (result.cycleCount != null) lines.push(`cycleCount: ${result.cycleCount}`);
|
||||
|
||||
if (result.filesChanged.length > 0) {
|
||||
lines.push('filesChanged:');
|
||||
for (const file of result.filesChanged) {
|
||||
lines.push(` - path: ${file.path}`);
|
||||
lines.push(` action: ${file.action}`);
|
||||
if (file.summary) lines.push(` summary: "${file.summary.replace(/"/g, '\\"')}"`);
|
||||
}
|
||||
} else {
|
||||
lines.push('filesChanged: []');
|
||||
}
|
||||
|
||||
if (result.artifacts && result.artifacts.length > 0) {
|
||||
lines.push('artifacts:');
|
||||
for (const artifact of result.artifacts) {
|
||||
lines.push(` - name: ${artifact.name}`);
|
||||
lines.push(` path: ${artifact.path}`);
|
||||
lines.push(` type: ${artifact.type}`);
|
||||
lines.push(` description: "${artifact.description.replace(/"/g, '\\"')}"`);
|
||||
}
|
||||
}
|
||||
|
||||
lines.push('---');
|
||||
lines.push('');
|
||||
lines.push(result.output);
|
||||
|
||||
return lines.join('\n');
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a result.md file into a SpawnResult.
|
||||
*/
|
||||
export function parseSpawnResult(content: string, agentId: string, fallbackDurationMs: number): SpawnResult | null {
|
||||
const parsed = parseYamlFrontmatter(content);
|
||||
if (!parsed) return null;
|
||||
|
||||
const { frontmatter, body } = parsed;
|
||||
|
||||
return {
|
||||
status: (['completed', 'failed', 'timeout', 'cancelled'].includes(String(frontmatter.status))
|
||||
? String(frontmatter.status) as SpawnResult['status']
|
||||
: 'completed'),
|
||||
error: frontmatter.error != null ? String(frontmatter.error) : undefined,
|
||||
durationMs: typeof frontmatter.durationMs === 'number' ? frontmatter.durationMs : fallbackDurationMs,
|
||||
tokens: {
|
||||
input: 0,
|
||||
output: 0,
|
||||
total: 0,
|
||||
},
|
||||
cost: typeof frontmatter.cost === 'number' ? frontmatter.cost : 0,
|
||||
summary: String(frontmatter.summary ?? 'No summary provided'),
|
||||
output: body,
|
||||
filesChanged: Array.isArray(frontmatter.filesChanged)
|
||||
? frontmatter.filesChanged.map((f: unknown) => {
|
||||
if (typeof f === 'object' && f !== null) {
|
||||
const obj = f as Record<string, unknown>;
|
||||
return {
|
||||
path: String(obj.path ?? ''),
|
||||
action: (['created', 'modified', 'deleted'].includes(String(obj.action)) ? String(obj.action) : 'modified') as 'created' | 'modified' | 'deleted',
|
||||
summary: obj.summary != null ? String(obj.summary) : undefined,
|
||||
};
|
||||
}
|
||||
return { path: String(f), action: 'modified' as const };
|
||||
})
|
||||
: [],
|
||||
artifacts: Array.isArray(frontmatter.artifacts)
|
||||
? frontmatter.artifacts.map((a: unknown) => {
|
||||
const obj = a as Record<string, unknown>;
|
||||
return {
|
||||
name: String(obj.name ?? ''),
|
||||
path: String(obj.path ?? ''),
|
||||
type: String(obj.type ?? 'unknown'),
|
||||
description: String(obj.description ?? ''),
|
||||
};
|
||||
})
|
||||
: undefined,
|
||||
agentId,
|
||||
completedAt: typeof frontmatter.completedAt === 'number' ? frontmatter.completedAt : Date.now(),
|
||||
cycleCount: typeof frontmatter.cycleCount === 'number' ? frontmatter.cycleCount : undefined,
|
||||
};
|
||||
}
|
||||
+8
-4
@@ -166,9 +166,9 @@ export class StateStore {
|
||||
JSON.parse(currentContent); // Validate
|
||||
writeFileSync(backupPath, currentContent, 'utf-8');
|
||||
}
|
||||
} catch {
|
||||
} catch (err) {
|
||||
// Backup failed - current file may be corrupt, continue with write
|
||||
console.warn('[StateStore] Could not create backup (current file may be corrupt)');
|
||||
console.warn('[StateStore] Could not create backup (current file may be corrupt):', err);
|
||||
}
|
||||
|
||||
// Step 3: Atomic write: write to temp file, then rename
|
||||
@@ -191,7 +191,9 @@ export class StateStore {
|
||||
if (existsSync(tempPath)) {
|
||||
unlinkSync(tempPath);
|
||||
}
|
||||
} catch { /* ignore cleanup errors */ }
|
||||
} catch (cleanupErr) {
|
||||
console.warn('[StateStore] Failed to cleanup temp file during save error:', cleanupErr);
|
||||
}
|
||||
|
||||
// Check circuit breaker threshold
|
||||
if (this.consecutiveSaveFailures >= MAX_CONSECUTIVE_FAILURES) {
|
||||
@@ -604,7 +606,9 @@ export class StateStore {
|
||||
if (existsSync(tempPath)) {
|
||||
unlinkSync(tempPath);
|
||||
}
|
||||
} catch { /* ignore cleanup errors */ }
|
||||
} catch (cleanupErr) {
|
||||
console.warn('[StateStore] Failed to cleanup temp file during Ralph state save error:', cleanupErr);
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
+54
-49
@@ -5,13 +5,14 @@
|
||||
* and emits structured events for tool calls, progress, and messages.
|
||||
*/
|
||||
|
||||
import { EventEmitter } from 'events';
|
||||
import { watch, statSync, readdirSync, existsSync, readFileSync, FSWatcher } from 'fs';
|
||||
import { createReadStream } from 'fs';
|
||||
import { createInterface } from 'readline';
|
||||
import { homedir } from 'os';
|
||||
import { join, basename } from 'path';
|
||||
import { execSync } from 'child_process';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { watch, statSync, readdirSync, existsSync, readFileSync, FSWatcher } from 'node:fs';
|
||||
import { createReadStream } from 'node:fs';
|
||||
import { createInterface } from 'node:readline';
|
||||
import { homedir } from 'node:os';
|
||||
import { join, basename } from 'node:path';
|
||||
import { execSync } from 'node:child_process';
|
||||
import { PENDING_TOOL_CALL_TTL_MS } from './config/map-limits.js';
|
||||
|
||||
// ========== Types ==========
|
||||
|
||||
@@ -160,8 +161,9 @@ export class SubagentWatcher extends EventEmitter {
|
||||
private livenessInterval: NodeJS.Timeout | null = null;
|
||||
private _isRunning = false;
|
||||
private knownSubagentDirs = new Set<string>();
|
||||
// Map of agentId -> Map of toolUseId -> tool name (for linking tool_result to tool_call)
|
||||
private pendingToolCalls = new Map<string, Map<string, string>>();
|
||||
// Map of agentId -> Map of toolUseId -> { toolName, timestamp } (for linking tool_result to tool_call)
|
||||
// Includes timestamp for TTL-based cleanup of orphaned entries
|
||||
private pendingToolCalls = new Map<string, Map<string, { toolName: string; timestamp: number }>>();
|
||||
|
||||
constructor() {
|
||||
super();
|
||||
@@ -312,6 +314,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
* - Completed agents: removed after STALE_COMPLETED_MAX_AGE_MS (1 hour)
|
||||
* - Idle agents: removed after STALE_IDLE_MAX_AGE_MS (4 hours)
|
||||
* Also enforces MAX_TRACKED_AGENTS limit with LRU eviction.
|
||||
* Also cleans up orphaned pending tool calls older than PENDING_TOOL_CALL_TTL_MS.
|
||||
*/
|
||||
private cleanupStaleAgents(): void {
|
||||
const now = Date.now();
|
||||
@@ -348,6 +351,24 @@ export class SubagentWatcher extends EventEmitter {
|
||||
for (const agentId of agentsToDelete) {
|
||||
this.removeAgent(agentId);
|
||||
}
|
||||
|
||||
// Clean up orphaned pending tool calls (older than TTL)
|
||||
// These can accumulate if tool_result is never received (e.g., agent crashed)
|
||||
for (const [agentId, agentCalls] of this.pendingToolCalls) {
|
||||
const idsToDelete: string[] = [];
|
||||
for (const [toolUseId, callInfo] of agentCalls) {
|
||||
if (now - callInfo.timestamp > PENDING_TOOL_CALL_TTL_MS) {
|
||||
idsToDelete.push(toolUseId);
|
||||
}
|
||||
}
|
||||
for (const id of idsToDelete) {
|
||||
agentCalls.delete(id);
|
||||
}
|
||||
// If agent has no more pending calls, remove the agent entry from the map
|
||||
if (agentCalls.size === 0) {
|
||||
this.pendingToolCalls.delete(agentId);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -719,12 +740,15 @@ export class SubagentWatcher extends EventEmitter {
|
||||
|
||||
/**
|
||||
* Extract the short description from the parent session's transcript.
|
||||
* This is the most reliable method because it reads the actual Task tool call
|
||||
* that spawned this agent, which contains the description parameter.
|
||||
* This is the most reliable method because it reads the actual Task tool result
|
||||
* that spawned this agent, which contains the description directly.
|
||||
*
|
||||
* The parent transcript contains:
|
||||
* 1. An assistant message with tool_use for Task, containing input.description
|
||||
* 2. A progress entry with data.agentId and parentToolUseID linking them
|
||||
* The parent transcript contains a 'user' entry with toolUseResult that has:
|
||||
* - agentId: the spawned agent's ID
|
||||
* - description: the Task description parameter
|
||||
*
|
||||
* We look for this format:
|
||||
* { "type": "user", "toolUseResult": { "agentId": "xxx", "description": "..." } }
|
||||
*/
|
||||
private extractDescriptionFromParentTranscript(
|
||||
projectHash: string,
|
||||
@@ -739,40 +763,16 @@ export class SubagentWatcher extends EventEmitter {
|
||||
const content = readFileSync(transcriptPath, 'utf8');
|
||||
const lines = content.split('\n').filter((l) => l.trim());
|
||||
|
||||
// First pass: Find the progress entry for this agent to get parentToolUseID
|
||||
let parentToolUseID: string | undefined;
|
||||
// Look for user entry with toolUseResult containing the agentId
|
||||
for (const line of lines) {
|
||||
try {
|
||||
const entry = JSON.parse(line);
|
||||
if (entry.type === 'progress' && entry.data?.agentId === agentId && entry.parentToolUseID) {
|
||||
parentToolUseID = entry.parentToolUseID;
|
||||
break;
|
||||
}
|
||||
} catch {
|
||||
// Skip malformed lines
|
||||
}
|
||||
}
|
||||
|
||||
if (!parentToolUseID) return undefined;
|
||||
|
||||
// Second pass: Find the Task tool_use with this ID and extract description
|
||||
for (const line of lines) {
|
||||
try {
|
||||
const entry = JSON.parse(line);
|
||||
if (entry.type === 'assistant' && entry.message?.content) {
|
||||
const content = entry.message.content;
|
||||
if (Array.isArray(content)) {
|
||||
for (const block of content) {
|
||||
if (
|
||||
block.type === 'tool_use' &&
|
||||
block.id === parentToolUseID &&
|
||||
block.name === 'Task' &&
|
||||
block.input?.description
|
||||
) {
|
||||
return block.input.description;
|
||||
}
|
||||
}
|
||||
}
|
||||
if (
|
||||
entry.type === 'user' &&
|
||||
entry.toolUseResult?.agentId === agentId &&
|
||||
entry.toolUseResult?.description
|
||||
) {
|
||||
return entry.toolUseResult.description;
|
||||
}
|
||||
} catch {
|
||||
// Skip malformed lines
|
||||
@@ -1169,12 +1169,15 @@ export class SubagentWatcher extends EventEmitter {
|
||||
} else {
|
||||
for (const content of entry.message.content) {
|
||||
if (content.type === 'tool_use' && content.name) {
|
||||
// Store toolUseId for linking to results
|
||||
// Store toolUseId for linking to results, with timestamp for TTL cleanup
|
||||
if (content.id) {
|
||||
if (!this.pendingToolCalls.has(agentId)) {
|
||||
this.pendingToolCalls.set(agentId, new Map());
|
||||
}
|
||||
this.pendingToolCalls.get(agentId)!.set(content.id, content.name);
|
||||
this.pendingToolCalls.get(agentId)!.set(content.id, {
|
||||
toolName: content.name,
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
}
|
||||
|
||||
const toolCall: SubagentToolCall = {
|
||||
@@ -1197,7 +1200,8 @@ export class SubagentWatcher extends EventEmitter {
|
||||
// Extract tool result
|
||||
const resultContent = this.extractToolResultContent(content.content);
|
||||
const agentPendingCalls = this.pendingToolCalls.get(agentId);
|
||||
const toolName = agentPendingCalls?.get(content.tool_use_id);
|
||||
const pendingCall = agentPendingCalls?.get(content.tool_use_id);
|
||||
const toolName = pendingCall?.toolName;
|
||||
// Delete after lookup to prevent memory leak
|
||||
agentPendingCalls?.delete(content.tool_use_id);
|
||||
|
||||
@@ -1247,7 +1251,8 @@ export class SubagentWatcher extends EventEmitter {
|
||||
if (content.type === 'tool_result' && content.tool_use_id) {
|
||||
const resultContent = this.extractToolResultContent(content.content);
|
||||
const agentPendingCalls = this.pendingToolCalls.get(agentId);
|
||||
const toolName = agentPendingCalls?.get(content.tool_use_id);
|
||||
const pendingCall = agentPendingCalls?.get(content.tool_use_id);
|
||||
const toolName = pendingCall?.toolName;
|
||||
// Delete after lookup to prevent memory leak
|
||||
agentPendingCalls?.delete(content.tool_use_id);
|
||||
|
||||
|
||||
@@ -10,10 +10,10 @@
|
||||
* The transcript path is provided by Claude Code hooks in the `transcript_path` field.
|
||||
*/
|
||||
|
||||
import { EventEmitter } from 'events';
|
||||
import { watch, statSync, existsSync, FSWatcher } from 'fs';
|
||||
import { createReadStream } from 'fs';
|
||||
import { createInterface } from 'readline';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { watch, statSync, existsSync, FSWatcher } from 'node:fs';
|
||||
import { createReadStream } from 'node:fs';
|
||||
import { createInterface } from 'node:readline';
|
||||
|
||||
// ========== Types ==========
|
||||
|
||||
|
||||
@@ -93,6 +93,12 @@ export type TaskStatus = 'pending' | 'running' | 'completed' | 'failed';
|
||||
/** Status of the Ralph Loop controller */
|
||||
export type RalphLoopStatus = 'stopped' | 'running' | 'paused';
|
||||
|
||||
/** Task execution status for plan tracking */
|
||||
export type PlanTaskStatus = 'pending' | 'in_progress' | 'completed' | 'failed' | 'blocked';
|
||||
|
||||
/** TDD phase categories */
|
||||
export type TddPhase = 'setup' | 'test' | 'impl' | 'verify' | 'review';
|
||||
|
||||
// ========== Session Types ==========
|
||||
|
||||
/**
|
||||
|
||||
+42
-4
@@ -42,6 +42,8 @@ export interface LRUMapOptions<K, V> {
|
||||
export class LRUMap<K, V> extends Map<K, V> {
|
||||
private readonly maxSize: number;
|
||||
private readonly onEvict?: (key: K, value: V) => void;
|
||||
/** Tracks the newest key for O(1) newest() access */
|
||||
private _newestKey: K | undefined = undefined;
|
||||
|
||||
/**
|
||||
* Creates a new LRUMap.
|
||||
@@ -71,6 +73,8 @@ export class LRUMap<K, V> extends Map<K, V> {
|
||||
|
||||
// Add the entry (will be at end = most recent)
|
||||
super.set(key, value);
|
||||
// Track newest key for O(1) newest() access
|
||||
this._newestKey = key;
|
||||
|
||||
// Evict oldest entries if over capacity
|
||||
while (super.size > this.maxSize) {
|
||||
@@ -100,6 +104,8 @@ export class LRUMap<K, V> extends Map<K, V> {
|
||||
const value = super.get(key)!;
|
||||
super.delete(key);
|
||||
super.set(key, value);
|
||||
// Track newest key for O(1) newest() access
|
||||
this._newestKey = key;
|
||||
return value;
|
||||
}
|
||||
|
||||
@@ -114,6 +120,36 @@ export class LRUMap<K, V> extends Map<K, V> {
|
||||
return super.has(key);
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete a key-value pair.
|
||||
* Updates _newestKey if the deleted key was the newest.
|
||||
*
|
||||
* @param key - Key to delete
|
||||
* @returns True if the key existed and was deleted
|
||||
*/
|
||||
override delete(key: K): boolean {
|
||||
const existed = super.delete(key);
|
||||
// If we deleted the newest key, we need to find the new newest
|
||||
// This is O(n) but delete is rare; set/get are the hot paths
|
||||
if (existed && this._newestKey === key) {
|
||||
this._newestKey = undefined;
|
||||
// Find the new newest by iterating (last entry)
|
||||
for (const k of super.keys()) {
|
||||
this._newestKey = k;
|
||||
}
|
||||
}
|
||||
return existed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Clear all entries.
|
||||
* Resets _newestKey to undefined.
|
||||
*/
|
||||
override clear(): void {
|
||||
super.clear();
|
||||
this._newestKey = undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* Peek at a value WITHOUT refreshing its position.
|
||||
* Use this when you want to read without affecting LRU order.
|
||||
@@ -138,15 +174,17 @@ export class LRUMap<K, V> extends Map<K, V> {
|
||||
|
||||
/**
|
||||
* Get the newest entry (most recently accessed).
|
||||
* O(1) operation using tracked newest key.
|
||||
*
|
||||
* @returns [key, value] of newest entry, or undefined if empty
|
||||
*/
|
||||
newest(): [K, V] | undefined {
|
||||
let last: [K, V] | undefined;
|
||||
for (const entry of super.entries()) {
|
||||
last = entry;
|
||||
if (this._newestKey === undefined || !super.has(this._newestKey)) {
|
||||
return undefined;
|
||||
}
|
||||
return last;
|
||||
// Use super.get to avoid refreshing the position
|
||||
const value = super.get(this._newestKey)!;
|
||||
return [this._newestKey, value];
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -11243,6 +11243,21 @@ class ClaudemanApp {
|
||||
this.planGenerationAbortController = null;
|
||||
}
|
||||
|
||||
// Clean up wizard-specific timers (leak fix: not cleared on SSE reconnect)
|
||||
if (this.wizardMinimizedTimer) {
|
||||
clearInterval(this.wizardMinimizedTimer);
|
||||
this.wizardMinimizedTimer = null;
|
||||
}
|
||||
|
||||
// Clean up wizard drag listeners (leak fix: document-level handlers)
|
||||
this.cleanupWizardDragging();
|
||||
|
||||
// Deactivate focus trap if wizard was open (leak fix: keydown listener)
|
||||
if (this.activeFocusTrap) {
|
||||
this.activeFocusTrap.deactivate();
|
||||
this.activeFocusTrap = null;
|
||||
}
|
||||
|
||||
// Clear minimized agents tracking
|
||||
this.minimizedSubagents.clear();
|
||||
|
||||
|
||||
@@ -1159,6 +1159,60 @@ describe('SubagentWatcher', () => {
|
||||
expect(updatedInfo.description).toBeDefined();
|
||||
expect(updatedInfo.description).toContain('Create unit tests');
|
||||
});
|
||||
|
||||
it('should extract description from parent transcript toolUseResult', async () => {
|
||||
// Create parent transcript with toolUseResult containing agentId and description
|
||||
const parentTranscript = JSON.stringify({
|
||||
type: 'user',
|
||||
timestamp: new Date().toISOString(),
|
||||
message: { role: 'user', content: [] },
|
||||
toolUseResult: {
|
||||
isAsync: true,
|
||||
status: 'async_launched',
|
||||
agentId: 'parentdesc',
|
||||
description: 'Research codebase cleanups',
|
||||
},
|
||||
});
|
||||
|
||||
const mockRl = new EventEmitter();
|
||||
mockCreateInterface.mockReturnValue(mockRl);
|
||||
mockCreateReadStream.mockReturnValue({});
|
||||
|
||||
mockExistsSync.mockReturnValue(true);
|
||||
mockReaddirSync.mockImplementation((path: string) => {
|
||||
if (path.includes('subagents')) return ['agent-parentdesc.jsonl'];
|
||||
if (path.includes('session1')) return ['subagents'];
|
||||
if (path.includes('project1')) return ['session1'];
|
||||
return ['project1'];
|
||||
});
|
||||
mockStatSync.mockReturnValue({
|
||||
isDirectory: () => true,
|
||||
birthtime: new Date(),
|
||||
mtime: new Date(),
|
||||
size: 100,
|
||||
});
|
||||
// Return parent transcript when reading the session transcript
|
||||
// Return empty for subagent file (will fall back, but we want to test parent extraction)
|
||||
mockReadFileSync.mockImplementation((filepath: string) => {
|
||||
if (filepath.includes('session1.jsonl')) {
|
||||
return parentTranscript;
|
||||
}
|
||||
return ''; // Empty subagent file
|
||||
});
|
||||
|
||||
const discoveredHandler = vi.fn();
|
||||
watcher.on('subagent:discovered', discoveredHandler);
|
||||
|
||||
watcher.start();
|
||||
mockRl.emit('close');
|
||||
|
||||
await vi.advanceTimersByTimeAsync(100);
|
||||
|
||||
expect(discoveredHandler).toHaveBeenCalled();
|
||||
const info = discoveredHandler.mock.calls[0][0] as SubagentInfo;
|
||||
// Should have extracted description from parent transcript
|
||||
expect(info.description).toBe('Research codebase cleanups');
|
||||
});
|
||||
});
|
||||
|
||||
describe('getRecentSubagents', () => {
|
||||
|
||||
Reference in New Issue
Block a user