From 1705b7a67d4eda2b8541fad1cd099e3432c9faf9 Mon Sep 17 00:00:00 2001 From: arkon Date: Fri, 23 Jan 2026 10:20:37 +0100 Subject: [PATCH] feat: implement spawn1337 autonomous agent protocol Add full lifecycle management for spawning autonomous Claude sessions as screen-based agents. Agents communicate via filesystem message bus, signal completion via RalphTracker mechanism, and enforce resource budgets (tokens, cost, timeout, depth limits). New files: - spawn-types.ts: Types, YAML parser, factory functions, serialization - spawn-detector.ts: Terminal pattern detection for spawn1337 tags - spawn-orchestrator.ts: Agent lifecycle (spawn, monitor, queue, cleanup) - spawn-claude-md.ts: CLAUDE.md generator for agent sessions Modified: - session.ts: SpawnDetector integration, parent/child tracking - server.ts: Orchestrator wiring, 11 API endpoints, SSE events - types.ts: Re-exports, SessionState additions Tests: 80 new tests across 3 test files (all passing) Co-Authored-By: Claude Opus 4.5 --- CLAUDE.md | 96 +++- src/session.ts | 78 +++ src/spawn-claude-md.ts | 154 ++++++ src/spawn-detector.ts | 292 ++++++++++ src/spawn-orchestrator.ts | 907 ++++++++++++++++++++++++++++++++ src/spawn-types.ts | 698 ++++++++++++++++++++++++ src/types.ts | 22 + src/web/server.ts | 206 ++++++++ test/spawn-detector.test.ts | 249 +++++++++ test/spawn-orchestrator.test.ts | 496 +++++++++++++++++ test/spawn-types.test.ts | 395 ++++++++++++++ 11 files changed, 3591 insertions(+), 2 deletions(-) create mode 100644 src/spawn-claude-md.ts create mode 100644 src/spawn-detector.ts create mode 100644 src/spawn-orchestrator.ts create mode 100644 src/spawn-types.ts create mode 100644 test/spawn-detector.test.ts create mode 100644 test/spawn-orchestrator.test.ts create mode 100644 test/spawn-types.test.ts diff --git a/CLAUDE.md b/CLAUDE.md index 69726762..981f7ab2 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -68,7 +68,7 @@ npx vitest run -t "should create session" # By pattern # 3115: integration-flows.test.ts # 3120: session-cleanup.test.ts # 3125: ralph-integration.test.ts -# Unit tests (no port needed): respawn-controller, ralph-tracker, pty-interactive, task-queue, task, ralph-loop, session-manager, state-store, types, templates, ralph-config +# Unit tests (no port needed): respawn-controller, ralph-tracker, pty-interactive, task-queue, task, ralph-loop, session-manager, state-store, types, templates, ralph-config, spawn-detector, spawn-types, spawn-orchestrator # Next available: 3127+ # Tests mock PTY - no real Claude CLI spawned @@ -148,6 +148,10 @@ claudeman reset # Reset all state | `src/types.ts` | All TypeScript interfaces | | `src/templates/claude-md.ts` | CLAUDE.md template generation with placeholder support | | `src/templates/case-template.md` | Default CLAUDE.md template for new cases (with placeholders) | +| `src/spawn-types.ts` | Types, YAML parser, factory functions for spawn1337 protocol | +| `src/spawn-detector.ts` | Detects `` tags in terminal output (like ralph-tracker.ts) | +| `src/spawn-orchestrator.ts` | Full agent lifecycle: spawn, monitor, budget, queue, cleanup | +| `src/spawn-claude-md.ts` | Generates CLAUDE.md for spawned agent sessions | ### Data Flow @@ -175,6 +179,72 @@ Steps can be skipped via config (`sendClear: false`, `sendInit: false`). Optiona **Step confirmation**: After sending each step (update, clear, init, kickstart), the controller waits for `completionConfirmMs` (5s) of output silence before proceeding to the next step. This prevents sending commands while Claude is still processing. +### Spawn1337 Protocol (Autonomous Agents) + +Spawned agents are full-power Claude sessions running in their own screen sessions. They communicate via a filesystem-based message bus and signal completion via RalphTracker's `` mechanism. + +**Protocol Flow:** +``` +Parent outputs task.md + → SpawnDetector parses tag + → SpawnOrchestrator reads + parses task spec file (YAML frontmatter) + → Creates agent directory: ~/claudeman-cases/spawn-/ + → Spawns interactive Claude session in screen + → Injects initial prompt via writeViaScreen() + → Agent works autonomously, writes progress to spawn-comms/ + → RalphTracker detects PHRASE on child + → Orchestrator reads result.md, notifies parent via SSE +``` + +**Tag Patterns** (detected by `SpawnDetector`): +- `path/to/task.md` - Spawn request +- `` - Status query +- `` - Cancel request +- `content` - Message to child + +**Agent Communication Directory:** +``` +~/claudeman-cases/spawn-/ +├── CLAUDE.md # Auto-generated agent instructions +├── spawn-comms/ +│ ├── task.md # Copy of original task spec +│ ├── progress.json # Agent updates periodically +│ ├── result.md # Final result (YAML frontmatter + body) +│ └── messages/ # Bidirectional message files (NNN-parent.md, NNN-agent.md) +└── workspace/ # Symlinked context files +``` + +**Task Spec Format** (YAML frontmatter in .md file): +```yaml +--- +agentId: my-agent-001 +name: My Agent +type: explore # explore|implement|test|review|refactor|research|generate|fix|general +priority: high # low|normal|high|critical +maxTokens: 150000 +maxCost: 0.50 +timeoutMinutes: 15 +canModifyParentFiles: false +contextFiles: [src/auth.ts, src/types.ts] +dependsOn: [other-agent-id] +completionPhrase: MY_AGENT_DONE +outputFormat: structured # markdown|json|code|structured|freeform +--- +Task instructions here... +``` + +**Resource Governance:** +- Budget warning at 80% of token/cost limit +- Graceful shutdown message at 100% +- Force kill at 110% (or timeout + 60s grace) +- Max concurrent agents: 5 (configurable) +- Max spawn depth: 3 (prevents infinite recursion) +- Default timeout: 30 minutes (max: 120) + +**Session Integration:** Each session has a `SpawnDetector` (alongside RalphTracker) that forwards terminal data. Spawn events (`spawnRequested`, `spawnStatusRequested`, `spawnCancelRequested`, `spawnMessageToChild`) are emitted on the session and wired to the orchestrator in server.ts. + +**Agent Tree:** Agents can spawn children (up to `maxSpawnDepth`). Sessions track `parentAgentId` and `childAgentIds`. Cancelling a parent cascades to all children. + ### Session Modes Sessions have a `mode` property (`SessionMode` type): @@ -370,7 +440,7 @@ Tab switch/new session fix: clear xterm → write buffer → resize PTY → Ctrl All events broadcast to `/api/events` with format: `{ type: string, sessionId?: string, data: any }`. -Event prefixes: `session:`, `task:`, `respawn:`, `scheduled:`, `case:`, `screen:`, `init`. +Event prefixes: `session:`, `task:`, `respawn:`, `spawn:`, `scheduled:`, `case:`, `screen:`, `init`. Key events for frontend handling (see `app.js:handleSSEEvent()`): - `session:idle`, `session:working` - Status indicator updates @@ -378,6 +448,8 @@ Key events for frontend handling (see `app.js:handleSSEEvent()`): - `session:completion`, `session:autoClear`, `session:autoCompact` - Lifecycle events - `session:ralphLoopUpdate`, `session:ralphTodoUpdate`, `session:ralphCompletionDetected` - Ralph tracking - `respawn:detectionUpdate` - Multi-layer idle detection status (confidence level, waiting state) +- `spawn:queued`, `spawn:started`, `spawn:completed`, `spawn:failed`, `spawn:timeout`, `spawn:cancelled` - Agent lifecycle +- `spawn:progress`, `spawn:message`, `spawn:budgetWarning`, `spawn:stateUpdate` - Agent monitoring ### Frontend (app.js) @@ -410,6 +482,7 @@ Writes debounced (500ms) to `~/.claudeman/state.json`. The web server persists f | `respawnEnabled` | Whether respawn controller is currently running | | `respawnConfig` | Full respawn config including `durationMinutes` | | `totalCost`, `inputTokens`, `outputTokens` | Token and cost tracking | +| `parentAgentId`, `childAgentIds` | Spawn agent tree relationships | **CLI visibility**: The `claudeman status` and `claudeman session list` commands read from `state.json` to display web server-managed sessions, even though they run in a separate process. @@ -433,6 +506,15 @@ Writes debounced (500ms) to `~/.claudeman/state.json`. The web server persists f | Respawn completion confirm | 5s | `RespawnConfig.completionConfirmMs` | | Respawn auto-accept delay | 8s | `RespawnConfig.autoAcceptDelayMs` | | Respawn no-output fallback | 30s | `RespawnConfig.noOutputTimeoutMs` | +| Spawn event debounce | 50ms | `spawn-detector.ts` | +| Spawn line buffer max | 64KB | `spawn-detector.ts` | +| Spawn progress poll | 5s | `spawn-orchestrator.ts` | +| Spawn default timeout | 30 min | `spawn-orchestrator.ts` | +| Spawn max timeout | 120 min | `spawn-orchestrator.ts` | +| Spawn max concurrent | 5 agents | `spawn-orchestrator.ts` | +| Spawn max depth | 3 levels | `spawn-orchestrator.ts` | +| Spawn budget warning | 80% | `spawn-orchestrator.ts` | +| Spawn budget grace | 60s | `spawn-orchestrator.ts` | ### TypeScript Config @@ -547,6 +629,16 @@ Long-running sessions are supported with automatic trimming: | GET | `/api/cases` | List available cases | | POST | `/api/cases` | Create new case | | GET | `/api/screens` | List screen sessions with stats | +| GET | `/api/spawn/agents` | List all spawn agents (active + completed + queued) | +| GET | `/api/spawn/agents/:agentId` | Detailed agent status + progress | +| GET | `/api/spawn/agents/:agentId/result` | Read agent's result.md | +| GET | `/api/spawn/agents/:agentId/messages` | List messages in channel | +| POST | `/api/spawn/agents/:agentId/message` | Send message to agent | +| POST | `/api/spawn/agents/:agentId/cancel` | Cancel agent (graceful stop) | +| DELETE | `/api/spawn/agents/:agentId` | Force kill + cleanup | +| GET | `/api/spawn/status` | Orchestrator status (counts, config) | +| PUT | `/api/spawn/config` | Update orchestrator config | +| POST | `/api/spawn/trigger` | Programmatic spawn (bypass terminal detection) | ## Keyboard Shortcuts (Web UI) diff --git a/src/session.ts b/src/session.ts index e9a4521b..813b8350 100644 --- a/src/session.ts +++ b/src/session.ts @@ -21,6 +21,7 @@ import * as pty from 'node-pty'; import { SessionState, SessionStatus, SessionConfig, ScreenSession, RalphTrackerState, RalphTodoItem } from './types.js'; import { TaskTracker, type BackgroundTask } from './task-tracker.js'; import { RalphTracker } from './ralph-tracker.js'; +import { SpawnDetector } from './spawn-detector.js'; import { ScreenManager } from './screen-manager.js'; export type { BackgroundTask } from './task-tracker.js'; @@ -225,6 +226,14 @@ export interface SessionEvents { ralphTodoUpdate: (todos: RalphTodoItem[]) => void; /** Ralph completion phrase detected */ ralphCompletionDetected: (phrase: string) => void; + /** Spawn1337 agent spawn requested */ + spawnRequested: (filePath: string, rawLine: string) => void; + /** Spawn1337 agent status query */ + spawnStatusRequested: (agentId: string) => void; + /** Spawn1337 agent cancel request */ + spawnCancelRequested: (agentId: string) => void; + /** Spawn1337 message to child agent */ + spawnMessageToChild: (agentId: string, content: string) => void; } /** @@ -325,6 +334,13 @@ export class Session extends EventEmitter { // Ralph tracking (Ralph Wiggum loops and todo lists inside Claude Code) private _ralphTracker: RalphTracker; + // Spawn1337 detection (agent spawning protocol) + private _spawnDetector: SpawnDetector; + + // Agent tree tracking + private _parentAgentId: string | null = null; + private _childAgentIds: string[] = []; + // Store handler references for cleanup (prevents memory leaks) private _taskTrackerHandlers: { taskCreated: (task: BackgroundTask) => void; @@ -339,6 +355,13 @@ export class Session extends EventEmitter { completionDetected: (phrase: string) => void; } | null = null; + private _spawnHandlers: { + spawnRequested: (filePath: string, rawLine: string) => void; + statusRequested: (agentId: string) => void; + cancelRequested: (agentId: string) => void; + messageToChild: (agentId: string, content: string) => void; + } | null = null; + constructor(config: Partial & { workingDir: string; mode?: SessionMode; @@ -381,6 +404,19 @@ export class Session extends EventEmitter { this._ralphTracker.on('loopUpdate', this._ralphHandlers.loopUpdate); this._ralphTracker.on('todoUpdate', this._ralphHandlers.todoUpdate); this._ralphTracker.on('completionDetected', this._ralphHandlers.completionDetected); + + // Initialize Spawn detector and forward events (store handlers for cleanup) + this._spawnDetector = new SpawnDetector(); + this._spawnHandlers = { + spawnRequested: (filePath, rawLine) => this.emit('spawnRequested', filePath, rawLine), + statusRequested: (agentId) => this.emit('spawnStatusRequested', agentId), + cancelRequested: (agentId) => this.emit('spawnCancelRequested', agentId), + messageToChild: (agentId, content) => this.emit('spawnMessageToChild', agentId, content), + }; + this._spawnDetector.on('spawnRequested', this._spawnHandlers.spawnRequested); + this._spawnDetector.on('statusRequested', this._spawnHandlers.statusRequested); + this._spawnDetector.on('cancelRequested', this._spawnHandlers.cancelRequested); + this._spawnDetector.on('messageToChild', this._spawnHandlers.messageToChild); } get status(): SessionStatus { @@ -468,6 +504,34 @@ export class Session extends EventEmitter { return this._ralphTracker.getTodoStats(); } + // Spawn1337 tracking getters + get spawnDetector(): SpawnDetector { + return this._spawnDetector; + } + + get parentAgentId(): string | null { + return this._parentAgentId; + } + + set parentAgentId(value: string | null) { + this._parentAgentId = value; + } + + get childAgentIds(): string[] { + return [...this._childAgentIds]; + } + + addChildAgentId(agentId: string): void { + if (!this._childAgentIds.includes(agentId)) { + this._childAgentIds.push(agentId); + } + } + + removeChildAgentId(agentId: string): void { + const idx = this._childAgentIds.indexOf(agentId); + if (idx >= 0) this._childAgentIds.splice(idx, 1); + } + // Token tracking getters and setters get totalTokens(): number { return this._totalInputTokens + this._totalOutputTokens; @@ -559,6 +623,8 @@ export class Session extends EventEmitter { outputTokens: this._totalOutputTokens, ralphEnabled: this._ralphTracker.enabled, ralphCompletionPhrase: this._ralphTracker.loopState.completionPhrase || undefined, + parentAgentId: this._parentAgentId || undefined, + childAgentIds: this._childAgentIds.length > 0 ? this._childAgentIds : undefined, }; } @@ -743,6 +809,9 @@ export class Session extends EventEmitter { // Forward to Ralph tracker to detect Ralph loops and todos this._ralphTracker.processTerminalData(data); + // Forward to Spawn detector to detect spawn1337 protocol tags + this._spawnDetector.processTerminalData(data); + // Parse token count from status line (e.g., "123.4k tokens" or "5234 tokens") this.parseTokensFromStatusLine(data); @@ -1390,6 +1459,15 @@ export class Session extends EventEmitter { this._ralphTracker.off('completionDetected', this._ralphHandlers.completionDetected); this._ralphHandlers = null; } + + // Remove SpawnDetector handlers + if (this._spawnHandlers) { + this._spawnDetector.off('spawnRequested', this._spawnHandlers.spawnRequested); + this._spawnDetector.off('statusRequested', this._spawnHandlers.statusRequested); + this._spawnDetector.off('cancelRequested', this._spawnHandlers.cancelRequested); + this._spawnDetector.off('messageToChild', this._spawnHandlers.messageToChild); + this._spawnHandlers = null; + } } /** diff --git a/src/spawn-claude-md.ts b/src/spawn-claude-md.ts new file mode 100644 index 00000000..04429702 --- /dev/null +++ b/src/spawn-claude-md.ts @@ -0,0 +1,154 @@ +/** + * @fileoverview Agent CLAUDE.md Generator for spawn1337 protocol. + * + * Generates a comprehensive CLAUDE.md for each spawned agent that tells it: + * - What its task is + * - How to communicate progress + * - How to signal completion + * - What constraints it has + * - How to read/write messages + * + * @module spawn-claude-md + */ + +import type { SpawnTask } from './spawn-types.js'; + +/** + * Generate a CLAUDE.md file for a spawned agent. + * + * This CLAUDE.md gives the agent full context about: + * - Its identity and task + * - Communication protocol (progress, messages, result) + * - Resource constraints (timeout, tokens, cost) + * - Working directory and available context files + * + * @param task - The full parsed task specification + * @param commsDir - Absolute path to the communication directory + * @param agentWorkingDir - Absolute path to the agent's working directory + * @returns The generated CLAUDE.md content + */ +export function generateAgentClaudeMd(task: SpawnTask, commsDir: string, agentWorkingDir: string): string { + const spec = task.spec; + + const constraintLines: string[] = []; + constraintLines.push(`- Timeout: ${spec.timeoutMinutes} minutes`); + if (spec.maxTokens) constraintLines.push(`- Token budget: ${spec.maxTokens.toLocaleString()} tokens`); + if (spec.maxCost) constraintLines.push(`- Cost budget: $${spec.maxCost.toFixed(2)}`); + if (!spec.canModifyParentFiles) { + constraintLines.push('- DO NOT modify files outside your workspace'); + } else { + constraintLines.push('- You MAY modify files in the parent project directory'); + } + constraintLines.push(`- Output format: ${spec.outputFormat}`); + + const contextSection = spec.contextFiles && spec.contextFiles.length > 0 + ? `\nContext files available in workspace:\n${spec.contextFiles.map(f => `- ${f}`).join('\n')}` + : ''; + + const progressSection = spec.progressIntervalSeconds > 0 + ? `### Progress Reporting + +Update \`${commsDir}/progress.json\` every ~${spec.progressIntervalSeconds} seconds with your current status: + +\`\`\`json +{ + "phase": "current phase description", + "percentComplete": 45, + "currentAction": "What you are doing right now", + "subtasks": [ + {"description": "Subtask 1", "status": "completed"}, + {"description": "Subtask 2", "status": "in_progress"} + ], + "filesModified": ["file1.ts", "file2.ts"], + "tokensUsed": 0, + "costSoFar": 0, + "updatedAt": ${Date.now()} +} +\`\`\`` + : '### Progress Reporting\n\nProgress reporting is disabled for this task.'; + + return `# Agent: ${spec.name} + +## Your Identity + +You are an autonomous agent (ID: \`${spec.agentId}\`) spawned by a parent Claude session. +You are running in your own screen session with full Claude Code capabilities. +Type: ${spec.type} | Priority: ${spec.priority} | Depth: ${task.depth} + +## Task + +${task.instructions} + +## Success Criteria + +${spec.successCriteria || 'Complete the task as described above.'} + +## Communication Protocol + +${progressSection} + +### Check for Messages + +Periodically check \`${commsDir}/messages/\` for instructions from the parent. +Files are named \`NNN-parent.md\` (from parent) or \`NNN-agent.md\` (from you). +Read any new \`*-parent.md\` files for additional instructions or clarifications. + +To send a message back to the parent, create a file like: +\`${commsDir}/messages/002-agent.md\` + +### Write Result + +When complete, write your final result to \`${commsDir}/result.md\` with YAML frontmatter: + +\`\`\`markdown +--- +status: completed +summary: "Brief 1-3 sentence summary of what you accomplished" +filesChanged: + - path: relative/path/to/file.ts + action: modified + summary: "What was changed" +--- + +## Detailed Output + +Your full output, analysis, or report here. +\`\`\` + +Valid status values: \`completed\`, \`failed\` + +### Signal Completion + +After writing result.md, output this EXACT phrase to signal you are done: + +${spec.completionPhrase} + +**IMPORTANT**: Only output the completion phrase AFTER you have written result.md. +The completion phrase triggers the orchestrator to read your result and clean up. + +## Constraints + +${constraintLines.join('\n')} + +## Working Directory + +Your workspace is: \`${agentWorkingDir}\` +${contextSection} + +## Important Notes + +- Work autonomously - do not ask for user input +- Focus exclusively on the task described above +- If you encounter errors, document them in result.md with status: failed +- Do not modify this CLAUDE.md file +- Stay within your resource constraints +`; +} + +/** + * Build the initial prompt injected into the agent session via writeViaScreen(). + * Intentionally brief - all detail is in the CLAUDE.md. + */ +export function buildInitialPrompt(task: SpawnTask): string { + return `Read your CLAUDE.md file for complete task instructions, communication protocol, and constraints. Begin working on the task immediately. Report progress to spawn-comms/progress.json periodically. When complete, write your result to spawn-comms/result.md and then output your completion phrase: ${task.spec.completionPhrase}`; +} diff --git a/src/spawn-detector.ts b/src/spawn-detector.ts new file mode 100644 index 00000000..402dd9df --- /dev/null +++ b/src/spawn-detector.ts @@ -0,0 +1,292 @@ +/** + * @fileoverview Spawn Detector - Detects spawn1337 tags in terminal output. + * + * Monitors terminal output for spawn protocol patterns: + * - filename.md - Agent spawn request + * - - Status query + * - - Cancel request + * - content - Message to child + * + * Same architecture as ralph-tracker.ts: line-buffered, auto-enabling, + * debounced events, pre-compiled patterns. + * + * @module spawn-detector + */ + +import { EventEmitter } from 'node:events'; +import { SpawnTrackerState, createInitialSpawnTrackerState } from './spawn-types.js'; + +// ========== Configuration Constants ========== + +/** Debounce interval for state event emissions (ms) */ +const EVENT_DEBOUNCE_MS = 50; + +/** Maximum line buffer size to prevent unbounded growth */ +const MAX_LINE_BUFFER_SIZE = 64 * 1024; + +// ========== Pre-compiled Regex Patterns ========== + +/** Matches spawn request tags: filename.md */ +const SPAWN_TAG_PATTERN = /([^<]+)<\/spawn1337>/; + +/** Quick check string before running regex */ +const SPAWN_QUICK_CHECK = 'spawn1337'; + +/** Matches status query: */ +const SPAWN_STATUS_PATTERN = //; + +/** Matches cancel request: */ +const SPAWN_CANCEL_PATTERN = //; + +/** Matches message to child: content */ +const SPAWN_MESSAGE_PATTERN = /([\s\S]*?)<\/spawn1337-message>/; + +/** Removes ANSI escape codes from terminal output */ +const ANSI_ESCAPE_PATTERN = /\x1b\[[0-9;]*[A-Za-z]/g; + +// ========== Event Types ========== + +/** + * Events emitted by SpawnDetector + */ +export interface SpawnDetectorEvents { + /** Emitted when a spawn request tag is detected */ + spawnRequested: (filePath: string, rawLine: string) => void; + /** Emitted when a status query is detected */ + statusRequested: (agentId: string) => void; + /** Emitted when a cancel request is detected */ + cancelRequested: (agentId: string) => void; + /** Emitted when a message to child is detected */ + messageToChild: (agentId: string, content: string) => void; + /** Emitted when tracker state changes */ + stateUpdate: (state: SpawnTrackerState) => void; +} + +/** + * SpawnDetector - Parses terminal output to detect spawn1337 protocol tags. + * + * This class monitors Claude Code session output to detect agent spawn requests + * and related communication patterns. It auto-enables when any spawn1337 pattern + * is first detected, reducing overhead for sessions not using the spawn protocol. + * + * ## Pattern Detection + * + * 1. **Spawn Request**: `path/to/task.md` + * 2. **Status Query**: `` + * 3. **Cancel Request**: `` + * 4. **Message**: `content` + * + * @extends EventEmitter + */ +export class SpawnDetector extends EventEmitter { + /** Whether the detector is actively monitoring output */ + private _enabled: boolean = false; + + /** Buffer for incomplete lines from terminal data */ + private _lineBuffer: string = ''; + + /** Current tracker state */ + private _state: SpawnTrackerState; + + /** Debounce timer for state events */ + private _stateUpdateTimer: NodeJS.Timeout | null = null; + + /** Flag indicating pending state update emission */ + private _stateUpdatePending: boolean = false; + + constructor() { + super(); + this._state = createInitialSpawnTrackerState(); + } + + /** + * Whether the detector is enabled and actively monitoring output. + */ + get enabled(): boolean { + return this._enabled; + } + + /** + * Get a copy of the current tracker state. + */ + get state(): SpawnTrackerState { + return { ...this._state }; + } + + /** + * Enable the detector to start monitoring terminal output. + */ + enable(): void { + if (!this._enabled) { + this._enabled = true; + this._state.enabled = true; + this.emitStateUpdateDebounced(); + } + } + + /** + * Disable the detector. + */ + disable(): void { + if (this._enabled) { + this._enabled = false; + this._state.enabled = false; + this.emitStateUpdateDebounced(); + } + } + + /** + * Reset all state. + */ + reset(): void { + this.clearDebounceTimers(); + this._enabled = false; + this._lineBuffer = ''; + this._state = createInitialSpawnTrackerState(); + this.emit('stateUpdate', this.state); + } + + /** + * Update state from orchestrator data. + * Called by server.ts when orchestrator state changes. + */ + updateState(state: Partial): void { + Object.assign(this._state, state); + this.emitStateUpdateDebounced(); + } + + /** + * Process raw terminal data to detect spawn patterns. + * + * @param data - Raw terminal data (may include ANSI codes) + */ + processTerminalData(data: string): void { + // Remove ANSI escape codes + const cleanData = data.replace(ANSI_ESCAPE_PATTERN, ''); + + // Buffer data for line-based processing + this._lineBuffer += cleanData; + + // Prevent unbounded line buffer growth + if (this._lineBuffer.length > MAX_LINE_BUFFER_SIZE) { + this._lineBuffer = this._lineBuffer.slice(-MAX_LINE_BUFFER_SIZE / 2); + } + + // Quick pre-check: if no spawn patterns in buffer, just drain lines + if (!this._lineBuffer.includes(SPAWN_QUICK_CHECK)) { + const lines = this._lineBuffer.split('\n'); + this._lineBuffer = lines.pop() || ''; + return; + } + + // Auto-enable on first spawn pattern detection + if (!this._enabled) { + this.enable(); + } + + // Process complete lines + const lines = this._lineBuffer.split('\n'); + this._lineBuffer = lines.pop() || ''; + + for (const line of lines) { + this.processLine(line); + } + + // Also check the full chunk for multi-line patterns (message tag can span lines) + this.checkMultiLinePatterns(cleanData); + } + + /** + * Process a single line for spawn patterns. + */ + private processLine(line: string): void { + const trimmed = line.trim(); + if (!trimmed || !trimmed.includes(SPAWN_QUICK_CHECK)) return; + + // Check spawn request: filename.md + const spawnMatch = trimmed.match(SPAWN_TAG_PATTERN); + if (spawnMatch) { + const filePath = spawnMatch[1].trim(); + this._state.totalSpawned++; + this.emit('spawnRequested', filePath, trimmed); + this.emitStateUpdateDebounced(); + return; + } + + // Check status query: + const statusMatch = trimmed.match(SPAWN_STATUS_PATTERN); + if (statusMatch) { + this.emit('statusRequested', statusMatch[1]); + return; + } + + // Check cancel request: + const cancelMatch = trimmed.match(SPAWN_CANCEL_PATTERN); + if (cancelMatch) { + this.emit('cancelRequested', cancelMatch[1]); + return; + } + + // Check message (single line): content + const msgMatch = trimmed.match(SPAWN_MESSAGE_PATTERN); + if (msgMatch) { + this.emit('messageToChild', msgMatch[1], msgMatch[2]); + return; + } + } + + /** + * Check for patterns that might span multiple lines. + * The message tag content can be multiline. + */ + private checkMultiLinePatterns(data: string): void { + if (!data.includes('spawn1337-message')) return; + + const msgMatch = data.match(SPAWN_MESSAGE_PATTERN); + if (msgMatch) { + this.emit('messageToChild', msgMatch[1], msgMatch[2]); + } + } + + /** + * Emit stateUpdate with debouncing. + */ + private emitStateUpdateDebounced(): void { + this._stateUpdatePending = true; + if (this._stateUpdateTimer) { + clearTimeout(this._stateUpdateTimer); + } + this._stateUpdateTimer = setTimeout(() => { + if (this._stateUpdatePending) { + this._stateUpdatePending = false; + this._stateUpdateTimer = null; + this.emit('stateUpdate', this.state); + } + }, EVENT_DEBOUNCE_MS); + } + + /** + * Flush any pending debounced events immediately. + */ + flushPendingEvents(): void { + if (this._stateUpdatePending) { + this._stateUpdatePending = false; + if (this._stateUpdateTimer) { + clearTimeout(this._stateUpdateTimer); + this._stateUpdateTimer = null; + } + this.emit('stateUpdate', this.state); + } + } + + /** + * Clear all debounce timers. + */ + private clearDebounceTimers(): void { + if (this._stateUpdateTimer) { + clearTimeout(this._stateUpdateTimer); + this._stateUpdateTimer = null; + } + this._stateUpdatePending = false; + } +} diff --git a/src/spawn-orchestrator.ts b/src/spawn-orchestrator.ts new file mode 100644 index 00000000..891c73df --- /dev/null +++ b/src/spawn-orchestrator.ts @@ -0,0 +1,907 @@ +/** + * @fileoverview Spawn Orchestrator - Full lifecycle management for spawned agents. + * + * Manages: + * - Agent creation from task spec files + * - Directory setup (CLAUDE.md, comms, workspace) + * - Session spawning via screen + * - Progress monitoring and timeout enforcement + * - Resource governance (tokens, cost, depth) + * - Bidirectional communication + * - Result collection and cleanup + * - Queue management with priority ordering + * + * @module spawn-orchestrator + */ + +import { EventEmitter } from 'node:events'; +import { join, resolve, isAbsolute } from 'node:path'; +import { existsSync, mkdirSync, writeFileSync, readFileSync, readdirSync, statSync, symlinkSync } from 'node:fs'; +import { v4 as uuidv4 } from 'uuid'; +import { + type SpawnOrchestratorConfig, + type SpawnTask, + type AgentContext, + type AgentProgress, + type AgentStatusReport, + type SpawnResult, + type SpawnTrackerState, + type SpawnMessage, + type SpawnPersistedState, + createDefaultOrchestratorConfig, + createEmptyAgentProgress, + parseTaskSpecFile, + parseSpawnResult, + MAX_TASK_FILE_SIZE, + MAX_CONTEXT_FILE_SIZE, + MAX_CONTEXT_FILES, + MAX_QUEUE_LENGTH, + BUDGET_WARNING_THRESHOLD, + MESSAGE_MAX_SIZE, + MAX_MESSAGES_PER_CHANNEL, + MAX_TRACKED_AGENTS, +} from './spawn-types.js'; +import { generateAgentClaudeMd, buildInitialPrompt } from './spawn-claude-md.js'; +import { getErrorMessage } from './types.js'; + +// ========== Types for integration ========== + +/** + * Interface for session creation callback. + * The orchestrator delegates session creation to the server to avoid circular deps. + */ +export interface SessionCreator { + createAgentSession(workingDir: string, name: string): Promise<{ sessionId: string }>; + writeToSession(sessionId: string, data: string): void; + getSessionTokens(sessionId: string): number; + getSessionCost(sessionId: string): number; + stopSession(sessionId: string): Promise; + onSessionCompletion(sessionId: string, handler: (phrase: string) => void): void; + removeSessionCompletionHandler(sessionId: string, handler: (phrase: string) => void): void; +} + +// ========== Events ========== + +export interface SpawnOrchestratorEvents { + /** Agent added to queue */ + queued: (data: { agentId: string; name: string; parentSessionId: string; position: number }) => void; + /** Agent directory being set up */ + initializing: (data: { agentId: string; name: string; workingDir: string }) => void; + /** Agent session started */ + started: (data: { agentId: string; name: string; sessionId: string }) => void; + /** Agent progress update */ + progress: (data: { agentId: string; progress: AgentProgress }) => void; + /** New message in channel */ + message: (data: { agentId: string; message: SpawnMessage }) => void; + /** Agent completed successfully */ + completed: (data: { agentId: string; result: SpawnResult }) => void; + /** Agent failed */ + failed: (data: { agentId: string; error: string; partialProgress: AgentProgress | null }) => void; + /** Agent timed out */ + timeout: (data: { agentId: string; elapsed: number; limit: number }) => void; + /** Agent cancelled */ + cancelled: (data: { agentId: string; reason: string }) => void; + /** Budget warning */ + budgetWarning: (data: { agentId: string; type: 'tokens' | 'cost'; used: number; limit: number }) => void; + /** Overall state changed */ + stateUpdate: (state: SpawnTrackerState) => void; +} + +/** + * SpawnOrchestrator - Manages the full lifecycle of spawned agents. + * + * Handles agent creation, monitoring, communication, resource governance, + * and cleanup. Integrates with Session, ScreenManager, and RalphTracker. + */ +export class SpawnOrchestrator extends EventEmitter { + private _agents: Map = new Map(); + private _completedAgents: Map = new Map(); + private _queue: SpawnTask[] = []; + private _config: SpawnOrchestratorConfig; + private _sessionCreator: SessionCreator | null = null; + private _totalSpawned: number = 0; + private _totalCompleted: number = 0; + private _totalFailed: number = 0; + private _maxDepthReached: number = 0; + private _completionHandlers: Map void> = new Map(); + + constructor(config?: Partial) { + super(); + this._config = { ...createDefaultOrchestratorConfig(), ...config }; + } + + /** + * Set the session creator callback. + * Must be called before any spawn requests can be processed. + */ + setSessionCreator(creator: SessionCreator): void { + this._sessionCreator = creator; + } + + /** + * Get current orchestrator configuration. + */ + get config(): SpawnOrchestratorConfig { + return { ...this._config }; + } + + /** + * Update orchestrator configuration. + */ + updateConfig(config: Partial): void { + Object.assign(this._config, config); + } + + /** + * Handle a spawn request detected from terminal output. + * + * @param filePath - Path to the task spec file (relative to parent's workingDir) + * @param parentSessionId - ID of the parent session + * @param parentWorkingDir - Working directory of the parent session + * @param parentDepth - Depth of the parent in the spawn tree + */ + async handleSpawnRequest( + filePath: string, + parentSessionId: string, + parentWorkingDir: string, + parentDepth: number = 0 + ): Promise { + if (!this._sessionCreator) { + console.error('[spawn-orchestrator] No session creator set, cannot spawn agent'); + return; + } + + // Resolve file path relative to parent's working directory + const resolvedPath = isAbsolute(filePath) ? filePath : join(parentWorkingDir, filePath); + + // Validate file exists and size + if (!existsSync(resolvedPath)) { + console.error(`[spawn-orchestrator] Task file not found: ${resolvedPath}`); + this.emit('failed', { agentId: 'unknown', error: `Task file not found: ${resolvedPath}`, partialProgress: null }); + return; + } + + const stat = statSync(resolvedPath); + if (stat.size > MAX_TASK_FILE_SIZE) { + console.error(`[spawn-orchestrator] Task file too large: ${stat.size} bytes (max ${MAX_TASK_FILE_SIZE})`); + this.emit('failed', { agentId: 'unknown', error: `Task file too large: ${stat.size} bytes`, partialProgress: null }); + return; + } + + // Parse task file + const content = readFileSync(resolvedPath, 'utf-8'); + const fallbackId = `agent-${uuidv4().slice(0, 8)}`; + const parsed = parseTaskSpecFile(content, fallbackId); + + if (!parsed) { + console.error(`[spawn-orchestrator] Failed to parse task file: ${resolvedPath}`); + this.emit('failed', { agentId: fallbackId, error: 'Failed to parse task spec YAML frontmatter', partialProgress: null }); + return; + } + + const childDepth = parentDepth + 1; + + // Depth check + if (childDepth > this._config.maxSpawnDepth) { + console.error(`[spawn-orchestrator] Max spawn depth (${this._config.maxSpawnDepth}) exceeded at depth ${childDepth}`); + this.emit('failed', { agentId: parsed.spec.agentId, error: `Max spawn depth exceeded (${this._config.maxSpawnDepth})`, partialProgress: null }); + return; + } + + // Enforce timeout limits + if (parsed.spec.timeoutMinutes > this._config.maxTimeoutMinutes) { + parsed.spec.timeoutMinutes = this._config.maxTimeoutMinutes; + } + + const task: SpawnTask = { + spec: parsed.spec, + instructions: parsed.instructions, + sourceFile: resolvedPath, + parentSessionId, + depth: childDepth, + }; + + // Check dependencies + if (task.spec.dependsOn && task.spec.dependsOn.length > 0) { + const unmetDeps = task.spec.dependsOn.filter(depId => { + const dep = this._completedAgents.get(depId); + return !dep || dep.status !== 'completed'; + }); + if (unmetDeps.length > 0) { + // Queue with dependency tracking + this.enqueueTask(task); + return; + } + } + + // Concurrency check + const activeCount = this.getActiveCount(); + if (activeCount >= this._config.maxConcurrentAgents) { + this.enqueueTask(task); + return; + } + + // Spawn immediately + await this.spawnAgent(task); + } + + /** + * Cancel an agent by ID. + */ + async cancelAgent(agentId: string, reason: string = 'Cancelled by parent'): Promise { + const agent = this._agents.get(agentId); + if (!agent) { + // Check queue + const queueIdx = this._queue.findIndex(t => t.spec.agentId === agentId); + if (queueIdx >= 0) { + this._queue.splice(queueIdx, 1); + this.emit('cancelled', { agentId, reason: 'Removed from queue' }); + this.emitStateUpdate(); + } + return; + } + + agent.status = 'cancelled'; + this.emit('cancelled', { agentId, reason }); + + await this.cleanupAgent(agentId); + } + + /** + * Send a message to an agent. + */ + async sendMessageToAgent(agentId: string, content: string): Promise { + const agent = this._agents.get(agentId); + if (!agent) return; + + if (content.length > MESSAGE_MAX_SIZE) { + content = content.slice(0, MESSAGE_MAX_SIZE); + } + + const messagesDir = join(agent.commsDir, 'messages'); + if (!existsSync(messagesDir)) { + mkdirSync(messagesDir, { recursive: true }); + } + + // Count existing messages + const existingMessages = readdirSync(messagesDir).filter(f => f.endsWith('.md')); + if (existingMessages.length >= MAX_MESSAGES_PER_CHANNEL) { + return; // Channel full + } + + const seq = existingMessages.length + 1; + const seqStr = String(seq).padStart(3, '0'); + const fileName = `${seqStr}-parent.md`; + + const message: SpawnMessage = { + sequence: seq, + sender: 'parent', + content, + sentAt: Date.now(), + read: false, + }; + + writeFileSync(join(messagesDir, fileName), content, 'utf-8'); + this.emit('message', { agentId, message }); + } + + /** + * Get status of a specific agent. + */ + getAgentStatus(agentId: string): AgentStatusReport | null { + const agent = this._agents.get(agentId) || this._completedAgents.get(agentId); + if (!agent) return null; + return this.buildStatusReport(agent); + } + + /** + * Get status of all agents (active + recently completed). + */ + getAllAgentStatuses(): AgentStatusReport[] { + const reports: AgentStatusReport[] = []; + for (const agent of this._agents.values()) { + reports.push(this.buildStatusReport(agent)); + } + for (const agent of this._completedAgents.values()) { + reports.push(this.buildStatusReport(agent)); + } + return reports; + } + + /** + * Get current orchestrator state. + */ + getState(): SpawnTrackerState { + return { + enabled: true, + activeCount: this.getActiveCount(), + queuedCount: this._queue.length, + totalSpawned: this._totalSpawned, + totalCompleted: this._totalCompleted, + totalFailed: this._totalFailed, + maxDepthReached: this._maxDepthReached, + agents: this.getAllAgentStatuses(), + }; + } + + /** + * Get state for persistence. + */ + getPersistedState(): SpawnPersistedState { + const agents: SpawnPersistedState['agents'] = {}; + for (const [id, agent] of this._agents) { + agents[id] = { + agentId: id, + status: agent.status, + parentSessionId: agent.parentSessionId, + childSessionId: agent.sessionId, + depth: agent.depth, + startedAt: agent.startedAt, + commsDir: agent.commsDir, + workingDir: agent.workingDir, + completionPhrase: agent.task.spec.completionPhrase, + timeoutMinutes: agent.task.spec.timeoutMinutes, + }; + } + return { config: this._config, agents }; + } + + /** + * Stop all agents. + */ + async stopAll(): Promise { + const agentIds = Array.from(this._agents.keys()); + for (const agentId of agentIds) { + await this.cancelAgent(agentId, 'Orchestrator shutdown'); + } + this._queue = []; + } + + /** + * Read an agent's result.md file. + */ + readAgentResult(agentId: string): SpawnResult | null { + const agent = this._agents.get(agentId) || this._completedAgents.get(agentId); + if (!agent) return null; + + const resultPath = join(agent.commsDir, 'result.md'); + if (!existsSync(resultPath)) return null; + + const content = readFileSync(resultPath, 'utf-8'); + const durationMs = agent.startedAt ? Date.now() - agent.startedAt : 0; + return parseSpawnResult(content, agentId, durationMs); + } + + /** + * Read an agent's progress.json file. + */ + readAgentProgress(agentId: string): AgentProgress | null { + const agent = this._agents.get(agentId) || this._completedAgents.get(agentId); + if (!agent) return null; + + const progressPath = join(agent.commsDir, 'progress.json'); + if (!existsSync(progressPath)) return null; + + try { + const content = readFileSync(progressPath, 'utf-8'); + return JSON.parse(content) as AgentProgress; + } catch { + return null; + } + } + + /** + * Read messages from an agent's communication channel. + */ + readAgentMessages(agentId: string): SpawnMessage[] { + const agent = this._agents.get(agentId) || this._completedAgents.get(agentId); + if (!agent) return []; + + const messagesDir = join(agent.commsDir, 'messages'); + if (!existsSync(messagesDir)) return []; + + const files = readdirSync(messagesDir) + .filter(f => f.endsWith('.md')) + .sort(); + + const messages: SpawnMessage[] = []; + for (const file of files) { + const match = file.match(/^(\d+)-(parent|agent)\.md$/); + if (!match) continue; + + const content = readFileSync(join(messagesDir, file), 'utf-8'); + messages.push({ + sequence: parseInt(match[1]), + sender: match[2] as 'parent' | 'agent', + content, + sentAt: statSync(join(messagesDir, file)).mtimeMs, + read: true, + }); + } + + return messages; + } + + /** + * Programmatically trigger a spawn without terminal detection. + */ + async triggerSpawn( + taskContent: string, + parentSessionId: string, + parentWorkingDir: string, + parentDepth: number = 0 + ): Promise { + const fallbackId = `agent-${uuidv4().slice(0, 8)}`; + const parsed = parseTaskSpecFile(taskContent, fallbackId); + if (!parsed) return null; + + // If the spec doesn't specify a workingDir, use the parent's + if (!parsed.spec.workingDir) { + parsed.spec.workingDir = parentWorkingDir; + } + + // Write task content to a temp file so setupAgentDirectory can read it + const tempDir = join(this._config.casesDir, '.spawn-tmp'); + mkdirSync(tempDir, { recursive: true }); + const tempFile = join(tempDir, `${parsed.spec.agentId}.md`); + writeFileSync(tempFile, taskContent, 'utf-8'); + + const task: SpawnTask = { + spec: parsed.spec, + instructions: parsed.instructions, + sourceFile: tempFile, + parentSessionId, + depth: parentDepth + 1, + }; + + await this.spawnAgent(task); + return task.spec.agentId; + } + + // ========== Internal Methods ========== + + private getActiveCount(): number { + let count = 0; + for (const agent of this._agents.values()) { + if (agent.status === 'initializing' || agent.status === 'running') { + count++; + } + } + return count; + } + + private enqueueTask(task: SpawnTask): void { + if (this._queue.length >= MAX_QUEUE_LENGTH) { + this.emit('failed', { + agentId: task.spec.agentId, + error: `Queue full (max ${MAX_QUEUE_LENGTH})`, + partialProgress: null, + }); + return; + } + + // Insert by priority (higher priority first) + const priorityOrder = { critical: 0, high: 1, normal: 2, low: 3 }; + const taskPriority = priorityOrder[task.spec.priority]; + let insertIdx = this._queue.length; + for (let i = 0; i < this._queue.length; i++) { + if (priorityOrder[this._queue[i].spec.priority] > taskPriority) { + insertIdx = i; + break; + } + } + this._queue.splice(insertIdx, 0, task); + + this.emit('queued', { + agentId: task.spec.agentId, + name: task.spec.name, + parentSessionId: task.parentSessionId, + position: insertIdx + 1, + }); + this.emitStateUpdate(); + } + + private async spawnAgent(task: SpawnTask): Promise { + if (!this._sessionCreator) return; + + const agentId = task.spec.agentId; + this._totalSpawned++; + if (task.depth > this._maxDepthReached) { + this._maxDepthReached = task.depth; + } + + // Create agent context + const workingDir = join(this._config.casesDir, `spawn-${agentId}`); + const commsDir = join(workingDir, 'spawn-comms'); + + const agent: AgentContext = { + task, + sessionId: null, + workingDir, + commsDir, + parentSessionId: task.parentSessionId, + depth: task.depth, + timeoutTimer: null, + progressTimer: null, + status: 'initializing', + startedAt: null, + tokenBudget: task.spec.maxTokens ?? null, + costBudget: task.spec.maxCost ?? null, + }; + + this._agents.set(agentId, agent); + this.emit('initializing', { agentId, name: task.spec.name, workingDir }); + this.emitStateUpdate(); + + try { + // Setup directory structure + this.setupAgentDirectory(task, workingDir, commsDir); + + // Create session + const { sessionId } = await this._sessionCreator.createAgentSession(workingDir, agentId); + agent.sessionId = sessionId; + agent.status = 'running'; + agent.startedAt = Date.now(); + + this.emit('started', { agentId, name: task.spec.name, sessionId }); + this.emitStateUpdate(); + + // Setup completion listener + this.setupCompletionListener(agent); + + // Setup progress monitor + this.setupProgressMonitor(agent); + + // Setup timeout + this.setupTimeout(agent); + + // Inject initial prompt (short delay to let session initialize) + setTimeout(() => { + if (agent.status === 'running' && this._sessionCreator) { + const prompt = buildInitialPrompt(task); + this._sessionCreator.writeToSession(sessionId, prompt + '\r'); + } + }, 3000); + + } catch (err) { + agent.status = 'failed'; + this._totalFailed++; + this.emit('failed', { agentId, error: getErrorMessage(err), partialProgress: null }); + await this.cleanupAgent(agentId); + } + } + + private setupAgentDirectory(task: SpawnTask, workingDir: string, commsDir: string): void { + // Create directory structure + mkdirSync(workingDir, { recursive: true }); + mkdirSync(commsDir, { recursive: true }); + mkdirSync(join(commsDir, 'messages'), { recursive: true }); + mkdirSync(join(commsDir, 'artifacts'), { recursive: true }); + mkdirSync(join(workingDir, 'workspace'), { recursive: true }); + + // Copy task.md to comms + writeFileSync(join(commsDir, 'task.md'), readFileSync(task.sourceFile, 'utf-8'), 'utf-8'); + + // Write initial progress.json + writeFileSync( + join(commsDir, 'progress.json'), + JSON.stringify(createEmptyAgentProgress(), null, 2), + 'utf-8' + ); + + // Generate and write CLAUDE.md + const claudeMd = generateAgentClaudeMd(task, commsDir, workingDir); + writeFileSync(join(workingDir, 'CLAUDE.md'), claudeMd, 'utf-8'); + + // Symlink context files into workspace + if (task.spec.contextFiles && task.spec.contextFiles.length > 0) { + const parentWorkingDir = this.resolveParentWorkingDir(task); + let fileCount = 0; + + for (const contextFile of task.spec.contextFiles) { + if (fileCount >= MAX_CONTEXT_FILES) break; + + const sourcePath = isAbsolute(contextFile) + ? contextFile + : join(parentWorkingDir, contextFile); + + if (!existsSync(sourcePath)) continue; + + const stat = statSync(sourcePath); + if (stat.size > MAX_CONTEXT_FILE_SIZE) continue; + + const destPath = join(workingDir, 'workspace', contextFile.split('/').pop() || contextFile); + try { + symlinkSync(sourcePath, destPath); + fileCount++; + } catch { + // Ignore symlink errors (e.g., dest already exists) + } + } + } + } + + private resolveParentWorkingDir(task: SpawnTask): string { + // If the task has a specified workingDir, resolve it + if (task.spec.workingDir) { + return isAbsolute(task.spec.workingDir) + ? task.spec.workingDir + : resolve(this._config.casesDir, task.spec.workingDir); + } + // Default: use casesDir + return this._config.casesDir; + } + + private setupCompletionListener(agent: AgentContext): void { + if (!this._sessionCreator || !agent.sessionId) return; + + const handler = (phrase: string) => { + if (phrase === agent.task.spec.completionPhrase) { + this.handleAgentCompletion(agent); + } + }; + + this._completionHandlers.set(agent.task.spec.agentId, handler); + this._sessionCreator.onSessionCompletion(agent.sessionId, handler); + } + + private setupProgressMonitor(agent: AgentContext): void { + if (this._config.progressPollIntervalMs <= 0) return; + + agent.progressTimer = setInterval(() => { + if (agent.status !== 'running') return; + + // Read progress + const progress = this.readAgentProgress(agent.task.spec.agentId); + if (progress) { + this.emit('progress', { agentId: agent.task.spec.agentId, progress }); + } + + // Check resource budgets + this.checkResourceBudgets(agent); + }, this._config.progressPollIntervalMs); + } + + private setupTimeout(agent: AgentContext): void { + const timeoutMs = agent.task.spec.timeoutMinutes * 60 * 1000; + + // Warning at 90% + const warningMs = timeoutMs * 0.9; + setTimeout(() => { + if (agent.status === 'running' && this._sessionCreator && agent.sessionId) { + this._sessionCreator.writeToSession( + agent.sessionId, + 'WARNING: You have less than 10% of your timeout remaining. Please wrap up and write your result.md soon.\r' + ); + } + }, warningMs); + + // Hard timeout + agent.timeoutTimer = setTimeout(() => { + if (agent.status === 'running') { + this.handleAgentTimeout(agent); + } + }, timeoutMs); + } + + private checkResourceBudgets(agent: AgentContext): void { + if (!this._sessionCreator || !agent.sessionId) return; + + // Token budget + if (agent.tokenBudget !== null) { + const tokensUsed = this._sessionCreator.getSessionTokens(agent.sessionId); + const ratio = tokensUsed / agent.tokenBudget; + + if (ratio >= 1.1) { + // Force kill at 110% + this.handleAgentTimeout(agent); + return; + } else if (ratio >= 1.0) { + // Graceful shutdown + this._sessionCreator.writeToSession( + agent.sessionId, + 'You have exceeded your token budget. Write your result.md NOW and output your completion phrase.\r' + ); + } else if (ratio >= BUDGET_WARNING_THRESHOLD) { + this.emit('budgetWarning', { + agentId: agent.task.spec.agentId, + type: 'tokens', + used: tokensUsed, + limit: agent.tokenBudget, + }); + } + } + + // Cost budget + if (agent.costBudget !== null) { + const costUsed = this._sessionCreator.getSessionCost(agent.sessionId); + const ratio = costUsed / agent.costBudget; + + if (ratio >= 1.1) { + this.handleAgentTimeout(agent); + return; + } else if (ratio >= 1.0) { + this._sessionCreator.writeToSession( + agent.sessionId, + 'You have exceeded your cost budget. Write your result.md NOW and output your completion phrase.\r' + ); + } else if (ratio >= BUDGET_WARNING_THRESHOLD) { + this.emit('budgetWarning', { + agentId: agent.task.spec.agentId, + type: 'cost', + used: costUsed, + limit: agent.costBudget, + }); + } + } + } + + private async handleAgentCompletion(agent: AgentContext): Promise { + if (agent.status !== 'running') return; + + agent.status = 'completing'; + this._totalCompleted++; + + // Read result + const result = this.readAgentResult(agent.task.spec.agentId); + if (result) { + // Update token/cost from session + if (this._sessionCreator && agent.sessionId) { + result.tokens.total = this._sessionCreator.getSessionTokens(agent.sessionId); + result.cost = this._sessionCreator.getSessionCost(agent.sessionId); + } + this.emit('completed', { agentId: agent.task.spec.agentId, result }); + } else { + // No result file found, create a minimal one + const minimalResult: SpawnResult = { + status: 'completed', + durationMs: agent.startedAt ? Date.now() - agent.startedAt : 0, + tokens: { input: 0, output: 0, total: 0 }, + cost: 0, + summary: 'Agent completed but no result.md was found', + output: '', + filesChanged: [], + agentId: agent.task.spec.agentId, + completedAt: Date.now(), + }; + this.emit('completed', { agentId: agent.task.spec.agentId, result: minimalResult }); + } + + agent.status = 'completed'; + await this.cleanupAgent(agent.task.spec.agentId); + this.processQueue(); + } + + private async handleAgentTimeout(agent: AgentContext): Promise { + if (agent.status !== 'running') return; + + agent.status = 'timeout'; + this._totalFailed++; + + const elapsed = agent.startedAt ? Date.now() - agent.startedAt : 0; + const limit = agent.task.spec.timeoutMinutes * 60 * 1000; + + this.emit('timeout', { agentId: agent.task.spec.agentId, elapsed, limit }); + await this.cleanupAgent(agent.task.spec.agentId); + this.processQueue(); + } + + private async cleanupAgent(agentId: string): Promise { + const agent = this._agents.get(agentId); + if (!agent) return; + + // Clear timers + if (agent.timeoutTimer) { + clearTimeout(agent.timeoutTimer); + agent.timeoutTimer = null; + } + if (agent.progressTimer) { + clearInterval(agent.progressTimer); + agent.progressTimer = null; + } + + // Remove completion handler + const handler = this._completionHandlers.get(agentId); + if (handler && this._sessionCreator && agent.sessionId) { + this._sessionCreator.removeSessionCompletionHandler(agent.sessionId, handler); + this._completionHandlers.delete(agentId); + } + + // Stop session + if (agent.sessionId && this._sessionCreator) { + try { + await this._sessionCreator.stopSession(agent.sessionId); + } catch { + // Ignore cleanup errors + } + } + + // Move to completed (LRU) + this._agents.delete(agentId); + this._completedAgents.set(agentId, agent); + + // LRU eviction for completed agents + if (this._completedAgents.size > MAX_TRACKED_AGENTS) { + const firstKey = this._completedAgents.keys().next().value; + if (firstKey) this._completedAgents.delete(firstKey); + } + + this.emitStateUpdate(); + } + + private processQueue(): void { + while (this._queue.length > 0 && this.getActiveCount() < this._config.maxConcurrentAgents) { + const task = this._queue.shift(); + if (!task) break; + + // Re-check dependencies + if (task.spec.dependsOn && task.spec.dependsOn.length > 0) { + const unmetDeps = task.spec.dependsOn.filter(depId => { + const dep = this._completedAgents.get(depId); + return !dep || dep.status !== 'completed'; + }); + if (unmetDeps.length > 0) { + // Put back in queue + this._queue.unshift(task); + break; + } + } + + // Spawn (async, don't await to allow multiple spawns) + this.spawnAgent(task).catch(err => { + console.error(`[spawn-orchestrator] Failed to spawn queued agent: ${getErrorMessage(err)}`); + }); + } + } + + private buildStatusReport(agent: AgentContext): AgentStatusReport { + const now = Date.now(); + const elapsed = agent.startedAt ? now - agent.startedAt : 0; + const timeoutMs = agent.task.spec.timeoutMinutes * 60 * 1000; + const timeRemaining = agent.startedAt ? Math.max(0, timeoutMs - elapsed) : timeoutMs; + + let tokensUsed = 0; + let costSoFar = 0; + if (agent.sessionId && this._sessionCreator) { + tokensUsed = this._sessionCreator.getSessionTokens(agent.sessionId); + costSoFar = this._sessionCreator.getSessionCost(agent.sessionId); + } + + // Check dependency status + let dependencyStatus: 'waiting' | 'ready' | 'n/a' = 'n/a'; + if (agent.task.spec.dependsOn && agent.task.spec.dependsOn.length > 0) { + const allMet = agent.task.spec.dependsOn.every(depId => { + const dep = this._completedAgents.get(depId); + return dep && dep.status === 'completed'; + }); + dependencyStatus = allMet ? 'ready' : 'waiting'; + } + + return { + agentId: agent.task.spec.agentId, + name: agent.task.spec.name, + type: agent.task.spec.type, + status: agent.status, + priority: agent.task.spec.priority, + parentSessionId: agent.parentSessionId, + childSessionId: agent.sessionId, + depth: agent.depth, + startedAt: agent.startedAt, + elapsedMs: elapsed, + progress: this.readAgentProgress(agent.task.spec.agentId), + tokensUsed, + costSoFar, + tokenBudget: agent.tokenBudget, + costBudget: agent.costBudget, + timeoutMinutes: agent.task.spec.timeoutMinutes, + timeRemainingMs: timeRemaining, + completionPhrase: agent.task.spec.completionPhrase, + dependsOn: agent.task.spec.dependsOn || [], + dependencyStatus, + }; + } + + private emitStateUpdate(): void { + this.emit('stateUpdate', this.getState()); + } +} diff --git a/src/spawn-types.ts b/src/spawn-types.ts new file mode 100644 index 00000000..f989fa71 --- /dev/null +++ b/src/spawn-types.ts @@ -0,0 +1,698 @@ +/** + * @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; + + // === 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 SpawnDetector) ========== + +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; + /** 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; +} + +// ========== 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; 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 = {}; + + let currentKey: string | null = null; + let currentArray: unknown[] | null = null; + let currentObject: Record | 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)[kvMatch[1]] = parseYamlValue(kvMatch[2]); + } else if (kvMatch && value.includes(':')) { + // New object in array + const obj: Record = {}; + 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 { + 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; + 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; + 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, + }; +} diff --git a/src/types.ts b/src/types.ts index c260fabb..7902b3c5 100644 --- a/src/types.ts +++ b/src/types.ts @@ -82,6 +82,10 @@ export interface SessionState { ralphEnabled?: boolean; /** Ralph completion phrase (if set) */ ralphCompletionPhrase?: string; + /** Parent agent ID if this session is a spawned agent */ + parentAgentId?: string; + /** Child agent IDs spawned by this session */ + childAgentIds?: string[]; } // ========== Task Types ========== @@ -722,3 +726,21 @@ export function getErrorMessage(error: unknown): string { } return 'An unknown error occurred'; } + +// ========== Spawn1337 Protocol Re-exports ========== + +export type { + SpawnPriority, + SpawnResultDelivery, + SpawnStatus, + SpawnTaskSpec, + SpawnTask, + AgentProgress, + SpawnResult, + SpawnMessage, + AgentStatusReport, + SpawnTrackerState, + SpawnOrchestratorConfig, + AgentContext, + SpawnPersistedState, +} from './spawn-types.js'; diff --git a/src/web/server.ts b/src/web/server.ts index d09d6426..fe096f13 100644 --- a/src/web/server.ts +++ b/src/web/server.ts @@ -19,6 +19,8 @@ import { homedir, totalmem, freemem, loadavg, cpus } from 'node:os'; import { EventEmitter } from 'node:events'; import { Session, ClaudeMessage, type BackgroundTask, type RalphTrackerState, type RalphTodoItem } from '../session.js'; import { RespawnController, RespawnConfig, RespawnState } from '../respawn-controller.js'; +import { SpawnOrchestrator, type SessionCreator } from '../spawn-orchestrator.js'; +import type { SpawnOrchestratorConfig } from '../spawn-types.js'; import { ScreenManager } from '../screen-manager.js'; import { getStore } from '../state-store.js'; import { generateClaudeMd } from '../templates/claude-md.js'; @@ -154,6 +156,8 @@ export class WebServer extends EventEmitter { private sseHealthCheckTimer: NodeJS.Timeout | null = null; // Flag to prevent new timers during shutdown private _isStopping: boolean = false; + // Spawn1337 agent orchestrator + private spawnOrchestrator: SpawnOrchestrator; constructor(port: number = 3000) { super(); @@ -174,6 +178,10 @@ export class WebServer extends EventEmitter { this.screenManager.on('statsUpdated', (screens) => { this.broadcast('screen:statsUpdated', screens); }); + + // Initialize spawn orchestrator + this.spawnOrchestrator = new SpawnOrchestrator(); + this.setupSpawnOrchestratorListeners(); } private async setupRoutes(): Promise { @@ -1278,6 +1286,93 @@ export class WebServer extends EventEmitter { this.app.get('/api/system/stats', async () => { return this.getSystemStats(); }); + + // ========== Spawn1337 Agent Protocol Endpoints ========== + + this.app.get('/api/spawn/agents', async () => { + return { success: true, data: this.spawnOrchestrator.getAllAgentStatuses() }; + }); + + this.app.get('/api/spawn/agents/:agentId', async (req) => { + const { agentId } = req.params as { agentId: string }; + const status = this.spawnOrchestrator.getAgentStatus(agentId); + if (!status) { + return createErrorResponse(ApiErrorCode.NOT_FOUND, `Agent ${agentId} not found`); + } + return { success: true, data: status }; + }); + + this.app.get('/api/spawn/agents/:agentId/result', async (req) => { + const { agentId } = req.params as { agentId: string }; + const result = this.spawnOrchestrator.readAgentResult(agentId); + if (!result) { + return createErrorResponse(ApiErrorCode.NOT_FOUND, `No result found for agent ${agentId}`); + } + return { success: true, data: result }; + }); + + this.app.get('/api/spawn/agents/:agentId/progress', async (req) => { + const { agentId } = req.params as { agentId: string }; + const progress = this.spawnOrchestrator.readAgentProgress(agentId); + return { success: true, data: progress }; + }); + + this.app.get('/api/spawn/agents/:agentId/messages', async (req) => { + const { agentId } = req.params as { agentId: string }; + const messages = this.spawnOrchestrator.readAgentMessages(agentId); + return { success: true, data: messages }; + }); + + this.app.post('/api/spawn/agents/:agentId/message', async (req) => { + const { agentId } = req.params as { agentId: string }; + const { content } = req.body as { content: string }; + if (!content) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Message content is required'); + } + await this.spawnOrchestrator.sendMessageToAgent(agentId, content); + return { success: true }; + }); + + this.app.post('/api/spawn/agents/:agentId/cancel', async (req) => { + const { agentId } = req.params as { agentId: string }; + const { reason } = (req.body as { reason?: string }) || {}; + await this.spawnOrchestrator.cancelAgent(agentId, reason || 'Cancelled via API'); + return { success: true }; + }); + + this.app.delete('/api/spawn/agents/:agentId', async (req) => { + const { agentId } = req.params as { agentId: string }; + await this.spawnOrchestrator.cancelAgent(agentId, 'Force killed via API'); + return { success: true }; + }); + + this.app.get('/api/spawn/status', async () => { + return { success: true, data: this.spawnOrchestrator.getState() }; + }); + + this.app.put('/api/spawn/config', async (req) => { + const config = req.body as Partial; + this.spawnOrchestrator.updateConfig(config); + return { success: true, data: this.spawnOrchestrator.config }; + }); + + this.app.post('/api/spawn/trigger', async (req) => { + const { taskContent, parentSessionId, parentWorkingDir } = req.body as { + taskContent: string; + parentSessionId: string; + parentWorkingDir?: string; + }; + if (!taskContent || !parentSessionId) { + return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'taskContent and parentSessionId are required'); + } + const session = this.sessions.get(parentSessionId); + const workingDir = parentWorkingDir || session?.workingDir || process.cwd(); + const agentId = await this.spawnOrchestrator.triggerSpawn(taskContent, parentSessionId, workingDir); + if (!agentId) { + return createErrorResponse(ApiErrorCode.OPERATION_FAILED, 'Failed to parse task spec'); + } + return { success: true, data: { agentId } }; + }); } /** Persists full session state including respawn config to state.json */ @@ -1514,6 +1609,35 @@ export class WebServer extends EventEmitter { session.on('ralphCompletionDetected', (phrase: string) => { this.broadcast('session:ralphCompletionDetected', { sessionId: session.id, phrase }); }); + + // Spawn1337 protocol events + session.on('spawnRequested', (filePath: string) => { + this.spawnOrchestrator.handleSpawnRequest( + filePath, + session.id, + session.workingDir, + session.parentAgentId ? 1 : 0 // Simple depth tracking + ).catch(err => { + console.error(`[Server] Spawn request failed for session ${session.id}:`, getErrorMessage(err)); + }); + }); + + session.on('spawnStatusRequested', (agentId: string) => { + const status = this.spawnOrchestrator.getAgentStatus(agentId); + this.broadcast('spawn:statusResponse', { sessionId: session.id, agentId, status }); + }); + + session.on('spawnCancelRequested', (agentId: string) => { + this.spawnOrchestrator.cancelAgent(agentId, 'Cancelled by parent session').catch(err => { + console.error(`[Server] Spawn cancel failed:`, getErrorMessage(err)); + }); + }); + + session.on('spawnMessageToChild', (agentId: string, content: string) => { + this.spawnOrchestrator.sendMessageToAgent(agentId, content).catch(err => { + console.error(`[Server] Spawn message failed:`, getErrorMessage(err)); + }); + }); } private setupRespawnListeners(sessionId: string, controller: RespawnController): void { @@ -1554,6 +1678,84 @@ export class WebServer extends EventEmitter { }); } + private setupSpawnOrchestratorListeners(): void { + const sessionCreator: SessionCreator = { + createAgentSession: async (workingDir: string, name: string) => { + const session = new Session({ + workingDir, + screenManager: this.screenManager, + useScreen: true, + mode: 'claude', + name: `spawn:${name}`, + }); + + this.sessions.set(session.id, session); + this.setupSessionListeners(session); + session.parentAgentId = name; + + await session.startInteractive(); + this.broadcast('session:created', session.toDetailedState()); + this.broadcast('session:interactive', { id: session.id }); + this.persistSessionState(session); + + // Configure ralph tracker for completion detection + session.ralphTracker.enable(); + + return { sessionId: session.id }; + }, + writeToSession: (sessionId: string, data: string) => { + const session = this.sessions.get(sessionId); + if (session) { + session.writeViaScreen(data); + } + }, + getSessionTokens: (sessionId: string) => { + const session = this.sessions.get(sessionId); + return session ? session.totalTokens : 0; + }, + getSessionCost: (sessionId: string) => { + const session = this.sessions.get(sessionId); + return session ? session.totalCost : 0; + }, + stopSession: async (sessionId: string) => { + const session = this.sessions.get(sessionId); + if (session) { + await session.stop(); + this.sessions.delete(sessionId); + this.broadcast('session:deleted', { id: sessionId }); + this.persistSessionState(session); + } + }, + onSessionCompletion: (sessionId: string, handler: (phrase: string) => void) => { + const session = this.sessions.get(sessionId); + if (session) { + session.on('ralphCompletionDetected', handler); + } + }, + removeSessionCompletionHandler: (sessionId: string, handler: (phrase: string) => void) => { + const session = this.sessions.get(sessionId); + if (session) { + session.off('ralphCompletionDetected', handler); + } + }, + }; + + this.spawnOrchestrator.setSessionCreator(sessionCreator); + + // Forward orchestrator events as SSE broadcasts + this.spawnOrchestrator.on('queued', (data) => this.broadcast('spawn:queued', data)); + this.spawnOrchestrator.on('initializing', (data) => this.broadcast('spawn:initializing', data)); + this.spawnOrchestrator.on('started', (data) => this.broadcast('spawn:started', data)); + this.spawnOrchestrator.on('progress', (data) => this.broadcast('spawn:progress', data)); + this.spawnOrchestrator.on('message', (data) => this.broadcast('spawn:message', data)); + this.spawnOrchestrator.on('completed', (data) => this.broadcast('spawn:completed', data)); + this.spawnOrchestrator.on('failed', (data) => this.broadcast('spawn:failed', data)); + this.spawnOrchestrator.on('timeout', (data) => this.broadcast('spawn:timeout', data)); + this.spawnOrchestrator.on('cancelled', (data) => this.broadcast('spawn:cancelled', data)); + this.spawnOrchestrator.on('budgetWarning', (data) => this.broadcast('spawn:budgetWarning', data)); + this.spawnOrchestrator.on('stateUpdate', (data) => this.broadcast('spawn:stateUpdate', data)); + } + private setupTimedRespawn(sessionId: string, durationMinutes: number): void { // Clear existing timer if any const existing = this.respawnTimers.get(sessionId); @@ -2156,6 +2358,10 @@ export class WebServer extends EventEmitter { } this.respawnControllers.clear(); + // Stop spawn orchestrator and all agents + await this.spawnOrchestrator.stopAll(); + this.spawnOrchestrator.removeAllListeners(); + // Stop all scheduled runs first (they have their own session cleanup) for (const [id] of this.scheduledRuns) { await this.stopScheduledRun(id); diff --git a/test/spawn-detector.test.ts b/test/spawn-detector.test.ts new file mode 100644 index 00000000..2fba749f --- /dev/null +++ b/test/spawn-detector.test.ts @@ -0,0 +1,249 @@ +import { describe, it, expect, beforeEach, vi } from 'vitest'; +import { SpawnDetector } from '../src/spawn-detector.js'; + +describe('SpawnDetector', () => { + let detector: SpawnDetector; + + beforeEach(() => { + detector = new SpawnDetector(); + }); + + describe('Initialization', () => { + it('should start disabled', () => { + expect(detector.enabled).toBe(false); + }); + + it('should start with initial state', () => { + const state = detector.state; + expect(state.enabled).toBe(false); + expect(state.activeCount).toBe(0); + expect(state.totalSpawned).toBe(0); + }); + }); + + describe('Auto-Enable', () => { + it('should auto-enable on spawn tag detection', () => { + detector.processTerminalData('task.md\n'); + expect(detector.enabled).toBe(true); + }); + + it('should not enable on unrelated data', () => { + detector.processTerminalData('Hello world\nsome output\n'); + expect(detector.enabled).toBe(false); + }); + + it('should auto-enable on status tag', () => { + detector.processTerminalData('\n'); + expect(detector.enabled).toBe(true); + }); + + it('should auto-enable on cancel tag', () => { + detector.processTerminalData('\n'); + expect(detector.enabled).toBe(true); + }); + }); + + describe('Spawn Request Detection', () => { + it('should detect spawn tag and emit event', () => { + const handler = vi.fn(); + detector.on('spawnRequested', handler); + + detector.processTerminalData('tasks/auth-explore.md\n'); + + expect(handler).toHaveBeenCalledWith('tasks/auth-explore.md', expect.any(String)); + }); + + it('should detect spawn tag with path containing slashes', () => { + const handler = vi.fn(); + detector.on('spawnRequested', handler); + + detector.processTerminalData('src/tasks/deep/nested.md\n'); + + expect(handler).toHaveBeenCalledWith('src/tasks/deep/nested.md', expect.any(String)); + }); + + it('should handle spawn tag with surrounding text', () => { + const handler = vi.fn(); + detector.on('spawnRequested', handler); + + detector.processTerminalData('Starting agent: task.md done\n'); + + expect(handler).toHaveBeenCalledWith('task.md', expect.any(String)); + }); + + it('should increment totalSpawned count', () => { + detector.processTerminalData('task1.md\n'); + detector.processTerminalData('task2.md\n'); + detector.flushPendingEvents(); + + expect(detector.state.totalSpawned).toBe(2); + }); + + it('should handle multiple spawn tags in one chunk', () => { + const handler = vi.fn(); + detector.on('spawnRequested', handler); + + detector.processTerminalData( + 'task1.md\ntask2.md\n' + ); + + expect(handler).toHaveBeenCalledTimes(2); + }); + + it('should strip ANSI codes before parsing', () => { + const handler = vi.fn(); + detector.on('spawnRequested', handler); + + detector.processTerminalData('\x1b[32mtask.md\x1b[0m\n'); + + expect(handler).toHaveBeenCalledWith('task.md', expect.any(String)); + }); + + it('should trim whitespace from file path', () => { + const handler = vi.fn(); + detector.on('spawnRequested', handler); + + detector.processTerminalData(' task.md \n'); + + expect(handler).toHaveBeenCalledWith('task.md', expect.any(String)); + }); + }); + + describe('Status Request Detection', () => { + it('should detect status query and emit event', () => { + const handler = vi.fn(); + detector.on('statusRequested', handler); + + detector.processTerminalData('\n'); + + expect(handler).toHaveBeenCalledWith('auth-001'); + }); + + it('should handle status with complex agent ID', () => { + const handler = vi.fn(); + detector.on('statusRequested', handler); + + detector.processTerminalData('\n'); + + expect(handler).toHaveBeenCalledWith('my-complex-agent-123'); + }); + }); + + describe('Cancel Request Detection', () => { + it('should detect cancel request and emit event', () => { + const handler = vi.fn(); + detector.on('cancelRequested', handler); + + detector.processTerminalData('\n'); + + expect(handler).toHaveBeenCalledWith('test-agent'); + }); + }); + + describe('Message Detection', () => { + it('should detect single-line message', () => { + const handler = vi.fn(); + detector.on('messageToChild', handler); + + detector.processTerminalData( + 'Focus on the JWT flow\n' + ); + + expect(handler).toHaveBeenCalledWith('agent-001', 'Focus on the JWT flow'); + }); + + it('should detect multi-line message via checkMultiLinePatterns', () => { + const handler = vi.fn(); + detector.on('messageToChild', handler); + + const multiline = 'Line 1\nLine 2\nLine 3'; + detector.processTerminalData(multiline + '\n'); + + expect(handler).toHaveBeenCalledWith('agent-001', 'Line 1\nLine 2\nLine 3'); + }); + }); + + describe('Enable/Disable', () => { + it('should emit stateUpdate when enabled', async () => { + const handler = vi.fn(); + detector.on('stateUpdate', handler); + + detector.enable(); + detector.flushPendingEvents(); + + expect(handler).toHaveBeenCalled(); + expect(handler.mock.calls[0][0].enabled).toBe(true); + }); + + it('should emit stateUpdate when disabled', () => { + detector.enable(); + detector.flushPendingEvents(); + + const handler = vi.fn(); + detector.on('stateUpdate', handler); + + detector.disable(); + detector.flushPendingEvents(); + + expect(handler).toHaveBeenCalled(); + expect(handler.mock.calls[0][0].enabled).toBe(false); + }); + }); + + describe('Reset', () => { + it('should reset all state', () => { + detector.processTerminalData('task.md\n'); + detector.flushPendingEvents(); + expect(detector.enabled).toBe(true); + + detector.reset(); + + expect(detector.enabled).toBe(false); + expect(detector.state.totalSpawned).toBe(0); + }); + }); + + describe('State Update', () => { + it('should allow external state updates', () => { + detector.updateState({ activeCount: 3, queuedCount: 2 }); + detector.flushPendingEvents(); + + expect(detector.state.activeCount).toBe(3); + expect(detector.state.queuedCount).toBe(2); + }); + }); + + describe('Line Buffer', () => { + it('should handle data split across chunks', () => { + const handler = vi.fn(); + detector.on('spawnRequested', handler); + + // Send in two chunks (split in the middle of the tag) + detector.processTerminalData('task.md\n'); + + expect(handler).toHaveBeenCalledWith('task.md', expect.any(String)); + }); + + it('should handle partial lines without emitting', () => { + const handler = vi.fn(); + detector.on('spawnRequested', handler); + + // No newline - data stays in buffer + detector.processTerminalData('task.md'); + + // Should not emit until newline + expect(handler).not.toHaveBeenCalled(); + }); + + it('should emit on subsequent newline', () => { + const handler = vi.fn(); + detector.on('spawnRequested', handler); + + detector.processTerminalData('task.md'); + detector.processTerminalData('\n'); + + expect(handler).toHaveBeenCalledWith('task.md', expect.any(String)); + }); + }); +}); diff --git a/test/spawn-orchestrator.test.ts b/test/spawn-orchestrator.test.ts new file mode 100644 index 00000000..8c6a3128 --- /dev/null +++ b/test/spawn-orchestrator.test.ts @@ -0,0 +1,496 @@ +import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'; +import { SpawnOrchestrator, type SessionCreator } from '../src/spawn-orchestrator.js'; +import { mkdirSync, writeFileSync, existsSync, rmSync, readFileSync } from 'node:fs'; +import { join } from 'node:path'; +import { tmpdir } from 'node:os'; + +/** + * SpawnOrchestrator Tests + * + * Tests the full lifecycle management of spawned agents. + * Uses a temporary directory and mock session creator. + */ + +describe('SpawnOrchestrator', () => { + let orchestrator: SpawnOrchestrator; + let testDir: string; + let mockSessionCreator: SessionCreator; + let completionHandlers: Map void>; + + beforeEach(() => { + testDir = join(tmpdir(), `spawn-test-${Date.now()}-${Math.random().toString(36).slice(2)}`); + mkdirSync(testDir, { recursive: true }); + + completionHandlers = new Map(); + + mockSessionCreator = { + createAgentSession: vi.fn().mockResolvedValue({ sessionId: `session-${Date.now()}` }), + writeToSession: vi.fn(), + getSessionTokens: vi.fn().mockReturnValue(0), + getSessionCost: vi.fn().mockReturnValue(0), + stopSession: vi.fn().mockResolvedValue(undefined), + onSessionCompletion: vi.fn().mockImplementation((sessionId, handler) => { + completionHandlers.set(sessionId, handler); + }), + removeSessionCompletionHandler: vi.fn().mockImplementation((sessionId) => { + completionHandlers.delete(sessionId); + }), + }; + + orchestrator = new SpawnOrchestrator({ + casesDir: testDir, + maxConcurrentAgents: 3, + maxSpawnDepth: 2, + defaultTimeoutMinutes: 5, + maxTimeoutMinutes: 10, + progressPollIntervalMs: 60000, // Long interval to avoid interference + }); + + orchestrator.setSessionCreator(mockSessionCreator); + }); + + afterEach(() => { + // Stop all agents and clear timers + orchestrator.stopAll().catch(() => {}); + orchestrator.removeAllListeners(); + // Clean up test directory + if (existsSync(testDir)) { + rmSync(testDir, { recursive: true, force: true }); + } + }); + + function createTaskFile(dir: string, filename: string, content: string): string { + const filePath = join(dir, filename); + mkdirSync(dir, { recursive: true }); + writeFileSync(filePath, content); + return filePath; + } + + const basicTaskContent = `--- +agentId: test-agent-001 +name: Test Agent +type: explore +priority: normal +timeoutMinutes: 5 +completionPhrase: TEST_DONE +canModifyParentFiles: false +--- + +# Test Task + +Do a simple test.`; + + describe('Configuration', () => { + it('should use provided config', () => { + expect(orchestrator.config.maxConcurrentAgents).toBe(3); + expect(orchestrator.config.maxSpawnDepth).toBe(2); + }); + + it('should update config', () => { + orchestrator.updateConfig({ maxConcurrentAgents: 10 }); + expect(orchestrator.config.maxConcurrentAgents).toBe(10); + }); + }); + + describe('handleSpawnRequest', () => { + it('should reject when no session creator is set', async () => { + const noCreator = new SpawnOrchestrator({ casesDir: testDir }); + const failHandler = vi.fn(); + noCreator.on('failed', failHandler); + + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await noCreator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + // Should not crash, just log error + expect(failHandler).not.toHaveBeenCalled(); // Silent failure with console.error + }); + + it('should fail when task file does not exist', async () => { + const failHandler = vi.fn(); + orchestrator.on('failed', failHandler); + + await orchestrator.handleSpawnRequest('nonexistent.md', 'parent-session', testDir); + + expect(failHandler).toHaveBeenCalledWith( + expect.objectContaining({ error: expect.stringContaining('not found') }) + ); + }); + + it('should fail when task file cannot be parsed', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'bad.md', 'No frontmatter here'); + const failHandler = vi.fn(); + orchestrator.on('failed', failHandler); + + await orchestrator.handleSpawnRequest('bad.md', 'parent-session', parentDir); + + expect(failHandler).toHaveBeenCalledWith( + expect.objectContaining({ error: expect.stringContaining('parse') }) + ); + }); + + it('should reject when max depth exceeded', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + const failHandler = vi.fn(); + orchestrator.on('failed', failHandler); + + // Max depth is 2, so parentDepth=2 means child would be 3 (exceeds) + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir, 2); + + expect(failHandler).toHaveBeenCalledWith( + expect.objectContaining({ error: expect.stringContaining('depth') }) + ); + }); + + it('should spawn agent and create directory structure', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + // Wait for async initialization + await vi.waitFor(() => { + expect(mockSessionCreator.createAgentSession).toHaveBeenCalled(); + }); + + // Check directory was created + const agentDir = join(testDir, 'spawn-test-agent-001'); + expect(existsSync(agentDir)).toBe(true); + expect(existsSync(join(agentDir, 'CLAUDE.md'))).toBe(true); + expect(existsSync(join(agentDir, 'spawn-comms'))).toBe(true); + expect(existsSync(join(agentDir, 'spawn-comms', 'task.md'))).toBe(true); + expect(existsSync(join(agentDir, 'spawn-comms', 'progress.json'))).toBe(true); + expect(existsSync(join(agentDir, 'spawn-comms', 'messages'))).toBe(true); + }); + + it('should generate proper CLAUDE.md for agent', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + expect(mockSessionCreator.createAgentSession).toHaveBeenCalled(); + }); + + const agentDir = join(testDir, 'spawn-test-agent-001'); + const claudeMd = readFileSync(join(agentDir, 'CLAUDE.md'), 'utf-8'); + expect(claudeMd).toContain('Agent: Test Agent'); + expect(claudeMd).toContain('test-agent-001'); + expect(claudeMd).toContain('TEST_DONE'); + expect(claudeMd).toContain('# Test Task'); + }); + + it('should emit initializing and started events', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + const initHandler = vi.fn(); + const startHandler = vi.fn(); + orchestrator.on('initializing', initHandler); + orchestrator.on('started', startHandler); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + expect(startHandler).toHaveBeenCalled(); + }); + + expect(initHandler).toHaveBeenCalledWith( + expect.objectContaining({ agentId: 'test-agent-001', name: 'Test Agent' }) + ); + expect(startHandler).toHaveBeenCalledWith( + expect.objectContaining({ agentId: 'test-agent-001', name: 'Test Agent' }) + ); + }); + + it('should enforce timeout limits', async () => { + const parentDir = join(testDir, 'parent'); + const content = basicTaskContent.replace('timeoutMinutes: 5', 'timeoutMinutes: 999'); + createTaskFile(parentDir, 'task.md', content); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + const status = orchestrator.getAgentStatus('test-agent-001'); + expect(status).not.toBeNull(); + expect(status!.timeoutMinutes).toBe(10); // Capped at maxTimeoutMinutes + }); + }); + }); + + describe('Queue Management', () => { + it('should queue agents when concurrency limit reached', async () => { + const parentDir = join(testDir, 'parent'); + + // Create 4 tasks (limit is 3) + for (let i = 1; i <= 4; i++) { + const content = basicTaskContent + .replace('test-agent-001', `agent-${i}`) + .replace('Test Agent', `Agent ${i}`); + createTaskFile(parentDir, `task${i}.md`, content); + } + + const queueHandler = vi.fn(); + orchestrator.on('queued', queueHandler); + + // Spawn 4 agents + for (let i = 1; i <= 4; i++) { + await orchestrator.handleSpawnRequest(`task${i}.md`, 'parent', parentDir); + } + + // Wait for first 3 to start + await vi.waitFor(() => { + expect(mockSessionCreator.createAgentSession).toHaveBeenCalledTimes(3); + }); + + // 4th should be queued + expect(queueHandler).toHaveBeenCalledWith( + expect.objectContaining({ agentId: 'agent-4' }) + ); + }); + + it('should order queue by priority', async () => { + const parentDir = join(testDir, 'parent'); + + // Fill concurrency first + for (let i = 1; i <= 3; i++) { + const content = basicTaskContent + .replace('test-agent-001', `filler-${i}`) + .replace('Test Agent', `Filler ${i}`); + createTaskFile(parentDir, `filler${i}.md`, content); + await orchestrator.handleSpawnRequest(`filler${i}.md`, 'parent', parentDir); + } + + await vi.waitFor(() => { + expect(mockSessionCreator.createAgentSession).toHaveBeenCalledTimes(3); + }); + + // Now add low and high priority + const lowContent = basicTaskContent + .replace('test-agent-001', 'low-agent') + .replace('priority: normal', 'priority: low'); + createTaskFile(parentDir, 'low.md', lowContent); + + const highContent = basicTaskContent + .replace('test-agent-001', 'high-agent') + .replace('priority: normal', 'priority: critical'); + createTaskFile(parentDir, 'high.md', highContent); + + await orchestrator.handleSpawnRequest('low.md', 'parent', parentDir); + await orchestrator.handleSpawnRequest('high.md', 'parent', parentDir); + + // State should show high priority first in queue + const state = orchestrator.getState(); + expect(state.queuedCount).toBe(2); + }); + }); + + describe('cancelAgent', () => { + it('should cancel a running agent', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + expect(orchestrator.getAgentStatus('test-agent-001')).not.toBeNull(); + }); + + const cancelHandler = vi.fn(); + orchestrator.on('cancelled', cancelHandler); + + await orchestrator.cancelAgent('test-agent-001', 'User cancelled'); + + expect(cancelHandler).toHaveBeenCalledWith( + expect.objectContaining({ agentId: 'test-agent-001', reason: 'User cancelled' }) + ); + }); + + it('should cancel a queued agent', async () => { + const parentDir = join(testDir, 'parent'); + + // Fill concurrency + for (let i = 1; i <= 3; i++) { + const content = basicTaskContent + .replace('test-agent-001', `filler-${i}`) + .replace('Test Agent', `Filler ${i}`); + createTaskFile(parentDir, `filler${i}.md`, content); + await orchestrator.handleSpawnRequest(`filler${i}.md`, 'parent', parentDir); + } + + // Add one more (queued) + const queuedContent = basicTaskContent.replace('test-agent-001', 'queued-agent'); + createTaskFile(parentDir, 'queued.md', queuedContent); + await orchestrator.handleSpawnRequest('queued.md', 'parent', parentDir); + + const cancelHandler = vi.fn(); + orchestrator.on('cancelled', cancelHandler); + + await orchestrator.cancelAgent('queued-agent'); + + expect(cancelHandler).toHaveBeenCalledWith( + expect.objectContaining({ agentId: 'queued-agent' }) + ); + }); + }); + + describe('sendMessageToAgent', () => { + it('should write message file to comms directory', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + expect(orchestrator.getAgentStatus('test-agent-001')).not.toBeNull(); + }); + + await orchestrator.sendMessageToAgent('test-agent-001', 'Focus on JWT'); + + const messagesDir = join(testDir, 'spawn-test-agent-001', 'spawn-comms', 'messages'); + expect(existsSync(join(messagesDir, '001-parent.md'))).toBe(true); + + const content = readFileSync(join(messagesDir, '001-parent.md'), 'utf-8'); + expect(content).toBe('Focus on JWT'); + }); + }); + + describe('getAgentStatus', () => { + it('should return null for unknown agent', () => { + expect(orchestrator.getAgentStatus('nonexistent')).toBeNull(); + }); + + it('should return status for active agent', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + const status = orchestrator.getAgentStatus('test-agent-001'); + expect(status).not.toBeNull(); + expect(status!.status).toBe('running'); + expect(status!.name).toBe('Test Agent'); + expect(status!.completionPhrase).toBe('TEST_DONE'); + }); + }); + }); + + describe('getState', () => { + it('should return complete orchestrator state', () => { + const state = orchestrator.getState(); + expect(state.enabled).toBe(true); + expect(state.activeCount).toBe(0); + expect(state.queuedCount).toBe(0); + expect(state.totalSpawned).toBe(0); + expect(state.agents).toEqual([]); + }); + + it('should update counts after spawn', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + const state = orchestrator.getState(); + expect(state.activeCount).toBe(1); + expect(state.totalSpawned).toBe(1); + }); + }); + }); + + describe('readAgentMessages', () => { + it('should return empty for unknown agent', () => { + expect(orchestrator.readAgentMessages('nonexistent')).toEqual([]); + }); + + it('should read messages from comms directory', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + expect(orchestrator.getAgentStatus('test-agent-001')).not.toBeNull(); + }); + + // Write a message + await orchestrator.sendMessageToAgent('test-agent-001', 'Hello agent'); + + const messages = orchestrator.readAgentMessages('test-agent-001'); + expect(messages).toHaveLength(1); + expect(messages[0].sender).toBe('parent'); + expect(messages[0].content).toBe('Hello agent'); + expect(messages[0].sequence).toBe(1); + }); + }); + + describe('triggerSpawn', () => { + it('should spawn from content string', async () => { + const agentId = await orchestrator.triggerSpawn( + basicTaskContent, + 'parent-session', + testDir + ); + + expect(agentId).toBe('test-agent-001'); + + await vi.waitFor(() => { + expect(mockSessionCreator.createAgentSession).toHaveBeenCalled(); + }); + }); + + it('should return null for unparseable content', async () => { + const agentId = await orchestrator.triggerSpawn( + 'Not valid YAML frontmatter', + 'parent-session', + testDir + ); + + expect(agentId).toBeNull(); + }); + }); + + describe('stopAll', () => { + it('should stop all active agents', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + expect(orchestrator.getState().activeCount).toBe(1); + }); + + await orchestrator.stopAll(); + + expect(orchestrator.getState().activeCount).toBe(0); + }); + }); + + describe('getPersistedState', () => { + it('should return serializable state', async () => { + const parentDir = join(testDir, 'parent'); + createTaskFile(parentDir, 'task.md', basicTaskContent); + + await orchestrator.handleSpawnRequest('task.md', 'parent-session', parentDir); + + await vi.waitFor(() => { + expect(orchestrator.getState().activeCount).toBe(1); + }); + + const persisted = orchestrator.getPersistedState(); + expect(persisted.config).toBeDefined(); + expect(persisted.agents['test-agent-001']).toBeDefined(); + expect(persisted.agents['test-agent-001'].completionPhrase).toBe('TEST_DONE'); + + // Should be JSON-serializable + expect(() => JSON.stringify(persisted)).not.toThrow(); + }); + }); +}); diff --git a/test/spawn-types.test.ts b/test/spawn-types.test.ts new file mode 100644 index 00000000..6c2daf6a --- /dev/null +++ b/test/spawn-types.test.ts @@ -0,0 +1,395 @@ +import { describe, it, expect } from 'vitest'; +import { + parseYamlFrontmatter, + parseTaskSpecFile, + createDefaultSpawnTaskSpec, + createEmptyAgentProgress, + createInitialSpawnTrackerState, + createDefaultOrchestratorConfig, + serializeSpawnResult, + parseSpawnResult, + AGENT_NAME_MAX_LENGTH, +} from '../src/spawn-types.js'; + +describe('spawn-types', () => { + describe('parseYamlFrontmatter', () => { + it('should parse basic frontmatter', () => { + const content = `--- +name: Test Agent +type: explore +priority: high +--- + +# Task Body + +Do something useful.`; + + const result = parseYamlFrontmatter(content); + expect(result).not.toBeNull(); + expect(result!.frontmatter.name).toBe('Test Agent'); + expect(result!.frontmatter.type).toBe('explore'); + expect(result!.frontmatter.priority).toBe('high'); + expect(result!.body).toContain('# Task Body'); + expect(result!.body).toContain('Do something useful.'); + }); + + it('should parse numbers', () => { + const content = `--- +timeoutMinutes: 30 +maxCost: 0.50 +maxTokens: 150000 +--- +body`; + + const result = parseYamlFrontmatter(content); + expect(result!.frontmatter.timeoutMinutes).toBe(30); + expect(result!.frontmatter.maxCost).toBe(0.5); + expect(result!.frontmatter.maxTokens).toBe(150000); + }); + + it('should parse booleans', () => { + const content = `--- +canModifyParentFiles: true +enabled: false +--- +body`; + + const result = parseYamlFrontmatter(content); + expect(result!.frontmatter.canModifyParentFiles).toBe(true); + expect(result!.frontmatter.enabled).toBe(false); + }); + + it('should parse inline arrays', () => { + const content = `--- +dependsOn: [agent-1, agent-2] +--- +body`; + + const result = parseYamlFrontmatter(content); + expect(result!.frontmatter.dependsOn).toEqual(['agent-1', 'agent-2']); + }); + + it('should parse block arrays', () => { + const content = `--- +contextFiles: + - src/auth.ts + - src/middleware.ts + - src/types.ts +--- +body`; + + const result = parseYamlFrontmatter(content); + expect(result!.frontmatter.contextFiles).toEqual(['src/auth.ts', 'src/middleware.ts', 'src/types.ts']); + }); + + it('should parse empty arrays', () => { + const content = `--- +dependsOn: [] +--- +body`; + + const result = parseYamlFrontmatter(content); + expect(result!.frontmatter.dependsOn).toEqual([]); + }); + + it('should parse quoted strings', () => { + const content = `--- +name: "Test Agent" +completionPhrase: 'AUTH_DONE' +--- +body`; + + const result = parseYamlFrontmatter(content); + expect(result!.frontmatter.name).toBe('Test Agent'); + expect(result!.frontmatter.completionPhrase).toBe('AUTH_DONE'); + }); + + it('should handle null values', () => { + const content = `--- +value1: null +value2: ~ +--- +body`; + + const result = parseYamlFrontmatter(content); + expect(result!.frontmatter.value1).toBeNull(); + expect(result!.frontmatter.value2).toBeNull(); + }); + + it('should handle comments', () => { + const content = `--- +# This is a comment +name: Test +# Another comment +type: explore +--- +body`; + + const result = parseYamlFrontmatter(content); + expect(result!.frontmatter.name).toBe('Test'); + expect(result!.frontmatter.type).toBe('explore'); + }); + + it('should return null if no frontmatter delimiters', () => { + const content = `No frontmatter here\nJust plain text`; + expect(parseYamlFrontmatter(content)).toBeNull(); + }); + + it('should return null if missing closing delimiter', () => { + const content = `---\nname: Test\nNo closing delimiter`; + expect(parseYamlFrontmatter(content)).toBeNull(); + }); + + it('should handle nested objects', () => { + const content = `--- +env: + NODE_ENV: production + DEBUG: true +--- +body`; + + const result = parseYamlFrontmatter(content); + expect(result!.frontmatter.env).toEqual({ NODE_ENV: 'production', DEBUG: true }); + }); + }); + + describe('parseTaskSpecFile', () => { + it('should parse a complete task spec', () => { + const content = `--- +agentId: auth-explorer-001 +name: Authentication Explorer +type: explore +priority: high +canModifyParentFiles: false +maxTokens: 150000 +maxCost: 0.50 +timeoutMinutes: 15 +resultDelivery: both +completionPhrase: AUTH_EXPLORE_DONE +progressIntervalSeconds: 30 +outputFormat: structured +successCriteria: "Document all auth patterns" +--- + +# Task: Explore Authentication + +Analyze the auth system.`; + + const result = parseTaskSpecFile(content, 'fallback-id'); + expect(result).not.toBeNull(); + expect(result!.spec.agentId).toBe('auth-explorer-001'); + expect(result!.spec.name).toBe('Authentication Explorer'); + expect(result!.spec.type).toBe('explore'); + expect(result!.spec.priority).toBe('high'); + expect(result!.spec.canModifyParentFiles).toBe(false); + expect(result!.spec.maxTokens).toBe(150000); + expect(result!.spec.maxCost).toBe(0.5); + expect(result!.spec.timeoutMinutes).toBe(15); + expect(result!.spec.completionPhrase).toBe('AUTH_EXPLORE_DONE'); + expect(result!.spec.outputFormat).toBe('structured'); + expect(result!.instructions).toContain('# Task: Explore Authentication'); + }); + + it('should use defaults for missing fields', () => { + const content = `--- +name: Simple Agent +--- +Do something.`; + + const result = parseTaskSpecFile(content, 'my-fallback'); + expect(result).not.toBeNull(); + expect(result!.spec.agentId).toBe('my-fallback'); + expect(result!.spec.type).toBe('general'); + expect(result!.spec.priority).toBe('normal'); + expect(result!.spec.timeoutMinutes).toBe(30); + expect(result!.spec.resultDelivery).toBe('both'); + expect(result!.spec.outputFormat).toBe('markdown'); + expect(result!.spec.canModifyParentFiles).toBe(false); + }); + + it('should truncate long names', () => { + const longName = 'A'.repeat(100); + const content = `--- +name: ${longName} +--- +body`; + + const result = parseTaskSpecFile(content, 'id'); + expect(result!.spec.name.length).toBe(AGENT_NAME_MAX_LENGTH); + }); + + it('should validate type values', () => { + const content = `--- +type: invalid_type +--- +body`; + + const result = parseTaskSpecFile(content, 'id'); + expect(result!.spec.type).toBe('general'); // Falls back to default + }); + + it('should validate priority values', () => { + const content = `--- +priority: super_high +--- +body`; + + const result = parseTaskSpecFile(content, 'id'); + expect(result!.spec.priority).toBe('normal'); // Falls back to default + }); + + it('should return null for content without frontmatter', () => { + const content = 'No frontmatter at all'; + expect(parseTaskSpecFile(content, 'id')).toBeNull(); + }); + + it('should parse contextFiles array', () => { + const content = `--- +contextFiles: + - src/auth.ts + - src/types.ts +--- +body`; + + const result = parseTaskSpecFile(content, 'id'); + expect(result!.spec.contextFiles).toEqual(['src/auth.ts', 'src/types.ts']); + }); + + it('should parse dependsOn array', () => { + const content = `--- +dependsOn: + - agent-1 + - agent-2 +--- +body`; + + const result = parseTaskSpecFile(content, 'id'); + expect(result!.spec.dependsOn).toEqual(['agent-1', 'agent-2']); + }); + }); + + describe('Factory Functions', () => { + it('createDefaultSpawnTaskSpec should generate valid defaults', () => { + const spec = createDefaultSpawnTaskSpec('my-agent'); + expect(spec.agentId).toBe('my-agent'); + expect(spec.name).toBe('my-agent'); + expect(spec.type).toBe('general'); + expect(spec.priority).toBe('normal'); + expect(spec.timeoutMinutes).toBe(30); + expect(spec.completionPhrase).toContain('MY_AGENT'); + expect(spec.completionPhrase).toContain('DONE'); + }); + + it('createDefaultSpawnTaskSpec should sanitize agent ID for completion phrase', () => { + const spec = createDefaultSpawnTaskSpec('my-agent-123'); + expect(spec.completionPhrase).toBe('AGENT_MY_AGENT_123_DONE'); + }); + + it('createEmptyAgentProgress should create valid progress', () => { + const progress = createEmptyAgentProgress(); + expect(progress.phase).toBe('initializing'); + expect(progress.percentComplete).toBe(0); + expect(progress.filesModified).toEqual([]); + expect(progress.tokensUsed).toBe(0); + }); + + it('createInitialSpawnTrackerState should create valid state', () => { + const state = createInitialSpawnTrackerState(); + expect(state.enabled).toBe(false); + expect(state.activeCount).toBe(0); + expect(state.agents).toEqual([]); + }); + + it('createDefaultOrchestratorConfig should create valid config', () => { + const config = createDefaultOrchestratorConfig(); + expect(config.maxConcurrentAgents).toBe(5); + expect(config.maxSpawnDepth).toBe(3); + expect(config.defaultTimeoutMinutes).toBe(30); + expect(config.maxTimeoutMinutes).toBe(120); + expect(config.progressPollIntervalMs).toBe(5000); + }); + }); + + describe('serializeSpawnResult', () => { + it('should serialize a completed result', () => { + const result = { + status: 'completed' as const, + durationMs: 60000, + tokens: { input: 1000, output: 500, total: 1500 }, + cost: 0.05, + summary: 'Task completed successfully', + output: '## Result\n\nDetailed output here.', + filesChanged: [ + { path: 'src/auth.ts', action: 'modified' as const, summary: 'Added validation' }, + ], + agentId: 'test-agent', + completedAt: 1700000000000, + }; + + const serialized = serializeSpawnResult(result); + expect(serialized).toContain('status: completed'); + expect(serialized).toContain('summary: "Task completed successfully"'); + expect(serialized).toContain('agentId: test-agent'); + expect(serialized).toContain('path: src/auth.ts'); + expect(serialized).toContain('## Result'); + }); + + it('should handle empty filesChanged', () => { + const result = { + status: 'failed' as const, + error: 'Something went wrong', + durationMs: 5000, + tokens: { input: 100, output: 50, total: 150 }, + cost: 0.01, + summary: 'Failed', + output: 'Error details', + filesChanged: [], + agentId: 'test', + completedAt: Date.now(), + }; + + const serialized = serializeSpawnResult(result); + expect(serialized).toContain('status: failed'); + expect(serialized).toContain('filesChanged: []'); + }); + }); + + describe('parseSpawnResult', () => { + it('should parse a result file', () => { + const content = `--- +status: completed +summary: "Found 3 auth patterns" +cost: 0.25 +--- + +## Analysis + +Detailed findings here.`; + + const result = parseSpawnResult(content, 'agent-001', 60000); + expect(result).not.toBeNull(); + expect(result!.status).toBe('completed'); + expect(result!.summary).toBe('Found 3 auth patterns'); + expect(result!.cost).toBe(0.25); + expect(result!.output).toContain('## Analysis'); + expect(result!.agentId).toBe('agent-001'); + }); + + it('should handle missing fields with defaults', () => { + const content = `--- +status: completed +--- +output`; + + const result = parseSpawnResult(content, 'agent', 30000); + expect(result!.durationMs).toBe(30000); + expect(result!.cost).toBe(0); + expect(result!.summary).toBe('No summary provided'); + expect(result!.filesChanged).toEqual([]); + }); + + it('should return null for invalid content', () => { + expect(parseSpawnResult('no frontmatter', 'id', 0)).toBeNull(); + }); + }); +});