From 55e9f9f0a50f97303a93e5d4c9dd7a80f8b93342 Mon Sep 17 00:00:00 2001 From: arkon Date: Wed, 28 Jan 2026 07:07:18 +0100 Subject: [PATCH] feat(execution): add ExecutionBridge for parallel task execution Implements the Execution Optimizer Rework plan that enables the system to actually use the optimizer metadata (parallelGroups, agentType, recommendedModel, requiresFreshContext, estimatedTokens) that was previously generated but ignored. New components: - ExecutionBridge: Central coordinator that loads optimized plans, manages parallel execution within groups, and coordinates with SpawnOrchestrator for session-based execution - ModelSelector: Routes tasks to appropriate models (opus/sonnet/haiku) based on user defaults and agent type overrides. Optimizer recommendations are advisory only - user preferences always win - GroupScheduler: Builds topologically ordered execution groups, manages dependencies, determines execution mode (session vs task-tool) - ContextManager: Handles fresh context requirements via /clear+/init or new session spawning Features: - Parallel task execution within groups (configurable limit) - Group-level dependency tracking (lower groups complete first) - Partial failure handling (continue with non-dependent tasks) - Model configuration in App Settings > Models tab - Agent type overrides (explore, implement, test, review) - Execution control API endpoints (start, pause, resume, cancel) - SSE events for real-time execution progress visibility - Execution history tracking API endpoints: - GET/POST/PUT /api/execution/* for execution control - GET/PUT /api/execution/model-config for model settings Co-Authored-By: Claude Opus 4.5 --- CLAUDE.md | 34 +- package.json | 2 +- src/config/execution-limits.ts | 93 ++++ src/context-manager.ts | 333 +++++++++++++ src/execution-bridge.ts | 824 +++++++++++++++++++++++++++++++++ src/group-scheduler.ts | 551 ++++++++++++++++++++++ src/model-selector.ts | 314 +++++++++++++ src/types.ts | 37 ++ src/web/public/app.js | 174 +++++++ src/web/public/index.html | 93 ++++ src/web/server.ts | 424 ++++++++++++++++- test/execution-bridge.test.ts | 270 +++++++++++ 12 files changed, 3140 insertions(+), 9 deletions(-) create mode 100644 src/config/execution-limits.ts create mode 100644 src/context-manager.ts create mode 100644 src/execution-bridge.ts create mode 100644 src/group-scheduler.ts create mode 100644 src/model-selector.ts create mode 100644 test/execution-bridge.test.ts diff --git a/CLAUDE.md b/CLAUDE.md index d6020adb..de25878a 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -16,7 +16,7 @@ When user says "COM": 1. Increment version in BOTH `package.json` AND `CLAUDE.md` 2. Run: `git add -A && git commit -m "chore: bump version to X.XXXX" && git push && npm run build && systemctl --user restart claudeman-web` -**Version**: 0.1415 (must match `package.json`) +**Version**: 0.1416 (must match `package.json`) ## Project Overview @@ -122,7 +122,7 @@ journalctl --user -u claudeman-web -f **Port allocation**: E2E tests use centralized ports in `test/e2e/e2e.config.ts`. Unit/integration tests pick unique ports manually. Search `const PORT =` or `TEST_PORT` in test files to find used ports before adding new tests. -**E2E tests**: Use Playwright. Run `npx playwright install chromium` first. See `test/e2e/fixtures/` for helpers. E2E config provides ports, timeouts, and helpers. +**E2E tests**: Use Playwright. Run `npx playwright install chromium` first. See `test/e2e/fixtures/` for helpers. E2E config (`test/e2e/e2e.config.ts`) provides ports (3183-3190), timeouts, and helpers. **Test config**: Vitest runs with `globals: true` (no imports needed for `describe`/`it`/`expect`) and `fileParallelism: false` (files run sequentially to respect screen limits). Unit test timeout is 30s, teardown timeout is 60s. E2E tests have longer timeouts defined in `test/e2e/e2e.config.ts` (90s test, 30s session creation). @@ -132,7 +132,7 @@ journalctl --user -u claudeman-web -f - Tracked resource cleanup (only kills screens/processes tests register) - Safe to run from within Claudeman-managed sessions -Respawn tests use MockSession to avoid spawning real Claude processes. +Respawn tests use MockSession to avoid spawning real Claude processes. See `test/respawn-test-utils.ts` for MockSession, MockAiIdleChecker, MockAiPlanChecker, state trackers, and terminal output generators. ## Debugging @@ -209,3 +209,31 @@ The `claudeman-mcp` binary provides Model Context Protocol integration for Claud ``` This enables Claude Desktop to spawn and manage agents via MCP tools. + +## Active Ralph Loop Task + +**Current Task**: Analyse the whole Start Ralph Loop Orchestrator from Claudeman! You have all files locally. +Does the Ralph Loop Orchestrator make sense like this? Should we force the agents to a pydantic json that all speak the same language or it good like this? +How could we make the agents more smarter, working better together? +Is the execution optimiser, really taking care of how it gets executed later? How is that managed? Recheck these logics and make it much better! +Also go through all prompts and optimise them heavily for better results! +Do the amount of agents make sense? Do we miss something, would it make sense to change its order? Would to make sense to limit their json output differently? Commit but don't push anything! + +**Case Folder**: `/home/arkon/claudeman-cases/claudeman` + +### Key Files +- **Plan Summary**: `/home/arkon/claudeman-cases/claudeman/ralph-wizard/summary.md` - Human-readable plan overview +- **Todo Items**: `/home/arkon/claudeman-cases/claudeman/ralph-wizard/final-result.json` - Contains `items` array with all todo tasks +- **Research**: `/home/arkon/claudeman-cases/claudeman/ralph-wizard/research/result.json` - External resources and codebase patterns + +### How to Work on This Task +1. Read the plan summary to understand the overall approach +2. Check `final-result.json` for the todo items array - each item has `id`, `title`, `description`, `priority` +3. Work through items in priority order (critical → high → medium → low) +4. Use `COMPLETION_PHRASE` when the entire task is complete + +### Research Insights +Check `/home/arkon/claudeman-cases/claudeman/ralph-wizard/research/result.json` for: +- External GitHub repos and documentation links to reference +- Existing codebase patterns to follow +- Technical recommendations from the research phase diff --git a/package.json b/package.json index eeb78394..e954f4ab 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "claudeman", - "version": "0.1415", + "version": "0.1416", "description": "The missing control plane for Claude Code - run 20 autonomous agents with real-time monitoring and session persistence", "type": "module", "main": "dist/index.js", diff --git a/src/config/execution-limits.ts b/src/config/execution-limits.ts new file mode 100644 index 00000000..e3a8ad6d --- /dev/null +++ b/src/config/execution-limits.ts @@ -0,0 +1,93 @@ +/** + * @fileoverview Centralized limits for the Execution Bridge system. + * + * These constants define resource limits for parallel task execution, + * group scheduling, and model selection. + * + * @module config/execution-limits + */ + +// ============================================================================ +// Parallel Execution Limits +// ============================================================================ + +/** + * Maximum tasks to execute in parallel within a single group. + * Higher values increase throughput but may strain system resources. + */ +export const MAX_PARALLEL_TASKS_PER_GROUP = 5; + +/** + * Default timeout for an execution group in milliseconds (30 minutes). + * If all tasks in a group don't complete within this time, stragglers are cancelled. + */ +export const GROUP_TIMEOUT_MS = 30 * 60 * 1000; + +/** + * Maximum retry attempts for a failed task before marking as permanently failed. + */ +export const MAX_TASK_RETRIES = 2; + +/** + * Delay between task retries in milliseconds (10 seconds). + */ +export const TASK_RETRY_DELAY_MS = 10 * 1000; + +// ============================================================================ +// Model Selection Limits +// ============================================================================ + +/** + * Default model when no recommendation is provided. + */ +export const DEFAULT_MODEL: 'opus' | 'sonnet' | 'haiku' = 'sonnet'; + +/** + * Token threshold for switching execution modes. + * Tasks with estimated tokens below this use task-tool mode. + * Tasks with estimated tokens above this use session mode. + */ +export const TOKEN_THRESHOLD_FOR_SESSION_MODE = 50000; + +/** + * Token threshold for low-complexity tasks (haiku-appropriate). + */ +export const TOKEN_THRESHOLD_HAIKU = 15000; + +// ============================================================================ +// Context Management Limits +// ============================================================================ + +/** + * Delay between /clear and /init commands in milliseconds. + */ +export const CONTEXT_REFRESH_DELAY_MS = 2000; + +/** + * Maximum pending context refresh operations. + */ +export const MAX_PENDING_CONTEXT_REFRESHES = 10; + +// ============================================================================ +// Execution Bridge Limits +// ============================================================================ + +/** + * Maximum execution groups to track in history. + */ +export const MAX_EXECUTION_HISTORY = 50; + +/** + * Polling interval for execution progress in milliseconds. + */ +export const EXECUTION_POLL_INTERVAL_MS = 1000; + +/** + * Maximum total tasks in a single execution plan. + */ +export const MAX_TASKS_PER_PLAN = 100; + +/** + * Grace period after group completion before cleanup in milliseconds. + */ +export const GROUP_CLEANUP_DELAY_MS = 5000; diff --git a/src/context-manager.ts b/src/context-manager.ts new file mode 100644 index 00000000..67dd6b54 --- /dev/null +++ b/src/context-manager.ts @@ -0,0 +1,333 @@ +/** + * @fileoverview Context Manager - Handles fresh context requirements. + * + * Manages context refresh operations for tasks that require starting + * with a clean context window. Supports both /clear + /init sequences + * and new session spawning. + * + * @module context-manager + */ + +import { EventEmitter } from 'node:events'; +import { + CONTEXT_REFRESH_DELAY_MS, + MAX_PENDING_CONTEXT_REFRESHES, +} from './config/execution-limits.js'; + +// ========== Types ========== + +/** Methods for refreshing context */ +export type ContextRefreshMethod = 'clear-init' | 'new-session'; + +/** Status of a context refresh operation */ +export type ContextRefreshStatus = 'pending' | 'clearing' | 'initializing' | 'completed' | 'failed'; + +/** + * Request for a context refresh. + */ +export interface ContextRefreshRequest { + /** Task ID requesting refresh */ + taskId: string; + /** Session ID to refresh */ + sessionId: string; + /** Preferred refresh method */ + method: ContextRefreshMethod; + /** Working directory (for new-session method) */ + workingDir?: string; + /** Optional init prompt to use */ + initPrompt?: string; +} + +/** + * Result of a context refresh operation. + */ +export interface ContextRefreshResult { + /** Request that was processed */ + request: ContextRefreshRequest; + /** Final status */ + status: ContextRefreshStatus; + /** New session ID (if method was new-session) */ + newSessionId?: string; + /** Error message if failed */ + error?: string; + /** Duration in milliseconds */ + durationMs: number; +} + +/** + * Tracked refresh operation. + */ +interface TrackedRefresh { + request: ContextRefreshRequest; + status: ContextRefreshStatus; + startedAt: number; + completedAt?: number; + error?: string; + newSessionId?: string; +} + +// ========== Events ========== + +export interface ContextManagerEvents { + /** Context refresh started */ + refreshStarted: (data: { taskId: string; sessionId: string; method: ContextRefreshMethod }) => void; + /** Context refresh completed */ + refreshCompleted: (result: ContextRefreshResult) => void; + /** Context refresh failed */ + refreshFailed: (data: { taskId: string; sessionId: string; error: string }) => void; +} + +// ========== Session Writer Interface ========== + +/** + * Interface for writing to sessions. + * Injected to avoid circular dependencies. + */ +export interface SessionWriter { + /** Write text to session */ + writeToSession(sessionId: string, text: string): void; + /** Create a new session */ + createSession?(workingDir: string, name?: string): Promise<{ sessionId: string }>; +} + +// ========== Context Manager ========== + +/** + * ContextManager - Handles context refresh operations. + * + * When a task requires fresh context (requiresFreshContext=true), + * this manager coordinates the refresh using one of two methods: + * + * 1. clear-init: Send /clear followed by /init to existing session + * 2. new-session: Spawn a completely new session (more expensive) + */ +export class ContextManager extends EventEmitter { + private _pending: Map = new Map(); + private _sessionWriter: SessionWriter | null = null; + private _refreshDelayMs: number; + + constructor(refreshDelayMs: number = CONTEXT_REFRESH_DELAY_MS) { + super(); + this._refreshDelayMs = refreshDelayMs; + } + + /** + * Set the session writer for performing actual operations. + */ + setSessionWriter(writer: SessionWriter): void { + this._sessionWriter = writer; + } + + /** + * Get number of pending refresh operations. + */ + get pendingCount(): number { + return this._pending.size; + } + + /** + * Check if a refresh is pending for a session. + */ + hasPendingRefresh(sessionId: string): boolean { + for (const tracked of this._pending.values()) { + if (tracked.request.sessionId === sessionId && + (tracked.status === 'pending' || tracked.status === 'clearing' || tracked.status === 'initializing')) { + return true; + } + } + return false; + } + + /** + * Request a context refresh for a task. + * + * @returns Promise that resolves when refresh completes + */ + async requestRefresh(request: ContextRefreshRequest): Promise { + if (!this._sessionWriter) { + return { + request, + status: 'failed', + error: 'No session writer configured', + durationMs: 0, + }; + } + + // Check pending limit + if (this._pending.size >= MAX_PENDING_CONTEXT_REFRESHES) { + return { + request, + status: 'failed', + error: `Max pending refreshes (${MAX_PENDING_CONTEXT_REFRESHES}) exceeded`, + durationMs: 0, + }; + } + + // Track the refresh + const tracked: TrackedRefresh = { + request, + status: 'pending', + startedAt: Date.now(), + }; + this._pending.set(request.taskId, tracked); + + this.emit('refreshStarted', { + taskId: request.taskId, + sessionId: request.sessionId, + method: request.method, + }); + + try { + if (request.method === 'clear-init') { + await this.performClearInit(tracked); + } else { + await this.performNewSession(tracked); + } + + tracked.status = 'completed'; + tracked.completedAt = Date.now(); + + const result: ContextRefreshResult = { + request, + status: 'completed', + newSessionId: tracked.newSessionId, + durationMs: tracked.completedAt - tracked.startedAt, + }; + + this.emit('refreshCompleted', result); + return result; + + } catch (err) { + tracked.status = 'failed'; + tracked.error = err instanceof Error ? err.message : String(err); + tracked.completedAt = Date.now(); + + const result: ContextRefreshResult = { + request, + status: 'failed', + error: tracked.error, + durationMs: tracked.completedAt - tracked.startedAt, + }; + + this.emit('refreshFailed', { + taskId: request.taskId, + sessionId: request.sessionId, + error: tracked.error, + }); + + return result; + + } finally { + // Clean up after a delay + setTimeout(() => { + this._pending.delete(request.taskId); + }, 5000); + } + } + + /** + * Perform /clear + /init sequence. + */ + private async performClearInit(tracked: TrackedRefresh): Promise { + if (!this._sessionWriter) throw new Error('No session writer'); + + const { sessionId, initPrompt } = tracked.request; + + // Send /clear + tracked.status = 'clearing'; + this._sessionWriter.writeToSession(sessionId, '/clear\r'); + + // Wait for clear to process + await this.delay(this._refreshDelayMs); + + // Send /init (or custom init prompt) + tracked.status = 'initializing'; + const prompt = initPrompt || '/init'; + this._sessionWriter.writeToSession(sessionId, prompt + '\r'); + + // Wait for init to complete + await this.delay(this._refreshDelayMs); + } + + /** + * Spawn a new session for complete context isolation. + */ + private async performNewSession(tracked: TrackedRefresh): Promise { + if (!this._sessionWriter?.createSession) { + throw new Error('Session creation not supported'); + } + + const { workingDir, taskId } = tracked.request; + if (!workingDir) { + throw new Error('Working directory required for new-session method'); + } + + tracked.status = 'initializing'; + + const result = await this._sessionWriter.createSession(workingDir, `task-${taskId}`); + tracked.newSessionId = result.sessionId; + } + + /** + * Cancel a pending refresh. + */ + cancelRefresh(taskId: string): boolean { + const tracked = this._pending.get(taskId); + if (tracked && tracked.status === 'pending') { + tracked.status = 'failed'; + tracked.error = 'Cancelled'; + tracked.completedAt = Date.now(); + this._pending.delete(taskId); + return true; + } + return false; + } + + /** + * Get status of all pending refreshes. + */ + getPendingStatus(): Array<{ taskId: string; sessionId: string; status: ContextRefreshStatus; elapsedMs: number }> { + const now = Date.now(); + return Array.from(this._pending.values()).map(tracked => ({ + taskId: tracked.request.taskId, + sessionId: tracked.request.sessionId, + status: tracked.status, + elapsedMs: now - tracked.startedAt, + })); + } + + /** + * Clear all pending operations (for cleanup). + */ + clearAll(): void { + this._pending.clear(); + } + + private delay(ms: number): Promise { + return new Promise(resolve => setTimeout(resolve, ms)); + } +} + +// ========== Singleton ========== + +let managerInstance: ContextManager | null = null; + +/** + * Get or create the singleton ContextManager instance. + */ +export function getContextManager(): ContextManager { + if (!managerInstance) { + managerInstance = new ContextManager(); + } + return managerInstance; +} + +/** + * Reset the singleton (for testing). + */ +export function resetContextManager(): void { + if (managerInstance) { + managerInstance.clearAll(); + } + managerInstance = null; +} diff --git a/src/execution-bridge.ts b/src/execution-bridge.ts new file mode 100644 index 00000000..b86da43c --- /dev/null +++ b/src/execution-bridge.ts @@ -0,0 +1,824 @@ +/** + * @fileoverview Execution Bridge - Coordinates parallel task execution. + * + * The central coordinator that: + * - Loads optimized plans from PlanOrchestrator + * - Converts plans to executable groups via GroupScheduler + * - Manages parallel execution within groups + * - Coordinates model selection via ModelSelector + * - Handles fresh context requirements via ContextManager + * - Tracks overall execution progress + * - Integrates with SpawnOrchestrator for session-based execution + * + * @module execution-bridge + */ + +import { EventEmitter } from 'node:events'; +import { + GroupScheduler, + getGroupScheduler, + type ExecutionGroup, + type GroupTask, + type ExecutionSchedule, +} from './group-scheduler.js'; +import { + ModelSelector, + getModelSelector, + type ModelConfig, + type ModelSelection, + type ExecutionMode, +} from './model-selector.js'; +import { + ContextManager, + getContextManager, + type SessionWriter, +} from './context-manager.js'; +import { + MAX_PARALLEL_TASKS_PER_GROUP, + GROUP_TIMEOUT_MS, + MAX_TASK_RETRIES, + TASK_RETRY_DELAY_MS, + EXECUTION_POLL_INTERVAL_MS, + MAX_EXECUTION_HISTORY, +} from './config/execution-limits.js'; + +// ========== Types ========== + +/** Overall execution status */ +export type ExecutionStatus = 'idle' | 'loading' | 'running' | 'paused' | 'completed' | 'partial' | 'failed' | 'cancelled'; + +/** + * Execution progress for UI display. + */ +export interface ExecutionProgress { + /** Current status */ + status: ExecutionStatus; + /** Current group being executed */ + currentGroup: number | null; + /** Total groups */ + totalGroups: number; + /** Completed groups */ + completedGroups: number; + /** Total tasks */ + totalTasks: number; + /** Completed tasks */ + completedTasks: number; + /** Failed tasks */ + failedTasks: number; + /** Tasks currently running */ + runningTasks: number; + /** Elapsed time in milliseconds */ + elapsedMs: number; + /** Estimated remaining time (if available) */ + estimatedRemainingMs?: number; +} + +/** + * Task assignment for spawning. + */ +export interface TaskAssignment { + /** Task from the schedule */ + task: GroupTask; + /** Selected model */ + model: ModelSelection; + /** Execution mode */ + executionMode: ExecutionMode; + /** Session ID (if session mode) */ + sessionId?: string; + /** Whether fresh context was requested */ + freshContextRequested: boolean; +} + +/** + * Plan item from PlanOrchestrator (input format). + */ +export interface PlanItem { + id: string; + title: string; + description: string; + parallelGroup?: number; + agentType?: string; + recommendedModel?: string; + requiresFreshContext?: boolean; + estimatedTokens?: number; + inputFiles?: string[]; + outputFiles?: string[]; + dependencies?: string[]; +} + +/** + * Execution history entry. + */ +export interface ExecutionHistoryEntry { + /** Unique execution ID */ + id: string; + /** When execution started */ + startedAt: number; + /** When execution ended */ + endedAt?: number; + /** Final status */ + status: ExecutionStatus; + /** Task counts */ + totalTasks: number; + completedTasks: number; + failedTasks: number; + /** Total cost estimate */ + estimatedCost?: number; +} + +// ========== Events ========== + +export interface ExecutionBridgeEvents { + /** Plan loaded and schedule built */ + planLoaded: (schedule: ExecutionSchedule) => void; + /** Execution started */ + started: () => void; + /** Execution paused */ + paused: () => void; + /** Execution resumed */ + resumed: () => void; + /** Execution completed */ + completed: (result: { status: ExecutionStatus; stats: ExecutionProgress }) => void; + /** Execution cancelled */ + cancelled: (reason: string) => void; + /** Group started */ + groupStarted: (data: { groupNumber: number; taskCount: number; executionMode: ExecutionMode }) => void; + /** Group completed */ + groupCompleted: (data: { groupNumber: number; status: string; completedCount: number; failedCount: number }) => void; + /** Task assigned to execution */ + taskAssigned: (assignment: TaskAssignment) => void; + /** Task completed */ + taskCompleted: (data: { taskId: string; groupNumber: number; durationMs: number }) => void; + /** Task failed */ + taskFailed: (data: { taskId: string; groupNumber: number; error: string; willRetry: boolean }) => void; + /** Fresh context triggered */ + freshContext: (data: { taskId: string; method: string; success: boolean }) => void; + /** Model selected for task */ + modelSelected: (data: { taskId: string; model: string; reason: string; optimizerSuggested?: string }) => void; + /** Progress update */ + progress: (progress: ExecutionProgress) => void; +} + +// ========== Spawn Interface ========== + +/** + * Interface for spawning agents. + * Injected to avoid circular dependencies. + */ +export interface AgentSpawner { + /** Spawn an agent with specific model */ + spawnAgentWithModel( + taskId: string, + workingDir: string, + prompt: string, + model: string, + options?: { requiresFreshContext?: boolean } + ): Promise<{ sessionId: string }>; + /** Use Task tool for lightweight execution */ + useTaskTool( + sessionId: string, + taskId: string, + prompt: string, + model: string + ): Promise; + /** Check if task is complete */ + isTaskComplete(taskId: string): boolean; + /** Get task result */ + getTaskResult(taskId: string): { success: boolean; output?: string; error?: string } | null; +} + +// ========== Execution Bridge ========== + +/** + * ExecutionBridge - Coordinates parallel task execution. + * + * This is the main entry point for the optimized execution system. + * It bridges the gap between PlanOrchestrator's output and actual + * task execution, making use of the optimizer's metadata. + */ +export class ExecutionBridge extends EventEmitter { + private _scheduler: GroupScheduler; + private _modelSelector: ModelSelector; + private _contextManager: ContextManager; + + private _status: ExecutionStatus = 'idle'; + private _executionId: string | null = null; + private _startedAt: number | null = null; + private _pausedAt: number | null = null; + + private _agentSpawner: AgentSpawner | null = null; + private _sessionWriter: SessionWriter | null = null; + + private _pollTimer: NodeJS.Timeout | null = null; + private _groupTimeoutTimers: Map = new Map(); + private _runningTasks: Map = new Map(); + + private _history: ExecutionHistoryEntry[] = []; + private _workingDir: string = process.cwd(); + + constructor(modelConfig?: Partial) { + super(); + this._scheduler = getGroupScheduler(); + this._modelSelector = getModelSelector(modelConfig); + this._contextManager = getContextManager(); + + // Forward scheduler events + this._scheduler.on('groupStarted', group => { + this.emit('groupStarted', { + groupNumber: group.groupNumber, + taskCount: group.tasks.length, + executionMode: group.executionMode, + }); + }); + + this._scheduler.on('groupCompleted', group => { + this.emit('groupCompleted', { + groupNumber: group.groupNumber, + status: group.status, + completedCount: group.completedCount, + failedCount: group.failedCount, + }); + }); + + this._scheduler.on('taskStatusChanged', data => { + if (data.newStatus === 'completed') { + const runningInfo = this._runningTasks.get(data.taskId); + const durationMs = runningInfo ? Date.now() - runningInfo.startedAt : 0; + this._runningTasks.delete(data.taskId); + this.emit('taskCompleted', { taskId: data.taskId, groupNumber: data.groupNumber, durationMs }); + } + }); + } + + /** + * Get current execution status. + */ + get status(): ExecutionStatus { + return this._status; + } + + /** + * Get current execution progress. + */ + getProgress(): ExecutionProgress { + const stats = this._scheduler.getStats(); + const schedule = this._scheduler.schedule; + + return { + status: this._status, + currentGroup: schedule?.currentGroupIndex ?? null, + totalGroups: stats.totalGroups, + completedGroups: stats.completedGroups, + totalTasks: stats.totalTasks, + completedTasks: stats.completedTasks, + failedTasks: stats.failedTasks, + runningTasks: this._runningTasks.size, + elapsedMs: this.calculateElapsedMs(), + }; + } + + /** + * Calculate elapsed time, accounting for pauses. + */ + private calculateElapsedMs(): number { + if (!this._startedAt) return 0; + if (this._pausedAt) { + return this._pausedAt - this._startedAt; + } + return Date.now() - this._startedAt; + } + + /** + * Set the agent spawner for session-based execution. + */ + setAgentSpawner(spawner: AgentSpawner): void { + this._agentSpawner = spawner; + } + + /** + * Set the session writer for context management. + */ + setSessionWriter(writer: SessionWriter): void { + this._sessionWriter = writer; + this._contextManager.setSessionWriter(writer); + } + + /** + * Set working directory for execution. + */ + setWorkingDir(dir: string): void { + this._workingDir = dir; + } + + /** + * Update model configuration. + */ + updateModelConfig(config: Partial): void { + this._modelSelector.updateConfig(config); + } + + /** + * Get current model configuration. + */ + getModelConfig(): ModelConfig { + return this._modelSelector.config; + } + + /** + * Load a plan and build execution schedule. + */ + loadPlan(items: PlanItem[]): ExecutionSchedule { + if (this._status === 'running') { + throw new Error('Cannot load plan while execution is running'); + } + + this._status = 'loading'; + const schedule = this._scheduler.buildSchedule(items); + this._status = 'idle'; + + this.emit('planLoaded', schedule); + return schedule; + } + + /** + * Start execution of the loaded plan. + */ + async start(): Promise { + const schedule = this._scheduler.schedule; + if (!schedule) { + throw new Error('No plan loaded'); + } + + if (this._status === 'running') { + return; + } + + if (!this._agentSpawner) { + throw new Error('No agent spawner configured'); + } + + this._status = 'running'; + this._executionId = `exec-${Date.now()}`; + this._startedAt = Date.now(); + + // Add to history + this._history.unshift({ + id: this._executionId, + startedAt: this._startedAt, + status: 'running', + totalTasks: schedule.totalTasks, + completedTasks: 0, + failedTasks: 0, + }); + + // Trim history + if (this._history.length > MAX_EXECUTION_HISTORY) { + this._history = this._history.slice(0, MAX_EXECUTION_HISTORY); + } + + this.emit('started'); + + // Start the execution loop + this.startExecutionLoop(); + } + + /** + * Pause execution. + */ + pause(): void { + if (this._status !== 'running') return; + + this._status = 'paused'; + this._pausedAt = Date.now(); + this.stopExecutionLoop(); + this.emit('paused'); + } + + /** + * Resume paused execution. + */ + resume(): void { + if (this._status !== 'paused') return; + + this._status = 'running'; + this._pausedAt = null; + this.startExecutionLoop(); + this.emit('resumed'); + } + + /** + * Cancel execution. + */ + async cancel(reason: string = 'User cancelled'): Promise { + if (this._status === 'idle' || this._status === 'completed' || this._status === 'cancelled') { + return; + } + + this._status = 'cancelled'; + this.stopExecutionLoop(); + + // Clear group timeouts + for (const timer of this._groupTimeoutTimers.values()) { + clearTimeout(timer); + } + this._groupTimeoutTimers.clear(); + + // Update history + this.updateHistoryEntry('cancelled'); + + this.emit('cancelled', reason); + } + + /** + * Get execution history. + */ + getHistory(): ExecutionHistoryEntry[] { + return [...this._history]; + } + + /** + * Start the main execution loop. + */ + private startExecutionLoop(): void { + if (this._pollTimer) return; + + this._pollTimer = setInterval(() => { + this.tick().catch(err => { + console.error('[execution-bridge] Tick error:', err); + }); + }, EXECUTION_POLL_INTERVAL_MS); + + // Run immediately + this.tick().catch(err => { + console.error('[execution-bridge] Initial tick error:', err); + }); + } + + /** + * Stop the execution loop. + */ + private stopExecutionLoop(): void { + if (this._pollTimer) { + clearInterval(this._pollTimer); + this._pollTimer = null; + } + } + + /** + * Main execution tick. + */ + private async tick(): Promise { + if (this._status !== 'running') return; + + const schedule = this._scheduler.schedule; + if (!schedule) return; + + // Check if we're done + if (schedule.status === 'completed' || schedule.status === 'partial' || schedule.status === 'failed') { + this.handleExecutionComplete(); + return; + } + + // Process current group or start next one + await this.processGroups(); + + // Emit progress + this.emit('progress', this.getProgress()); + } + + /** + * Process groups - start new ones or continue existing. + */ + private async processGroups(): Promise { + const schedule = this._scheduler.schedule; + if (!schedule) return; + + // Find current running group + const runningGroup = schedule.groups.find(g => g.status === 'running'); + + if (runningGroup) { + // Continue processing current group + await this.processGroup(runningGroup); + } else { + // Try to start next group + const nextGroup = this._scheduler.getNextReadyGroup(); + if (nextGroup) { + await this.startGroup(nextGroup); + } + } + } + + /** + * Start executing a group. + */ + private async startGroup(group: ExecutionGroup): Promise { + this._scheduler.startGroup(group.groupNumber); + + // Set up group timeout + const timeoutTimer = setTimeout(() => { + this.handleGroupTimeout(group.groupNumber); + }, GROUP_TIMEOUT_MS); + this._groupTimeoutTimers.set(group.groupNumber, timeoutTimer); + + // Start initial tasks + await this.processGroup(group); + } + + /** + * Process tasks within a group. + */ + private async processGroup(group: ExecutionGroup): Promise { + if (!this._agentSpawner) return; + + // Get ready tasks + const readyTasks = this._scheduler.getReadyTasksInGroup(group.groupNumber); + if (readyTasks.length === 0) return; + + // Limit parallel tasks + const slotsAvailable = MAX_PARALLEL_TASKS_PER_GROUP - this._runningTasks.size; + if (slotsAvailable <= 0) return; + + const tasksToStart = readyTasks.slice(0, slotsAvailable); + + for (const task of tasksToStart) { + await this.assignTask(task, group); + } + } + + /** + * Assign a task for execution. + */ + private async assignTask(task: GroupTask, group: ExecutionGroup): Promise { + if (!this._agentSpawner) return; + + // Select model + const modelSelection = this._modelSelector.selectModel(task.id, { + estimatedTokens: task.estimatedTokens, + agentType: task.agentType, + recommendedModel: task.recommendedModel, + outputFiles: task.outputFiles, + inputFiles: task.inputFiles, + }); + + this.emit('modelSelected', { + taskId: task.id, + model: modelSelection.model, + reason: modelSelection.reason, + optimizerSuggested: modelSelection.optimizerRecommendation, + }); + + // Handle fresh context if required + let freshContextRequested = false; + if (task.requiresFreshContext && this._sessionWriter) { + freshContextRequested = true; + // Context refresh will be handled by the spawner + } + + // Mark task as running + this._scheduler.updateTaskStatus(task.id, 'running'); + this._runningTasks.set(task.id, { startedAt: Date.now() }); + + const assignment: TaskAssignment = { + task, + model: modelSelection, + executionMode: group.executionMode, + freshContextRequested, + }; + + this.emit('taskAssigned', assignment); + + // Execute based on mode + try { + if (group.executionMode === 'session') { + const result = await this._agentSpawner.spawnAgentWithModel( + task.id, + this._workingDir, + task.description, + modelSelection.model, + { requiresFreshContext: task.requiresFreshContext } + ); + assignment.sessionId = result.sessionId; + this._runningTasks.set(task.id, { startedAt: Date.now(), sessionId: result.sessionId }); + } else { + // task-tool mode - would use Task tool in main session + // For now, fall back to session mode + const result = await this._agentSpawner.spawnAgentWithModel( + task.id, + this._workingDir, + task.description, + modelSelection.model + ); + assignment.sessionId = result.sessionId; + this._runningTasks.set(task.id, { startedAt: Date.now(), sessionId: result.sessionId }); + } + + if (task.requiresFreshContext) { + this.emit('freshContext', { + taskId: task.id, + method: 'session', + success: true, + }); + } + + } catch (err) { + const error = err instanceof Error ? err.message : String(err); + await this.handleTaskFailure(task, group.groupNumber, error); + } + } + + /** + * Handle task failure. + */ + private async handleTaskFailure(task: GroupTask, groupNumber: number, error: string): Promise { + task.retryCount++; + const willRetry = task.retryCount < MAX_TASK_RETRIES; + + this.emit('taskFailed', { + taskId: task.id, + groupNumber, + error, + willRetry, + }); + + if (willRetry) { + // Schedule retry + setTimeout(() => { + task.status = 'pending'; + task.error = undefined; + this._runningTasks.delete(task.id); + }, TASK_RETRY_DELAY_MS); + } else { + // Mark as permanently failed + this._scheduler.updateTaskStatus(task.id, 'failed', error); + this._runningTasks.delete(task.id); + + // Mark dependent tasks as blocked + this._scheduler.markDependentTasksBlocked(task.id); + } + } + + /** + * Mark a task as complete (called externally when agent finishes). + */ + markTaskComplete(taskId: string): void { + const schedule = this._scheduler.schedule; + if (!schedule) return; + + const groupNum = this.findTaskGroup(taskId); + if (groupNum === null) return; + + this._scheduler.updateTaskStatus(taskId, 'completed'); + this._runningTasks.delete(taskId); + } + + /** + * Mark a task as failed (called externally when agent fails). + */ + markTaskFailed(taskId: string, error: string): void { + const schedule = this._scheduler.schedule; + if (!schedule) return; + + const groupNum = this.findTaskGroup(taskId); + if (groupNum === null) return; + + // Get the task to check retry count + for (const group of schedule.groups) { + const task = group.tasks.find(t => t.id === taskId); + if (task) { + this.handleTaskFailure(task, group.groupNumber, error); + return; + } + } + } + + /** + * Find which group a task belongs to. + */ + private findTaskGroup(taskId: string): number | null { + const schedule = this._scheduler.schedule; + if (!schedule) return null; + + for (const group of schedule.groups) { + if (group.tasks.some(t => t.id === taskId)) { + return group.groupNumber; + } + } + return null; + } + + /** + * Handle group timeout. + */ + private handleGroupTimeout(groupNumber: number): void { + const schedule = this._scheduler.schedule; + if (!schedule) return; + + const group = schedule.groups.find(g => g.groupNumber === groupNumber); + if (!group || group.status !== 'running') return; + + console.warn(`[execution-bridge] Group ${groupNumber} timed out`); + + // Mark running tasks as failed + for (const task of group.tasks) { + if (task.status === 'running') { + this._scheduler.updateTaskStatus(task.id, 'failed', 'Group timeout'); + this._runningTasks.delete(task.id); + } else if (task.status === 'pending') { + this._scheduler.updateTaskStatus(task.id, 'skipped', 'Group timeout'); + } + } + + this._groupTimeoutTimers.delete(groupNumber); + } + + /** + * Handle execution complete. + */ + private handleExecutionComplete(): void { + this.stopExecutionLoop(); + + // Clear timers + for (const timer of this._groupTimeoutTimers.values()) { + clearTimeout(timer); + } + this._groupTimeoutTimers.clear(); + + const schedule = this._scheduler.schedule; + if (!schedule) return; + + // Map schedule status to execution status + if (schedule.status === 'completed') { + this._status = 'completed'; + } else if (schedule.status === 'partial') { + this._status = 'partial'; + } else { + this._status = 'failed'; + } + + // Update history + this.updateHistoryEntry(this._status); + + this.emit('completed', { + status: this._status, + stats: this.getProgress(), + }); + } + + /** + * Update history entry for current execution. + */ + private updateHistoryEntry(status: ExecutionStatus): void { + if (!this._executionId) return; + + const entry = this._history.find(e => e.id === this._executionId); + if (entry) { + entry.status = status; + entry.endedAt = Date.now(); + entry.completedTasks = this._scheduler.getStats().completedTasks; + entry.failedTasks = this._scheduler.getStats().failedTasks; + } + } + + /** + * Reset the bridge for a new execution. + */ + reset(): void { + this.stopExecutionLoop(); + + for (const timer of this._groupTimeoutTimers.values()) { + clearTimeout(timer); + } + this._groupTimeoutTimers.clear(); + + this._scheduler.reset(); + this._runningTasks.clear(); + this._status = 'idle'; + this._executionId = null; + this._startedAt = null; + this._pausedAt = null; + } +} + +// ========== Singleton ========== + +let bridgeInstance: ExecutionBridge | null = null; + +/** + * Get or create the singleton ExecutionBridge instance. + */ +export function getExecutionBridge(modelConfig?: Partial): ExecutionBridge { + if (!bridgeInstance) { + bridgeInstance = new ExecutionBridge(modelConfig); + } + return bridgeInstance; +} + +/** + * Reset the singleton (for testing). + */ +export function resetExecutionBridge(): void { + if (bridgeInstance) { + bridgeInstance.reset(); + } + bridgeInstance = null; +} diff --git a/src/group-scheduler.ts b/src/group-scheduler.ts new file mode 100644 index 00000000..23c74c95 --- /dev/null +++ b/src/group-scheduler.ts @@ -0,0 +1,551 @@ +/** + * @fileoverview Group Scheduler - Topological ordering and dependency management. + * + * Builds execution schedules from parallel groups, manages dependencies + * between groups, and tracks group-level completion status. + * + * @module group-scheduler + */ + +import { EventEmitter } from 'node:events'; +import type { ExecutionMode } from './model-selector.js'; +import type { ModelTier, AgentType } from './model-selector.js'; + +// ========== Types ========== + +/** Status of a task within a group */ +export type GroupTaskStatus = 'pending' | 'running' | 'completed' | 'failed' | 'blocked' | 'skipped'; + +/** Status of an execution group */ +export type ExecutionGroupStatus = 'pending' | 'ready' | 'running' | 'completed' | 'partial' | 'failed'; + +/** + * A task within an execution group. + */ +export interface GroupTask { + /** Task ID from the plan */ + id: string; + /** Task title/description */ + title: string; + /** Full task description */ + description: string; + /** Parallel group number */ + parallelGroup: number; + /** Agent type for model selection */ + agentType: AgentType; + /** Recommended model */ + recommendedModel?: ModelTier; + /** Whether task requires fresh context */ + requiresFreshContext: boolean; + /** Estimated token usage */ + estimatedTokens?: number; + /** Input files (read-only) */ + inputFiles?: string[]; + /** Output files (will be modified) */ + outputFiles?: string[]; + /** Current status */ + status: GroupTaskStatus; + /** Task dependencies (other task IDs) */ + dependencies: string[]; + /** Error message if failed */ + error?: string; + /** Retry count */ + retryCount: number; +} + +/** + * An execution group containing tasks that can run in parallel. + */ +export interface ExecutionGroup { + /** Group number (from parallelGroup) */ + groupNumber: number; + /** Tasks in this group */ + tasks: GroupTask[]; + /** Group status */ + status: ExecutionGroupStatus; + /** Execution mode for this group */ + executionMode: ExecutionMode; + /** Rationale for execution mode choice */ + executionModeRationale: string; + /** Groups that must complete before this one */ + dependsOnGroups: number[]; + /** When group started executing */ + startedAt?: number; + /** When group completed */ + completedAt?: number; + /** Number of tasks completed */ + completedCount: number; + /** Number of tasks failed */ + failedCount: number; + /** Number of tasks skipped due to dependency failures */ + skippedCount: number; +} + +/** + * Full execution schedule. + */ +export interface ExecutionSchedule { + /** Ordered groups (by dependency order) */ + groups: ExecutionGroup[]; + /** Total task count */ + totalTasks: number; + /** Completed task count */ + completedTasks: number; + /** Failed task count */ + failedTasks: number; + /** Current executing group index (-1 if not started) */ + currentGroupIndex: number; + /** Overall status */ + status: 'pending' | 'running' | 'completed' | 'partial' | 'failed'; +} + +// ========== Events ========== + +export interface GroupSchedulerEvents { + /** Schedule built from plan */ + scheduleBuilt: (schedule: ExecutionSchedule) => void; + /** Group started executing */ + groupStarted: (group: ExecutionGroup) => void; + /** Group completed (fully or partially) */ + groupCompleted: (group: ExecutionGroup) => void; + /** Task status changed */ + taskStatusChanged: (data: { taskId: string; groupNumber: number; oldStatus: GroupTaskStatus; newStatus: GroupTaskStatus }) => void; +} + +// ========== Group Scheduler ========== + +/** + * GroupScheduler - Manages execution order and dependencies. + * + * Responsibilities: + * - Build topologically ordered groups from plan items + * - Track group dependencies (lower groups must complete first) + * - Handle partial failures (continue with independent tasks) + * - Determine group execution mode (session vs task-tool) + */ +export class GroupScheduler extends EventEmitter { + private _schedule: ExecutionSchedule | null = null; + private _taskToGroup: Map = new Map(); + + constructor() { + super(); + } + + /** + * Get current schedule. + */ + get schedule(): ExecutionSchedule | null { + return this._schedule; + } + + /** + * Build execution schedule from plan items. + * + * @param items - Array of plan items with parallelGroup assignments + * @returns Built execution schedule + */ + buildSchedule(items: Array<{ + id: string; + title: string; + description: string; + parallelGroup?: number; + agentType?: string; + recommendedModel?: string; + requiresFreshContext?: boolean; + estimatedTokens?: number; + inputFiles?: string[]; + outputFiles?: string[]; + dependencies?: string[]; + }>): ExecutionSchedule { + // Group items by parallel group + const groupMap = new Map(); + + for (const item of items) { + const groupNum = item.parallelGroup ?? 0; + const task: GroupTask = { + id: item.id, + title: item.title, + description: item.description, + parallelGroup: groupNum, + agentType: (item.agentType as AgentType) ?? 'general', + recommendedModel: item.recommendedModel as ModelTier | undefined, + requiresFreshContext: item.requiresFreshContext ?? false, + estimatedTokens: item.estimatedTokens, + inputFiles: item.inputFiles, + outputFiles: item.outputFiles, + status: 'pending', + dependencies: item.dependencies ?? [], + retryCount: 0, + }; + + if (!groupMap.has(groupNum)) { + groupMap.set(groupNum, []); + } + groupMap.get(groupNum)!.push(task); + this._taskToGroup.set(item.id, groupNum); + } + + // Sort groups by number + const sortedGroupNums = Array.from(groupMap.keys()).sort((a, b) => a - b); + + // Build execution groups + const groups: ExecutionGroup[] = sortedGroupNums.map(groupNum => { + const tasks = groupMap.get(groupNum)!; + + // Determine which groups this one depends on + const dependsOnGroups = new Set(); + for (const task of tasks) { + for (const depId of task.dependencies) { + const depGroup = this._taskToGroup.get(depId); + if (depGroup !== undefined && depGroup !== groupNum && depGroup < groupNum) { + dependsOnGroups.add(depGroup); + } + } + } + + // Determine execution mode based on task characteristics + const { mode, rationale } = this.determineGroupExecutionMode(tasks); + + return { + groupNumber: groupNum, + tasks, + status: 'pending', + executionMode: mode, + executionModeRationale: rationale, + dependsOnGroups: Array.from(dependsOnGroups).sort((a, b) => a - b), + completedCount: 0, + failedCount: 0, + skippedCount: 0, + }; + }); + + this._schedule = { + groups, + totalTasks: items.length, + completedTasks: 0, + failedTasks: 0, + currentGroupIndex: -1, + status: 'pending', + }; + + this.emit('scheduleBuilt', this._schedule); + return this._schedule; + } + + /** + * Determine execution mode for a group based on task characteristics. + */ + private determineGroupExecutionMode(tasks: GroupTask[]): { mode: ExecutionMode; rationale: string } { + // High token estimate → session mode + const highTokenTask = tasks.find(t => t.estimatedTokens && t.estimatedTokens > 50000); + if (highTokenTask) { + return { + mode: 'session', + rationale: `Task ${highTokenTask.id} has high token estimate (${highTokenTask.estimatedTokens})`, + }; + } + + // Complex agent types → session mode + const complexTask = tasks.find(t => t.agentType === 'implement' || t.agentType === 'review'); + if (complexTask) { + return { + mode: 'session', + rationale: `Task ${complexTask.id} has complex agent type (${complexTask.agentType})`, + }; + } + + // Multiple output files in any task → session mode + const multiOutputTask = tasks.find(t => t.outputFiles && t.outputFiles.length > 2); + if (multiOutputTask) { + return { + mode: 'session', + rationale: `Task ${multiOutputTask.id} has multiple output files (${multiOutputTask.outputFiles!.length})`, + }; + } + + // Fresh context required → session mode + const freshContextTask = tasks.find(t => t.requiresFreshContext); + if (freshContextTask) { + return { + mode: 'session', + rationale: `Task ${freshContextTask.id} requires fresh context`, + }; + } + + // All low-token explore tasks → task-tool mode + const allLowToken = tasks.every(t => !t.estimatedTokens || t.estimatedTokens < 15000); + const allExplore = tasks.every(t => t.agentType === 'explore' || t.agentType === 'general'); + if (allLowToken && allExplore) { + return { + mode: 'task-tool', + rationale: 'All tasks are low-token explore/general tasks', + }; + } + + // Default to session mode for reliability + return { + mode: 'session', + rationale: 'Default to session mode for reliability', + }; + } + + /** + * Get the next group ready for execution. + */ + getNextReadyGroup(): ExecutionGroup | null { + if (!this._schedule) return null; + + for (const group of this._schedule.groups) { + if (group.status === 'pending' && this.areGroupDependenciesSatisfied(group)) { + group.status = 'ready'; + return group; + } + } + + return null; + } + + /** + * Check if a group's dependencies are satisfied. + */ + areGroupDependenciesSatisfied(group: ExecutionGroup): boolean { + if (!this._schedule) return false; + + for (const depGroupNum of group.dependsOnGroups) { + const depGroup = this._schedule.groups.find(g => g.groupNumber === depGroupNum); + if (!depGroup) continue; + + // Dependency must be completed (fully or partially) + if (depGroup.status !== 'completed' && depGroup.status !== 'partial') { + return false; + } + } + + return true; + } + + /** + * Mark a group as started. + */ + startGroup(groupNumber: number): void { + if (!this._schedule) return; + + const group = this._schedule.groups.find(g => g.groupNumber === groupNumber); + if (!group) return; + + group.status = 'running'; + group.startedAt = Date.now(); + this._schedule.status = 'running'; + this._schedule.currentGroupIndex = this._schedule.groups.indexOf(group); + + this.emit('groupStarted', group); + } + + /** + * Update task status within a group. + */ + updateTaskStatus(taskId: string, status: GroupTaskStatus, error?: string): void { + if (!this._schedule) return; + + const groupNum = this._taskToGroup.get(taskId); + if (groupNum === undefined) return; + + const group = this._schedule.groups.find(g => g.groupNumber === groupNum); + if (!group) return; + + const task = group.tasks.find(t => t.id === taskId); + if (!task) return; + + const oldStatus = task.status; + task.status = status; + if (error) task.error = error; + + // Update group counters + if (status === 'completed') { + group.completedCount++; + this._schedule.completedTasks++; + } else if (status === 'failed') { + group.failedCount++; + this._schedule.failedTasks++; + } else if (status === 'skipped') { + group.skippedCount++; + } + + this.emit('taskStatusChanged', { taskId, groupNumber: groupNum, oldStatus, newStatus: status }); + + // Check if group is complete + this.checkGroupCompletion(group); + } + + /** + * Mark tasks blocked by a failed dependency. + */ + markDependentTasksBlocked(failedTaskId: string): void { + if (!this._schedule) return; + + for (const group of this._schedule.groups) { + for (const task of group.tasks) { + if (task.dependencies.includes(failedTaskId) && task.status === 'pending') { + this.updateTaskStatus(task.id, 'skipped', `Blocked by failed task ${failedTaskId}`); + } + } + } + } + + /** + * Get tasks ready to execute in a group. + */ + getReadyTasksInGroup(groupNumber: number): GroupTask[] { + if (!this._schedule) return []; + + const group = this._schedule.groups.find(g => g.groupNumber === groupNumber); + if (!group) return []; + + return group.tasks.filter(task => { + if (task.status !== 'pending') return false; + + // Check if task dependencies are satisfied (within and across groups) + for (const depId of task.dependencies) { + // Check if dependency is in the same group + const sameGroupDep = group.tasks.find(t => t.id === depId); + if (sameGroupDep && sameGroupDep.status !== 'completed') { + return false; + } + + // Check if dependency is in a different group + const depGroupNum = this._taskToGroup.get(depId); + if (depGroupNum !== undefined && depGroupNum !== groupNumber) { + const depGroup = this._schedule!.groups.find(g => g.groupNumber === depGroupNum); + const depTask = depGroup?.tasks.find(t => t.id === depId); + if (depTask && depTask.status !== 'completed') { + return false; + } + } + } + + return true; + }); + } + + /** + * Check if a group has completed (all tasks done, failed, or skipped). + */ + private checkGroupCompletion(group: ExecutionGroup): void { + const pendingOrRunning = group.tasks.filter( + t => t.status === 'pending' || t.status === 'running' + ); + + if (pendingOrRunning.length > 0) return; + + group.completedAt = Date.now(); + + // Determine final status + if (group.failedCount === 0 && group.skippedCount === 0) { + group.status = 'completed'; + } else if (group.completedCount > 0) { + group.status = 'partial'; + } else { + group.status = 'failed'; + } + + this.emit('groupCompleted', group); + + // Check if all groups are done + this.checkScheduleCompletion(); + } + + /** + * Check if entire schedule has completed. + */ + private checkScheduleCompletion(): void { + if (!this._schedule) return; + + const pendingOrRunning = this._schedule.groups.filter( + g => g.status === 'pending' || g.status === 'ready' || g.status === 'running' + ); + + if (pendingOrRunning.length > 0) return; + + // Determine final status + const failedGroups = this._schedule.groups.filter(g => g.status === 'failed'); + const partialGroups = this._schedule.groups.filter(g => g.status === 'partial'); + + if (failedGroups.length === this._schedule.groups.length) { + this._schedule.status = 'failed'; + } else if (failedGroups.length > 0 || partialGroups.length > 0) { + this._schedule.status = 'partial'; + } else { + this._schedule.status = 'completed'; + } + } + + /** + * Get schedule statistics. + */ + getStats(): { + totalGroups: number; + completedGroups: number; + failedGroups: number; + partialGroups: number; + totalTasks: number; + completedTasks: number; + failedTasks: number; + skippedTasks: number; + } { + if (!this._schedule) { + return { + totalGroups: 0, + completedGroups: 0, + failedGroups: 0, + partialGroups: 0, + totalTasks: 0, + completedTasks: 0, + failedTasks: 0, + skippedTasks: 0, + }; + } + + return { + totalGroups: this._schedule.groups.length, + completedGroups: this._schedule.groups.filter(g => g.status === 'completed').length, + failedGroups: this._schedule.groups.filter(g => g.status === 'failed').length, + partialGroups: this._schedule.groups.filter(g => g.status === 'partial').length, + totalTasks: this._schedule.totalTasks, + completedTasks: this._schedule.completedTasks, + failedTasks: this._schedule.failedTasks, + skippedTasks: this._schedule.groups.reduce((sum, g) => sum + g.skippedCount, 0), + }; + } + + /** + * Reset the scheduler. + */ + reset(): void { + this._schedule = null; + this._taskToGroup.clear(); + } +} + +// ========== Singleton ========== + +let schedulerInstance: GroupScheduler | null = null; + +/** + * Get or create the singleton GroupScheduler instance. + */ +export function getGroupScheduler(): GroupScheduler { + if (!schedulerInstance) { + schedulerInstance = new GroupScheduler(); + } + return schedulerInstance; +} + +/** + * Reset the singleton (for testing). + */ +export function resetGroupScheduler(): void { + if (schedulerInstance) { + schedulerInstance.reset(); + } + schedulerInstance = null; +} diff --git a/src/model-selector.ts b/src/model-selector.ts new file mode 100644 index 00000000..b734ff2d --- /dev/null +++ b/src/model-selector.ts @@ -0,0 +1,314 @@ +/** + * @fileoverview Model Selector - Routes tasks to appropriate Claude models. + * + * Handles model selection based on: + * - User-configured default model + * - Optimizer recommendations (advisory) + * - Agent type mappings + * - Per-task overrides + * + * @module model-selector + */ + +import { EventEmitter } from 'node:events'; +import { + DEFAULT_MODEL, + TOKEN_THRESHOLD_HAIKU, + TOKEN_THRESHOLD_FOR_SESSION_MODE, +} from './config/execution-limits.js'; + +// ========== Types ========== + +/** Supported Claude model tiers */ +export type ModelTier = 'opus' | 'sonnet' | 'haiku'; + +/** Agent types that influence model selection */ +export type AgentType = 'explore' | 'implement' | 'test' | 'review' | 'general'; + +/** Execution mode for tasks */ +export type ExecutionMode = 'session' | 'task-tool'; + +/** + * User-configurable model settings. + * Stored in settings.json and editable via App Settings. + */ +export interface ModelConfig { + /** User's preferred default model */ + defaultModel: ModelTier; + /** Whether to show optimizer recommendations in UI (advisory only) */ + showRecommendations: boolean; + /** Override map for specific agent types */ + agentTypeOverrides: Partial>; +} + +/** + * Model selection result with reasoning. + */ +export interface ModelSelection { + /** The model to use */ + model: ModelTier; + /** Why this model was selected */ + reason: string; + /** What the optimizer recommended (if different) */ + optimizerRecommendation?: ModelTier; + /** Was user default used? */ + usedUserDefault: boolean; +} + +/** + * Execution mode selection result. + */ +export interface ExecutionModeSelection { + /** How to execute this task */ + mode: ExecutionMode; + /** Why this mode was selected */ + rationale: string; +} + +/** + * Task characteristics for selection decisions. + */ +export interface TaskCharacteristics { + /** Estimated token usage */ + estimatedTokens?: number; + /** Agent type */ + agentType?: AgentType; + /** Optimizer's recommended model */ + recommendedModel?: ModelTier; + /** Files task will modify */ + outputFiles?: string[]; + /** Files task will read */ + inputFiles?: string[]; + /** Task complexity hint */ + complexity?: 'low' | 'medium' | 'high'; +} + +// ========== Events ========== + +export interface ModelSelectorEvents { + /** Emitted when model is selected */ + modelSelected: (data: { taskId: string; selection: ModelSelection }) => void; + /** Emitted when config changes */ + configUpdated: (config: ModelConfig) => void; +} + +// ========== Default Configuration ========== + +/** + * Default optimizer recommendations by agent type. + * These are advisory - user default always wins. + */ +const DEFAULT_RECOMMENDATIONS: Record = { + explore: 'haiku', + implement: 'sonnet', + test: 'sonnet', + review: 'opus', + general: 'sonnet', +}; + +/** + * Creates default model configuration. + */ +export function createDefaultModelConfig(): ModelConfig { + return { + defaultModel: DEFAULT_MODEL, + showRecommendations: true, + agentTypeOverrides: {}, + }; +} + +// ========== Model Selector ========== + +/** + * ModelSelector - Manages model selection for task execution. + * + * User preferences always take precedence. Optimizer recommendations + * are shown in the UI for awareness but don't override user settings. + */ +export class ModelSelector extends EventEmitter { + private _config: ModelConfig; + + constructor(config?: Partial) { + super(); + this._config = { ...createDefaultModelConfig(), ...config }; + } + + /** + * Get current configuration. + */ + get config(): ModelConfig { + return { ...this._config }; + } + + /** + * Update configuration. + */ + updateConfig(config: Partial): void { + Object.assign(this._config, config); + this.emit('configUpdated', this._config); + } + + /** + * Select model for a task. + * + * Priority order: + * 1. User's agent type override (if set) + * 2. User's default model + * 3. Optimizer recommendation (only if no user preference) + */ + selectModel(taskId: string, characteristics: TaskCharacteristics): ModelSelection { + const { agentType, recommendedModel } = characteristics; + + let model: ModelTier; + let reason: string; + let usedUserDefault = false; + + // Check for agent type override + if (agentType && this._config.agentTypeOverrides[agentType]) { + model = this._config.agentTypeOverrides[agentType]!; + reason = `User override for ${agentType} tasks`; + } else { + // Use user's default model + model = this._config.defaultModel; + reason = 'User default model'; + usedUserDefault = true; + } + + // Determine what optimizer would have recommended + const optimizerRecommendation = recommendedModel || + (agentType ? DEFAULT_RECOMMENDATIONS[agentType] : undefined); + + const selection: ModelSelection = { + model, + reason, + usedUserDefault, + }; + + // Include optimizer recommendation if different (for UI display) + if (optimizerRecommendation && optimizerRecommendation !== model) { + selection.optimizerRecommendation = optimizerRecommendation; + if (this._config.showRecommendations) { + selection.reason += ` (optimizer suggested ${optimizerRecommendation})`; + } + } + + this.emit('modelSelected', { taskId, selection }); + return selection; + } + + /** + * Select execution mode for a task. + * + * Decision based on task characteristics: + * - High token estimate → session mode (needs full context) + * - Complex agent types → session mode (dedicated Claude) + * - Low token estimate → task-tool mode (efficient) + * - Multiple output files → session mode (avoid conflicts) + * - Read-only tasks → task-tool mode (no side effects) + */ + selectExecutionMode(characteristics: TaskCharacteristics): ExecutionModeSelection { + const { estimatedTokens, agentType, outputFiles, inputFiles, complexity } = characteristics; + + // High token estimate needs session mode + if (estimatedTokens && estimatedTokens > TOKEN_THRESHOLD_FOR_SESSION_MODE) { + return { + mode: 'session', + rationale: `High token estimate (${estimatedTokens} > ${TOKEN_THRESHOLD_FOR_SESSION_MODE})`, + }; + } + + // Complex agent types benefit from session mode + if (agentType === 'implement' || agentType === 'review') { + return { + mode: 'session', + rationale: `Complex agent type (${agentType}) benefits from dedicated context`, + }; + } + + // Multiple output files need session mode to avoid conflicts + if (outputFiles && outputFiles.length > 2) { + return { + mode: 'session', + rationale: `Multiple output files (${outputFiles.length}) - avoid conflicts`, + }; + } + + // High complexity tasks need session mode + if (complexity === 'high') { + return { + mode: 'session', + rationale: 'High complexity task needs dedicated context', + }; + } + + // Low token tasks can use task-tool mode + if (estimatedTokens && estimatedTokens < TOKEN_THRESHOLD_HAIKU) { + return { + mode: 'task-tool', + rationale: `Low token estimate (${estimatedTokens} < ${TOKEN_THRESHOLD_HAIKU})`, + }; + } + + // Explore tasks typically work well with task-tool + if (agentType === 'explore') { + return { + mode: 'task-tool', + rationale: 'Explore tasks efficient with shared context', + }; + } + + // Read-only tasks (input files only) can use task-tool + if (inputFiles && inputFiles.length > 0 && (!outputFiles || outputFiles.length === 0)) { + return { + mode: 'task-tool', + rationale: 'Read-only task (no output files)', + }; + } + + // Default to session mode for safety + return { + mode: 'session', + rationale: 'Default to session mode for reliability', + }; + } + + /** + * Get optimizer's recommendation for an agent type (for UI display). + */ + getOptimizerRecommendation(agentType: AgentType): ModelTier { + return DEFAULT_RECOMMENDATIONS[agentType]; + } + + /** + * Get model cost multiplier (relative to sonnet). + * Used for cost estimation in UI. + */ + getModelCostMultiplier(model: ModelTier): number { + switch (model) { + case 'opus': return 5.0; // ~5x more expensive + case 'sonnet': return 1.0; // baseline + case 'haiku': return 0.04; // ~25x cheaper + } + } +} + +// ========== Singleton ========== + +let selectorInstance: ModelSelector | null = null; + +/** + * Get or create the singleton ModelSelector instance. + */ +export function getModelSelector(config?: Partial): ModelSelector { + if (!selectorInstance) { + selectorInstance = new ModelSelector(config); + } + return selectorInstance; +} + +/** + * Reset the singleton (for testing). + */ +export function resetModelSelector(): void { + selectorInstance = null; +} diff --git a/src/types.ts b/src/types.ts index 6cfc949a..d5430f68 100644 --- a/src/types.ts +++ b/src/types.ts @@ -139,6 +139,8 @@ export interface SessionState { autoCompactThreshold?: number; /** Auto-compact prompt */ autoCompactPrompt?: string; + /** Image watcher enabled for this session */ + imageWatcherEnabled?: boolean; /** Total cost in USD */ totalCost?: number; /** Input tokens used */ @@ -1270,3 +1272,38 @@ export type { AgentContext, SpawnPersistedState, } from './spawn-types.js'; + +// ========== Execution Bridge Re-exports ========== + +export type { + ExecutionStatus, + ExecutionProgress, + TaskAssignment as ExecutionTaskAssignment, + PlanItem, + ExecutionHistoryEntry, +} from './execution-bridge.js'; + +export type { + ModelTier, + AgentType, + ExecutionMode, + ModelConfig, + ModelSelection, + ExecutionModeSelection, + TaskCharacteristics, +} from './model-selector.js'; + +export type { + GroupTaskStatus, + ExecutionGroupStatus, + GroupTask, + ExecutionGroup, + ExecutionSchedule, +} from './group-scheduler.js'; + +export type { + ContextRefreshMethod, + ContextRefreshStatus, + ContextRefreshRequest, + ContextRefreshResult, +} from './context-manager.js'; diff --git a/src/web/public/app.js b/src/web/public/app.js index b2915377..9e5190e0 100644 --- a/src/web/public/app.js +++ b/src/web/public/app.js @@ -6193,6 +6193,7 @@ class ClaudemanApp { document.getElementById('modalAutoCompactPrompt').value = session.autoCompactPrompt ?? ''; document.getElementById('modalAutoClearEnabled').checked = session.autoClearEnabled ?? false; document.getElementById('modalAutoClearThreshold').value = session.autoClearThreshold ?? 140000; + document.getElementById('modalImageWatcherEnabled').checked = session.imageWatcherEnabled ?? true; // Populate Ralph Wiggum form with current session values const ralphState = this.ralphStates.get(sessionId); @@ -6234,6 +6235,26 @@ class ClaudemanApp { } catch { /* silent */ } } + async toggleSessionImageWatcher() { + if (!this.editingSessionId) return; + const enabled = document.getElementById('modalImageWatcherEnabled').checked; + try { + await fetch(`/api/sessions/${this.editingSessionId}/image-watcher`, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ enabled }) + }); + // Update local session state + const session = this.sessions.get(this.editingSessionId); + if (session) { + session.imageWatcherEnabled = enabled; + } + this.showToast(`Image watcher ${enabled ? 'enabled' : 'disabled'}`, 'success'); + } catch (err) { + this.showToast('Failed to toggle image watcher', 'error'); + } + } + async autoSaveRespawnConfig() { if (!this.editingSessionId) return; const config = { @@ -6834,6 +6855,7 @@ class ClaudemanApp { document.getElementById('appSettingsShowSubagents').checked = settings.showSubagents ?? true; document.getElementById('appSettingsSubagentTracking').checked = settings.subagentTrackingEnabled ?? true; document.getElementById('appSettingsSubagentActiveTabOnly').checked = settings.subagentActiveTabOnly ?? false; + document.getElementById('appSettingsImageWatcherEnabled').checked = settings.imageWatcherEnabled ?? true; // Claude CLI settings const claudeModeSelect = document.getElementById('appSettingsClaudeMode'); const allowedToolsRow = document.getElementById('allowedToolsRow'); @@ -6848,6 +6870,8 @@ class ClaudemanApp { const niceSettings = settings.nice || {}; document.getElementById('appSettingsNiceEnabled').checked = niceSettings.enabled ?? false; document.getElementById('appSettingsNiceValue').value = niceSettings.niceValue ?? 10; + // Model configuration (loaded from server) + this.loadModelConfigForSettings(); // Notification settings const notifPrefs = this.notificationManager?.preferences || {}; document.getElementById('appSettingsNotifEnabled').checked = notifPrefs.enabled ?? true; @@ -6906,6 +6930,7 @@ class ClaudemanApp { showSubagents: document.getElementById('appSettingsShowSubagents').checked, subagentTrackingEnabled: document.getElementById('appSettingsSubagentTracking').checked, subagentActiveTabOnly: document.getElementById('appSettingsSubagentActiveTabOnly').checked, + imageWatcherEnabled: document.getElementById('appSettingsImageWatcherEnabled').checked, // Claude CLI settings claudeMode: document.getElementById('appSettingsClaudeMode').value, allowedTools: document.getElementById('appSettingsAllowedTools').value.trim(), @@ -6947,6 +6972,10 @@ class ClaudemanApp { headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ ...settings, notificationPreferences: notifPrefsToSave }) }); + + // Save model configuration separately + await this.saveModelConfigFromSettings(); + this.showToast('Settings saved', 'success'); } catch (err) { // Server save failed but localStorage succeeded @@ -6956,6 +6985,71 @@ class ClaudemanApp { this.closeAppSettings(); } + // Load model configuration from server for the settings modal + async loadModelConfigForSettings() { + try { + const res = await fetch('/api/execution/model-config'); + const data = await res.json(); + if (data.success && data.data) { + const config = data.data; + // Default model + const defaultModelEl = document.getElementById('appSettingsDefaultModel'); + if (defaultModelEl) { + defaultModelEl.value = config.defaultModel || 'sonnet'; + } + // Show recommendations + const showRecsEl = document.getElementById('appSettingsShowModelRecommendations'); + if (showRecsEl) { + showRecsEl.checked = config.showRecommendations ?? true; + } + // Agent type overrides + const overrides = config.agentTypeOverrides || {}; + const exploreEl = document.getElementById('appSettingsModelExplore'); + const implementEl = document.getElementById('appSettingsModelImplement'); + const testEl = document.getElementById('appSettingsModelTest'); + const reviewEl = document.getElementById('appSettingsModelReview'); + if (exploreEl) exploreEl.value = overrides.explore || ''; + if (implementEl) implementEl.value = overrides.implement || ''; + if (testEl) testEl.value = overrides.test || ''; + if (reviewEl) reviewEl.value = overrides.review || ''; + } + } catch (err) { + console.warn('Failed to load model config:', err); + } + } + + // Save model configuration from settings modal to server + async saveModelConfigFromSettings() { + const defaultModelEl = document.getElementById('appSettingsDefaultModel'); + const showRecsEl = document.getElementById('appSettingsShowModelRecommendations'); + const exploreEl = document.getElementById('appSettingsModelExplore'); + const implementEl = document.getElementById('appSettingsModelImplement'); + const testEl = document.getElementById('appSettingsModelTest'); + const reviewEl = document.getElementById('appSettingsModelReview'); + + const agentTypeOverrides = {}; + if (exploreEl?.value) agentTypeOverrides.explore = exploreEl.value; + if (implementEl?.value) agentTypeOverrides.implement = implementEl.value; + if (testEl?.value) agentTypeOverrides.test = testEl.value; + if (reviewEl?.value) agentTypeOverrides.review = reviewEl.value; + + const config = { + defaultModel: defaultModelEl?.value || 'sonnet', + showRecommendations: showRecsEl?.checked ?? true, + agentTypeOverrides, + }; + + try { + await fetch('/api/execution/model-config', { + method: 'PUT', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify(config) + }); + } catch (err) { + console.warn('Failed to save model config:', err); + } + } + // Get the global Ralph tracker enabled setting isRalphTrackerEnabledByDefault() { const settings = this.loadAppSettingsFromStorage(); @@ -7067,6 +7161,40 @@ class ClaudemanApp { localStorage.setItem('claudeman-app-settings', JSON.stringify(settings)); } + async clearAllSubagents() { + const count = this.subagents.size; + if (count === 0) { + this.showToast('No subagents to clear', 'info'); + return; + } + + if (!confirm(`Clear all ${count} tracked subagent(s)? This removes them from the UI but does not affect running processes.`)) { + return; + } + + try { + const res = await fetch('/api/subagents', { method: 'DELETE' }); + const data = await res.json(); + if (data.success) { + // Clear local state + this.subagents.clear(); + this.subagentActivity.clear(); + this.subagentToolResults.clear(); + // Close any open subagent windows + this.cleanupAllFloatingWindows(); + // Update UI + this.renderSubagentPanel(); + this.renderMonitorSubagents(); + this.updateSubagentBadge(); + this.showToast(`Cleared ${data.data.cleared} subagent(s)`, 'success'); + } else { + this.showToast('Failed to clear subagents: ' + data.error, 'error'); + } + } catch (err) { + this.showToast('Failed to clear subagents', 'error'); + } + } + toggleSubagentsPanel() { const panel = document.getElementById('subagentsPanel'); const toggleBtn = document.getElementById('subagentsToggleBtn'); @@ -8713,6 +8841,9 @@ class ClaudemanApp { // Always update badge count this.updateSubagentBadge(); + // Always update monitor panel (even if subagent panel is hidden) + this.renderMonitorSubagents(); + // If panel is not visible, don't render content if (!this.subagentPanelVisible) { return; @@ -11002,6 +11133,49 @@ class ClaudemanApp { body.innerHTML = html; } + renderMonitorSubagents() { + const body = document.getElementById('monitorSubagentsBody'); + const stats = document.getElementById('monitorSubagentStats'); + if (!body) return; + + const subagents = Array.from(this.subagents.values()); + const activeCount = subagents.filter(s => s.status === 'active' || s.status === 'idle').length; + + if (stats) { + stats.textContent = `${subagents.length} tracked` + (activeCount > 0 ? `, ${activeCount} active` : ''); + } + + if (subagents.length === 0) { + body.innerHTML = '
No background agents
'; + return; + } + + let html = ''; + for (const agent of subagents) { + const statusClass = agent.status === 'active' ? 'active' : agent.status === 'idle' ? 'idle' : 'completed'; + const modelBadge = agent.modelShort ? `${agent.modelShort}` : ''; + const desc = agent.description ? this.escapeHtml(agent.description.substring(0, 40)) : agent.agentId; + + html += ` +
+ ${agent.status} +
+
${modelBadge} ${desc}
+
+ ID: ${agent.agentId} + ${agent.toolCallCount || 0} tools +
+
+
+ ${agent.status !== 'completed' ? `` : ''} +
+
+ `; + } + + body.innerHTML = html; + } + async killScreen(sessionId) { if (!confirm('Kill this screen session?')) return; diff --git a/src/web/public/index.html b/src/web/public/index.html index 2a88fd1c..7be37bc9 100644 --- a/src/web/public/index.html +++ b/src/web/public/index.html @@ -331,6 +331,16 @@
No screen sessions
+
+
+ Background Agents + 0 tracked + +
+
+
No background agents
+
+
+ +
+ + + Auto-popup new images (screenshots, etc.) in this session's directory +
@@ -685,6 +705,7 @@ @@ -763,6 +784,16 @@ + + +
Image Watcher
+
+ Enable Globally + +
@@ -805,6 +836,68 @@ Process priority (-20 to 19, higher = lower priority, default: 10) + +