chore: bump version to 0.1431

This commit is contained in:
arkon
2026-01-29 13:22:02 +01:00
parent 8ed1c218f4
commit 7b866c1998
14 changed files with 29 additions and 3512 deletions
-370
View File
@@ -1,370 +0,0 @@
#!/usr/bin/env node
/**
* @fileoverview MCP Server for Claudeman Spawn1337 protocol.
*
* Exposes spawn capabilities as native MCP tools that Claude Code can call
* directly, replacing the terminal-tag-parsing approach (SpawnDetector).
*
* Tools:
* - spawn_agent: Spawn a new autonomous agent
* - list_agents: List all agents (active + completed)
* - get_agent_status: Get detailed agent status
* - get_agent_result: Read a completed agent's result
* - send_agent_message: Send a message to a running agent
* - cancel_agent: Cancel a running agent
*
* Environment:
* - CLAUDEMAN_API_URL: Base URL for the Claudeman API (default: http://localhost:3000)
* - CLAUDEMAN_SESSION_ID: Session ID of the calling Claude session
*
* @module mcp-server
*/
import { McpServer } from '@modelcontextprotocol/sdk/server/mcp.js';
import { StdioServerTransport } from '@modelcontextprotocol/sdk/server/stdio.js';
import { z } from 'zod';
// ========== Configuration ==========
const API_URL = process.env.CLAUDEMAN_API_URL || 'http://localhost:3000';
const SESSION_ID = process.env.CLAUDEMAN_SESSION_ID || '';
// ========== API Helper ==========
/**
* Make an HTTP request to the Claudeman API.
*/
async function apiRequest(method: string, path: string, body?: unknown): Promise<{ status: number; data: unknown }> {
const url = `${API_URL}${path}`;
const options: RequestInit = {
method,
headers: { 'Content-Type': 'application/json' },
};
if (body !== undefined) {
options.body = JSON.stringify(body);
}
const response = await fetch(url, options);
const data = await response.json();
return { status: response.status, data };
}
// ========== YAML Construction ==========
/**
* Build a YAML frontmatter + body task spec from structured parameters.
*/
export function buildTaskSpec(params: {
agentId: string;
name: string;
instructions: string;
type?: string;
priority?: string;
maxTokens?: number;
maxCost?: number;
timeoutMinutes?: number;
canModifyParentFiles?: boolean;
contextFiles?: string[];
dependsOn?: string[];
completionPhrase?: string;
outputFormat?: string;
successCriteria?: string;
workingDir?: string;
}): string {
const lines: string[] = ['---'];
lines.push(`agentId: ${params.agentId}`);
lines.push(`name: ${params.name}`);
if (params.type) lines.push(`type: ${params.type}`);
if (params.priority) lines.push(`priority: ${params.priority}`);
if (params.maxTokens != null) lines.push(`maxTokens: ${params.maxTokens}`);
if (params.maxCost != null) lines.push(`maxCost: ${params.maxCost}`);
if (params.timeoutMinutes != null) lines.push(`timeoutMinutes: ${params.timeoutMinutes}`);
if (params.canModifyParentFiles != null) lines.push(`canModifyParentFiles: ${params.canModifyParentFiles}`);
if (params.completionPhrase) lines.push(`completionPhrase: ${params.completionPhrase}`);
if (params.outputFormat) lines.push(`outputFormat: ${params.outputFormat}`);
if (params.successCriteria) lines.push(`successCriteria: "${params.successCriteria.replace(/"/g, '\\"')}"`);
if (params.workingDir) lines.push(`workingDir: ${params.workingDir}`);
if (params.contextFiles && params.contextFiles.length > 0) {
lines.push(`contextFiles: [${params.contextFiles.join(', ')}]`);
}
if (params.dependsOn && params.dependsOn.length > 0) {
lines.push(`dependsOn: [${params.dependsOn.join(', ')}]`);
}
lines.push('---');
lines.push('');
lines.push(params.instructions);
return lines.join('\n');
}
// ========== MCP Server Setup ==========
const server = new McpServer({
name: 'claudeman-spawn',
version: '1.0.0',
});
// ---- spawn_agent ----
server.tool(
'spawn_agent',
'Spawn a new autonomous Claude agent to handle a subtask. The agent runs in its own session with full capabilities.',
{
agentId: z.string().describe('Unique identifier for the agent (e.g., "research-auth-001")'),
name: z.string().describe('Human-readable name for the agent'),
instructions: z.string().describe('Detailed task instructions for the agent (markdown)'),
type: z.enum(['explore', 'implement', 'test', 'review', 'refactor', 'research', 'generate', 'fix', 'general']).optional().describe('Task type/category'),
priority: z.enum(['low', 'normal', 'high', 'critical']).optional().describe('Priority for queue ordering'),
maxTokens: z.number().optional().describe('Maximum token budget (input + output combined)'),
maxCost: z.number().optional().describe('Maximum cost in USD'),
timeoutMinutes: z.number().optional().describe('Maximum runtime in minutes (max: 120)'),
canModifyParentFiles: z.boolean().optional().describe('Whether agent can modify files in parent project directory'),
contextFiles: z.array(z.string()).optional().describe('Files to symlink into agent workspace as context'),
dependsOn: z.array(z.string()).optional().describe('Agent IDs that must complete before this one starts'),
completionPhrase: z.string().optional().describe('Phrase agent outputs when finished (default: auto-generated)'),
outputFormat: z.enum(['markdown', 'json', 'code', 'structured', 'freeform']).optional().describe('Expected output format'),
successCriteria: z.string().optional().describe('Success criteria included in agent instructions'),
workingDir: z.string().optional().describe('Working directory (relative to parent, or absolute)'),
},
async (params) => {
if (!SESSION_ID) {
return {
content: [{ type: 'text', text: 'Error: CLAUDEMAN_SESSION_ID not set. This tool must be run within a Claudeman-managed session.' }],
isError: true,
};
}
const taskSpec = buildTaskSpec(params);
try {
const { status, data } = await apiRequest('POST', '/api/spawn/trigger', {
parentSessionId: SESSION_ID,
taskContent: taskSpec,
});
if (status >= 400) {
const errorData = data as { error?: { code?: string; details?: string } };
const errorMsg = errorData.error?.details || errorData.error?.code || 'Unknown error';
return {
content: [{ type: 'text', text: `Error spawning agent: ${errorMsg}` }],
isError: true,
};
}
const result = data as { success: boolean; data?: { agentId: string } };
const agentId = result.data?.agentId || params.agentId;
return {
content: [{ type: 'text', text: `Agent spawned successfully.\n\nAgent ID: ${agentId}` }],
};
} catch (err) {
return {
content: [{ type: 'text', text: `Error: Could not connect to Claudeman API at ${API_URL}. Is the web server running?\n\n${err instanceof Error ? err.message : String(err)}` }],
isError: true,
};
}
}
);
// ---- list_agents ----
server.tool(
'list_agents',
'List all spawn agents (active, queued, and completed).',
{},
async () => {
try {
const { status, data } = await apiRequest('GET', '/api/spawn/agents');
if (status >= 400) {
return {
content: [{ type: 'text', text: `Error listing agents: ${JSON.stringify(data)}` }],
isError: true,
};
}
const agents = data as Array<{ agentId: string; name: string; status: string; type: string; priority: string }>;
if (agents.length === 0) {
return { content: [{ type: 'text', text: 'No agents found.' }] };
}
const summary = agents.map(a => `- ${a.agentId} (${a.name}): ${a.status} [${a.type}, ${a.priority}]`).join('\n');
return { content: [{ type: 'text', text: `Agents (${agents.length}):\n\n${summary}` }] };
} catch (err) {
return {
content: [{ type: 'text', text: `Error: Could not connect to Claudeman API at ${API_URL}.\n\n${err instanceof Error ? err.message : String(err)}` }],
isError: true,
};
}
}
);
// ---- get_agent_status ----
server.tool(
'get_agent_status',
'Get detailed status and progress of a specific agent.',
{
agentId: z.string().describe('The agent ID to query'),
},
async ({ agentId }) => {
try {
const { status, data } = await apiRequest('GET', `/api/spawn/agents/${encodeURIComponent(agentId)}`);
if (status === 404) {
return {
content: [{ type: 'text', text: `Agent not found: ${agentId}` }],
isError: true,
};
}
if (status >= 400) {
return {
content: [{ type: 'text', text: `Error getting agent status: ${JSON.stringify(data)}` }],
isError: true,
};
}
return { content: [{ type: 'text', text: JSON.stringify(data, null, 2) }] };
} catch (err) {
return {
content: [{ type: 'text', text: `Error: Could not connect to Claudeman API at ${API_URL}.\n\n${err instanceof Error ? err.message : String(err)}` }],
isError: true,
};
}
}
);
// ---- get_agent_result ----
server.tool(
'get_agent_result',
'Read the result of a completed agent. Returns the agent\'s output and metadata.',
{
agentId: z.string().describe('The agent ID whose result to read'),
},
async ({ agentId }) => {
try {
const { status, data } = await apiRequest('GET', `/api/spawn/agents/${encodeURIComponent(agentId)}/result`);
if (status === 404) {
return {
content: [{ type: 'text', text: `Agent or result not found: ${agentId}. The agent may not have completed yet.` }],
isError: true,
};
}
if (status >= 400) {
return {
content: [{ type: 'text', text: `Error getting agent result: ${JSON.stringify(data)}` }],
isError: true,
};
}
// Result could be a string (raw markdown) or an object
const resultText = typeof data === 'string' ? data : JSON.stringify(data, null, 2);
return { content: [{ type: 'text', text: resultText }] };
} catch (err) {
return {
content: [{ type: 'text', text: `Error: Could not connect to Claudeman API at ${API_URL}.\n\n${err instanceof Error ? err.message : String(err)}` }],
isError: true,
};
}
}
);
// ---- send_agent_message ----
server.tool(
'send_agent_message',
'Send a message to a running agent. The message is written to the agent\'s communication channel.',
{
agentId: z.string().describe('The agent ID to message'),
message: z.string().describe('The message content (markdown)'),
},
async ({ agentId, message }) => {
try {
const { status, data } = await apiRequest('POST', `/api/spawn/agents/${encodeURIComponent(agentId)}/message`, {
content: message,
sender: 'parent',
});
if (status === 404) {
return {
content: [{ type: 'text', text: `Agent not found: ${agentId}` }],
isError: true,
};
}
if (status >= 400) {
return {
content: [{ type: 'text', text: `Error sending message: ${JSON.stringify(data)}` }],
isError: true,
};
}
return { content: [{ type: 'text', text: `Message sent to agent ${agentId}.` }] };
} catch (err) {
return {
content: [{ type: 'text', text: `Error: Could not connect to Claudeman API at ${API_URL}.\n\n${err instanceof Error ? err.message : String(err)}` }],
isError: true,
};
}
}
);
// ---- cancel_agent ----
server.tool(
'cancel_agent',
'Cancel a running agent. Sends a graceful shutdown signal.',
{
agentId: z.string().describe('The agent ID to cancel'),
reason: z.string().optional().describe('Reason for cancellation'),
},
async ({ agentId, reason }) => {
try {
const { status, data } = await apiRequest('POST', `/api/spawn/agents/${encodeURIComponent(agentId)}/cancel`, {
reason: reason || 'Cancelled by parent session',
});
if (status === 404) {
return {
content: [{ type: 'text', text: `Agent not found: ${agentId}` }],
isError: true,
};
}
if (status >= 400) {
return {
content: [{ type: 'text', text: `Error cancelling agent: ${JSON.stringify(data)}` }],
isError: true,
};
}
return { content: [{ type: 'text', text: `Agent ${agentId} cancel request sent.` }] };
} catch (err) {
return {
content: [{ type: 'text', text: `Error: Could not connect to Claudeman API at ${API_URL}.\n\n${err instanceof Error ? err.message : String(err)}` }],
isError: true,
};
}
}
);
// ========== Start Server ==========
async function main(): Promise<void> {
const transport = new StdioServerTransport();
await server.connect(transport);
}
main().catch((err) => {
console.error('MCP server failed to start:', err);
process.exit(1);
});
-154
View File
@@ -1,154 +0,0 @@
/**
* @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:
<promise>${spec.completionPhrase}</promise>
**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: <promise>${task.spec.completionPhrase}</promise>`;
}
-964
View File
@@ -1,964 +0,0 @@
/**
* @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';
// ========== Local Constants ==========
/** UUID truncation length for fallback agent IDs */
const UUID_TRUNCATE_LENGTH = 8;
/** Message sequence number padding length */
const MESSAGE_SEQUENCE_PAD_LENGTH = 3;
/** Timeout warning threshold (90% of timeout) */
const TIMEOUT_WARNING_RATIO = 0.9;
/** Budget hard limit ratio (110% - force stop) */
const BUDGET_HARD_LIMIT_RATIO = 1.1;
/** Budget soft limit ratio (100% - warning) */
const BUDGET_SOFT_LIMIT_RATIO = 1.0;
// ========== 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<void>;
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<string, AgentContext> = new Map();
private _completedAgents: Map<string, AgentContext> = 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<string, (phrase: string) => void> = new Map();
constructor(config?: Partial<SpawnOrchestratorConfig>) {
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<SpawnOrchestratorConfig>): 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<void> {
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, UUID_TRUNCATE_LENGTH)}`;
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.
* Cascades cancellation to all child agents before cleaning up the parent.
*/
async cancelAgent(agentId: string, reason: string = 'Cancelled by parent'): Promise<void> {
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;
}
// Cancel all child agents first (cascade)
// Child agents have their parentSessionId set to this agent's sessionId
if (agent.sessionId) {
for (const [childId, childAgent] of this._agents) {
if (childAgent.parentSessionId === agent.sessionId && childAgent.status !== 'cancelled') {
await this.cancelAgent(childId, `Parent ${agentId} cancelled`);
}
}
}
// Also remove any queued tasks that depend on this agent's session
if (agent.sessionId) {
const queuedChildren = this._queue.filter(t => t.parentSessionId === agent.sessionId);
for (const task of queuedChildren) {
const idx = this._queue.indexOf(task);
if (idx >= 0) {
this._queue.splice(idx, 1);
this.emit('cancelled', { agentId: task.spec.agentId, reason: `Parent ${agentId} cancelled` });
}
}
}
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<void> {
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(MESSAGE_SEQUENCE_PAD_LENGTH, '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<void> {
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 (err) {
console.warn(`[spawn-orchestrator] Failed to read progress for agent ${agentId}: ${getErrorMessage(err)}`);
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 [];
try {
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;
try {
const filePath = join(messagesDir, file);
const content = readFileSync(filePath, 'utf-8');
messages.push({
sequence: parseInt(match[1]),
sender: match[2] as 'parent' | 'agent',
content,
sentAt: statSync(filePath).mtimeMs,
read: true,
});
} catch (err) {
console.warn(`[spawn-orchestrator] Failed to read message file ${file}: ${getErrorMessage(err)}`);
// Continue processing other messages
}
}
return messages;
} catch (err) {
console.warn(`[spawn-orchestrator] Failed to read messages for agent ${agentId}: ${getErrorMessage(err)}`);
return [];
}
}
/**
* Programmatically trigger a spawn without terminal detection.
*/
async triggerSpawn(
taskContent: string,
parentSessionId: string,
parentWorkingDir: string,
parentDepth: number = 0
): Promise<string | null> {
const fallbackId = `agent-${uuidv4().slice(0, UUID_TRUNCATE_LENGTH)}`;
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<void> {
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,
warningTimer: 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% - store timer for cleanup
const warningMs = timeoutMs * TIMEOUT_WARNING_RATIO;
agent.warningTimer = 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 >= BUDGET_HARD_LIMIT_RATIO) {
// Force kill at 110%
this.handleAgentTimeout(agent);
return;
} else if (ratio >= BUDGET_SOFT_LIMIT_RATIO) {
// 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 >= BUDGET_HARD_LIMIT_RATIO) {
this.handleAgentTimeout(agent);
return;
} else if (ratio >= BUDGET_SOFT_LIMIT_RATIO) {
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<void> {
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<void> {
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<void> {
const agent = this._agents.get(agentId);
if (!agent) return;
// Clear timers
if (agent.timeoutTimer) {
clearTimeout(agent.timeoutTimer);
agent.timeoutTimer = null;
}
if (agent.warningTimer) {
clearTimeout(agent.warningTimer);
agent.warningTimer = 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 (err) {
console.warn(`[spawn-orchestrator] Failed to stop session for agent ${agentId}: ${getErrorMessage(err)}`);
}
}
// 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());
}
}
-18
View File
@@ -1255,24 +1255,6 @@ export interface ImageDetectedEvent {
size: number;
}
// ========== Spawn1337 Protocol Re-exports ==========
export type {
SpawnPriority,
SpawnResultDelivery,
SpawnStatus,
SpawnTaskSpec,
SpawnTask,
AgentProgress,
SpawnResult,
SpawnMessage,
AgentStatusReport,
SpawnTrackerState,
SpawnOrchestratorConfig,
AgentContext,
SpawnPersistedState,
} from './spawn-types.js';
// ========== Execution Bridge Re-exports ==========
export type {
+25 -51
View File
@@ -466,7 +466,7 @@ class ClaudemanApp {
this._subagentHideTimeout = null; // Timeout for hover-based dropdown hide
this.ralphStatePanelCollapsed = true; // Default to collapsed
// Plan subagent windows (visible Opus agents during plan generation)
// Plan subagent windows (visible agents during plan generation)
this.planSubagents = new Map(); // Map<agentId, { type, model, status, startTime, element, relativePos }>
this.planSubagentWindowZIndex = 1100;
this.planGenerationStopped = false; // Flag to ignore SSE events after Stop
@@ -641,7 +641,7 @@ class ClaudemanApp {
fontFamily: '"Fira Code", "Cascadia Code", "JetBrains Mono", "SF Mono", Monaco, monospace',
fontSize: 14,
lineHeight: 1.2,
cursorBlink: true,
cursorBlink: false,
cursorStyle: 'block',
scrollback: scrollback,
allowTransparency: true,
@@ -1146,6 +1146,10 @@ class ClaudemanApp {
// This connects subagents that were waiting for the session to identify itself
if (claudeSessionIdJustSet) {
this.recheckOrphanSubagents();
// Update connection lines after DOM settles (ensure tabs are rendered)
requestAnimationFrame(() => {
this.updateConnectionLines();
});
}
});
@@ -1690,55 +1694,6 @@ class ClaudemanApp {
this.handleBashToolsUpdate(data.sessionId, data.tools);
});
// Spawn agent notification events
addListener('spawn:failed', (e) => {
const data = JSON.parse(e.data);
this.notificationManager?.notify({
urgency: 'critical',
category: 'spawn-failed',
sessionId: data.sessionId,
sessionName: data.agentId || 'agent',
title: 'Agent Failed',
message: `Agent "${data.agentId}" failed: ${data.reason || 'unknown'}`,
});
});
addListener('spawn:timeout', (e) => {
const data = JSON.parse(e.data);
this.notificationManager?.notify({
urgency: 'critical',
category: 'spawn-timeout',
sessionId: data.sessionId,
sessionName: data.agentId || 'agent',
title: 'Agent Timeout',
message: `Agent "${data.agentId}" exceeded time limit`,
});
});
addListener('spawn:budgetWarning', (e) => {
const data = JSON.parse(e.data);
this.notificationManager?.notify({
urgency: 'warning',
category: 'spawn-budget',
sessionId: data.sessionId,
sessionName: data.agentId || 'agent',
title: 'Budget Warning',
message: `Agent "${data.agentId}" at ${data.percent || 80}% budget`,
});
});
addListener('spawn:completed', (e) => {
const data = JSON.parse(e.data);
this.notificationManager?.notify({
urgency: 'info',
category: 'spawn-completed',
sessionId: data.sessionId,
sessionName: data.agentId || 'agent',
title: 'Agent Complete',
message: `Agent "${data.agentId}" finished successfully`,
});
});
// Hook events (from Claude Code hooks system)
// Use pendingHooks state machine to track hook events and derive tab alerts.
// This ensures alerts persist even when session:working events fire.
@@ -1832,6 +1787,11 @@ class ClaudemanApp {
if (data.status === 'active') {
this.openSubagentWindow(data.agentId);
}
// Ensure connection lines are updated after window is created and DOM settles
requestAnimationFrame(() => {
this.updateConnectionLines();
});
});
addListener('subagent:updated', (e) => {
@@ -7658,6 +7618,11 @@ class ClaudemanApp {
this.renderSessionTabs(); // Update tab badges
this.saveSubagentWindowStates(); // Persist corrected mappings
// Update connection lines after all windows are restored (use rAF to ensure DOM is ready)
requestAnimationFrame(() => {
this.updateConnectionLines();
});
}
// ========== Help Modal ==========
@@ -9425,11 +9390,20 @@ class ClaudemanApp {
* Called when session:updated fires, in case claudeSessionId was just set.
*/
recheckOrphanSubagents() {
let anyFound = false;
for (const [agentId, agent] of this.subagents) {
if (!agent.parentSessionId && agent.sessionId) {
const hadParent = agent.parentSessionId;
this.findParentSessionForSubagent(agentId);
if (!hadParent && this.subagents.get(agentId)?.parentSessionId) {
anyFound = true;
}
}
}
// Ensure connection lines are updated after all orphans are processed
if (anyFound) {
this.updateConnectionLines();
}
}
/**
-207
View File
@@ -22,8 +22,6 @@ import { EventEmitter } from 'node:events';
import { Session, ClaudeMessage, type BackgroundTask, type RalphTrackerState, type RalphTodoItem, type ActiveBashTool } from '../session.js';
import { fileStreamManager } from '../file-stream-manager.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';
@@ -307,8 +305,6 @@ 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;
// Token recording for daily stats (track what's been recorded to avoid double-counting)
private lastRecordedTokens: Map<string, { input: number; output: number }> = new Map();
private tokenRecordingTimer: NodeJS.Timeout | null = null;
@@ -350,10 +346,6 @@ export class WebServer extends EventEmitter {
this.broadcast('screen:statsUpdated', screens);
});
// Initialize spawn orchestrator
this.spawnOrchestrator = new SpawnOrchestrator();
this.setupSpawnOrchestratorListeners();
// Initialize execution bridge with model config from settings
this.executionBridge = getExecutionBridge(this.loadModelConfig());
this.setupExecutionBridgeListeners();
@@ -510,7 +502,6 @@ export class WebServer extends EventEmitter {
// Returns comprehensive memory metrics for debugging memory leaks
this.app.get('/api/debug/memory', async () => {
const mem = process.memoryUsage();
const spawnState = this.spawnOrchestrator.getState();
const subagentStats = subagentWatcher.getStats();
// Calculate total Map entries for memory estimation
@@ -569,13 +560,6 @@ export class WebServer extends EventEmitter {
subagentIdleTimers: subagentStats.idleTimerCount,
total: this.respawnTimers.size + this.pendingRespawnStarts.size + subagentStats.idleTimerCount,
},
spawn: {
activeAgents: spawnState.activeCount,
queuedAgents: spawnState.queuedCount,
totalSpawned: spawnState.totalSpawned,
totalCompleted: spawnState.totalCompleted,
totalFailed: spawnState.totalFailed,
},
uptime: {
seconds: Math.round(process.uptime()),
formatted: formatUptime(process.uptime()),
@@ -1987,9 +1971,6 @@ export class WebServer extends EventEmitter {
const claudeMd = generateClaudeMd(name, description || '', templatePath);
writeFileSync(join(casePath, 'CLAUDE.md'), claudeMd);
// Write .mcp.json for Claude Code to discover spawn tools
this.writeMcpConfig(casePath);
// Write .claude/settings.local.json with hooks for desktop notifications
writeHooksConfig(casePath);
@@ -2236,9 +2217,6 @@ export class WebServer extends EventEmitter {
const claudeMd = generateClaudeMd(caseName, '', templatePath);
writeFileSync(join(casePath, 'CLAUDE.md'), claudeMd);
// Write .mcp.json for Claude Code to discover spawn tools
this.writeMcpConfig(casePath);
// Write .claude/settings.local.json with hooks for desktop notifications
writeHooksConfig(casePath);
@@ -3019,93 +2997,6 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
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<SpawnOrchestratorConfig>;
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 } };
});
// ========== Execution Bridge Endpoints ==========
// Get execution status
@@ -3971,83 +3862,6 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
});
}
private setupSpawnOrchestratorListeners(): void {
const sessionCreator: SessionCreator = {
createAgentSession: async (workingDir: string, name: string) => {
const globalNice = this.getGlobalNiceConfig();
const session = new Session({
workingDir,
screenManager: this.screenManager,
useScreen: true,
mode: 'claude',
name: `spawn:${name}`,
niceConfig: globalNice,
});
this.sessions.set(session.id, session);
this.store.incrementSessionsCreated();
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) => {
// Use cleanupSession to properly clean up all resources (respawn controllers,
// run summary trackers, file streams, Ralph state, etc.)
await this.cleanupSession(sessionId);
},
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));
}
/**
* Load model configuration from settings file.
*/
@@ -4373,23 +4187,6 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
return undefined;
}
/**
* Write .mcp.json to a case directory for Claude Code to discover spawn MCP tools.
*/
private writeMcpConfig(casePath: string): void {
const projectRoot = join(__dirname, '..', '..');
const mcpServerPath = join(projectRoot, 'dist', 'mcp-server.js');
const mcpConfig = {
mcpServers: {
'claudeman-spawn': {
command: 'node',
args: [mcpServerPath],
},
},
};
writeFileSync(join(casePath, '.mcp.json'), JSON.stringify(mcpConfig, null, 2) + '\n');
}
private async startScheduledRun(prompt: string, workingDir: string, durationMinutes: number): Promise<ScheduledRun> {
const id = uuidv4();
const now = Date.now();
@@ -5147,10 +4944,6 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
}
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);