From 61b5ec095c78430811058e102f38c1c7a01a7c8d Mon Sep 17 00:00:00 2001 From: arkon Date: Sat, 21 Mar 2026 07:20:18 +0100 Subject: [PATCH] =?UTF-8?q?feat:=20add=20Orchestrator=20Loop=20=E2=80=94?= =?UTF-8?q?=20phased=20plan=20execution=20with=20team=20agents?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds a new autonomous loop that accepts high-level goals, generates phased execution plans via AI, and executes them step-by-step with verification gates between phases. Core components: - OrchestratorLoop: state machine (idle→planning→approval→executing→verifying→completed) - OrchestratorPlanner: plan generation via PlanOrchestrator, Kahn's algorithm phase grouping - OrchestratorVerifier: phase verification (strict/moderate/lenient modes) - Prompt templates for phase execution, team delegation, verification, replanning API (10 endpoints): - POST start/approve/reject/pause/resume/stop - GET status/plan - POST phase/:id/skip, phase/:id/retry Frontend: orchestrator-panel.js with SSE-driven state, phase progress, task tracking Tests: 22 tests (18 route + 4 unit), all passing. Typecheck/lint/format clean. Co-Authored-By: Claude Opus 4.6 --- docs/orchestrator-loop-architecture.md | 367 ++++++++++ docs/orchestrator-loop-plan.md | 633 ++++++++++++++++ docs/orchestrator-loop-research.md | 157 ++++ src/orchestrator-loop.ts | 912 ++++++++++++++++++++++++ src/orchestrator-planner.ts | 412 +++++++++++ src/orchestrator-verifier.ts | 292 ++++++++ src/prompts/index.ts | 7 + src/prompts/orchestrator.ts | 116 +++ src/state-store.ts | 24 + src/types/app-state.ts | 2 + src/types/index.ts | 2 + src/types/orchestrator.ts | 285 ++++++++ src/web/ports/index.ts | 1 + src/web/ports/orchestrator-port.ts | 11 + src/web/public/app.js | 17 + src/web/public/constants.js | 13 + src/web/public/index.html | 12 + src/web/public/orchestrator-panel.js | 456 ++++++++++++ src/web/public/styles.css | 219 ++++++ src/web/routes/index.ts | 1 + src/web/routes/orchestrator-routes.ts | 250 +++++++ src/web/schemas.ts | 28 + src/web/server.ts | 21 + src/web/sse-events.ts | 40 +- test/mocks/mock-route-context.ts | 4 + test/orchestrator-planner.test.ts | 102 +++ test/routes/orchestrator-routes.test.ts | 335 +++++++++ 27 files changed, 4718 insertions(+), 1 deletion(-) create mode 100644 docs/orchestrator-loop-architecture.md create mode 100644 docs/orchestrator-loop-plan.md create mode 100644 docs/orchestrator-loop-research.md create mode 100644 src/orchestrator-loop.ts create mode 100644 src/orchestrator-planner.ts create mode 100644 src/orchestrator-verifier.ts create mode 100644 src/prompts/orchestrator.ts create mode 100644 src/types/orchestrator.ts create mode 100644 src/web/ports/orchestrator-port.ts create mode 100644 src/web/public/orchestrator-panel.js create mode 100644 src/web/routes/orchestrator-routes.ts create mode 100644 test/orchestrator-planner.test.ts create mode 100644 test/routes/orchestrator-routes.test.ts diff --git a/docs/orchestrator-loop-architecture.md b/docs/orchestrator-loop-architecture.md new file mode 100644 index 00000000..a81d99b6 --- /dev/null +++ b/docs/orchestrator-loop-architecture.md @@ -0,0 +1,367 @@ +# Orchestrator Loop — Architecture & Data Flow + +> Technical architecture document. Not for GitHub. + +## System Overview + +``` +┌─────────────────────────────────────────────────────────────────────┐ +│ CODEMAN WEB UI │ +│ ┌──────────────────────────────────────────────────────────────┐ │ +│ │ Orchestrator Dashboard │ │ +│ │ [Goal Input] [Plan View] [Phase Progress] [Agent Activity] │ │ +│ └───────────────────────────┬──────────────────────────────────┘ │ +│ │ SSE Events │ +│ ▼ │ +│ ┌──────────────────────────────────────────────────────────────┐ │ +│ │ Orchestrator API Routes (/api/orchestrator/*) │ │ +│ └───────────────────────────┬──────────────────────────────────┘ │ +└───────────────────────────────┼─────────────────────────────────────┘ + ▼ +┌─────────────────────────────────────────────────────────────────────┐ +│ ORCHESTRATOR LOOP │ +│ │ +│ ┌──────────────┐ ┌──────────────┐ ┌──────────────────────┐ │ +│ │ Orchestrator │ │ Orchestrator │ │ Orchestrator │ │ +│ │ Planner │ │ Loop (state │ │ Verifier │ │ +│ │ │ │ machine) │ │ │ │ +│ │ • Research │◄──►│ • Phase mgmt │◄──►│ • Test runner │ │ +│ │ • Plan gen │ │ • Task queue │ │ • AI review │ │ +│ │ • Phasing │ │ • Event loop │ │ • Output checks │ │ +│ └──────┬───────┘ └──────┬───────┘ └──────────┬───────────┘ │ +│ │ │ │ │ +│ ▼ ▼ ▼ │ +│ ┌──────────────────────────────────────────────────────────────┐ │ +│ │ EXISTING CODEMAN INFRASTRUCTURE │ │ +│ │ │ │ +│ │ SessionManager ←→ Sessions ←→ PTY (Claude CLI) │ │ +│ │ ↑ ↑ ↑ │ │ +│ │ │ │ │ │ │ +│ │ TaskQueue RalphTracker RespawnController │ │ +│ │ StateStore HooksConfig TeamWatcher │ │ +│ │ Auto-Ops SubagentWatcher SSE Broadcast │ │ +│ └──────────────────────────────────────────────────────────────┘ │ +└─────────────────────────────────────────────────────────────────────┘ +``` + +## Data Flow: Complete Lifecycle + +### 1. User Submits Goal + +``` +User → POST /api/orchestrator/start { goal: "Build a REST API...", config: {...} } + → OrchestratorLoop.start(goal) + → state = PLANNING + → emit('stateChanged', 'planning') + → SSE: orchestrator:stateChanged +``` + +### 2. Planning Phase + +``` +OrchestratorPlanner.generatePlan(goal) + → PlanOrchestrator.generateDetailedPlan(goal) + → [Research Agent] → enriched task description + → [Planner Agent] → PlanItem[] + → groupIntoPhases(planItems) + → topological sort by dependencies + → group into layers + → assign team strategies + → OrchestratorPlan { phases: [...] } + → state = APPROVAL + → emit('planReady', plan) + → SSE: orchestrator:planReady +``` + +### 3. User Approves Plan + +``` +User → POST /api/orchestrator/approve + → OrchestratorLoop.approvePlan() + → state = EXECUTING + → executePhase(phases[0]) +``` + +### 4. Phase Execution + +``` +executePhase(phase) + → For each task in phase: + → Convert to CreateTaskOptions + → Add to TaskQueue with completion phrase "PHASE_{N}_TASK_{M}_DONE" + → If phase.teamStrategy.type === 'team': + → Start session with AGENT_TEAMS enabled + → Send team orchestration prompt to lead + → Else: + → Assign tasks to available sessions (same as RalphLoop) + + → Listen for task completion events: + → TaskQueue emits taskCompleted + → Check: all phase tasks done? + → Yes → state = VERIFYING → verifyPhase(phase) + → No → wait for more completions +``` + +### 5. Verification + +``` +verifyPhase(phase) + → OrchestratorVerifier.verify(phase, session) + → Run test commands via session + → Check file existence + → AI review (optional) + → If passed: + → phase.status = 'passed' + → emit('phaseCompleted', phase) + → If more phases: executePhase(nextPhase) + → If last phase: state = COMPLETED + → If failed: + → phase.attempts++ + → If attempts < maxAttempts: + → state = REPLANNING + → Generate recovery tasks + → state = EXECUTING (retry) + → Else: + → state = FAILED + → emit('phaseFailed', phase, reason) +``` + +### 6. Context Management Between Phases + +``` +After phase completion: + → If config.compactBetweenPhases: + → session.sendInput('/compact') + → Wait for compact to complete + → If config.respawnBetweenMilestones && phase is a milestone: + → Save orchestrator state to StateStore + → Respawn session (kill + recreate) + → Send resume prompt with phase context +``` + +## File Layout + +``` +src/ +├── orchestrator-loop.ts # Main state machine (~400 lines) +├── orchestrator-planner.ts # Plan generation + phase grouping (~300 lines) +├── orchestrator-verifier.ts # Phase verification (~200 lines) +├── types/ +│ └── orchestrator.ts # All orchestrator types (~150 lines) +├── prompts/ +│ └── orchestrator.ts # Prompt templates (~200 lines) +├── web/ +│ ├── routes/ +│ │ └── orchestrator-routes.ts # API endpoints (~250 lines) +│ └── public/ +│ └── orchestrator-ui.js # Frontend panel (~500 lines) +``` + +## Integration Points with Existing Code + +### StateStore (`src/state-store.ts`) +```typescript +// Add to AppState interface +orchestrator?: OrchestratorPersistState; + +// Add methods +getOrchestratorState(): OrchestratorPersistState; +setOrchestratorState(state: Partial): void; +``` + +### SSE Events (`src/web/sse-events.ts`) +```typescript +// Add ~8 new events +export const SseEvent = { + // ... existing + ORCHESTRATOR_STATE_CHANGED: 'orchestrator:stateChanged', + ORCHESTRATOR_PLAN_READY: 'orchestrator:planReady', + ORCHESTRATOR_PHASE_STARTED: 'orchestrator:phaseStarted', + ORCHESTRATOR_PHASE_COMPLETED: 'orchestrator:phaseCompleted', + ORCHESTRATOR_PHASE_FAILED: 'orchestrator:phaseFailed', + ORCHESTRATOR_VERIFICATION: 'orchestrator:verificationResult', + ORCHESTRATOR_COMPLETED: 'orchestrator:completed', + ORCHESTRATOR_ERROR: 'orchestrator:error', +} as const; +``` + +### Frontend Constants (`src/web/public/constants.js`) +```javascript +// Mirror SSE events +SSE_EVENTS.ORCHESTRATOR_STATE_CHANGED = 'orchestrator:stateChanged'; +// ... etc +``` + +### Route Registration (`src/web/routes/index.ts`) +```typescript +import { registerOrchestratorRoutes } from './orchestrator-routes.js'; +// Add to barrel export +``` + +### Server (`src/web/server.ts`) +```typescript +// Initialize OrchestratorLoop alongside RalphLoop +const orchestratorLoop = new OrchestratorLoop(config); + +// Register routes +registerOrchestratorRoutes(app, { ...ctx, orchestrator: orchestratorLoop }); +``` + +### Port Interface (`src/web/ports/`) +```typescript +// New port +export interface OrchestratorPort { + orchestrator: OrchestratorLoop; +} +``` + +## Prompt Flow Through System + +The key insight is how prompts flow from Orchestrator → Session → Claude: + +``` +OrchestratorLoop decides to execute Phase 3, Task 2 + │ + ▼ +Converts OrchestratorTask to CreateTaskOptions: + { + prompt: "Implement the rate limiter middleware. Read src/middleware/auth.ts + for the pattern. Add to src/middleware/rate-limiter.ts. Must export + a Fastify plugin. When done: PHASE_3_TASK_2_DONE", + priority: 100, + dependencies: ["phase-3-task-1"], // Must finish auth middleware first + completionPhrase: "PHASE_3_TASK_2_DONE", + timeoutMs: 600000 // 10 minutes + } + │ + ▼ +TaskQueue.addTask(options) + │ + ▼ +RalphLoop.tick() → assignTasks() // OR OrchestratorLoop does its own assignment + │ + ▼ +session.sendInput(task.prompt) + │ + ▼ +writeViaMux() → tmux send-keys -l "prompt..." + Enter + │ + ▼ +Claude CLI receives prompt, executes, outputs results + │ + ▼ +RalphTracker.processData() → detects "PHASE_3_TASK_2_DONE" + │ + ▼ +emit('completionDetected') → OrchestratorLoop.handleTaskCompleted() + │ + ▼ +Check: all tasks in Phase 3 done? → If yes → verifyPhase(phase3) +``` + +## Team Agent Flow (When Enabled) + +``` +Phase has teamStrategy.type === 'team' + │ + ▼ +OrchestratorLoop creates/reuses a session with: + env: { CLAUDE_CODE_EXPERIMENTAL_AGENT_TEAMS: '1' } + │ + ▼ +Sends team orchestration prompt: + "You're the team lead for Phase 3: Core Implementation. + + Your team should work on these tasks in parallel: + 1. Rate limiter middleware (teammate 1) + 2. Error handling middleware (teammate 2) + 3. Validation layer (teammate 3) + + Context files to read first: [...] + Each teammate should output their task's completion phrase when done. + When ALL tasks are complete, output: PHASE_3_COMPLETE" + │ + ▼ +Claude Code team-lead spawns teammates + │ + ▼ +TeamWatcher detects new team in ~/.claude/teams/ + → Matches to session via leadSessionId + → Tracks teammate activity + │ + ▼ +Teammates work in parallel (in-process threads) + │ + ▼ +hook: teammate_idle → POST /api/hook-event + → OrchestratorLoop notes teammate finished + │ + ▼ +hook: task_completed → POST /api/hook-event + → Or: RalphTracker detects PHASE_3_COMPLETE + → OrchestratorLoop → phase complete → verify +``` + +## Error Recovery Strategy + +``` +Task fails (timeout, error, session crash) + │ + ├─ Task-level retry (up to 2 retries per task) + │ → Reset task to pending + │ → Re-queue with modified prompt: "Previous attempt failed: {error}. Try again..." + │ + ├─ Phase-level retry (up to 3 retries per phase) + │ → Respawn session (fresh context) + │ → Re-execute entire phase with learnings from failure + │ → Modified prompt includes what went wrong + │ + └─ Orchestration-level failure + → All retries exhausted + → state = FAILED + → Notify user with detailed failure report + → User can: modify plan → retry, skip phase → continue, or stop +``` + +## Interaction with Ralph Loop + +Ralph Loop and Orchestrator Loop are **mutually exclusive** on the same sessions: + +``` +if (orchestratorLoop.isRunning()) { + // Orchestrator controls task assignment + // Ralph Loop should not interfere + // Respawn Controller uses 'orchestrator' preset +} + +if (ralphLoop.isRunning()) { + // Ralph controls task assignment + // Orchestrator should not start +} +``` + +The Orchestrator can optionally USE the Ralph Loop internally for phase execution (delegate phase tasks to Ralph's queue), or manage task assignment directly. Decision: **manage directly** — gives more control over phase boundaries and verification timing. + +## Summary of What Touches What + +| Existing File | Change | +|---|---| +| `src/types/index.ts` | Export orchestrator types | +| `src/state-store.ts` | Add orchestrator state persistence | +| `src/web/sse-events.ts` | Add ~8 orchestrator events | +| `src/web/routes/index.ts` | Register orchestrator routes | +| `src/web/server.ts` | Initialize OrchestratorLoop | +| `src/web/public/constants.js` | Mirror SSE events | +| `src/web/public/app.js` | Add orchestrator event listeners, panel toggle | +| `src/web/route-helpers.ts` | Add 'orchestrator' respawn preset | + +| New File | Purpose | +|---|---| +| `src/orchestrator-loop.ts` | Core state machine | +| `src/orchestrator-planner.ts` | Plan generation + phasing | +| `src/orchestrator-verifier.ts` | Phase verification | +| `src/types/orchestrator.ts` | Type definitions | +| `src/prompts/orchestrator.ts` | Prompt templates | +| `src/web/routes/orchestrator-routes.ts` | API endpoints | +| `src/web/public/orchestrator-ui.js` | Frontend panel | +| `src/web/ports/orchestrator-port.ts` | Port interface | diff --git a/docs/orchestrator-loop-plan.md b/docs/orchestrator-loop-plan.md new file mode 100644 index 00000000..9d067a47 --- /dev/null +++ b/docs/orchestrator-loop-plan.md @@ -0,0 +1,633 @@ +# Orchestrator Loop — Detailed Implementation Plan (v2) + +> Internal research/planning document. Not for GitHub. + +## Vision + +The **Orchestrator Loop** is a new autonomous execution mode that transforms high-level user goals into phased, verified, team-coordinated implementations. Unlike Ralph Loop (flat task queue → idle sessions), the Orchestrator manages the full lifecycle: **plan → approve → execute → verify → adapt → complete**. + +``` +USER: "Add OAuth2 login with Google/GitHub, role-based access control, and API key management" + +ORCHESTRATOR: + Phase 1: Research & Setup ✅ (3m) — scaffold, deps, config + Phase 2: Auth Core ✅ (8m) — OAuth2 flow, session mgmt + Phase 3: Provider Integration 🔄 (12m) — Google + GitHub (parallel via team agents) + Phase 4: RBAC ⏳ — roles, permissions, middleware + Phase 5: API Keys ⏳ — generation, validation, rate limits + Phase 6: Testing & Review ⏳ — integration tests, security review + + Progress: ━━━━━━━━━━━━━━━━━━━━ 40% | Agents: 3 active | Time: 23m +``` + +## Architecture + +``` +┌─────────────────────────────────────────────────────────────────┐ +│ OrchestratorLoop │ +│ │ +│ ┌────────────────┐ ┌────────────────┐ ┌──────────────────┐ │ +│ │ Orchestrator │ │ Orchestrator │ │ Orchestrator │ │ +│ │ Planner │ │ Executor │ │ Verifier │ │ +│ │ │ │ │ │ │ │ +│ │ PlanOrchestrator│ │ TaskQueue │ │ AI review │ │ +│ │ + phase grouper│ │ SessionManager │ │ Test commands │ │ +│ │ + team strategy│ │ Team prompts │ │ File checks │ │ +│ └───────┬────────┘ └───────┬────────┘ └─────────┬────────┘ │ +│ │ │ │ │ +│ └───────────────────┼──────────────────────┘ │ +│ │ │ +│ ┌─────────▼─────────┐ │ +│ │ Existing Codeman │ │ +│ │ Infrastructure │ │ +│ │ │ │ +│ │ SessionManager │ │ +│ │ TaskQueue │ │ +│ │ RespawnController │ │ +│ │ TeamWatcher │ │ +│ │ PlanOrchestrator │ │ +│ │ StateStore │ │ +│ │ Hooks + SSE │ │ +│ └────────────────────┘ │ +└─────────────────────────────────────────────────────────────────┘ +``` + +## State Machine + +``` + ┌─────────┐ + │ IDLE │ + └────┬────┘ + │ start(goal) + ▼ + ┌─────────┐ + ┌────────│PLANNING │────────┐ + │ fail └────┬────┘ │ + ▼ │ plan ready │ user cancels + ┌────────┐ ▼ ▼ + │ FAILED │ ┌─────────┐ ┌────────┐ + └────────┘ │APPROVAL │ │ IDLE │ + ▲ └────┬────┘ └────────┘ + │ │ approve + │ ▼ + │ ┌──────────┐ + │ ┌───►│EXECUTING │◄────────────────────┐ + │ │ └────┬─────┘ │ + │ │ │ all tasks in phase done │ + │ │ ▼ │ + │ │ ┌──────────┐ │ + │ │ │VERIFYING │ │ + │ │ └────┬─────┘ │ + │ │ pass │ │ fail │ + │ │ ▼ ▼ │ + │ │ more ┌──────────┐ │ + │ │ phases?│REPLANNING│── retry ────────┘ + │ │ │ └────┬─────┘ + │ │ │ │ max retries + │ │ │ ▼ + │ │ │ ┌────────┐ + │ └────┘ │ FAILED │ + │ next └────────┘ + │ phase + │ │ + │ ▼ + │ ┌───────────┐ + └─│ COMPLETED │ + └───────────┘ +``` + +**States:** `idle` | `planning` | `approval` | `executing` | `verifying` | `replanning` | `completed` | `failed` | `paused` + +Transitions are event-driven. The state machine is the single source of truth — all methods check `this.state` before acting. + +## Type Definitions + +### `src/types/orchestrator.ts` + +```typescript +// ═══════════════════════════════════════════════════════════════ +// State Machine +// ═══════════════════════════════════════════════════════════════ + +export type OrchestratorState = + | 'idle' + | 'planning' + | 'approval' + | 'executing' + | 'verifying' + | 'replanning' + | 'completed' + | 'failed' + | 'paused'; + +// ═══════════════════════════════════════════════════════════════ +// Plan Structure +// ═══════════════════════════════════════════════════════════════ + +export interface OrchestratorPlan { + id: string; + goal: string; + createdAt: number; + phases: OrchestratorPhase[]; + metadata: { + totalTasks: number; + estimatedComplexity: 'low' | 'medium' | 'high'; + modelUsed: string; + planDurationMs: number; + }; +} + +export interface OrchestratorPhase { + id: string; // "phase-1", "phase-2" + name: string; // Human-readable name + description: string; + order: number; + status: PhaseStatus; + tasks: OrchestratorTask[]; + verificationCriteria: string[]; + testCommands: string[]; + maxAttempts: number; // Default: 3 + attempts: number; // Current attempt count + startedAt: number | null; + completedAt: number | null; + durationMs: number | null; + teamStrategy: TeamStrategy; +} + +export type PhaseStatus = + | 'pending' + | 'executing' + | 'verifying' + | 'passed' + | 'failed' + | 'skipped'; + +export interface OrchestratorTask { + id: string; // "phase-1-task-1" + phaseId: string; + prompt: string; // Single-line prompt for Claude + status: 'pending' | 'running' | 'completed' | 'failed'; + assignedSessionId: string | null; + queueTaskId: string | null; // Links to TaskQueue task + parallel: boolean; // Can run in parallel with sibling tasks + completionPhrase: string; // Unique phrase for completion detection + timeoutMs: number; + startedAt: number | null; + completedAt: number | null; + error: string | null; + retries: number; +} + +// ═══════════════════════════════════════════════════════════════ +// Team Strategy +// ═══════════════════════════════════════════════════════════════ + +export type TeamStrategy = + | { type: 'single' } // One session handles all + | { type: 'parallel'; maxSessions: number } // Multiple sessions + | { type: 'team'; config: TeamSetup } // Agent teams + +export interface TeamSetup { + leadPrompt: string; + suggestedTeammates: string[]; // Role descriptions + maxTeammates: number; +} + +// ═══════════════════════════════════════════════════════════════ +// Verification +// ═══════════════════════════════════════════════════════════════ + +export interface VerificationResult { + passed: boolean; + checks: VerificationCheck[]; + summary: string; + suggestions: string[]; // Recovery hints for replanning +} + +export interface VerificationCheck { + type: 'test_command' | 'ai_review' | 'file_check'; + description: string; + passed: boolean; + output?: string; +} + +// ═══════════════════════════════════════════════════════════════ +// Configuration +// ═══════════════════════════════════════════════════════════════ + +export interface OrchestratorConfig { + plannerModel: string; // Default: 'opus' + researchEnabled: boolean; // Default: true + autoApprove: boolean; // Default: false + maxPhaseRetries: number; // Default: 3 + phaseTimeoutMs: number; // Default: 1800000 (30min) + enableTeamAgents: boolean; // Default: true + maxParallelSessions: number; // Default: 3 + verificationMode: 'strict' | 'moderate' | 'lenient'; + compactBetweenPhases: boolean; // Default: true +} + +// ═══════════════════════════════════════════════════════════════ +// Persistence (saved to ~/.codeman/state.json) +// ═══════════════════════════════════════════════════════════════ + +export interface OrchestratorPersistState { + state: OrchestratorState; + plan: OrchestratorPlan | null; + currentPhaseIndex: number; + startedAt: number | null; + completedAt: number | null; + config: OrchestratorConfig; + stats: OrchestratorStats; +} + +export interface OrchestratorStats { + phasesCompleted: number; + phasesFailed: number; + totalTasksCompleted: number; + totalTasksFailed: number; + totalDurationMs: number; + replanCount: number; +} +``` + +## New Files (Implementation Order) + +### Step 1: `src/types/orchestrator.ts` — Type definitions +All interfaces above. No dependencies. ~120 lines. + +### Step 2: `src/orchestrator-planner.ts` — Plan generation + phase grouping +~300 lines. Wraps existing PlanOrchestrator. + +```typescript +/** + * @fileoverview Orchestrator plan generation — converts goals into phased plans. + * + * Uses PlanOrchestrator for AI plan generation, then groups PlanItems into + * sequential phases with team strategies and verification criteria. + * + * @module orchestrator-planner + */ + +export class OrchestratorPlanner { + constructor(mux: TerminalMultiplexer, workingDir: string, config: OrchestratorConfig); + + /** Generate plan from goal. Uses PlanOrchestrator internally. */ + async generatePlan(goal: string, onProgress?: ProgressCallback): Promise; + + /** Cancel in-progress plan generation. */ + async cancel(): Promise; + + // Internal + private groupIntoPhases(items: PlanItem[], goal: string): OrchestratorPhase[]; + private assignTeamStrategies(phases: OrchestratorPhase[]): void; + private generateCompletionPhrases(plan: OrchestratorPlan): void; +} +``` + +**Phase grouping algorithm:** +1. Topological sort by `PlanItem.dependencies` +2. Group into dependency layers (Kahn's algorithm) +3. Within each layer, sub-group by `tddPhase` (setup → test → impl → verify → review) +4. Merge adjacent small phases (< 2 tasks) if they share the same tddPhase +5. Assign team strategies: + - 1-2 tasks → `{ type: 'single' }` + - 3+ independent tasks → `{ type: 'parallel', maxSessions: Math.min(taskCount, config.maxParallelSessions) }` + - 4+ tasks with high complexity → `{ type: 'team', config: { ... } }` +6. Generate unique completion phrases per task: `ORCH_P{phaseOrder}_T{taskIndex}` + +### Step 3: `src/orchestrator-verifier.ts` — Phase verification +~200 lines. + +```typescript +/** + * @fileoverview Orchestrator phase verification. + * + * Runs verification checks after each phase completes: + * test commands, AI review, and file existence checks. + * + * @module orchestrator-verifier + */ + +export class OrchestratorVerifier { + constructor(config: OrchestratorConfig); + + /** Run all verification checks for a completed phase. */ + async verifyPhase( + phase: OrchestratorPhase, + session: Session, + mode: 'strict' | 'moderate' | 'lenient' + ): Promise; + + // Verification strategies + private async runTestCommands(commands: string[], session: Session): Promise; + private async aiReview(phase: OrchestratorPhase, session: Session): Promise; +} +``` + +**Verification modes:** +- `strict`: ALL test commands must pass AND AI review must approve +- `moderate`: Test commands must pass, AI review is advisory +- `lenient`: At least one test command passes, AI review skipped + +**AI review prompt (sent as a task to the session):** +``` +Review Phase "{phase.name}" completion. Check: +1. Expected functionality works +2. No obvious regressions +3. Code quality is acceptable + +Criteria: {phase.verificationCriteria.join('\n')} + +If ALL criteria are met, respond: ORCH_VERIFY_PASS +If ANY criteria fail, respond: ORCH_VERIFY_FAIL and explain what failed. +``` + +### Step 4: `src/orchestrator-loop.ts` — Core state machine +~500 lines. Main orchestrator engine. + +```typescript +/** + * @fileoverview Orchestrator Loop — phased plan execution with team agents. + * + * State machine that generates plans from user goals, executes them + * phase-by-phase with verification gates, and adapts on failure. + * + * @module orchestrator-loop + */ + +export interface OrchestratorLoopEvents { + stateChanged: (state: OrchestratorState, prevState: OrchestratorState) => void; + planReady: (plan: OrchestratorPlan) => void; + phaseStarted: (phase: OrchestratorPhase) => void; + phaseCompleted: (phase: OrchestratorPhase) => void; + phaseFailed: (phase: OrchestratorPhase, reason: string) => void; + taskAssigned: (task: OrchestratorTask, sessionId: string) => void; + taskCompleted: (task: OrchestratorTask) => void; + taskFailed: (task: OrchestratorTask, error: string) => void; + verificationResult: (phase: OrchestratorPhase, result: VerificationResult) => void; + completed: (stats: OrchestratorStats) => void; + error: (error: Error) => void; +} + +export class OrchestratorLoop extends EventEmitter { + private state: OrchestratorState = 'idle'; + private plan: OrchestratorPlan | null = null; + private currentPhaseIndex = 0; + private config: OrchestratorConfig; + private planner: OrchestratorPlanner; + private verifier: OrchestratorVerifier; + private sessionManager: SessionManager; + private taskQueue: TaskQueue; + private store: StateStore; + private stats: OrchestratorStats; + private cleanup: CleanupManager; + private pausedState: OrchestratorState | null = null; // State before pause + + // ── Lifecycle ────────────────────────────────────────────── + + constructor(mux: TerminalMultiplexer, workingDir: string, config?: Partial); + + /** Start orchestration with a goal. Transitions: idle → planning */ + async start(goal: string): Promise; + + /** Approve the generated plan. Transitions: approval → executing */ + async approve(): Promise; + + /** Reject plan with feedback. Transitions: approval → planning (regenerate) */ + async reject(feedback: string): Promise; + + /** Pause execution. Saves current state. */ + pause(): void; + + /** Resume from pause. */ + resume(): void; + + /** Stop everything and clean up. → idle */ + async stop(): Promise; + + /** Skip current phase. → executing (next phase) or completed */ + async skipPhase(phaseId: string): Promise; + + /** Retry a failed phase. → executing */ + async retryPhase(phaseId: string): Promise; + + // ── Getters ──────────────────────────────────────────────── + + getState(): OrchestratorState; + getPlan(): OrchestratorPlan | null; + getCurrentPhase(): OrchestratorPhase | null; + getStats(): OrchestratorStats; + getStatus(): OrchestratorPersistState; + + // ── Internal: Phase Execution ────────────────────────────── + + private async executeCurrentPhase(): Promise; + private async executePhase(phase: OrchestratorPhase): Promise; + private async assignPhaseTasks(phase: OrchestratorPhase): Promise; + private handleTaskCompleted(taskId: string): void; + private handleTaskFailed(taskId: string, error: string): void; + private async onPhaseTasksComplete(phase: OrchestratorPhase): Promise; + + // ── Internal: Verification ───────────────────────────────── + + private async verifyCurrentPhase(): Promise; + private async handleVerificationResult(phase: OrchestratorPhase, result: VerificationResult): Promise; + + // ── Internal: Replanning ─────────────────────────────────── + + private async replanPhase(phase: OrchestratorPhase, failures: string[]): Promise; + + // ── Internal: State Machine ──────────────────────────────── + + private setState(newState: OrchestratorState): void; + private advanceToNextPhase(): Promise; + private persist(): void; + private restore(): void; +} +``` + +**Key execution flow in `executePhase()`:** +1. Mark phase as `executing`, emit `phaseStarted` +2. For each task in phase: + - Create a `CreateTaskOptions` from `OrchestratorTask` + - Add to `TaskQueue` with proper dependencies + completion phrase + - Store the TaskQueue task ID in `OrchestratorTask.queueTaskId` +3. Poll task completion (listen to TaskQueue events) +4. When all tasks complete → call `onPhaseTasksComplete()` +5. `onPhaseTasksComplete()` triggers verification + +**How tasks get assigned to sessions:** +The OrchestratorLoop does NOT manage session assignment directly. It adds tasks to the existing TaskQueue and starts a mini poll loop that assigns pending tasks to idle sessions — the same pattern as RalphLoop's `assignTasks()`. This reuses existing session management. + +**Team agent flow:** +For phases with `teamStrategy.type === 'team'`: +- Start a single session with `CLAUDE_CODE_EXPERIMENTAL_AGENT_TEAMS=1` +- Instead of adding individual tasks to TaskQueue, send ONE comprehensive prompt to the lead +- The prompt instructs the lead to create teammates and delegate +- Monitor via TeamWatcher for team task completion + hook events +- Phase completion is detected via the lead's completion phrase + +### Step 5: `src/web/routes/orchestrator-routes.ts` — API endpoints +~300 lines. + +``` +POST /api/orchestrator/start — { goal, config? } → start planning +POST /api/orchestrator/approve — approve generated plan +POST /api/orchestrator/reject — { feedback } → reject + replan +POST /api/orchestrator/pause — pause execution +POST /api/orchestrator/resume — resume execution +POST /api/orchestrator/stop — stop orchestration +GET /api/orchestrator/status — full state + plan + stats +GET /api/orchestrator/plan — plan details only +POST /api/orchestrator/phase/:id/skip — skip a phase +POST /api/orchestrator/phase/:id/retry — retry a failed phase +``` + +Port dependency: `SessionPort & EventPort & RespawnPort & ConfigPort & InfraPort` + +The route module receives the OrchestratorLoop instance via the InfraPort (added to `createRouteContext()`). + +### Step 6: SSE Events — `src/web/sse-events.ts` additions + +```typescript +// ─── Orchestrator ──────────────────────────────────────────────────────────── + +/** Orchestrator state machine transitioned. */ +export const OrchestratorStateChanged = 'orchestrator:stateChanged' as const; +/** Orchestrator plan generated and ready for approval. */ +export const OrchestratorPlanReady = 'orchestrator:planReady' as const; +/** Orchestrator phase started executing. */ +export const OrchestratorPhaseStarted = 'orchestrator:phaseStarted' as const; +/** Orchestrator phase completed successfully. */ +export const OrchestratorPhaseCompleted = 'orchestrator:phaseCompleted' as const; +/** Orchestrator phase failed. */ +export const OrchestratorPhaseFailed = 'orchestrator:phaseFailed' as const; +/** Orchestrator verification result for a phase. */ +export const OrchestratorVerification = 'orchestrator:verification' as const; +/** Orchestrator task assigned to session. */ +export const OrchestratorTaskAssigned = 'orchestrator:taskAssigned' as const; +/** Orchestrator task completed. */ +export const OrchestratorTaskCompleted = 'orchestrator:taskCompleted' as const; +/** Orchestrator task failed. */ +export const OrchestratorTaskFailed = 'orchestrator:taskFailed' as const; +/** All phases completed successfully. */ +export const OrchestratorCompleted = 'orchestrator:completed' as const; +/** Orchestrator error. */ +export const OrchestratorError = 'orchestrator:error' as const; +``` + +11 new events. Add to `SseEvent` namespace object + mirror in `constants.js`. + +### Step 7: State persistence — `src/state-store.ts` additions + +Add to `AppState`: +```typescript +orchestrator?: OrchestratorPersistState; +``` + +Add methods: +```typescript +getOrchestratorState(): OrchestratorPersistState | null; +setOrchestratorState(state: Partial): void; +clearOrchestratorState(): void; +``` + +### Step 8: Server integration — `src/web/server.ts` modifications + +1. Import `OrchestratorLoop` and `registerOrchestratorRoutes` +2. Add `private orchestratorLoop: OrchestratorLoop` field +3. Initialize in constructor (lazy — created on first start, not at boot) +4. Add to `createRouteContext()` InfraPort: `orchestratorLoop: this.orchestratorLoop` +5. Wire up OrchestratorLoop events → SSE broadcasts +6. Register routes: `registerOrchestratorRoutes(this.app, ctx)` +7. Clean up in `stop()` + +### Step 9: `src/web/public/orchestrator-ui.js` — Frontend panel +~500 lines. New frontend module. + +**Load order**: After `panels-ui.js` (11), before `ralph-wizard.js` (13). So load order = 11.5. + +**UI elements:** +- Goal input form (text area + config toggles) +- Plan approval view (phase list, task details, approve/reject buttons) +- Execution dashboard (progress bar, phase cards, task status indicators) +- Agent activity panel (session count, team status) +- Controls (pause, resume, stop, skip phase, retry phase) + +**SSE listeners:** +- All 11 orchestrator events → update UI state +- Reuses existing session/respawn/team event handlers for agent monitoring + +### Step 10: `src/prompts/orchestrator.ts` — Prompt templates +~200 lines. + +Templates for: +- Phase execution prompt (tells Claude what to do in this phase) +- Team lead delegation prompt (instructs lead to create and coordinate teammates) +- Verification prompt (asks Claude to verify phase output) +- Replan prompt (gives failure context, asks for recovery steps) + +### Step 11: Constants, schemas, route barrel updates + +- `src/web/public/constants.js` — Add 11 SSE event mirrors +- `src/web/schemas.ts` — Add Zod schemas for orchestrator API input validation +- `src/web/routes/index.ts` — Export `registerOrchestratorRoutes` +- `src/web/ports/infra-port.ts` — Add `orchestratorLoop` to InfraPort +- `src/types/index.ts` — Export orchestrator types + +## Existing File Modifications Summary + +| File | Change | Lines | +|------|--------|-------| +| `src/types/index.ts` | Add orchestrator barrel export | +1 | +| `src/web/sse-events.ts` | Add 11 orchestrator events + SseEvent entries | +30 | +| `src/web/public/constants.js` | Mirror 11 SSE events | +15 | +| `src/web/routes/index.ts` | Export registerOrchestratorRoutes | +1 | +| `src/web/ports/infra-port.ts` | Add orchestratorLoop to InfraPort | +3 | +| `src/web/server.ts` | Initialize OrchestratorLoop, wire events, register routes | +40 | +| `src/web/schemas.ts` | Add orchestrator Zod schemas | +20 | +| `src/state-store.ts` | Add orchestrator state persistence | +20 | +| `src/web/public/app.js` | Add orchestrator SSE listeners + panel toggle | +30 | +| `src/web/public/index.html` | Add orchestrator-ui.js script tag | +1 | + +**Total new code**: ~2,300 lines across 6 new files +**Total modifications**: ~160 lines across 10 existing files + +## Implementation Execution Order + +This is the actual build order — each step is a commit checkpoint: + +1. **Types** — `src/types/orchestrator.ts` + barrel export. Zero risk, pure types. +2. **SSE events** — Add all 11 events to both `sse-events.ts` and `constants.js`. Wire in SseEvent namespace. +3. **State persistence** — Add orchestrator state to StateStore. Small, isolated change. +4. **Schemas** — Add Zod validation schemas for API input. +5. **Planner** — `src/orchestrator-planner.ts`. Can test in isolation. +6. **Verifier** — `src/orchestrator-verifier.ts`. Can test in isolation. +7. **Core loop** — `src/orchestrator-loop.ts`. The big one. Depends on planner + verifier. +8. **Prompts** — `src/prompts/orchestrator.ts`. Templates used by core loop. +9. **Port + routes** — `src/web/ports/infra-port.ts` update + `src/web/routes/orchestrator-routes.ts`. +10. **Server integration** — Wire OrchestratorLoop into WebServer. Routes become live. +11. **Frontend** — `src/web/public/orchestrator-ui.js` + app.js listeners + index.html script tag. +12. **Tests** — `test/orchestrator-*.test.ts`. +13. **Typecheck + lint** — Fix all issues, ensure CI passes. + +## Edge Cases & Error Handling + +- **Session limit reached**: Queue tasks and wait for sessions to free up (existing SessionManager handles this) +- **All sessions crash during phase**: Mark phase as failed, attempt replan +- **Verification flaky**: `moderate` mode allows test retries; `lenient` skips AI review +- **Plan too large**: Cap at 10 phases, 50 total tasks. Warn user. +- **Context overflow**: Auto-compact between phases. Respawn if needed (orchestrator state is external). +- **User pauses mid-phase**: Pause task assignment, don't cancel running tasks. Resume picks up where it left off. +- **Network/API errors during planning**: Retry plan generation up to 2 times, then fail with clear message. +- **Orchestrator vs Ralph conflict**: Mutually exclusive. Starting orchestrator stops Ralph if running. Starting Ralph stops orchestrator. + +## Testing Strategy + +- **Unit tests**: `test/orchestrator-planner.test.ts` — phase grouping algorithm, team strategy assignment +- **Unit tests**: `test/orchestrator-verifier.test.ts` — verification logic with mocked sessions +- **Integration tests**: `test/orchestrator-loop.test.ts` — state machine transitions, task lifecycle +- **Route tests**: `test/routes/orchestrator-routes.test.ts` — API validation, status responses + +All tests use `MockSession` pattern from existing test infrastructure. No real tmux needed. diff --git a/docs/orchestrator-loop-research.md b/docs/orchestrator-loop-research.md new file mode 100644 index 00000000..c9c8fa0f --- /dev/null +++ b/docs/orchestrator-loop-research.md @@ -0,0 +1,157 @@ +# Orchestrator Loop — Research Findings + +> Research doc for the new "Orchestrator Loop" feature. Not for GitHub. + +## What We're Building + +A new autonomous loop variant — **Orchestrator Loop** — that takes high-level user tasks, decomposes them into a detailed plan using team agents, and executes the plan step-by-step with quality gates. Unlike Ralph Loop (which executes a flat task queue), the Orchestrator coordinates **planning, delegation, and verification** as a continuous cycle. + +**Core idea**: User inputs a goal → Orchestrator creates a detailed plan → spins up team agents for parallel execution → validates each step → adapts the plan based on results → delivers polished output. + +## Existing Infrastructure Analysis + +### What We Can Reuse + +#### 1. Ralph Loop (`src/ralph-loop.ts`) +- **Pattern**: Poll loop with `start() → tick() → stop()` lifecycle +- **Reusable**: Event-driven task assignment, session completion handling, timeout management +- **Limitation**: Flat task queue — no concept of phases, dependencies between task groups, or adaptive replanning +- **Key insight**: `assignTaskToSession()` uses `session.sendInput(task.prompt)` — simple prompt injection into PTY + +#### 2. Task Queue (`src/task-queue.ts`) + Task (`src/task.ts`) +- **Already has**: Priority ordering, dependency tracking between tasks, completion phrase detection +- **Limitation**: No task *groups* or *phases*. Dependencies are task-to-task, not phase-to-phase +- **Key insight**: Tasks support `completionPhrase` — a string the task watches for in output. This is how Ralph knows a task is done + +#### 3. Plan Orchestrator (`src/plan-orchestrator.ts`) +- **Already has**: 2-agent plan generation (Research Agent → Planner Agent), TDD-aware plan items with P0/P1/P2 priorities +- **Output**: `PlanItem[]` with dependencies, verification criteria, TDD phases, complexity ratings +- **Limitation**: Plan generation only — no execution. Plans are generated then sit in state/UI for human review +- **Key insight**: Uses `Session` directly to run Claude subagent instances for research and planning. Returns structured JSON + +#### 4. Team Agents (`src/team-watcher.ts`, `~/.claude/teams/`) +- **Already has**: Team creation, member tracking, filesystem inbox messaging, task management via `~/.claude/tasks/{team-name}/` +- **Limitation**: Codeman can only *observe* teams (TeamWatcher is read-only polling), not *create* or *orchestrate* them +- **Key insight**: Teams are a Claude Code feature. Codeman monitors them but doesn't control them. We can't programmatically create teammates — Claude Code does that when you use `CLAUDE_CODE_EXPERIMENTAL_AGENT_TEAMS=1` + +#### 5. Respawn Controller (`src/respawn-controller.ts`) +- **Already has**: Preset-based automation (ralph-todo, overnight-autonomous), circuit breaker, health scoring +- **Key insight**: The `ralph-todo` preset (8s idle, 480min max) is designed for autonomous task execution. We'd need a new preset or make Orchestrator Loop set its own timing + +#### 6. Session Auto-Ops (`src/session-auto-ops.ts`) +- **Already has**: Auto-compact at token thresholds, auto-clear for context management +- **Key insight**: Critical for long Orchestrator runs — prevents context overflow during multi-step execution + +#### 7. Hooks (`src/hooks-config.ts`) +- **Already has**: `idle_prompt`, `stop`, `teammate_idle`, `task_completed` hook events +- **Key insight**: Hooks fire POST to `/api/hook-event` — this is how Codeman knows when Claude is idle, stopped, or completed a task. The Orchestrator Loop can listen to these same events + +### What We Need to Build New + +1. **Plan → Task decomposition**: Convert PlanOrchestrator output (PlanItem[]) into executable task groups with phase ordering +2. **Multi-phase execution engine**: Execute plan phases sequentially, tasks within phases in parallel +3. **Verification gates**: After each phase, run verification (test commands, AI review) before proceeding +4. **Adaptive replanning**: When a task fails or verification fails, generate a recovery plan +5. **Team agent orchestration**: Leverage Claude Code's agent teams for parallel execution within phases +6. **Progress tracking & UI**: Real-time dashboard showing plan progress, phase status, agent activity + +## How Teams Actually Work (Important Constraint) + +After deep research, here's the reality of agent teams: + +``` +User starts session with CLAUDE_CODE_EXPERIMENTAL_AGENT_TEAMS=1 + → Claude Code creates a team-lead + → Team-lead spawns teammates (in-process threads) + → Teammates appear as subagents (detected by SubagentWatcher) + → Communication via ~/.claude/teams/{name}/inboxes/{member}.json + → Tasks tracked in ~/.claude/tasks/{team-name}/{N}.json +``` + +**Codeman cannot programmatically create team members.** This is a Claude Code internal feature. However, Codeman CAN: +- Start a session that has teams enabled +- Send a prompt to the lead that instructs it to use agent teams +- Monitor team activity via TeamWatcher +- React to teammate_idle and task_completed hook events +- Read team task status from the filesystem + +**This means**: The Orchestrator Loop orchestrates at the *session prompt* level, not the *team member* level. We tell the lead what to do, and the lead decides how to use its team. + +## Architecture Decision: Prompt-Level Orchestration + +Given the team constraint, the Orchestrator Loop works by: + +1. **Planning phase**: Use PlanOrchestrator to generate a detailed plan from user input +2. **Execution phase**: Feed plan steps as prompts to sessions, one phase at a time +3. **Verification phase**: After each phase, run verification prompts and check results +4. **Adaptation phase**: If verification fails, generate recovery prompts + +The "team agents" aspect works by: +- Starting sessions with `CLAUDE_CODE_EXPERIMENTAL_AGENT_TEAMS=1` +- Crafting prompts that *instruct the lead to delegate* to teammates +- Monitoring team activity to track parallel progress +- The lead agent is smart enough to decompose work across its team + +## Key Technical Findings + +### Session Input Mechanics +```typescript +// From session.ts - how we send prompts +await session.sendInput(task.prompt); // Uses writeViaMux() internally +// writeViaMux() does: tmux send-keys -l "prompt text" + tmux send-keys Enter +// CRITICAL: Single-line only! Multi-line breaks Ink rendering +``` + +### Completion Detection Chain +``` +PTY output → RalphTracker.processData() → completion phrase fuzzy match + → CompletionConfidence scoring (multi-signal: promise tag + todos + exit signal) + → If confident → emit 'completionDetected' + → RalphLoop listens → marks task complete → assigns next +``` + +### How Plan Items Map to Tasks +```typescript +// PlanItem has: +interface PlanItem { + id: string; // "P0-001" + content: string; // "Implement error handling for API endpoints" + priority: 'P0' | 'P1' | 'P2'; + dependencies: string[]; // ["P0-000"] — other PlanItem IDs + verificationCriteria: string; + testCommand: string; + tddPhase: 'setup' | 'test' | 'impl' | 'verify' | 'review'; + complexity: 'low' | 'medium' | 'high'; +} + +// Task has: +interface CreateTaskOptions { + prompt: string; + priority: number; + dependencies: string[]; // Task IDs + completionPhrase: string; + timeoutMs: number; +} + +// Natural mapping: PlanItem.content → Task.prompt +// PlanItem.dependencies → Task.dependencies +// PlanItem.priority → Task.priority (P0=100, P1=50, P2=10) +// PlanItem.verificationCriteria → verification task prompt +``` + +### Context Management for Long Runs +- Auto-compact at ~110k tokens (configurable) +- Auto-clear at ~140k tokens (configurable) +- Respawn cycling: kill + restart session to reset context entirely +- For Orchestrator: we want compact between phases, respawn between major milestones + +## Risk Assessment + +| Risk | Severity | Mitigation | +|------|----------|------------| +| Context overflow during complex phases | High | Auto-compact between tasks, respawn between phases | +| Team agents not predictable | Medium | Orchestrate at session level, let Claude decide team delegation | +| Plan too ambitious → infinite loop | High | Phase budgets (max attempts per phase), circuit breaker | +| Verification too strict → blocks progress | Medium | Configurable strictness, human override via UI | +| Single-line prompt limit | Medium | Use CLAUDE.md file for complex instructions, prompt references file | +| Long planning phase delays execution | Low | Show plan for approval before execution | diff --git a/src/orchestrator-loop.ts b/src/orchestrator-loop.ts new file mode 100644 index 00000000..f5113780 --- /dev/null +++ b/src/orchestrator-loop.ts @@ -0,0 +1,912 @@ +/** + * @fileoverview Orchestrator Loop — phased plan execution with team agents. + * + * State machine that generates plans from user goals, executes them + * phase-by-phase with verification gates, and adapts on failure. + * + * States: idle → planning → approval → executing → verifying → (replanning) → completed/failed + * + * Key exports: + * - `OrchestratorLoop` class — main engine, extends EventEmitter + * - `OrchestratorLoopEvents` interface — typed event map + * + * Lifecycle: `start(goal)` → plan → approve → execute phases → verify → complete + * + * @dependencies orchestrator-planner (plan generation), orchestrator-verifier (phase verification), + * session-manager (sessions), task-queue (task execution), state-store (persistence), + * prompts/orchestrator (prompt templates) + * @consumedby web/server (orchestrator routes, SSE) + * @emits stateChanged, planReady, phaseStarted, phaseCompleted, phaseFailed, + * taskAssigned, taskCompleted, taskFailed, verificationResult, completed, error + * @persistence Orchestrator state saved to `~/.codeman/state.json` (orchestrator key) + * + * @module orchestrator-loop + */ + +import { EventEmitter } from 'node:events'; +import { getSessionManager, SessionManager } from './session-manager.js'; +import { getTaskQueue, TaskQueue } from './task-queue.js'; +import { getStore, StateStore } from './state-store.js'; +import { OrchestratorPlanner } from './orchestrator-planner.js'; +import { OrchestratorVerifier } from './orchestrator-verifier.js'; +import { PHASE_EXECUTION_PROMPT, REPLAN_PROMPT, SINGLE_TASK_PROMPT, TEAM_LEAD_PROMPT } from './prompts/index.js'; +import type { TerminalMultiplexer } from './mux-interface.js'; +import type { CreateTaskOptions } from './task.js'; +import { + type OrchestratorState, + type OrchestratorPlan, + type OrchestratorPhase, + type OrchestratorTask, + type OrchestratorConfig, + type OrchestratorStats, + type OrchestratorPersistState, + type VerificationResult, + DEFAULT_ORCHESTRATOR_CONFIG, + createInitialOrchestratorStats, + getErrorMessage, +} from './types.js'; + +// ═══════════════════════════════════════════════════════════════ +// Constants +// ═══════════════════════════════════════════════════════════════ + +/** Poll interval for checking task completion within a phase (2 seconds) */ +const PHASE_POLL_INTERVAL_MS = 2000; + +/** Delay between phase completion and verification (1 second) */ +const POST_PHASE_DELAY_MS = 1000; + +// ═══════════════════════════════════════════════════════════════ +// Events +// ═══════════════════════════════════════════════════════════════ + +export interface OrchestratorLoopEvents { + stateChanged: (state: OrchestratorState, prevState: OrchestratorState) => void; + planReady: (plan: OrchestratorPlan) => void; + phaseStarted: (phase: OrchestratorPhase) => void; + phaseCompleted: (phase: OrchestratorPhase) => void; + phaseFailed: (phase: OrchestratorPhase, reason: string) => void; + taskAssigned: (task: OrchestratorTask, sessionId: string) => void; + taskCompleted: (task: OrchestratorTask) => void; + taskFailed: (task: OrchestratorTask, error: string) => void; + verificationResult: (phase: OrchestratorPhase, result: VerificationResult) => void; + completed: (stats: OrchestratorStats) => void; + error: (error: Error) => void; +} + +// ═══════════════════════════════════════════════════════════════ +// OrchestratorLoop +// ═══════════════════════════════════════════════════════════════ + +export class OrchestratorLoop extends EventEmitter { + private _state: OrchestratorState = 'idle'; + private plan: OrchestratorPlan | null = null; + private currentPhaseIndex = 0; + private config: OrchestratorConfig; + private stats: OrchestratorStats; + private startedAt: number | null = null; + private completedAt: number | null = null; + + private workingDir: string; + private planner: OrchestratorPlanner; + private verifier: OrchestratorVerifier; + private sessionManager: SessionManager; + private taskQueue: TaskQueue; + private store: StateStore; + + /** State before pause (to resume to correct state) */ + private pausedState: OrchestratorState | null = null; + + /** Phase poll timer for checking task completion */ + private phasePollTimer: NodeJS.Timeout | null = null; + + /** Session completion listener (bound for cleanup) */ + private sessionCompletionListener: ((sessionId: string, phrase: string) => void) | null = null; + + /** Active sessions assigned to current phase */ + private phaseSessionIds: Set = new Set(); + + constructor(mux: TerminalMultiplexer, workingDir: string, config?: Partial) { + super(); + this.workingDir = workingDir; + this.config = { ...DEFAULT_ORCHESTRATOR_CONFIG, ...config }; + this.stats = createInitialOrchestratorStats(); + this.sessionManager = getSessionManager(); + this.taskQueue = getTaskQueue(); + this.store = getStore(); + this.planner = new OrchestratorPlanner(mux, workingDir, this.config); + this.verifier = new OrchestratorVerifier(this.config); + + // Restore state if crashed while running + this.restore(); + } + + // ═══════════════════════════════════════════════════════════════ + // Public API — Lifecycle + // ═══════════════════════════════════════════════════════════════ + + /** Start orchestration with a goal. Transitions: idle → planning */ + async start(goal: string): Promise { + if (this._state !== 'idle' && this._state !== 'failed' && this._state !== 'completed') { + throw new Error(`Cannot start from state "${this._state}"`); + } + + this.reset(); + this.startedAt = Date.now(); + this.setState('planning'); + + try { + const plan = await this.planner.generatePlan(goal, (phase, detail) => { + // Forward planning progress — could emit an event here + console.log(`[Orchestrator] Planning: ${phase} — ${detail}`); + }); + + if ((this._state as OrchestratorState) !== 'planning') { + // Cancelled during planning + return; + } + + this.plan = plan; + this.persist(); + + if (this.config.autoApprove) { + this.setState('executing'); + await this.executeCurrentPhase(); + } else { + this.setState('approval'); + this.emit('planReady', plan); + } + } catch (err) { + this.handleError(err); + } + } + + /** Approve the generated plan. Transitions: approval → executing */ + async approve(): Promise { + if ((this._state as OrchestratorState) !== 'approval') { + throw new Error(`Cannot approve from state "${this._state}"`); + } + if (!this.plan) { + throw new Error('No plan to approve'); + } + + this.setState('executing'); + await this.executeCurrentPhase(); + } + + /** Reject plan with feedback. Transitions: approval → planning (regenerate) */ + async reject(feedback: string): Promise { + if ((this._state as OrchestratorState) !== 'approval') { + throw new Error(`Cannot reject from state "${this._state}"`); + } + if (!this.plan) { + throw new Error('No plan to reject'); + } + + const goal = this.plan.goal + '\n\nFeedback on previous plan: ' + feedback; + this.plan = null; + this.setState('planning'); + + try { + const plan = await this.planner.generatePlan(goal); + + if ((this._state as OrchestratorState) !== 'planning') return; + + this.plan = plan; + this.persist(); + this.setState('approval'); + this.emit('planReady', plan); + } catch (err) { + this.handleError(err); + } + } + + /** Pause execution. Saves current state. */ + pause(): void { + if (this._state === 'idle' || this._state === 'paused' || this._state === 'completed' || this._state === 'failed') { + return; + } + this.pausedState = this._state; + this.clearPhasePoll(); + this.setState('paused'); + } + + /** Resume from pause. */ + async resume(): Promise { + if (this._state !== 'paused' || !this.pausedState) { + throw new Error('Not paused'); + } + + const resumeTo = this.pausedState; + this.pausedState = null; + this.setState(resumeTo); + + // Re-enter the appropriate phase of execution + if (resumeTo === 'executing') { + await this.executeCurrentPhase(); + } else if (resumeTo === 'verifying') { + await this.verifyCurrentPhase(); + } + } + + /** Stop everything and clean up. */ + async stop(): Promise { + this.clearPhasePoll(); + this.cleanupTaskHandlers(); + await this.planner.cancel(); + this.setState('idle'); + this.store.clearOrchestratorState(); + } + + /** Skip a specific phase. */ + async skipPhase(phaseId: string): Promise { + if (!this.plan) return; + + const phase = this.plan.phases.find((p) => p.id === phaseId); + if (!phase) throw new Error(`Phase "${phaseId}" not found`); + + phase.status = 'skipped'; + phase.completedAt = Date.now(); + this.persist(); + + // If this is the current phase, advance + if (this.plan.phases[this.currentPhaseIndex]?.id === phaseId) { + await this.advanceToNextPhase(); + } + } + + /** Retry a failed phase. */ + async retryPhase(phaseId: string): Promise { + if (!this.plan) return; + if (this._state !== 'executing' && this._state !== 'failed') { + throw new Error(`Cannot retry from state "${this._state}"`); + } + + const phaseIndex = this.plan.phases.findIndex((p) => p.id === phaseId); + if (phaseIndex === -1) throw new Error(`Phase "${phaseId}" not found`); + + const phase = this.plan.phases[phaseIndex]; + phase.status = 'pending'; + phase.attempts = 0; + for (const task of phase.tasks) { + task.status = 'pending'; + task.error = null; + task.assignedSessionId = null; + task.queueTaskId = null; + } + + this.currentPhaseIndex = phaseIndex; + this.setState('executing'); + await this.executeCurrentPhase(); + } + + // ═══════════════════════════════════════════════════════════════ + // Public API — Getters + // ═══════════════════════════════════════════════════════════════ + + get state(): OrchestratorState { + return this._state; + } + + getPlan(): OrchestratorPlan | null { + return this.plan; + } + + getCurrentPhase(): OrchestratorPhase | null { + if (!this.plan) return null; + return this.plan.phases[this.currentPhaseIndex] ?? null; + } + + getStats(): OrchestratorStats { + return { ...this.stats }; + } + + getStatus(): OrchestratorPersistState { + return { + state: this._state, + plan: this.plan, + currentPhaseIndex: this.currentPhaseIndex, + startedAt: this.startedAt, + completedAt: this.completedAt, + config: this.config, + stats: this.stats, + }; + } + + isRunning(): boolean { + return this._state !== 'idle' && this._state !== 'completed' && this._state !== 'failed'; + } + + // ═══════════════════════════════════════════════════════════════ + // Internal — Phase Execution + // ═══════════════════════════════════════════════════════════════ + + private async executeCurrentPhase(): Promise { + if (!this.plan || this._state !== 'executing') return; + + const phase = this.plan.phases[this.currentPhaseIndex]; + if (!phase) { + // All phases done + await this.handleCompletion(); + return; + } + + // Skip already completed/skipped phases + if (phase.status === 'passed' || phase.status === 'skipped') { + await this.advanceToNextPhase(); + return; + } + + phase.status = 'executing'; + phase.startedAt = Date.now(); + phase.attempts++; + this.persist(); + this.emit('phaseStarted', phase); + + try { + await this.assignPhaseTasks(phase); + this.startPhasePoll(phase); + } catch (err) { + this.handlePhaseError(phase, getErrorMessage(err)); + } + } + + private async assignPhaseTasks(phase: OrchestratorPhase): Promise { + // For team strategy, send a single comprehensive prompt to a lead session + if (phase.teamStrategy.type === 'team') { + await this.assignTeamPhase(phase); + return; + } + + // For single/parallel strategy, add individual tasks to TaskQueue + for (const task of phase.tasks) { + if (task.status !== 'pending') continue; + + const prompt = this.buildTaskPrompt(task, phase); + const taskOptions: CreateTaskOptions = { + prompt, + workingDir: this.workingDir, + priority: 100 - phase.order, // Earlier phases get higher priority + completionPhrase: task.completionPhrase, + timeoutMs: Math.min(task.timeoutMs, this.config.phaseTimeoutMs), + }; + + const queueTask = this.taskQueue.addTask(taskOptions); + task.queueTaskId = queueTask.id; + task.status = 'running'; + } + + this.persist(); + this.setupTaskHandlers(); + + // Manually assign tasks to idle sessions + await this.assignQueuedTasksToSessions(); + } + + private async assignTeamPhase(phase: OrchestratorPhase): Promise { + const teamConfig = phase.teamStrategy.type === 'team' ? phase.teamStrategy.config : null; + if (!teamConfig) return; + + // Find or use an idle session + const sessions = this.sessionManager.getIdleSessions(); + if (sessions.length === 0) { + throw new Error('No idle sessions available for team phase execution'); + } + + const session = sessions[0]; + this.phaseSessionIds.add(session.id); + + // Mark all tasks as running under this session + for (const task of phase.tasks) { + task.status = 'running'; + task.assignedSessionId = session.id; + } + + // Build and send the team lead prompt + const prompt = TEAM_LEAD_PROMPT.replace('{PHASE_NAME}', phase.name) + .replace('{TASK_LIST}', phase.tasks.map((t, i) => `${i + 1}. ${t.prompt}`).join('\n')) + .replace('{TEAMMATE_HINTS}', teamConfig.suggestedTeammates.map((h, i) => `${i + 1}. ${h}`).join('\n')) + .replace('{COMPLETION_PHRASE}', `${phase.id.toUpperCase()}_COMPLETE`); + + // Create a TaskQueue task for the entire phase + const queueTask = this.taskQueue.addTask({ + prompt, + workingDir: this.workingDir, + priority: 100 - phase.order, + completionPhrase: `${phase.id.toUpperCase()}_COMPLETE`, + timeoutMs: this.config.phaseTimeoutMs, + }); + + // Link all phase tasks to this single queue task + for (const task of phase.tasks) { + task.queueTaskId = queueTask.id; + } + + this.persist(); + this.setupTaskHandlers(); + + // Assign the task to the session + try { + queueTask.assign(session.id); + session.assignTask(queueTask.id); + this.taskQueue.updateTask(queueTask); + await session.sendInput(prompt); + } catch (err) { + queueTask.fail(getErrorMessage(err)); + this.taskQueue.updateTask(queueTask); + throw err; + } + } + + private async assignQueuedTasksToSessions(): Promise { + const idleSessions = this.sessionManager.getIdleSessions(); + const maxSessions = + this.getCurrentPhase()?.teamStrategy.type === 'parallel' + ? (this.getCurrentPhase()?.teamStrategy as { type: 'parallel'; maxSessions: number }).maxSessions + : 1; + + const sessionsToUse = idleSessions.slice(0, maxSessions); + + for (const session of sessionsToUse) { + const task = this.taskQueue.next(); + if (!task) break; + + try { + task.assign(session.id); + session.assignTask(task.id); + this.taskQueue.updateTask(task); + await session.sendInput(task.prompt); + + this.phaseSessionIds.add(session.id); + + // Find the orchestrator task linked to this queue task + const orchTask = this.findOrchestratorTaskByQueueId(task.id); + if (orchTask) { + orchTask.assignedSessionId = session.id; + orchTask.startedAt = Date.now(); + this.emit('taskAssigned', orchTask, session.id); + } + } catch (err) { + task.fail(getErrorMessage(err)); + session.clearTask(); + this.taskQueue.updateTask(task); + } + } + } + + // ═══════════════════════════════════════════════════════════════ + // Internal — Task Completion Tracking + // ═══════════════════════════════════════════════════════════════ + + private setupTaskHandlers(): void { + this.cleanupTaskHandlers(); + + this.sessionCompletionListener = (_sessionId: string, _phrase: string) => { + // Session completion — check if it's related to our phase tasks + this.checkPhaseCompletion(); + }; + + this.sessionManager.on('sessionCompletion', this.sessionCompletionListener); + } + + private cleanupTaskHandlers(): void { + if (this.sessionCompletionListener) { + this.sessionManager.off('sessionCompletion', this.sessionCompletionListener); + this.sessionCompletionListener = null; + } + } + + private handleTaskCompleted(queueTaskId: string): void { + const orchTask = this.findOrchestratorTaskByQueueId(queueTaskId); + if (!orchTask) return; + + orchTask.status = 'completed'; + orchTask.completedAt = Date.now(); + this.stats.totalTasksCompleted++; + this.persist(); + this.emit('taskCompleted', orchTask); + + this.checkPhaseCompletion(); + } + + private handleTaskFailed(queueTaskId: string, error: string): void { + const orchTask = this.findOrchestratorTaskByQueueId(queueTaskId); + if (!orchTask) return; + + orchTask.status = 'failed'; + orchTask.error = error; + this.stats.totalTasksFailed++; + this.persist(); + this.emit('taskFailed', orchTask, error); + + // Check if we should retry the task or fail the phase + if (orchTask.retries < 2) { + orchTask.retries++; + orchTask.status = 'pending'; + orchTask.error = null; + orchTask.queueTaskId = null; + // Will be re-queued on next poll + } else { + this.checkPhaseCompletion(); + } + } + + private startPhasePoll(phase: OrchestratorPhase): void { + this.clearPhasePoll(); + this.phasePollTimer = setInterval(() => { + if (this._state !== 'executing') { + this.clearPhasePoll(); + return; + } + this.pollPhaseStatus(phase); + }, PHASE_POLL_INTERVAL_MS); + } + + private clearPhasePoll(): void { + if (this.phasePollTimer) { + clearInterval(this.phasePollTimer); + this.phasePollTimer = null; + } + } + + private pollPhaseStatus(phase: OrchestratorPhase): void { + // Check for queued tasks that need assignment + const pendingTasks = phase.tasks.filter((t) => t.status === 'pending' && !t.queueTaskId); + if (pendingTasks.length > 0) { + // Re-queue pending tasks + for (const task of pendingTasks) { + const prompt = this.buildTaskPrompt(task, phase); + const queueTask = this.taskQueue.addTask({ + prompt, + workingDir: this.workingDir, + priority: 100 - phase.order, + completionPhrase: task.completionPhrase, + timeoutMs: Math.min(task.timeoutMs, this.config.phaseTimeoutMs), + }); + task.queueTaskId = queueTask.id; + task.status = 'running'; + } + this.assignQueuedTasksToSessions().catch(() => {}); // Best effort + } + + // Check completion status of queue tasks + for (const task of phase.tasks) { + if (task.status === 'running' && task.queueTaskId) { + const queueTask = this.taskQueue.getTask(task.queueTaskId); + if (queueTask) { + if (queueTask.isCompleted()) { + this.handleTaskCompleted(task.queueTaskId); + } else if (queueTask.isFailed()) { + this.handleTaskFailed(task.queueTaskId, queueTask.error || 'Task failed'); + } + } + } + } + + this.checkPhaseCompletion(); + } + + private checkPhaseCompletion(): void { + if (this._state !== 'executing') return; + + const phase = this.getCurrentPhase(); + if (!phase) return; + + const allDone = phase.tasks.every((t) => t.status === 'completed' || t.status === 'failed'); + if (!allDone) return; + + const anyFailed = phase.tasks.some((t) => t.status === 'failed'); + + this.clearPhasePoll(); + + if (anyFailed) { + // Phase has failed tasks + this.handlePhaseError(phase, 'One or more tasks failed'); + } else { + // All tasks completed — run verification + setTimeout(() => { + this.verifyCurrentPhase().catch((err) => this.handleError(err)); + }, POST_PHASE_DELAY_MS); + } + } + + // ═══════════════════════════════════════════════════════════════ + // Internal — Verification + // ═══════════════════════════════════════════════════════════════ + + private async verifyCurrentPhase(): Promise { + if (!this.plan) return; + + const phase = this.plan.phases[this.currentPhaseIndex]; + if (!phase) return; + + // Skip verification if no criteria defined + if (phase.verificationCriteria.length === 0 && phase.testCommands.length === 0) { + phase.status = 'passed'; + phase.completedAt = Date.now(); + phase.durationMs = phase.startedAt ? Date.now() - phase.startedAt : null; + this.stats.phasesCompleted++; + this.persist(); + this.emit('phaseCompleted', phase); + await this.advanceToNextPhase(); + return; + } + + this.setState('verifying'); + + // Get a session for verification + const sessions = this.sessionManager.getIdleSessions(); + if (sessions.length === 0) { + // No idle sessions — mark as passed (can't verify) + phase.status = 'passed'; + phase.completedAt = Date.now(); + phase.durationMs = phase.startedAt ? Date.now() - phase.startedAt : null; + this.stats.phasesCompleted++; + this.persist(); + this.emit('phaseCompleted', phase); + this.setState('executing'); + await this.advanceToNextPhase(); + return; + } + + try { + const result = await this.verifier.verifyPhase(phase, sessions[0]); + this.emit('verificationResult', phase, result); + + if (result.passed) { + phase.status = 'passed'; + phase.completedAt = Date.now(); + phase.durationMs = phase.startedAt ? Date.now() - phase.startedAt : null; + this.stats.phasesCompleted++; + this.persist(); + this.emit('phaseCompleted', phase); + this.setState('executing'); + await this.advanceToNextPhase(); + } else { + // Verification failed — attempt replan + await this.handleVerificationFailure(phase, result); + } + } catch (err) { + // Verification error — treat as pass (don't block on verification bugs) + console.warn('[Orchestrator] Verification error, treating as pass:', err); + phase.status = 'passed'; + phase.completedAt = Date.now(); + phase.durationMs = phase.startedAt ? Date.now() - phase.startedAt : null; + this.stats.phasesCompleted++; + this.persist(); + this.emit('phaseCompleted', phase); + this.setState('executing'); + await this.advanceToNextPhase(); + } + } + + private async handleVerificationFailure(phase: OrchestratorPhase, result: VerificationResult): Promise { + if (phase.attempts >= phase.maxAttempts) { + // Max retries exceeded + phase.status = 'failed'; + phase.completedAt = Date.now(); + phase.durationMs = phase.startedAt ? Date.now() - phase.startedAt : null; + this.stats.phasesFailed++; + this.persist(); + this.emit('phaseFailed', phase, `Verification failed after ${phase.attempts} attempts: ${result.summary}`); + this.setState('failed'); + return; + } + + // Replan and retry + this.stats.replanCount++; + this.setState('replanning'); + + try { + await this.replanPhase(phase, result); + // Reset task states for retry + for (const task of phase.tasks) { + task.status = 'pending'; + task.error = null; + task.assignedSessionId = null; + task.queueTaskId = null; + task.completedAt = null; + task.startedAt = null; + } + phase.status = 'pending'; + phase.startedAt = null; + this.persist(); + + this.setState('executing'); + await this.executeCurrentPhase(); + } catch (err) { + this.handleError(err); + } + } + + private async replanPhase(phase: OrchestratorPhase, result: VerificationResult): Promise { + const sessions = this.sessionManager.getIdleSessions(); + if (sessions.length === 0) return; + + const prompt = REPLAN_PROMPT.replace('{PHASE_NAME}', phase.name) + .replace('{ATTEMPT_NUMBER}', String(phase.attempts)) + .replace('{MAX_ATTEMPTS}', String(phase.maxAttempts)) + .replace('{FAILURE_SUMMARY}', result.summary) + .replace('{SUGGESTIONS}', result.suggestions.join('\n')) + .replace('{ORIGINAL_TASKS}', phase.tasks.map((t, i) => `${i + 1}. ${t.prompt}`).join('\n')) + .replace('{COMPLETION_PHRASE}', phase.tasks[0]?.completionPhrase || `${phase.id.toUpperCase()}_FIXED`); + + // Send the replan prompt to an idle session + await sessions[0].sendInput(prompt); + } + + // ═══════════════════════════════════════════════════════════════ + // Internal — State Machine + // ═══════════════════════════════════════════════════════════════ + + private setState(newState: OrchestratorState): void { + const prev = this._state; + if (prev === newState) return; + this._state = newState; + this.persist(); + this.emit('stateChanged', newState, prev); + } + + private async advanceToNextPhase(): Promise { + this.currentPhaseIndex++; + this.phaseSessionIds.clear(); + this.persist(); + + if (!this.plan || this.currentPhaseIndex >= this.plan.phases.length) { + await this.handleCompletion(); + } else { + // Compact between phases if configured + if (this.config.compactBetweenPhases) { + const sessions = this.sessionManager.getIdleSessions(); + for (const session of sessions) { + try { + await session.writeViaMux('/compact'); + } catch { + // Best effort + } + } + // Brief delay for compact to take effect + await new Promise((resolve) => setTimeout(resolve, 2000)); + } + + await this.executeCurrentPhase(); + } + } + + private async handleCompletion(): Promise { + this.completedAt = Date.now(); + this.stats.totalDurationMs = this.startedAt ? this.completedAt - this.startedAt : 0; + this.clearPhasePoll(); + this.cleanupTaskHandlers(); + this.setState('completed'); + this.emit('completed', this.stats); + } + + private handlePhaseError(phase: OrchestratorPhase, error: string): void { + if (phase.attempts >= phase.maxAttempts) { + phase.status = 'failed'; + phase.completedAt = Date.now(); + phase.durationMs = phase.startedAt ? Date.now() - phase.startedAt : null; + this.stats.phasesFailed++; + this.persist(); + this.emit('phaseFailed', phase, error); + this.setState('failed'); + } else { + // Retry the phase + for (const task of phase.tasks) { + if (task.status === 'failed') { + task.status = 'pending'; + task.error = null; + task.queueTaskId = null; + task.assignedSessionId = null; + } + } + phase.status = 'pending'; + this.persist(); + this.executeCurrentPhase().catch((err) => this.handleError(err)); + } + } + + private handleError(err: unknown): void { + const error = err instanceof Error ? err : new Error(getErrorMessage(err)); + console.error('[Orchestrator] Error:', error.message); + this.setState('failed'); + this.emit('error', error); + } + + // ═══════════════════════════════════════════════════════════════ + // Internal — Persistence + // ═══════════════════════════════════════════════════════════════ + + private persist(): void { + this.store.setOrchestratorState(this.getStatus()); + } + + private restore(): void { + const saved = this.store.getOrchestratorState(); + if (!saved) return; + + // If we crashed while running, reset to failed + if (saved.state === 'executing' || saved.state === 'verifying' || saved.state === 'replanning') { + this._state = 'failed'; + this.plan = saved.plan; + this.currentPhaseIndex = saved.currentPhaseIndex; + this.startedAt = saved.startedAt; + this.config = saved.config; + this.stats = saved.stats; + this.store.setOrchestratorState({ ...saved, state: 'failed' }); + } else if (saved.state === 'planning' || saved.state === 'approval') { + // Planning/approval — reset to idle (plan is lost) + this.store.clearOrchestratorState(); + } else if (saved.state === 'completed' || saved.state === 'failed') { + // Preserve completed/failed state for UI display + this._state = saved.state; + this.plan = saved.plan; + this.currentPhaseIndex = saved.currentPhaseIndex; + this.startedAt = saved.startedAt; + this.completedAt = saved.completedAt; + this.config = saved.config; + this.stats = saved.stats; + } + } + + private reset(): void { + this._state = 'idle'; + this.plan = null; + this.currentPhaseIndex = 0; + this.startedAt = null; + this.completedAt = null; + this.stats = createInitialOrchestratorStats(); + this.pausedState = null; + this.phaseSessionIds.clear(); + this.clearPhasePoll(); + this.cleanupTaskHandlers(); + } + + // ═══════════════════════════════════════════════════════════════ + // Internal — Helpers + // ═══════════════════════════════════════════════════════════════ + + private buildTaskPrompt(task: OrchestratorTask, phase: OrchestratorPhase): string { + if (phase.tasks.length === 1) { + // Single task — use simpler prompt + const completedPhases = this.getCompletedPhasesSummary(); + return SINGLE_TASK_PROMPT.replace('{TASK}', task.prompt) + .replace('{GOAL}', this.plan?.goal || '') + .replace('{CONTEXT}', completedPhases ? `Previous phases completed: ${completedPhases}` : '') + .replace('{COMPLETION_PHRASE}', task.completionPhrase); + } + + // Multi-task phase — use full prompt + return PHASE_EXECUTION_PROMPT.replace('{PHASE_NAME}', phase.name) + .replace('{GOAL}', this.plan?.goal || '') + .replace('{COMPLETED_PHASES}', this.getCompletedPhasesSummary() || 'None yet') + .replace('{TASK_LIST}', phase.tasks.map((t, i) => `${i + 1}. ${t.prompt}`).join('\n')) + .replace('{VERIFICATION_CRITERIA}', phase.verificationCriteria.join('\n') || 'No specific criteria') + .replace('{COMPLETION_PHRASE}', task.completionPhrase); + } + + private getCompletedPhasesSummary(): string { + if (!this.plan) return ''; + return this.plan.phases + .filter((p) => p.status === 'passed' || p.status === 'skipped') + .map((p) => `${p.name}: ${p.status}`) + .join(', '); + } + + private findOrchestratorTaskByQueueId(queueTaskId: string): OrchestratorTask | null { + if (!this.plan) return null; + for (const phase of this.plan.phases) { + for (const task of phase.tasks) { + if (task.queueTaskId === queueTaskId) return task; + } + } + return null; + } + + /** Clean up resources when the loop is being destroyed. */ + destroy(): void { + this.clearPhasePoll(); + this.cleanupTaskHandlers(); + } +} diff --git a/src/orchestrator-planner.ts b/src/orchestrator-planner.ts new file mode 100644 index 00000000..748926b9 --- /dev/null +++ b/src/orchestrator-planner.ts @@ -0,0 +1,412 @@ +/** + * @fileoverview Orchestrator plan generation — converts goals into phased plans. + * + * Wraps PlanOrchestrator for AI-powered plan generation, then groups the + * resulting PlanItems into sequential phases with team strategies and + * verification criteria. + * + * Phase grouping algorithm: + * 1. Topological sort by dependencies (Kahn's algorithm) + * 2. Group into dependency layers + * 3. Sub-group by TDD phase within layers + * 4. Merge small adjacent phases + * 5. Assign team strategies based on parallelism potential + * + * Key exports: + * - `OrchestratorPlanner` class — plan generation + phase grouping + * + * @dependencies plan-orchestrator (AI plan generation), types (OrchestratorPlan, PlanItem) + * @consumedby orchestrator-loop + * + * @module orchestrator-planner + */ + +import { v4 as uuidv4 } from 'uuid'; +import { PlanOrchestrator, type DetailedPlanResult, type ProgressCallback } from './plan-orchestrator.js'; +import type { TerminalMultiplexer } from './mux-interface.js'; +import type { + PlanItem, + TddPhase, + OrchestratorPlan, + OrchestratorPhase, + OrchestratorTask, + OrchestratorConfig, + TeamStrategy, + PhaseStatus, +} from './types.js'; + +// ═══════════════════════════════════════════════════════════════ +// Constants +// ═══════════════════════════════════════════════════════════════ + +/** Maximum number of phases (prevents runaway plans) */ +const MAX_PHASES = 10; + +/** Maximum total tasks across all phases */ +const MAX_TOTAL_TASKS = 50; + +/** Default task timeout (10 minutes) */ +const DEFAULT_TASK_TIMEOUT_MS = 10 * 60 * 1000; + +/** Minimum tasks in a phase before it gets merged with adjacent */ +const MIN_PHASE_TASKS = 2; + +/** TDD phase ordering for grouping */ +const TDD_PHASE_ORDER: Record = { + setup: 0, + test: 1, + impl: 2, + verify: 3, + review: 4, +}; + +// ═══════════════════════════════════════════════════════════════ +// OrchestratorPlanner +// ═══════════════════════════════════════════════════════════════ + +export class OrchestratorPlanner { + private mux: TerminalMultiplexer; + private workingDir: string; + private config: OrchestratorConfig; + private orchestrator: PlanOrchestrator | null = null; + + constructor(mux: TerminalMultiplexer, workingDir: string, config: OrchestratorConfig) { + this.mux = mux; + this.workingDir = workingDir; + this.config = config; + } + + /** + * Generate a phased plan from a user goal. + * + * Uses PlanOrchestrator for AI plan generation, then groups results into phases. + */ + async generatePlan(goal: string, onProgress?: ProgressCallback): Promise { + const startTime = Date.now(); + + // Create a PlanOrchestrator for this plan generation + this.orchestrator = new PlanOrchestrator(this.mux, this.workingDir, undefined, { + defaultModel: this.config.plannerModel, + }); + + try { + onProgress?.('planning', 'Generating detailed plan...'); + + const result: DetailedPlanResult = await this.orchestrator.generateDetailedPlan(goal, onProgress); + + if (!result.success || !result.items || result.items.length === 0) { + throw new Error(result.error || 'Plan generation returned no items'); + } + + // Cap total tasks + const items = result.items.slice(0, MAX_TOTAL_TASKS); + + onProgress?.('grouping', 'Organizing plan into phases...'); + + // Group items into phases + const phases = this.groupIntoPhases(items, goal); + + // Assign team strategies + this.assignTeamStrategies(phases); + + // Generate unique completion phrases + this.generateCompletionPhrases(phases); + + const plan: OrchestratorPlan = { + id: uuidv4(), + goal, + createdAt: Date.now(), + phases, + metadata: { + totalTasks: phases.reduce((sum, p) => sum + p.tasks.length, 0), + estimatedComplexity: this.estimateComplexity(items), + modelUsed: this.config.plannerModel, + planDurationMs: Date.now() - startTime, + }, + }; + + return plan; + } finally { + this.orchestrator = null; + } + } + + /** Cancel in-progress plan generation. */ + async cancel(): Promise { + if (this.orchestrator) { + await this.orchestrator.cancel(); + this.orchestrator = null; + } + } + + // ═══════════════════════════════════════════════════════════════ + // Phase Grouping + // ═══════════════════════════════════════════════════════════════ + + /** + * Group PlanItems into sequential phases. + * + * Algorithm: + * 1. Build dependency graph and assign IDs to items without them + * 2. Topological sort into dependency layers (Kahn's algorithm) + * 3. Sub-group within each layer by TDD phase + * 4. Merge small phases with their neighbors + */ + private groupIntoPhases(items: PlanItem[], _goal: string): OrchestratorPhase[] { + // Ensure all items have IDs + const indexedItems = items.map((item, i) => ({ + ...item, + id: item.id || `task-${i}`, + })); + + // Build adjacency and in-degree for Kahn's algorithm + const idSet = new Set(indexedItems.map((item) => item.id!)); + const inDegree = new Map(); + const dependents = new Map(); // id → items that depend on it + + for (const item of indexedItems) { + inDegree.set(item.id!, 0); + dependents.set(item.id!, []); + } + + for (const item of indexedItems) { + const deps = (item.dependencies || []).filter((d) => idSet.has(d)); + inDegree.set(item.id!, deps.length); + for (const dep of deps) { + dependents.get(dep)!.push(item.id!); + } + } + + // Kahn's algorithm — produce dependency layers + const layers: PlanItem[][] = []; + const remaining = new Set(indexedItems.map((item) => item.id!)); + + while (remaining.size > 0) { + // Find items with no remaining dependencies (in-degree 0) + const layer: PlanItem[] = []; + for (const id of remaining) { + if (inDegree.get(id)! === 0) { + layer.push(indexedItems.find((item) => item.id === id)!); + } + } + + if (layer.length === 0) { + // Circular dependency — add all remaining items as a single layer + for (const id of remaining) { + layer.push(indexedItems.find((item) => item.id === id)!); + } + } + + layers.push(layer); + + // Remove this layer's items and update in-degrees + for (const item of layer) { + remaining.delete(item.id!); + for (const dep of dependents.get(item.id!) || []) { + if (remaining.has(dep)) { + inDegree.set(dep, Math.max(0, inDegree.get(dep)! - 1)); + } + } + } + } + + // Sub-group each layer by TDD phase + const rawPhases: PlanItem[][] = []; + for (const layer of layers) { + const byPhase = new Map(); + for (const item of layer) { + const phase = item.tddPhase || 'impl'; + if (!byPhase.has(phase)) byPhase.set(phase, []); + byPhase.get(phase)!.push(item); + } + + // Sort sub-groups by TDD phase order + const sorted = [...byPhase.entries()].sort( + ([a], [b]) => (TDD_PHASE_ORDER[a as TddPhase] ?? 2) - (TDD_PHASE_ORDER[b as TddPhase] ?? 2) + ); + + for (const [, items] of sorted) { + rawPhases.push(items); + } + } + + // Merge small phases with their previous neighbor + const mergedPhases: PlanItem[][] = []; + for (const phase of rawPhases) { + if (mergedPhases.length > 0 && phase.length < MIN_PHASE_TASKS) { + const prev = mergedPhases[mergedPhases.length - 1]; + if (prev.length < MIN_PHASE_TASKS) { + // Merge with previous + prev.push(...phase); + continue; + } + } + mergedPhases.push([...phase]); + } + + // Cap at MAX_PHASES by merging tail phases + while (mergedPhases.length > MAX_PHASES) { + const last = mergedPhases.pop()!; + mergedPhases[mergedPhases.length - 1].push(...last); + } + + // Convert to OrchestratorPhase objects + return mergedPhases.map((phaseItems, index) => this.createPhase(phaseItems, index)); + } + + private createPhase(items: PlanItem[], order: number): OrchestratorPhase { + // Derive phase name from TDD phases and priorities + const tddPhases = [...new Set(items.map((i) => i.tddPhase).filter(Boolean))]; + const name = this.generatePhaseName(items, tddPhases as TddPhase[], order); + const description = items.map((i) => i.content).join('; '); + + const tasks: OrchestratorTask[] = items.map((item, i) => ({ + id: `phase-${order + 1}-task-${i + 1}`, + phaseId: `phase-${order + 1}`, + prompt: item.content, + status: 'pending' as const, + assignedSessionId: null, + queueTaskId: null, + parallel: items.length > 1, // Tasks within a phase are parallel by default + completionPhrase: '', // Assigned later + timeoutMs: DEFAULT_TASK_TIMEOUT_MS, + startedAt: null, + completedAt: null, + error: null, + retries: 0, + })); + + // Extract verification criteria and test commands from items + const verificationCriteria = items + .map((i) => i.verificationCriteria) + .filter((v): v is string => v != null && v.length > 0); + + const testCommands = items.map((i) => i.testCommand).filter((t): t is string => t != null && t.length > 0); + + return { + id: `phase-${order + 1}`, + name, + description, + order, + status: 'pending' as PhaseStatus, + tasks, + verificationCriteria, + testCommands, + maxAttempts: this.config.maxPhaseRetries, + attempts: 0, + startedAt: null, + completedAt: null, + durationMs: null, + teamStrategy: { type: 'single' }, // Assigned later + }; + } + + private generatePhaseName(items: PlanItem[], tddPhases: TddPhase[], order: number): string { + // Try to create a meaningful name based on content + const priorities = [...new Set(items.map((i) => i.priority).filter(Boolean))]; + + if (tddPhases.length === 1) { + const phaseNames: Record = { + setup: 'Setup & Configuration', + test: 'Test Definition', + impl: 'Implementation', + verify: 'Verification', + review: 'Review & Polish', + }; + return `Phase ${order + 1}: ${phaseNames[tddPhases[0]]}`; + } + + if (priorities.includes('P0') && priorities.length === 1) { + return `Phase ${order + 1}: Critical Foundation`; + } + + return `Phase ${order + 1}: ${items.length > 1 ? 'Parallel Tasks' : items[0].content.slice(0, 50)}`; + } + + // ═══════════════════════════════════════════════════════════════ + // Team Strategy Assignment + // ═══════════════════════════════════════════════════════════════ + + private assignTeamStrategies(phases: OrchestratorPhase[]): void { + for (const phase of phases) { + phase.teamStrategy = this.computeTeamStrategy(phase); + } + } + + private computeTeamStrategy(phase: OrchestratorPhase): TeamStrategy { + const taskCount = phase.tasks.length; + const parallelTasks = phase.tasks.filter((t) => t.parallel).length; + + // Single task or no parallel potential → single session + if (taskCount <= 2 || parallelTasks <= 1) { + return { type: 'single' }; + } + + // If team agents are disabled, use parallel sessions instead + if (!this.config.enableTeamAgents) { + return { + type: 'parallel', + maxSessions: Math.min(parallelTasks, this.config.maxParallelSessions), + }; + } + + // 4+ parallel tasks with team agents enabled → team mode + if (parallelTasks >= 4) { + return { + type: 'team', + config: { + leadPrompt: this.buildTeamLeadPrompt(phase), + suggestedTeammates: phase.tasks.slice(0, 4).map((t) => `Specialist for: ${t.prompt.slice(0, 80)}`), + maxTeammates: Math.min(parallelTasks, 4), + }, + }; + } + + // 3 parallel tasks → parallel sessions + return { + type: 'parallel', + maxSessions: Math.min(parallelTasks, this.config.maxParallelSessions), + }; + } + + private buildTeamLeadPrompt(phase: OrchestratorPhase): string { + const taskList = phase.tasks.map((t, i) => `${i + 1}. ${t.prompt}`).join('\n'); + + return [ + `You are the team lead for "${phase.name}".`, + `Create teammates and delegate the following tasks for parallel execution:`, + '', + taskList, + '', + `Each teammate should focus on one task area.`, + `When all tasks are complete, verify the results and output: ${phase.id.toUpperCase()}_COMPLETE`, + ].join('\n'); + } + + // ═══════════════════════════════════════════════════════════════ + // Completion Phrases + // ═══════════════════════════════════════════════════════════════ + + private generateCompletionPhrases(phases: OrchestratorPhase[]): void { + for (const phase of phases) { + for (const task of phase.tasks) { + // Generate a unique, deterministic completion phrase per task + task.completionPhrase = `ORCH_P${phase.order + 1}_T${phase.tasks.indexOf(task) + 1}`; + } + } + } + + // ═══════════════════════════════════════════════════════════════ + // Helpers + // ═══════════════════════════════════════════════════════════════ + + private estimateComplexity(items: PlanItem[]): 'low' | 'medium' | 'high' { + const total = items.length; + const highComplexity = items.filter((i) => i.complexity === 'high').length; + const p0Count = items.filter((i) => i.priority === 'P0').length; + + if (total > 20 || highComplexity > 5 || p0Count > 8) return 'high'; + if (total > 10 || highComplexity > 2 || p0Count > 4) return 'medium'; + return 'low'; + } +} diff --git a/src/orchestrator-verifier.ts b/src/orchestrator-verifier.ts new file mode 100644 index 00000000..8005fad0 --- /dev/null +++ b/src/orchestrator-verifier.ts @@ -0,0 +1,292 @@ +/** + * @fileoverview Orchestrator phase verification. + * + * Runs verification checks after each phase completes: + * - Test commands (shell commands via session) + * - AI review (ask Claude to evaluate phase results) + * + * Three verification modes: + * - strict: ALL test commands must pass AND AI review must approve + * - moderate: Test commands must pass, AI review is advisory + * - lenient: At least one test command passes, AI review skipped + * + * Key exports: + * - `OrchestratorVerifier` class — phase verification engine + * + * @dependencies types (OrchestratorPhase, VerificationResult, VerificationCheck, OrchestratorConfig) + * @consumedby orchestrator-loop + * + * @module orchestrator-verifier + */ + +import type { Session } from './session.js'; +import type { OrchestratorPhase, OrchestratorConfig, VerificationResult, VerificationCheck } from './types.js'; + +// ═══════════════════════════════════════════════════════════════ +// Constants +// ═══════════════════════════════════════════════════════════════ + +/** Timeout for individual test command execution (2 minutes) */ +const TEST_COMMAND_TIMEOUT_MS = 2 * 60 * 1000; + +/** Timeout for AI review (3 minutes) */ +const AI_REVIEW_TIMEOUT_MS = 3 * 60 * 1000; + +/** Completion phrase for AI verification pass */ +const VERIFY_PASS_PHRASE = 'ORCH_VERIFY_PASS'; + +/** Completion phrase for AI verification fail */ +const VERIFY_FAIL_PHRASE = 'ORCH_VERIFY_FAIL'; + +// ═══════════════════════════════════════════════════════════════ +// OrchestratorVerifier +// ═══════════════════════════════════════════════════════════════ + +export class OrchestratorVerifier { + private config: OrchestratorConfig; + + constructor(config: OrchestratorConfig) { + this.config = config; + } + + /** + * Run all verification checks for a completed phase. + * + * @param phase - The phase to verify + * @param session - Session to use for running commands/reviews + * @returns Verification result with pass/fail and suggestions + */ + async verifyPhase(phase: OrchestratorPhase, session: Session): Promise { + const checks: VerificationCheck[] = []; + const mode = this.config.verificationMode; + + // Skip verification entirely in lenient mode with no test commands + if (mode === 'lenient' && phase.testCommands.length === 0 && phase.verificationCriteria.length === 0) { + return { + passed: true, + checks: [], + summary: 'Verification skipped (lenient mode, no checks defined)', + suggestions: [], + }; + } + + // Run test commands if any are defined + if (phase.testCommands.length > 0) { + const testChecks = await this.runTestCommands(phase.testCommands, session); + checks.push(...testChecks); + } + + // Run AI review in strict and moderate modes + if (mode !== 'lenient' && phase.verificationCriteria.length > 0) { + const aiCheck = await this.aiReview(phase, session); + checks.push(aiCheck); + } + + // Determine pass/fail based on mode + const passed = this.evaluateChecks(checks, mode); + + // Generate suggestions for failed checks + const suggestions = this.generateSuggestions(checks, phase); + + const passedCount = checks.filter((c) => c.passed).length; + const summary = + checks.length === 0 ? 'No verification checks defined' : `${passedCount}/${checks.length} checks passed`; + + return { passed, checks, summary, suggestions }; + } + + // ═══════════════════════════════════════════════════════════════ + // Test Command Execution + // ═══════════════════════════════════════════════════════════════ + + private async runTestCommands(commands: string[], session: Session): Promise { + const checks: VerificationCheck[] = []; + + for (const command of commands) { + try { + const check = await this.runSingleTestCommand(command, session); + checks.push(check); + } catch (err) { + checks.push({ + type: 'test_command', + description: `Run: ${command}`, + passed: false, + output: err instanceof Error ? err.message : String(err), + }); + } + } + + return checks; + } + + private async runSingleTestCommand(command: string, session: Session): Promise { + // Send the test command to the session and wait for completion + // We use a unique marker to detect when the command finishes + const marker = `ORCH_TEST_${Date.now()}`; + const wrappedCommand = `${command} && echo ${marker}_PASS || echo ${marker}_FAIL`; + + const result = await this.sendAndWaitForMarker(session, wrappedCommand, marker, TEST_COMMAND_TIMEOUT_MS); + + return { + type: 'test_command', + description: `Run: ${command}`, + passed: result.includes(`${marker}_PASS`), + output: result.slice(0, 2000), // Truncate output + }; + } + + // ═══════════════════════════════════════════════════════════════ + // AI Review + // ═══════════════════════════════════════════════════════════════ + + private async aiReview(phase: OrchestratorPhase, session: Session): Promise { + const prompt = this.buildVerificationPrompt(phase); + + try { + const result = await this.sendAndWaitForMarker( + session, + prompt, + VERIFY_PASS_PHRASE, + AI_REVIEW_TIMEOUT_MS, + VERIFY_FAIL_PHRASE + ); + + const passed = result.includes(VERIFY_PASS_PHRASE); + + return { + type: 'ai_review', + description: `AI review of "${phase.name}"`, + passed, + output: result.slice(0, 3000), + }; + } catch (err) { + return { + type: 'ai_review', + description: `AI review of "${phase.name}"`, + passed: false, + output: `AI review timed out or failed: ${err instanceof Error ? err.message : String(err)}`, + }; + } + } + + private buildVerificationPrompt(phase: OrchestratorPhase): string { + const criteria = phase.verificationCriteria.map((c, i) => `${i + 1}. ${c}`).join('\n'); + + return [ + `Review the work done in "${phase.name}". Check these criteria:`, + '', + criteria, + '', + `If ALL criteria are met, respond with: ${VERIFY_PASS_PHRASE}`, + `If ANY criteria fail, respond with: ${VERIFY_FAIL_PHRASE} and explain what failed.`, + ].join('\n'); + } + + // ═══════════════════════════════════════════════════════════════ + // Evaluation + // ═══════════════════════════════════════════════════════════════ + + private evaluateChecks(checks: VerificationCheck[], mode: OrchestratorConfig['verificationMode']): boolean { + if (checks.length === 0) return true; + + const testChecks = checks.filter((c) => c.type === 'test_command'); + const aiChecks = checks.filter((c) => c.type === 'ai_review'); + + switch (mode) { + case 'strict': + // ALL checks must pass + return checks.every((c) => c.passed); + + case 'moderate': + // All test commands must pass; AI review is advisory + return testChecks.length === 0 || testChecks.every((c) => c.passed); + + case 'lenient': + // At least one test passes (AI review skipped in lenient mode) + return testChecks.length === 0 || testChecks.some((c) => c.passed); + + default: + return aiChecks.every((c) => c.passed) && testChecks.every((c) => c.passed); + } + } + + private generateSuggestions(checks: VerificationCheck[], phase: OrchestratorPhase): string[] { + const suggestions: string[] = []; + const failedChecks = checks.filter((c) => !c.passed); + + if (failedChecks.length === 0) return suggestions; + + for (const check of failedChecks) { + if (check.type === 'test_command') { + suggestions.push(`Fix failing test: ${check.description}`); + } else if (check.type === 'ai_review' && check.output) { + // Extract failure reasons from AI review output + suggestions.push(`Address AI review feedback for "${phase.name}"`); + } + } + + return suggestions; + } + + // ═══════════════════════════════════════════════════════════════ + // Session Communication + // ═══════════════════════════════════════════════════════════════ + + /** + * Send a prompt to a session and wait for a marker phrase in the output. + * + * @param session - Session to send to + * @param input - Prompt/command to send + * @param marker - Primary marker to watch for + * @param timeoutMs - Maximum wait time + * @param altMarker - Alternative marker (for pass/fail detection) + * @returns Captured output containing the marker + */ + private sendAndWaitForMarker( + session: Session, + input: string, + marker: string, + timeoutMs: number, + altMarker?: string + ): Promise { + return new Promise((resolve, reject) => { + let output = ''; + let resolved = false; + + const timer = setTimeout(() => { + if (!resolved) { + resolved = true; + cleanup(); + reject(new Error(`Timeout waiting for marker "${marker}" after ${timeoutMs}ms`)); + } + }, timeoutMs); + + const handler = (data: string) => { + if (resolved) return; + output += data; + + if (output.includes(marker) || (altMarker && output.includes(altMarker))) { + resolved = true; + cleanup(); + resolve(output); + } + }; + + const cleanup = () => { + clearTimeout(timer); + session.off('terminal', handler); + }; + + session.on('terminal', handler); + + // Send the input + session.sendInput(input).catch((err) => { + if (!resolved) { + resolved = true; + cleanup(); + reject(err); + } + }); + }); + } +} diff --git a/src/prompts/index.ts b/src/prompts/index.ts index 0c91a154..281e48c3 100644 --- a/src/prompts/index.ts +++ b/src/prompts/index.ts @@ -7,3 +7,10 @@ export { RESEARCH_AGENT_PROMPT } from './research-agent.js'; export { PLANNER_PROMPT } from './planner.js'; +export { + PHASE_EXECUTION_PROMPT, + TEAM_LEAD_PROMPT, + VERIFICATION_PROMPT, + REPLAN_PROMPT, + SINGLE_TASK_PROMPT, +} from './orchestrator.js'; diff --git a/src/prompts/orchestrator.ts b/src/prompts/orchestrator.ts new file mode 100644 index 00000000..dafcf161 --- /dev/null +++ b/src/prompts/orchestrator.ts @@ -0,0 +1,116 @@ +/** + * @fileoverview Orchestrator Loop prompt templates. + * + * Templates for phase execution, team delegation, verification, and replanning. + * Placeholders use {VARIABLE} syntax and are replaced at runtime. + * + * @module prompts/orchestrator + */ + +/** + * Phase execution prompt — tells Claude what to accomplish in this phase. + * + * Placeholders: + * - {PHASE_NUMBER}: Phase index (1-based) + * - {PHASE_NAME}: Human-readable phase name + * - {GOAL}: Original user goal + * - {COMPLETED_PHASES}: Summary of previously completed phases + * - {TASK_LIST}: Numbered task list for this phase + * - {VERIFICATION_CRITERIA}: What will be checked after this phase + * - {COMPLETION_PHRASE}: The phrase to output when done + */ +export const PHASE_EXECUTION_PROMPT = `You are executing {PHASE_NAME} of a larger project. + +OVERALL GOAL: {GOAL} + +COMPLETED SO FAR: +{COMPLETED_PHASES} + +YOUR TASKS FOR THIS PHASE: +{TASK_LIST} + +Complete each task thoroughly. Run tests after each change to catch issues early. + +VERIFICATION (will be checked after you finish): +{VERIFICATION_CRITERIA} + +When ALL tasks in this phase are complete and verified, output: {COMPLETION_PHRASE}`; + +/** + * Team lead delegation prompt — instructs a lead to coordinate teammates. + * + * Placeholders: + * - {PHASE_NAME}: Phase name + * - {TASK_LIST}: Numbered task list + * - {TEAMMATE_HINTS}: Suggested teammate specializations + * - {COMPLETION_PHRASE}: Phrase for when all work is done + */ +export const TEAM_LEAD_PROMPT = `You are the team lead for {PHASE_NAME}. + +Create teammates and delegate the following tasks for parallel execution: + +{TASK_LIST} + +Suggested teammate roles: +{TEAMMATE_HINTS} + +Each teammate should focus on their assigned task area. Monitor their progress. +When ALL tasks are complete and you've verified the results, output: {COMPLETION_PHRASE}`; + +/** + * Verification prompt — asks Claude to verify phase completion. + * + * Placeholders: + * - {PHASE_NAME}: Phase name + * - {CRITERIA}: Numbered verification criteria + * - {PASS_PHRASE}: Phrase to output on success + * - {FAIL_PHRASE}: Phrase to output on failure + */ +export const VERIFICATION_PROMPT = `Review the work done in "{PHASE_NAME}". Check these criteria: + +{CRITERIA} + +If ALL criteria are met, respond with: {PASS_PHRASE} +If ANY criteria fail, respond with: {FAIL_PHRASE} and explain what failed.`; + +/** + * Replan prompt — gives failure context and asks for recovery. + * + * Placeholders: + * - {PHASE_NAME}: Phase name + * - {ATTEMPT_NUMBER}: Current retry attempt + * - {MAX_ATTEMPTS}: Maximum attempts allowed + * - {FAILURE_SUMMARY}: What went wrong + * - {SUGGESTIONS}: Recovery suggestions from verification + * - {ORIGINAL_TASKS}: The original task list + * - {COMPLETION_PHRASE}: Phrase for when recovery is done + */ +export const REPLAN_PROMPT = `Phase "{PHASE_NAME}" verification failed (attempt {ATTEMPT_NUMBER}/{MAX_ATTEMPTS}). + +WHAT WENT WRONG: +{FAILURE_SUMMARY} + +SUGGESTIONS: +{SUGGESTIONS} + +ORIGINAL TASKS: +{ORIGINAL_TASKS} + +Fix the issues identified above. Focus on making the verification criteria pass. +When the fixes are complete, output: {COMPLETION_PHRASE}`; + +/** + * Single-task execution prompt — for phases with a single task. + * + * Placeholders: + * - {TASK}: The task description + * - {GOAL}: Original user goal + * - {CONTEXT}: Any relevant context + * - {COMPLETION_PHRASE}: Phrase for when done + */ +export const SINGLE_TASK_PROMPT = `{TASK} + +Context: This is part of a larger project — {GOAL} +{CONTEXT} + +When done, output: {COMPLETION_PHRASE}`; diff --git a/src/state-store.ts b/src/state-store.ts index 1867f12d..f238776d 100644 --- a/src/state-store.ts +++ b/src/state-store.ts @@ -547,6 +547,30 @@ export class StateStore { this.save(); } + // ========== Orchestrator Loop State Methods ========== + + /** Returns the orchestrator loop state, or null if never initialized. */ + getOrchestratorState() { + return this.state.orchestrator ?? null; + } + + /** Updates orchestrator loop state (partial merge) and triggers a debounced save. */ + setOrchestratorState(orchestrator: Partial>) { + if (this.state.orchestrator) { + this.state.orchestrator = { ...this.state.orchestrator, ...orchestrator }; + } else { + // First initialization — caller must provide full state + this.state.orchestrator = orchestrator as NonNullable; + } + this.save(); + } + + /** Clears orchestrator state and triggers a debounced save. */ + clearOrchestratorState() { + this.state.orchestrator = undefined; + this.save(); + } + /** Returns the application configuration. */ getConfig() { return this.state.config; diff --git a/src/types/app-state.ts b/src/types/app-state.ts index 210a95aa..56809a6b 100644 --- a/src/types/app-state.ts +++ b/src/types/app-state.ts @@ -109,6 +109,8 @@ export interface AppState { globalStats?: GlobalStats; /** Daily token usage statistics */ tokenStats?: TokenStats; + /** Orchestrator Loop state (phased plan execution) */ + orchestrator?: import('./orchestrator.js').OrchestratorPersistState; } // ========== Default Configuration ========== diff --git a/src/types/index.ts b/src/types/index.ts index e0f54a27..6420817a 100644 --- a/src/types/index.ts +++ b/src/types/index.ts @@ -26,6 +26,7 @@ * | teams | TeamConfig, TeamMember, TeamTask, InboxMessage, PaneInfo | `~/.claude/teams/`, `~/.claude/tasks/` → `GET /api/teams` | * | push | PushSubscriptionRecord, VapidKeys | `~/.codeman/push-keys.json`, `~/.codeman/push-subscriptions.json` | * | plan | PlanItem, PlanTaskStatus, TddPhase | In-memory → `GET /api/sessions/:id/plan/tasks` | + * | orchestrator | OrchestratorState, OrchestratorPlan, OrchestratorConfig, OrchestratorPersistState | `~/.codeman/state.json` → `GET /api/orchestrator/status` | * * ## Cross-domain relationship map * @@ -64,3 +65,4 @@ export * from './tools.js'; export * from './teams.js'; export * from './push.js'; export * from './plan.js'; +export * from './orchestrator.js'; diff --git a/src/types/orchestrator.ts b/src/types/orchestrator.ts new file mode 100644 index 00000000..b2ec40f3 --- /dev/null +++ b/src/types/orchestrator.ts @@ -0,0 +1,285 @@ +/** + * @fileoverview Orchestrator Loop type definitions. + * + * Types for the phased plan execution system: state machine, plan structure, + * phase grouping, team strategies, verification, configuration, and persistence. + * + * Key exports: + * - OrchestratorState — state machine states (idle → planning → approval → executing → verifying → ...) + * - OrchestratorPlan / OrchestratorPhase / OrchestratorTask — hierarchical plan structure + * - TeamStrategy — how agents are coordinated per phase (single, parallel, team) + * - VerificationResult / VerificationCheck — phase verification output + * - OrchestratorConfig — user-configurable options + * - OrchestratorPersistState / OrchestratorStats — persistence and metrics + * + * Cross-domain relationships: + * - OrchestratorTask.queueTaskId links to TaskState.id (task domain) + * - OrchestratorTask.assignedSessionId links to SessionState.id (session domain) + * - OrchestratorPersistState is embedded in AppState.orchestrator (app-state domain) + * + * Served at `GET /api/orchestrator/status` and `GET /api/orchestrator/plan`. + * No dependencies on other domain modules. + */ + +// ═══════════════════════════════════════════════════════════════ +// State Machine +// ═══════════════════════════════════════════════════════════════ + +/** Orchestrator loop states */ +export type OrchestratorState = + | 'idle' + | 'planning' + | 'approval' + | 'executing' + | 'verifying' + | 'replanning' + | 'completed' + | 'failed' + | 'paused'; + +// ═══════════════════════════════════════════════════════════════ +// Plan Structure +// ═══════════════════════════════════════════════════════════════ + +/** Top-level orchestrator plan generated from a user goal */ +export interface OrchestratorPlan { + /** Unique plan identifier */ + id: string; + /** Original user goal/task description */ + goal: string; + /** When the plan was generated */ + createdAt: number; + /** Ordered list of execution phases */ + phases: OrchestratorPhase[]; + /** Plan generation metadata */ + metadata: OrchestratorPlanMetadata; +} + +/** Metadata from plan generation */ +export interface OrchestratorPlanMetadata { + /** Total tasks across all phases */ + totalTasks: number; + /** Estimated overall complexity */ + estimatedComplexity: 'low' | 'medium' | 'high'; + /** Model used for plan generation */ + modelUsed: string; + /** Time taken to generate the plan */ + planDurationMs: number; +} + +/** A sequential execution phase containing parallel tasks */ +export interface OrchestratorPhase { + /** Phase identifier (e.g., "phase-1") */ + id: string; + /** Human-readable phase name */ + name: string; + /** Detailed description of what this phase accomplishes */ + description: string; + /** Execution order (0-based) */ + order: number; + /** Current phase status */ + status: PhaseStatus; + /** Tasks within this phase */ + tasks: OrchestratorTask[]; + /** Criteria to verify after phase completion */ + verificationCriteria: string[]; + /** Shell commands to run for verification */ + testCommands: string[]; + /** Maximum retry attempts for this phase */ + maxAttempts: number; + /** Current attempt count */ + attempts: number; + /** When execution started */ + startedAt: number | null; + /** When phase completed (passed or failed) */ + completedAt: number | null; + /** Total execution duration */ + durationMs: number | null; + /** How to coordinate agents for this phase */ + teamStrategy: TeamStrategy; +} + +/** Phase execution status */ +export type PhaseStatus = 'pending' | 'executing' | 'verifying' | 'passed' | 'failed' | 'skipped'; + +/** A single executable task within a phase */ +export interface OrchestratorTask { + /** Task identifier (e.g., "phase-1-task-1") */ + id: string; + /** Parent phase identifier */ + phaseId: string; + /** Single-line prompt to send to Claude */ + prompt: string; + /** Current task status */ + status: 'pending' | 'running' | 'completed' | 'failed'; + /** Session running this task */ + assignedSessionId: string | null; + /** Links to TaskQueue task ID (for completion tracking) */ + queueTaskId: string | null; + /** Whether this task can run in parallel with siblings */ + parallel: boolean; + /** Unique phrase for completion detection */ + completionPhrase: string; + /** Timeout in milliseconds */ + timeoutMs: number; + /** When task started executing */ + startedAt: number | null; + /** When task completed */ + completedAt: number | null; + /** Error message if task failed */ + error: string | null; + /** Number of retry attempts */ + retries: number; +} + +// ═══════════════════════════════════════════════════════════════ +// Team Strategy +// ═══════════════════════════════════════════════════════════════ + +/** How agents are coordinated for a phase */ +export type TeamStrategy = + | { type: 'single' } + | { type: 'parallel'; maxSessions: number } + | { type: 'team'; config: TeamSetup }; + +/** Configuration for team-based phase execution */ +export interface TeamSetup { + /** Prompt to send to the team lead */ + leadPrompt: string; + /** Suggested teammate role descriptions */ + suggestedTeammates: string[]; + /** Maximum number of teammates to create */ + maxTeammates: number; +} + +// ═══════════════════════════════════════════════════════════════ +// Verification +// ═══════════════════════════════════════════════════════════════ + +/** Result of phase verification */ +export interface VerificationResult { + /** Whether all checks passed */ + passed: boolean; + /** Individual verification checks */ + checks: VerificationCheck[]; + /** Human-readable summary */ + summary: string; + /** Suggestions for replanning if verification failed */ + suggestions: string[]; +} + +/** A single verification check result */ +export interface VerificationCheck { + /** Type of check performed */ + type: 'test_command' | 'ai_review' | 'file_check'; + /** What was checked */ + description: string; + /** Whether this check passed */ + passed: boolean; + /** Command output or review text */ + output?: string; +} + +// ═══════════════════════════════════════════════════════════════ +// Configuration +// ═══════════════════════════════════════════════════════════════ + +/** User-configurable orchestrator options */ +export interface OrchestratorConfig { + /** Model to use for plan generation (default: 'opus') */ + plannerModel: string; + /** Whether to run research agent before planning (default: true) */ + researchEnabled: boolean; + /** Auto-approve generated plans without user review (default: false) */ + autoApprove: boolean; + /** Maximum retry attempts per phase (default: 3) */ + maxPhaseRetries: number; + /** Phase execution timeout in ms (default: 1800000 = 30min) */ + phaseTimeoutMs: number; + /** Enable Claude Code agent teams for parallel phases (default: true) */ + enableTeamAgents: boolean; + /** Maximum parallel sessions for task execution (default: 3) */ + maxParallelSessions: number; + /** Verification strictness (default: 'moderate') */ + verificationMode: 'strict' | 'moderate' | 'lenient'; + /** Run /compact between phases to manage context (default: true) */ + compactBetweenPhases: boolean; +} + +/** Default orchestrator configuration */ +export const DEFAULT_ORCHESTRATOR_CONFIG: OrchestratorConfig = { + plannerModel: 'opus', + researchEnabled: true, + autoApprove: false, + maxPhaseRetries: 3, + phaseTimeoutMs: 30 * 60 * 1000, // 30 minutes + enableTeamAgents: true, + maxParallelSessions: 3, + verificationMode: 'moderate', + compactBetweenPhases: true, +}; + +// ═══════════════════════════════════════════════════════════════ +// Persistence +// ═══════════════════════════════════════════════════════════════ + +/** Orchestrator state persisted to ~/.codeman/state.json */ +export interface OrchestratorPersistState { + /** Current state machine state */ + state: OrchestratorState; + /** Generated plan (null before planning) */ + plan: OrchestratorPlan | null; + /** Index of currently executing phase */ + currentPhaseIndex: number; + /** When orchestration started */ + startedAt: number | null; + /** When orchestration completed */ + completedAt: number | null; + /** User configuration */ + config: OrchestratorConfig; + /** Execution statistics */ + stats: OrchestratorStats; +} + +/** Orchestrator execution statistics */ +export interface OrchestratorStats { + /** Number of phases completed successfully */ + phasesCompleted: number; + /** Number of phases that failed (after all retries) */ + phasesFailed: number; + /** Total individual tasks completed */ + totalTasksCompleted: number; + /** Total individual tasks failed */ + totalTasksFailed: number; + /** Total time spent executing (ms) */ + totalDurationMs: number; + /** Number of times replanning was triggered */ + replanCount: number; +} + +/** Factory function for initial orchestrator stats */ +export function createInitialOrchestratorStats(): OrchestratorStats { + return { + phasesCompleted: 0, + phasesFailed: 0, + totalTasksCompleted: 0, + totalTasksFailed: 0, + totalDurationMs: 0, + replanCount: 0, + }; +} + +/** Factory function for initial orchestrator persist state */ +export function createInitialOrchestratorPersistState( + config: OrchestratorConfig = DEFAULT_ORCHESTRATOR_CONFIG +): OrchestratorPersistState { + return { + state: 'idle', + plan: null, + currentPhaseIndex: 0, + startedAt: null, + completedAt: null, + config, + stats: createInitialOrchestratorStats(), + }; +} diff --git a/src/web/ports/index.ts b/src/web/ports/index.ts index 6186476c..57d22f74 100644 --- a/src/web/ports/index.ts +++ b/src/web/ports/index.ts @@ -12,3 +12,4 @@ export type { RespawnPort } from './respawn-port.js'; export type { ConfigPort } from './config-port.js'; export type { InfraPort, ScheduledRun } from './infra-port.js'; export type { AuthPort, AuthSessionRecord } from './auth-port.js'; +export type { OrchestratorPort } from './orchestrator-port.js'; diff --git a/src/web/ports/orchestrator-port.ts b/src/web/ports/orchestrator-port.ts new file mode 100644 index 00000000..cc218089 --- /dev/null +++ b/src/web/ports/orchestrator-port.ts @@ -0,0 +1,11 @@ +/** + * @fileoverview Orchestrator port — capabilities for orchestrator loop management. + * Route modules that interact with the orchestrator depend on this port. + */ + +import type { OrchestratorLoop } from '../../orchestrator-loop.js'; + +export interface OrchestratorPort { + readonly orchestratorLoop: OrchestratorLoop | null; + initOrchestratorLoop(): OrchestratorLoop; +} diff --git a/src/web/public/app.js b/src/web/public/app.js index a501242c..f7fb35cd 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -235,6 +235,19 @@ const _SSE_HANDLER_MAP = [ [SSE_EVENTS.PLAN_STARTED, '_onPlanStarted'], [SSE_EVENTS.PLAN_CANCELLED, '_onPlanCancelled'], [SSE_EVENTS.PLAN_COMPLETED, '_onPlanCompleted'], + + // Orchestrator loop + [SSE_EVENTS.ORCHESTRATOR_STATE_CHANGED, '_onOrchestratorStateChanged'], + [SSE_EVENTS.ORCHESTRATOR_PLAN_READY, '_onOrchestratorPlanReady'], + [SSE_EVENTS.ORCHESTRATOR_PHASE_STARTED, '_onOrchestratorPhaseStarted'], + [SSE_EVENTS.ORCHESTRATOR_PHASE_COMPLETED, '_onOrchestratorPhaseCompleted'], + [SSE_EVENTS.ORCHESTRATOR_PHASE_FAILED, '_onOrchestratorPhaseFailed'], + [SSE_EVENTS.ORCHESTRATOR_VERIFICATION, '_onOrchestratorVerification'], + [SSE_EVENTS.ORCHESTRATOR_TASK_ASSIGNED, '_onOrchestratorTaskAssigned'], + [SSE_EVENTS.ORCHESTRATOR_TASK_COMPLETED, '_onOrchestratorTaskCompleted'], + [SSE_EVENTS.ORCHESTRATOR_TASK_FAILED, '_onOrchestratorTaskFailed'], + [SSE_EVENTS.ORCHESTRATOR_COMPLETED, '_onOrchestratorCompleted'], + [SSE_EVENTS.ORCHESTRATOR_ERROR, '_onOrchestratorError'], ]; // ═══════════════════════════════════════════════════════════════ @@ -317,6 +330,10 @@ class CodemanApp { this.logViewerWindows = new Map(); // Map this.logViewerWindowZIndex = ZINDEX_LOG_VIEWER_BASE; this.projectInsightsPanelVisible = false; + + // Orchestrator loop state + this.orchestratorState = null; // { state, plan, currentPhaseIndex, stats } + this.orchestratorPanelVisible = false; this.currentSessionWorkingDir = null; // Track current session's working dir for path normalization // Image popup windows (auto-open for detected screenshots/images) diff --git a/src/web/public/constants.js b/src/web/public/constants.js index 28388fb5..a26962c0 100644 --- a/src/web/public/constants.js +++ b/src/web/public/constants.js @@ -277,6 +277,19 @@ const SSE_EVENTS = { PLAN_STARTED: 'plan:started', PLAN_CANCELLED: 'plan:cancelled', PLAN_COMPLETED: 'plan:completed', + + // Orchestrator Loop + ORCHESTRATOR_STATE_CHANGED: 'orchestrator:stateChanged', + ORCHESTRATOR_PLAN_READY: 'orchestrator:planReady', + ORCHESTRATOR_PHASE_STARTED: 'orchestrator:phaseStarted', + ORCHESTRATOR_PHASE_COMPLETED: 'orchestrator:phaseCompleted', + ORCHESTRATOR_PHASE_FAILED: 'orchestrator:phaseFailed', + ORCHESTRATOR_VERIFICATION: 'orchestrator:verification', + ORCHESTRATOR_TASK_ASSIGNED: 'orchestrator:taskAssigned', + ORCHESTRATOR_TASK_COMPLETED: 'orchestrator:taskCompleted', + ORCHESTRATOR_TASK_FAILED: 'orchestrator:taskFailed', + ORCHESTRATOR_COMPLETED: 'orchestrator:completed', + ORCHESTRATOR_ERROR: 'orchestrator:error', }; // ═══════════════════════════════════════════════════════════════ diff --git a/src/web/public/index.html b/src/web/public/index.html index 81151a12..83499f77 100644 --- a/src/web/public/index.html +++ b/src/web/public/index.html @@ -137,6 +137,17 @@ + + +