mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 12:39:42 +02:00
refactor: simplify ralph wizard from 9 agents to 2
Remove execution layer that never actually controlled execution: - execution-bridge.ts (model param was ignored) - group-scheduler.ts (task-tool mode never used) - model-selector.ts (recommendations were display-only) - context-manager.ts (never invoked) - execution-limits.ts Remove redundant agent prompts (overlapping outputs): - requirements-analyst, architecture-planner, risk-analyst - testing-specialist, verification (merged into planner.ts) - execution-optimizer (output was ignored) - final-review (scores were cosmetic) Simplify plan-orchestrator.ts from 2400 LOC to 520 LOC: - Before: 9 agents, 6 phases, ~40-60 minutes - After: 2 agents (research + planner), ~18 minutes Remove /api/execution/* endpoints and related server code. Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
@@ -1,93 +0,0 @@
|
||||
/**
|
||||
* @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;
|
||||
@@ -1,333 +0,0 @@
|
||||
/**
|
||||
* @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<string, TrackedRefresh> = 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<ContextRefreshResult> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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;
|
||||
}
|
||||
@@ -1,846 +0,0 @@
|
||||
/**
|
||||
* @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<void>;
|
||||
/** 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<number, NodeJS.Timeout> = new Map();
|
||||
private _retryTimers: Map<string, NodeJS.Timeout> = new Map();
|
||||
private _runningTasks: Map<string, { startedAt: number; sessionId?: string }> = new Map();
|
||||
|
||||
private _history: ExecutionHistoryEntry[] = [];
|
||||
private _workingDir: string = process.cwd();
|
||||
|
||||
constructor(modelConfig?: Partial<ModelConfig>) {
|
||||
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 => {
|
||||
// Clear the group timeout timer to prevent memory leak
|
||||
const timer = this._groupTimeoutTimers.get(group.groupNumber);
|
||||
if (timer) {
|
||||
clearTimeout(timer);
|
||||
this._groupTimeoutTimers.delete(group.groupNumber);
|
||||
}
|
||||
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<ModelConfig>): 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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
task.retryCount++;
|
||||
const willRetry = task.retryCount < MAX_TASK_RETRIES;
|
||||
|
||||
this.emit('taskFailed', {
|
||||
taskId: task.id,
|
||||
groupNumber,
|
||||
error,
|
||||
willRetry,
|
||||
});
|
||||
|
||||
if (willRetry) {
|
||||
// Schedule retry with tracked timer
|
||||
const retryTimer = setTimeout(() => {
|
||||
this._retryTimers.delete(task.id);
|
||||
task.status = 'pending';
|
||||
task.error = undefined;
|
||||
this._runningTasks.delete(task.id);
|
||||
}, TASK_RETRY_DELAY_MS);
|
||||
this._retryTimers.set(task.id, retryTimer);
|
||||
} 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;
|
||||
|
||||
// Clear any pending retry timer for this task
|
||||
const retryTimer = this._retryTimers.get(taskId);
|
||||
if (retryTimer) {
|
||||
clearTimeout(retryTimer);
|
||||
this._retryTimers.delete(taskId);
|
||||
}
|
||||
|
||||
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();
|
||||
|
||||
// Clear any pending retry timers
|
||||
for (const timer of this._retryTimers.values()) {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
this._retryTimers.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<ModelConfig>): ExecutionBridge {
|
||||
if (!bridgeInstance) {
|
||||
bridgeInstance = new ExecutionBridge(modelConfig);
|
||||
}
|
||||
return bridgeInstance;
|
||||
}
|
||||
|
||||
/**
|
||||
* Reset the singleton (for testing).
|
||||
*/
|
||||
export function resetExecutionBridge(): void {
|
||||
if (bridgeInstance) {
|
||||
bridgeInstance.reset();
|
||||
}
|
||||
bridgeInstance = null;
|
||||
}
|
||||
@@ -1,551 +0,0 @@
|
||||
/**
|
||||
* @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<string, number> = 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<number, GroupTask[]>();
|
||||
|
||||
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<number>();
|
||||
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;
|
||||
}
|
||||
@@ -1,320 +0,0 @@
|
||||
/**
|
||||
* @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<Record<AgentType, ModelTier>>;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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<AgentType, ModelTier> = {
|
||||
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<ModelConfig>) {
|
||||
super();
|
||||
this._config = { ...createDefaultModelConfig(), ...config };
|
||||
}
|
||||
|
||||
/**
|
||||
* Get current configuration.
|
||||
*/
|
||||
get config(): ModelConfig {
|
||||
return { ...this._config };
|
||||
}
|
||||
|
||||
/**
|
||||
* Update configuration.
|
||||
*/
|
||||
updateConfig(config: Partial<ModelConfig>): 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
|
||||
default: {
|
||||
// Exhaustive check - if new ModelTier added, this will catch it
|
||||
const _exhaustive: never = model;
|
||||
console.warn(`[ModelSelector] Unknown model tier: ${_exhaustive}, using sonnet multiplier`);
|
||||
return 1.0;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ========== Singleton ==========
|
||||
|
||||
let selectorInstance: ModelSelector | null = null;
|
||||
|
||||
/**
|
||||
* Get or create the singleton ModelSelector instance.
|
||||
*/
|
||||
export function getModelSelector(config?: Partial<ModelConfig>): ModelSelector {
|
||||
if (!selectorInstance) {
|
||||
selectorInstance = new ModelSelector(config);
|
||||
}
|
||||
return selectorInstance;
|
||||
}
|
||||
|
||||
/**
|
||||
* Reset the singleton (for testing).
|
||||
*/
|
||||
export function resetModelSelector(): void {
|
||||
selectorInstance = null;
|
||||
}
|
||||
+243
-2172
File diff suppressed because it is too large
Load Diff
@@ -1,32 +0,0 @@
|
||||
/**
|
||||
* Architecture Planner Prompt
|
||||
*
|
||||
* Designs software component architecture for the task.
|
||||
*
|
||||
* Placeholders: {TASK}, {RESEARCH_CONTEXT}
|
||||
*/
|
||||
|
||||
export const ARCHITECTURE_PLANNER_PROMPT = `You are an Architecture Planner specializing in software component design.
|
||||
|
||||
## YOUR TASK
|
||||
Design the architecture for implementing this task:
|
||||
|
||||
## TASK DESCRIPTION
|
||||
{TASK}
|
||||
|
||||
{RESEARCH_CONTEXT}
|
||||
|
||||
## INSTRUCTIONS
|
||||
1. Identify all modules/components needed
|
||||
2. Define interfaces between components
|
||||
3. Specify data structures and types
|
||||
4. Note configuration and setup requirements
|
||||
5. Consider separation of concerns
|
||||
|
||||
## OUTPUT FORMAT
|
||||
Return ONLY a JSON array:
|
||||
[
|
||||
{"category": "module|interface|type|config|infrastructure", "content": "component description", "rationale": "why needed"}
|
||||
]
|
||||
|
||||
Generate 10-20 items. Think about the complete system architecture.`;
|
||||
@@ -1,134 +0,0 @@
|
||||
/**
|
||||
* Execution Optimizer Prompt
|
||||
*
|
||||
* Optimizes the plan for Claude Code execution with parallel groups,
|
||||
* agent types, model recommendations, and token estimates.
|
||||
*
|
||||
* Placeholders: {TASK}, {PLAN}
|
||||
*/
|
||||
|
||||
export const EXECUTION_OPTIMIZER_PROMPT = `You are a Claude Code Execution Optimizer. Your job is to analyze an implementation plan and optimize it for efficient execution using Claude Code's agent system.
|
||||
|
||||
## ORIGINAL TASK
|
||||
{TASK}
|
||||
|
||||
## CURRENT PLAN
|
||||
{PLAN}
|
||||
|
||||
## YOUR MISSION
|
||||
Analyze and enhance this plan for optimal Claude Code execution:
|
||||
|
||||
### 1. PARALLEL EXECUTION GROUPS
|
||||
Identify tasks that can run simultaneously in separate agents:
|
||||
- Tasks with NO dependencies between them
|
||||
- Tasks that modify DIFFERENT files
|
||||
- Tasks that read-only operations (exploration, analysis)
|
||||
- Assign a parallelGroup ID (e.g., "parallel-1", "parallel-2") to related tasks
|
||||
|
||||
### 2. AGENT TYPE RECOMMENDATIONS
|
||||
For each task, recommend the optimal Claude Code agent type:
|
||||
- "explore": For codebase exploration, finding files, understanding patterns
|
||||
- "implement": For writing new code, features, modifications
|
||||
- "test": For writing and running tests
|
||||
- "review": For code review, security analysis, best practices
|
||||
- "general": For mixed or unclear tasks
|
||||
|
||||
### 3. FRESH CONTEXT RECOMMENDATIONS
|
||||
Mark tasks that benefit from a fresh context (new conversation):
|
||||
- After large file modifications (>500 lines changed)
|
||||
- When switching between unrelated features
|
||||
- After test failures that need fresh analysis
|
||||
- When accumulated context might cause confusion
|
||||
|
||||
### 4. MODEL RECOMMENDATIONS
|
||||
Suggest the optimal model for each task:
|
||||
- "opus": Complex architecture, critical decisions, security review
|
||||
- "sonnet": Standard implementation, most coding tasks
|
||||
- "haiku": Quick exploration, simple searches, routine checks
|
||||
|
||||
### 5. FILE SCOPE ANALYSIS
|
||||
For each task, identify:
|
||||
- inputFiles: Files the task will need to READ
|
||||
- outputFiles: Files the task will CREATE or MODIFY
|
||||
|
||||
### 6. TOKEN ESTIMATION
|
||||
Estimate token usage for each task:
|
||||
- Small (exploration, simple changes): 5000-15000
|
||||
- Medium (feature implementation): 15000-50000
|
||||
- Large (complex features, refactoring): 50000-100000
|
||||
|
||||
## OUTPUT FORMAT
|
||||
Return ONLY a JSON object:
|
||||
{
|
||||
"optimizedPlan": [
|
||||
{
|
||||
"id": "P0-001",
|
||||
"content": "Explore existing auth patterns in codebase",
|
||||
"priority": "P0",
|
||||
"tddPhase": "setup",
|
||||
"verificationCriteria": "Documented auth patterns with file locations",
|
||||
"dependencies": [],
|
||||
"parallelGroup": "parallel-1",
|
||||
"agentType": "explore",
|
||||
"recommendedModel": "haiku",
|
||||
"requiresFreshContext": false,
|
||||
"estimatedTokens": 8000,
|
||||
"inputFiles": ["src/auth/**/*.ts", "src/middleware/*.ts"],
|
||||
"outputFiles": [],
|
||||
"executionNotes": "Quick exploration, can run alongside P0-002"
|
||||
},
|
||||
{
|
||||
"id": "P0-002",
|
||||
"content": "Explore test patterns and fixtures",
|
||||
"priority": "P0",
|
||||
"tddPhase": "setup",
|
||||
"verificationCriteria": "Understood test setup and conventions",
|
||||
"dependencies": [],
|
||||
"parallelGroup": "parallel-1",
|
||||
"agentType": "explore",
|
||||
"recommendedModel": "haiku",
|
||||
"requiresFreshContext": false,
|
||||
"estimatedTokens": 6000,
|
||||
"inputFiles": ["test/**/*.test.ts", "test/fixtures/**/*"],
|
||||
"outputFiles": [],
|
||||
"executionNotes": "Parallel with P0-001, different file scope"
|
||||
}
|
||||
],
|
||||
"parallelGroups": [
|
||||
{
|
||||
"id": "parallel-1",
|
||||
"tasks": ["P0-001", "P0-002"],
|
||||
"rationale": "Independent exploration tasks with no file overlap",
|
||||
"estimatedDuration": "2-3 minutes",
|
||||
"totalTokens": 14000
|
||||
}
|
||||
],
|
||||
"executionStrategy": {
|
||||
"totalParallelGroups": 3,
|
||||
"sequentialBlockers": ["P0-005 blocks all P1 tasks"],
|
||||
"freshContextPoints": ["After P0-005 (large refactor)", "After P1-003 (test failures)"],
|
||||
"estimatedTotalTokens": 150000,
|
||||
"estimatedAgentSpawns": 8,
|
||||
"criticalPath": ["P0-001", "P0-003", "P0-005", "P1-001"],
|
||||
"optimizationNotes": [
|
||||
"Group 1 saves ~3 min by parallelizing exploration",
|
||||
"Use haiku for 4 exploration tasks to reduce cost",
|
||||
"Fresh context after auth refactor prevents confusion"
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
CRITICAL REQUIREMENTS:
|
||||
1. Every task MUST have parallelGroup, agentType, recommendedModel
|
||||
2. Parallel groups MUST NOT have overlapping outputFiles
|
||||
3. Tasks in same parallelGroup MUST NOT depend on each other
|
||||
4. Preserve all existing task fields (id, content, priority, etc.)
|
||||
5. Add executionNotes explaining WHY this optimization
|
||||
|
||||
PARALLELIZATION GUIDELINES (BE CONSERVATIVE):
|
||||
- Only parallelize tasks when you are CERTAIN they have no file conflicts
|
||||
- Prefer sequential execution for complex or risky tasks
|
||||
- Limit parallel groups to 2-3 tasks maximum per group
|
||||
- When in doubt, keep tasks sequential - correctness over speed
|
||||
- Focus parallelization on exploration/read-only tasks, not implementations
|
||||
- Never parallelize tasks that might share state or side effects`;
|
||||
@@ -1,100 +0,0 @@
|
||||
/**
|
||||
* Final Review Expert Prompt
|
||||
*
|
||||
* Provides holistic analysis of the complete implementation plan
|
||||
* with scoring and improvement suggestions.
|
||||
*
|
||||
* Placeholders: {TASK}, {PLAN}
|
||||
*/
|
||||
|
||||
export const FINAL_REVIEW_PROMPT = `You are a Final Review Expert providing a holistic analysis of an implementation plan.
|
||||
|
||||
## ORIGINAL TASK
|
||||
{TASK}
|
||||
|
||||
## COMPLETE PLAN
|
||||
{PLAN}
|
||||
|
||||
## YOUR MISSION
|
||||
Review the ENTIRE plan from a high-level perspective. You have the bird's eye view.
|
||||
|
||||
### 1. LOGICAL FLOW ANALYSIS
|
||||
Check if the plan makes logical sense:
|
||||
- Does the order of tasks make sense?
|
||||
- Are there circular dependencies or impossible orderings?
|
||||
- Is there a clear progression from setup → implementation → testing → review?
|
||||
- Are foundation tasks (types, configs, setup) done before dependent tasks?
|
||||
|
||||
### 2. COMPLETENESS CHECK
|
||||
Verify nothing is missing:
|
||||
- Every implementation has a corresponding test?
|
||||
- Every test has clear verification criteria?
|
||||
- Error handling and edge cases are covered?
|
||||
- Setup and teardown steps are included?
|
||||
- Documentation tasks if needed?
|
||||
|
||||
### 3. COHERENCE VALIDATION
|
||||
Ensure the plan is internally consistent:
|
||||
- Do task descriptions match their dependencies?
|
||||
- Are file references consistent across tasks?
|
||||
- Do parallel groups actually make sense together?
|
||||
- Are priority levels justified?
|
||||
|
||||
### 4. FEASIBILITY ASSESSMENT
|
||||
Is this plan actually achievable?
|
||||
- Are any tasks too vague to execute?
|
||||
- Are there unrealistic expectations?
|
||||
- Are there hidden complexities not addressed?
|
||||
- Is the scope creep under control?
|
||||
|
||||
### 5. SUGGESTED IMPROVEMENTS
|
||||
Provide actionable fixes:
|
||||
- Tasks to add if missing
|
||||
- Tasks to split if too large
|
||||
- Tasks to merge if redundant
|
||||
- Order changes if needed
|
||||
- Clarifications needed
|
||||
|
||||
## OUTPUT FORMAT
|
||||
Return ONLY a JSON object:
|
||||
{
|
||||
"overallAssessment": "ready|needs-revision|major-issues",
|
||||
"logicScore": 0.85,
|
||||
"completenessScore": 0.90,
|
||||
"coherenceScore": 0.88,
|
||||
"feasibilityScore": 0.82,
|
||||
"overallScore": 0.86,
|
||||
"summary": "Brief 2-3 sentence summary of the plan quality",
|
||||
"logicIssues": [
|
||||
{"severity": "warning|error", "issue": "Description", "affectedTasks": ["P0-001"], "suggestion": "How to fix"}
|
||||
],
|
||||
"missingTasks": [
|
||||
{"content": "Add database migration script", "reason": "Schema changes require migration", "insertAfter": "P0-002", "priority": "P0"}
|
||||
],
|
||||
"tasksToSplit": [
|
||||
{"taskId": "P1-005", "reason": "Too complex", "splitInto": ["Implement auth logic", "Add session management"]}
|
||||
],
|
||||
"tasksToMerge": [
|
||||
{"taskIds": ["P2-001", "P2-002"], "reason": "Redundant", "mergedContent": "Combined task description"}
|
||||
],
|
||||
"orderChanges": [
|
||||
{"taskId": "P0-003", "currentPosition": 3, "suggestedPosition": 1, "reason": "Should run earlier"}
|
||||
],
|
||||
"clarificationsNeeded": [
|
||||
{"taskId": "P1-002", "issue": "Unclear which API endpoint", "question": "Is this REST or GraphQL?"}
|
||||
],
|
||||
"finalRecommendations": [
|
||||
"Start with P0 tasks in sequence for stable foundation",
|
||||
"Consider adding integration tests after P1-004",
|
||||
"Review security implications of auth changes"
|
||||
]
|
||||
}
|
||||
|
||||
SCORING GUIDELINES:
|
||||
- 0.9+: Excellent, ready to execute
|
||||
- 0.8-0.9: Good, minor tweaks recommended
|
||||
- 0.7-0.8: Acceptable, some issues to address
|
||||
- 0.6-0.7: Needs revision before execution
|
||||
- <0.6: Major issues, significant rework needed
|
||||
|
||||
Be thorough but constructive. The goal is to catch issues before execution, not to criticize.`;
|
||||
@@ -6,11 +6,5 @@
|
||||
*/
|
||||
|
||||
export { RESEARCH_AGENT_PROMPT } from './research-agent.js';
|
||||
export { REQUIREMENTS_ANALYST_PROMPT } from './requirements-analyst.js';
|
||||
export { ARCHITECTURE_PLANNER_PROMPT } from './architecture-planner.js';
|
||||
export { TESTING_SPECIALIST_PROMPT } from './testing-specialist.js';
|
||||
export { RISK_ANALYST_PROMPT } from './risk-analyst.js';
|
||||
export { PLANNER_PROMPT } from './planner.js';
|
||||
export { CODE_REVIEWER_PROMPT } from './code-reviewer.js';
|
||||
export { VERIFICATION_PROMPT } from './verification.js';
|
||||
export { EXECUTION_OPTIMIZER_PROMPT } from './execution-optimizer.js';
|
||||
export { FINAL_REVIEW_PROMPT } from './final-review.js';
|
||||
|
||||
@@ -0,0 +1,83 @@
|
||||
/**
|
||||
* Planner Prompt - Single agent for TDD plan generation
|
||||
*
|
||||
* Combines what was previously 5 separate agents:
|
||||
* - Requirements Analyst (redundant)
|
||||
* - Architecture Planner (redundant)
|
||||
* - Testing Specialist (kept - TDD focus)
|
||||
* - Risk Analyst (redundant)
|
||||
* - Verification Expert (kept - structure)
|
||||
*
|
||||
* Placeholders: {TASK}, {RESEARCH_CONTEXT}
|
||||
*/
|
||||
|
||||
export const PLANNER_PROMPT = `You are a TDD Plan Generator. Create a complete implementation plan with test-first approach.
|
||||
|
||||
## TASK DESCRIPTION
|
||||
{TASK}
|
||||
|
||||
{RESEARCH_CONTEXT}
|
||||
|
||||
## YOUR MISSION
|
||||
Generate a complete TDD implementation plan with:
|
||||
1. Tests BEFORE implementations (red-green-refactor)
|
||||
2. Review tasks AFTER implementations
|
||||
3. Clear priorities (P0=blocking, P1=required, P2=polish)
|
||||
4. Dependencies between tasks
|
||||
|
||||
## TDD CYCLE
|
||||
For each feature:
|
||||
1. Write failing test first
|
||||
2. Implement to make test pass
|
||||
3. Review implementation
|
||||
|
||||
## PRIORITY GUIDELINES
|
||||
- P0: Foundation, types, project setup, blocking dependencies
|
||||
- P1: Core features, main implementation, error handling
|
||||
- P2: Polish, optimization, documentation
|
||||
|
||||
## OUTPUT FORMAT
|
||||
Return ONLY a JSON object:
|
||||
{
|
||||
"items": [
|
||||
{
|
||||
"id": "P0-001",
|
||||
"content": "Write failing test for user authentication endpoint",
|
||||
"priority": "P0",
|
||||
"tddPhase": "test",
|
||||
"verificationCriteria": "Test file exists, test fails with 'not implemented'",
|
||||
"testCommand": "npm test -- --grep='auth'",
|
||||
"dependencies": []
|
||||
},
|
||||
{
|
||||
"id": "P0-002",
|
||||
"content": "Implement user authentication handler",
|
||||
"priority": "P0",
|
||||
"tddPhase": "impl",
|
||||
"verificationCriteria": "npm test -- --grep='auth' passes",
|
||||
"pairedWith": "P0-001",
|
||||
"dependencies": ["P0-001"]
|
||||
},
|
||||
{
|
||||
"id": "P0-003",
|
||||
"content": "Review auth implementation for security",
|
||||
"priority": "P0",
|
||||
"tddPhase": "review",
|
||||
"verificationCriteria": "No security issues, follows best practices",
|
||||
"reviewChecklist": ["Input validation", "XSS prevention", "Error handling"],
|
||||
"pairedWith": "P0-002",
|
||||
"dependencies": ["P0-002"]
|
||||
}
|
||||
],
|
||||
"gaps": ["any missing requirements noted"],
|
||||
"warnings": ["any concerns or risks identified"]
|
||||
}
|
||||
|
||||
CRITICAL REQUIREMENTS:
|
||||
1. Every implementation MUST have a paired test task that comes BEFORE it
|
||||
2. Every implementation MUST have a review task that comes AFTER it
|
||||
3. Use sequential IDs: P0-001, P0-002, P1-001, etc.
|
||||
4. verificationCriteria must be SPECIFIC and observable
|
||||
5. Dependencies must form a valid DAG (no cycles)
|
||||
|
||||
Generate 15-40 items covering the complete implementation.`;
|
||||
@@ -1,31 +0,0 @@
|
||||
/**
|
||||
* Requirements Analyst Prompt
|
||||
*
|
||||
* Extracts explicit and implicit requirements from task descriptions.
|
||||
*
|
||||
* Placeholders: {TASK}, {RESEARCH_CONTEXT}
|
||||
*/
|
||||
|
||||
export const REQUIREMENTS_ANALYST_PROMPT = `You are a Requirements Analyst specializing in extracting all requirements from task descriptions.
|
||||
|
||||
## YOUR TASK
|
||||
Analyze the following task and extract ALL requirements (explicit and implicit):
|
||||
|
||||
## TASK DESCRIPTION
|
||||
{TASK}
|
||||
|
||||
{RESEARCH_CONTEXT}
|
||||
|
||||
## INSTRUCTIONS
|
||||
1. Identify explicit requirements (directly stated)
|
||||
2. Infer implicit requirements (unstated but necessary)
|
||||
3. Note any assumptions that should be validated
|
||||
4. Consider non-functional requirements (performance, security, usability)
|
||||
|
||||
## OUTPUT FORMAT
|
||||
Return ONLY a JSON array:
|
||||
[
|
||||
{"category": "functional|non-functional|constraint|assumption", "content": "requirement description", "rationale": "why this is needed"}
|
||||
]
|
||||
|
||||
Generate 8-15 items. Be thorough - missing requirements cause project failures.`;
|
||||
@@ -1,32 +0,0 @@
|
||||
/**
|
||||
* Risk Analyst Prompt
|
||||
*
|
||||
* Identifies potential issues, edge cases, and blockers.
|
||||
*
|
||||
* Placeholders: {TASK}, {RESEARCH_CONTEXT}
|
||||
*/
|
||||
|
||||
export const RISK_ANALYST_PROMPT = `You are a Risk Analyst identifying potential issues and blockers.
|
||||
|
||||
## YOUR TASK
|
||||
Identify risks and edge cases for this task:
|
||||
|
||||
## TASK DESCRIPTION
|
||||
{TASK}
|
||||
|
||||
{RESEARCH_CONTEXT}
|
||||
|
||||
## INSTRUCTIONS
|
||||
1. Identify potential failure points
|
||||
2. Note edge cases that could cause bugs
|
||||
3. Consider security vulnerabilities
|
||||
4. Flag performance concerns
|
||||
5. Identify dependencies that could block progress
|
||||
|
||||
## OUTPUT FORMAT
|
||||
Return ONLY a JSON array:
|
||||
[
|
||||
{"category": "failure|edge-case|security|performance|dependency", "content": "risk description", "rationale": "mitigation approach"}
|
||||
]
|
||||
|
||||
Generate 8-15 items. Being proactive about risks prevents surprises.`;
|
||||
@@ -1,76 +0,0 @@
|
||||
/**
|
||||
* Testing Specialist Prompt
|
||||
*
|
||||
* Designs comprehensive, realistic test coverage with TDD approach.
|
||||
*
|
||||
* Placeholders: {TASK}, {RESEARCH_CONTEXT}
|
||||
*/
|
||||
|
||||
export const TESTING_SPECIALIST_PROMPT = `You are a TDD Specialist designing a comprehensive, REALISTIC test strategy.
|
||||
|
||||
## YOUR TASK
|
||||
Design detailed, executable test coverage for this task:
|
||||
|
||||
## TASK DESCRIPTION
|
||||
{TASK}
|
||||
|
||||
{RESEARCH_CONTEXT}
|
||||
|
||||
## INSTRUCTIONS
|
||||
Create REALISTIC tests that would actually run in a real codebase:
|
||||
|
||||
### 1. Unit Tests (test individual functions/methods in isolation)
|
||||
- Mock external dependencies (databases, APIs, file system)
|
||||
- Test pure logic with specific input/output examples
|
||||
- Include exact assertion values, not placeholders
|
||||
|
||||
### 2. Integration Tests (test component interactions)
|
||||
- Test API endpoints with realistic request/response bodies
|
||||
- Test database operations with actual schema
|
||||
- Test service-to-service communication
|
||||
|
||||
### 3. Edge Cases & Boundary Tests
|
||||
- Empty inputs, null values, undefined
|
||||
- Maximum/minimum values, overflow conditions
|
||||
- Unicode, special characters, injection attempts
|
||||
- Concurrent access, race conditions
|
||||
|
||||
### 4. Error Scenario Tests
|
||||
- Network failures, timeouts, connection refused
|
||||
- Invalid input validation with specific error messages
|
||||
- Authorization failures, permission denied
|
||||
- Resource not found, conflict states
|
||||
|
||||
### 5. Performance & Load Tests (where applicable)
|
||||
- Response time thresholds
|
||||
- Memory usage limits
|
||||
- Concurrent user handling
|
||||
|
||||
## REALISTIC TEST EXAMPLE
|
||||
BAD: "Test user login" (too vague)
|
||||
GOOD: "Test POST /api/auth/login with valid email 'test@example.com' and password 'ValidPass123!' returns 200 with JWT token containing userId and exp claims, sets httpOnly cookie 'session'"
|
||||
|
||||
## OUTPUT FORMAT
|
||||
Return ONLY a JSON array:
|
||||
[
|
||||
{
|
||||
"category": "unit|integration|edge-case|error|e2e|performance",
|
||||
"content": "Test POST /api/users with email 'new@test.com' creates user and returns 201 with {id, email, createdAt}",
|
||||
"rationale": "Validates user creation happy path with all required response fields",
|
||||
"verificationCriteria": "Response status 201, body contains id (uuid), email matches input, createdAt is valid ISO timestamp",
|
||||
"testCommand": "npm test -- --grep='POST /api/users creates user'",
|
||||
"testSetup": "Clear users table, seed with test data",
|
||||
"testTeardown": "Delete created test user",
|
||||
"pairedImpl": "Implement POST /api/users endpoint with validation and database insert",
|
||||
"mockDependencies": ["database connection", "email service"],
|
||||
"assertionDetails": ["status === 201", "body.id matches UUID regex", "body.email === 'new@test.com'"]
|
||||
}
|
||||
]
|
||||
|
||||
CRITICAL REQUIREMENTS:
|
||||
- verificationCriteria: SPECIFIC observable outcomes with exact values
|
||||
- testCommand: Actual runnable command (npm test, pytest, vitest, etc.)
|
||||
- pairedImpl: The exact implementation step this test validates
|
||||
- assertionDetails: List of specific assertions to make
|
||||
|
||||
Generate 15-30 detailed test items. Tests MUST be specific enough to implement directly.`;
|
||||
@@ -1,96 +0,0 @@
|
||||
/**
|
||||
* Verification Expert Prompt
|
||||
*
|
||||
* Reviews and enhances the synthesized plan with priorities,
|
||||
* verification criteria, and TDD pairing.
|
||||
*
|
||||
* Placeholders: {TASK}, {PLAN}
|
||||
*/
|
||||
|
||||
export const VERIFICATION_PROMPT = `You are a Plan Verification Expert reviewing an implementation plan for completeness and quality.
|
||||
|
||||
## ORIGINAL TASK
|
||||
{TASK}
|
||||
|
||||
## SYNTHESIZED PLAN (from multiple analysis subagents)
|
||||
{PLAN}
|
||||
|
||||
## YOUR MISSION
|
||||
Review and enhance this plan:
|
||||
1. Assign priorities (P0=critical/blocking, P1=required, P2=enhancement)
|
||||
2. Add verification criteria to EVERY task (how to know it's done)
|
||||
3. Pair test tasks with implementation tasks (TDD cycle)
|
||||
4. Add dependencies where one task blocks another
|
||||
5. Identify gaps and calculate quality score
|
||||
|
||||
## PRIORITY GUIDELINES
|
||||
- P0: Foundation tasks, type definitions, project setup, blocking dependencies
|
||||
- P1: Core implementation, tests, main features, error handling
|
||||
- P2: Polish, optimization, documentation, nice-to-have features
|
||||
|
||||
## TDD + REVIEW CYCLE RULES
|
||||
The complete cycle is: test → impl → review
|
||||
- Every implementation task should have a corresponding test task AND review task
|
||||
- Test task comes BEFORE its paired implementation task
|
||||
- Review task comes AFTER the implementation it reviews
|
||||
- Use "pairedWith" to link test ↔ implementation ↔ review
|
||||
- Verification criteria should reference test results where applicable
|
||||
|
||||
## REVIEW TASK REQUIREMENTS
|
||||
After EVERY implementation task, add a review task that checks:
|
||||
- Best practices for the language/framework
|
||||
- Security vulnerabilities (OWASP top 10)
|
||||
- Performance concerns
|
||||
- Error handling completeness
|
||||
- Code quality (DRY, SOLID, readability)
|
||||
|
||||
## OUTPUT FORMAT
|
||||
Return ONLY a JSON object:
|
||||
{
|
||||
"validatedPlan": [
|
||||
{
|
||||
"id": "P0-001",
|
||||
"content": "Write failing test for user authentication",
|
||||
"priority": "P0",
|
||||
"tddPhase": "test",
|
||||
"verificationCriteria": "Test file exists, test fails with 'not implemented'",
|
||||
"testCommand": "npm test -- --grep='auth'",
|
||||
"pairedWith": "P0-002",
|
||||
"dependencies": [],
|
||||
"complexity": "low"
|
||||
},
|
||||
{
|
||||
"id": "P0-002",
|
||||
"content": "Implement user authentication handler",
|
||||
"priority": "P0",
|
||||
"tddPhase": "impl",
|
||||
"verificationCriteria": "npm test -- --grep='auth' passes",
|
||||
"pairedWith": "P0-001",
|
||||
"dependencies": ["P0-001"],
|
||||
"complexity": "medium"
|
||||
},
|
||||
{
|
||||
"id": "P0-003",
|
||||
"content": "Review auth implementation for security and best practices",
|
||||
"priority": "P0",
|
||||
"tddPhase": "review",
|
||||
"verificationCriteria": "No security issues found, follows TypeScript best practices",
|
||||
"reviewChecklist": ["Input validation", "XSS prevention", "Session security", "Error handling"],
|
||||
"pairedWith": "P0-002",
|
||||
"dependencies": ["P0-002"],
|
||||
"complexity": "low"
|
||||
}
|
||||
],
|
||||
"gaps": ["missing requirement 1", "missing test coverage for X"],
|
||||
"warnings": ["consider Y before Z", "potential issue with..."],
|
||||
"qualityScore": 0.85
|
||||
}
|
||||
|
||||
CRITICAL REQUIREMENTS:
|
||||
1. EVERY task MUST have verificationCriteria (how to verify completion)
|
||||
2. Implementation tasks MUST have a paired test task AND a review task
|
||||
3. Review tasks MUST have a reviewChecklist with specific items to check
|
||||
4. Dependencies must form a valid DAG (no cycles)
|
||||
5. Use sequential IDs: P0-001, P0-002, P0-003, P1-001, etc.
|
||||
|
||||
Be critical but constructive. A thorough review catches issues that tests miss.`;
|
||||
+2
-33
@@ -1255,37 +1255,6 @@ export interface ImageDetectedEvent {
|
||||
size: number;
|
||||
}
|
||||
|
||||
// ========== Execution Bridge Re-exports ==========
|
||||
// ========== Plan Orchestrator 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';
|
||||
export type { PlanItem } from './plan-orchestrator.js';
|
||||
|
||||
@@ -34,13 +34,6 @@ import { v4 as uuidv4 } from 'uuid';
|
||||
import { createRequire } from 'node:module';
|
||||
import { RunSummaryTracker } from '../run-summary.js';
|
||||
import { PlanOrchestrator, type DetailedPlanResult } from '../plan-orchestrator.js';
|
||||
import {
|
||||
ExecutionBridge,
|
||||
getExecutionBridge,
|
||||
type PlanItem as ExecutionPlanItem,
|
||||
type AgentSpawner,
|
||||
} from '../execution-bridge.js';
|
||||
import { type ModelConfig } from '../model-selector.js';
|
||||
|
||||
// Load version from package.json
|
||||
const require = createRequire(import.meta.url);
|
||||
@@ -314,8 +307,6 @@ export class WebServer extends EventEmitter {
|
||||
private pendingRespawnStarts: Map<string, NodeJS.Timeout> = new Map();
|
||||
// Active plan orchestrators (for cancellation via API)
|
||||
private activePlanOrchestrators: Map<string, PlanOrchestrator> = new Map();
|
||||
// Execution bridge for parallel task execution
|
||||
private executionBridge: ExecutionBridge;
|
||||
// Grace period before starting restored respawn controllers (2 minutes)
|
||||
private static readonly RESPAWN_RESTORE_GRACE_PERIOD_MS = 2 * 60 * 1000;
|
||||
|
||||
@@ -346,10 +337,6 @@ export class WebServer extends EventEmitter {
|
||||
this.broadcast('screen:statsUpdated', screens);
|
||||
});
|
||||
|
||||
// Initialize execution bridge with model config from settings
|
||||
this.executionBridge = getExecutionBridge(this.loadModelConfig());
|
||||
this.setupExecutionBridgeListeners();
|
||||
|
||||
// Set up subagent watcher listeners
|
||||
this.setupSubagentWatcherListeners();
|
||||
|
||||
@@ -3045,148 +3032,6 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
|
||||
return this.getSystemStats();
|
||||
});
|
||||
|
||||
// ========== Execution Bridge Endpoints ==========
|
||||
|
||||
// Get execution status
|
||||
this.app.get('/api/execution/status', async () => {
|
||||
return {
|
||||
success: true,
|
||||
data: {
|
||||
status: this.executionBridge.status,
|
||||
progress: this.executionBridge.getProgress(),
|
||||
},
|
||||
};
|
||||
});
|
||||
|
||||
// Get execution schedule (the current plan)
|
||||
this.app.get('/api/execution/schedule', async () => {
|
||||
const schedule = this.executionBridge['_scheduler'].schedule;
|
||||
return { success: true, data: schedule };
|
||||
});
|
||||
|
||||
// Load a plan for execution
|
||||
// Accepts either ExecutionPlanItem format (id, title, description) or
|
||||
// PlanOrchestrator PlanItem format (id, content) and converts as needed
|
||||
this.app.post('/api/execution/load', async (req) => {
|
||||
const { items, workingDir } = req.body as {
|
||||
items: Array<{
|
||||
id?: string;
|
||||
title?: string;
|
||||
description?: string;
|
||||
content?: string;
|
||||
parallelGroup?: number;
|
||||
agentType?: string;
|
||||
recommendedModel?: string;
|
||||
requiresFreshContext?: boolean;
|
||||
estimatedTokens?: number;
|
||||
inputFiles?: string[];
|
||||
outputFiles?: string[];
|
||||
dependencies?: string[];
|
||||
}>;
|
||||
workingDir?: string;
|
||||
};
|
||||
if (!items || !Array.isArray(items)) {
|
||||
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'items array is required');
|
||||
}
|
||||
try {
|
||||
if (workingDir) {
|
||||
this.executionBridge.setWorkingDir(workingDir);
|
||||
}
|
||||
// Convert to ExecutionPlanItem format
|
||||
const execItems: ExecutionPlanItem[] = items.map((item, idx) => ({
|
||||
id: item.id || `task-${idx}`,
|
||||
title: item.title || item.content?.slice(0, 50) || `Task ${idx + 1}`,
|
||||
description: item.description || item.content || '',
|
||||
parallelGroup: item.parallelGroup,
|
||||
agentType: item.agentType,
|
||||
recommendedModel: item.recommendedModel,
|
||||
requiresFreshContext: item.requiresFreshContext,
|
||||
estimatedTokens: item.estimatedTokens,
|
||||
inputFiles: item.inputFiles,
|
||||
outputFiles: item.outputFiles,
|
||||
dependencies: item.dependencies,
|
||||
}));
|
||||
const schedule = this.executionBridge.loadPlan(execItems);
|
||||
return { success: true, data: schedule };
|
||||
} catch (err) {
|
||||
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, getErrorMessage(err));
|
||||
}
|
||||
});
|
||||
|
||||
// Start execution
|
||||
this.app.post('/api/execution/start', async () => {
|
||||
try {
|
||||
await this.executionBridge.start();
|
||||
return { success: true };
|
||||
} catch (err) {
|
||||
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, getErrorMessage(err));
|
||||
}
|
||||
});
|
||||
|
||||
// Pause execution
|
||||
this.app.post('/api/execution/pause', async () => {
|
||||
this.executionBridge.pause();
|
||||
return { success: true };
|
||||
});
|
||||
|
||||
// Resume execution
|
||||
this.app.post('/api/execution/resume', async () => {
|
||||
this.executionBridge.resume();
|
||||
return { success: true };
|
||||
});
|
||||
|
||||
// Cancel execution
|
||||
this.app.post('/api/execution/cancel', async (req) => {
|
||||
const { reason } = (req.body as { reason?: string }) || {};
|
||||
await this.executionBridge.cancel(reason || 'Cancelled via API');
|
||||
return { success: true };
|
||||
});
|
||||
|
||||
// Reset execution bridge
|
||||
this.app.post('/api/execution/reset', async () => {
|
||||
this.executionBridge.reset();
|
||||
return { success: true };
|
||||
});
|
||||
|
||||
// Get execution history
|
||||
this.app.get('/api/execution/history', async () => {
|
||||
return { success: true, data: this.executionBridge.getHistory() };
|
||||
});
|
||||
|
||||
// Get model configuration
|
||||
this.app.get('/api/execution/model-config', async () => {
|
||||
return { success: true, data: this.executionBridge.getModelConfig() };
|
||||
});
|
||||
|
||||
// Update model configuration
|
||||
this.app.put('/api/execution/model-config', async (req) => {
|
||||
const config = req.body as Partial<ModelConfig>;
|
||||
try {
|
||||
this.executionBridge.updateModelConfig(config);
|
||||
const fullConfig = this.executionBridge.getModelConfig();
|
||||
this.saveModelConfig(fullConfig);
|
||||
this.broadcast('execution:modelConfigUpdated', fullConfig);
|
||||
return { success: true, data: fullConfig };
|
||||
} catch (err) {
|
||||
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, getErrorMessage(err));
|
||||
}
|
||||
});
|
||||
|
||||
// Mark task complete (external signal)
|
||||
this.app.post('/api/execution/tasks/:taskId/complete', async (req) => {
|
||||
const { taskId } = req.params as { taskId: string };
|
||||
this.executionBridge.markTaskComplete(taskId);
|
||||
return { success: true };
|
||||
});
|
||||
|
||||
// Mark task failed (external signal)
|
||||
this.app.post('/api/execution/tasks/:taskId/failed', async (req) => {
|
||||
const { taskId } = req.params as { taskId: string };
|
||||
const { error } = (req.body as { error?: string }) || {};
|
||||
this.executionBridge.markTaskFailed(taskId, error || 'Failed via API');
|
||||
return { success: true };
|
||||
});
|
||||
|
||||
// ========== Subagent Monitoring (Claude Code Background Agents) ==========
|
||||
|
||||
// List all known subagents
|
||||
@@ -3910,193 +3755,6 @@ NOW: Generate the implementation plan for the task above. Think step by step.`;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Load model configuration from settings file.
|
||||
*/
|
||||
private loadModelConfig(): Partial<ModelConfig> {
|
||||
const settingsPath = join(homedir(), '.claudeman', 'settings.json');
|
||||
try {
|
||||
if (existsSync(settingsPath)) {
|
||||
const content = readFileSync(settingsPath, 'utf-8');
|
||||
const settings = JSON.parse(content);
|
||||
return settings.modelConfig || {};
|
||||
}
|
||||
} catch (err) {
|
||||
console.error('Failed to load model config:', getErrorMessage(err));
|
||||
}
|
||||
return {};
|
||||
}
|
||||
|
||||
/**
|
||||
* Save model configuration to settings file.
|
||||
*/
|
||||
private saveModelConfig(config: ModelConfig): void {
|
||||
const settingsPath = join(homedir(), '.claudeman', 'settings.json');
|
||||
try {
|
||||
const dir = dirname(settingsPath);
|
||||
if (!existsSync(dir)) {
|
||||
mkdirSync(dir, { recursive: true });
|
||||
}
|
||||
let settings: Record<string, unknown> = {};
|
||||
if (existsSync(settingsPath)) {
|
||||
settings = JSON.parse(readFileSync(settingsPath, 'utf-8'));
|
||||
}
|
||||
settings.modelConfig = config;
|
||||
writeFileSync(settingsPath, JSON.stringify(settings, null, 2));
|
||||
} catch (err) {
|
||||
console.error('Failed to save model config:', getErrorMessage(err));
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Set up event listeners for execution bridge.
|
||||
* Broadcasts execution progress to SSE clients.
|
||||
*/
|
||||
private setupExecutionBridgeListeners(): void {
|
||||
// Create agent spawner interface
|
||||
const spawner: AgentSpawner = {
|
||||
spawnAgentWithModel: async (taskId, workingDir, prompt, _model, _options) => {
|
||||
// Note: model and options are reserved for future model selection integration
|
||||
const globalNice = this.getGlobalNiceConfig();
|
||||
const session = new Session({
|
||||
workingDir,
|
||||
screenManager: this.screenManager,
|
||||
useScreen: true,
|
||||
mode: 'claude',
|
||||
name: `exec:${taskId}`,
|
||||
niceConfig: globalNice,
|
||||
});
|
||||
|
||||
this.sessions.set(session.id, session);
|
||||
this.store.incrementSessionsCreated();
|
||||
this.setupSessionListeners(session);
|
||||
|
||||
await session.startInteractive();
|
||||
this.broadcast('session:created', session.toDetailedState());
|
||||
this.broadcast('session:interactive', { id: session.id });
|
||||
this.persistSessionState(session);
|
||||
|
||||
// Configure ralph tracker for completion detection
|
||||
session.ralphTracker.enable();
|
||||
|
||||
// Set up completion listener for execution bridge
|
||||
session.on('ralphCompletionDetected', () => {
|
||||
this.executionBridge.markTaskComplete(taskId);
|
||||
});
|
||||
|
||||
// Send the prompt
|
||||
setTimeout(() => {
|
||||
session.writeViaScreen(prompt + '\r');
|
||||
}, 2000);
|
||||
|
||||
return { sessionId: session.id };
|
||||
},
|
||||
useTaskTool: async (_sessionId, _taskId, _prompt, _model) => {
|
||||
// Task tool mode - not yet implemented, falls back to session mode
|
||||
throw new Error('Task tool mode not yet implemented');
|
||||
},
|
||||
isTaskComplete: (taskId) => {
|
||||
// Check if task is marked complete in scheduler
|
||||
const schedule = this.executionBridge['_scheduler'].schedule;
|
||||
if (!schedule) return false;
|
||||
for (const group of schedule.groups) {
|
||||
const task = group.tasks.find(t => t.id === taskId);
|
||||
if (task) return task.status === 'completed';
|
||||
}
|
||||
return false;
|
||||
},
|
||||
getTaskResult: (taskId) => {
|
||||
const schedule = this.executionBridge['_scheduler'].schedule;
|
||||
if (!schedule) return null;
|
||||
for (const group of schedule.groups) {
|
||||
const task = group.tasks.find(t => t.id === taskId);
|
||||
if (task) {
|
||||
if (task.status === 'completed') {
|
||||
return { success: true };
|
||||
} else if (task.status === 'failed') {
|
||||
return { success: false, error: task.error };
|
||||
}
|
||||
}
|
||||
}
|
||||
return null;
|
||||
},
|
||||
};
|
||||
|
||||
this.executionBridge.setAgentSpawner(spawner);
|
||||
|
||||
// Set up session writer for context management
|
||||
this.executionBridge.setSessionWriter({
|
||||
writeToSession: (sessionId, text) => {
|
||||
const session = this.sessions.get(sessionId);
|
||||
if (session) {
|
||||
session.writeViaScreen(text);
|
||||
}
|
||||
},
|
||||
createSession: async (workingDir, name) => {
|
||||
const globalNice = this.getGlobalNiceConfig();
|
||||
const session = new Session({
|
||||
workingDir,
|
||||
screenManager: this.screenManager,
|
||||
useScreen: true,
|
||||
mode: 'claude',
|
||||
name: name || 'exec-context',
|
||||
niceConfig: globalNice,
|
||||
});
|
||||
this.sessions.set(session.id, session);
|
||||
this.store.incrementSessionsCreated();
|
||||
this.setupSessionListeners(session);
|
||||
await session.startInteractive();
|
||||
this.broadcast('session:created', session.toDetailedState());
|
||||
this.persistSessionState(session);
|
||||
return { sessionId: session.id };
|
||||
},
|
||||
});
|
||||
|
||||
// Forward execution bridge events as SSE broadcasts
|
||||
this.executionBridge.on('planLoaded', (schedule) => {
|
||||
this.broadcast('execution:planLoaded', { schedule });
|
||||
});
|
||||
this.executionBridge.on('started', () => {
|
||||
this.broadcast('execution:started', {});
|
||||
});
|
||||
this.executionBridge.on('paused', () => {
|
||||
this.broadcast('execution:paused', {});
|
||||
});
|
||||
this.executionBridge.on('resumed', () => {
|
||||
this.broadcast('execution:resumed', {});
|
||||
});
|
||||
this.executionBridge.on('completed', (data) => {
|
||||
this.broadcast('execution:completed', data);
|
||||
});
|
||||
this.executionBridge.on('cancelled', (reason) => {
|
||||
this.broadcast('execution:cancelled', { reason });
|
||||
});
|
||||
this.executionBridge.on('groupStarted', (data) => {
|
||||
this.broadcast('execution:groupStarted', data);
|
||||
});
|
||||
this.executionBridge.on('groupCompleted', (data) => {
|
||||
this.broadcast('execution:groupCompleted', data);
|
||||
});
|
||||
this.executionBridge.on('taskAssigned', (data) => {
|
||||
this.broadcast('execution:taskAssigned', data);
|
||||
});
|
||||
this.executionBridge.on('taskCompleted', (data) => {
|
||||
this.broadcast('execution:taskCompleted', data);
|
||||
});
|
||||
this.executionBridge.on('taskFailed', (data) => {
|
||||
this.broadcast('execution:taskFailed', data);
|
||||
});
|
||||
this.executionBridge.on('freshContext', (data) => {
|
||||
this.broadcast('execution:freshContext', data);
|
||||
});
|
||||
this.executionBridge.on('modelSelected', (data) => {
|
||||
this.broadcast('execution:modelSelected', data);
|
||||
});
|
||||
this.executionBridge.on('progress', (progress) => {
|
||||
this.broadcast('execution:progress', progress);
|
||||
});
|
||||
}
|
||||
|
||||
private setupTimedRespawn(sessionId: string, durationMinutes: number): void {
|
||||
// Clear existing timer if any
|
||||
const existing = this.respawnTimers.get(sessionId);
|
||||
|
||||
@@ -1,270 +0,0 @@
|
||||
/**
|
||||
* @fileoverview Tests for the Execution Bridge system.
|
||||
*
|
||||
* Tests the core functionality of:
|
||||
* - ModelSelector
|
||||
* - GroupScheduler
|
||||
* - ExecutionBridge
|
||||
*
|
||||
* Uses mocks to avoid spawning real Claude sessions.
|
||||
*
|
||||
* Port: none (unit tests, no server)
|
||||
*/
|
||||
|
||||
import { describe, it, expect, beforeEach, afterEach } from 'vitest';
|
||||
import {
|
||||
ModelSelector,
|
||||
resetModelSelector,
|
||||
createDefaultModelConfig,
|
||||
type ModelConfig,
|
||||
type AgentType,
|
||||
} from '../src/model-selector.js';
|
||||
import {
|
||||
GroupScheduler,
|
||||
resetGroupScheduler,
|
||||
type GroupTask,
|
||||
} from '../src/group-scheduler.js';
|
||||
import {
|
||||
ExecutionBridge,
|
||||
resetExecutionBridge,
|
||||
type PlanItem,
|
||||
} from '../src/execution-bridge.js';
|
||||
|
||||
describe('ModelSelector', () => {
|
||||
let selector: ModelSelector;
|
||||
|
||||
beforeEach(() => {
|
||||
resetModelSelector();
|
||||
selector = new ModelSelector();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
resetModelSelector();
|
||||
});
|
||||
|
||||
it('should use default model when no overrides', () => {
|
||||
const selection = selector.selectModel('task-1', { agentType: 'explore' });
|
||||
expect(selection.model).toBe('sonnet');
|
||||
expect(selection.usedUserDefault).toBe(true);
|
||||
});
|
||||
|
||||
it('should respect user default model', () => {
|
||||
selector.updateConfig({ defaultModel: 'opus' });
|
||||
const selection = selector.selectModel('task-1', { agentType: 'explore' });
|
||||
expect(selection.model).toBe('opus');
|
||||
expect(selection.usedUserDefault).toBe(true);
|
||||
});
|
||||
|
||||
it('should use agent type override when set', () => {
|
||||
selector.updateConfig({
|
||||
agentTypeOverrides: { explore: 'haiku' },
|
||||
});
|
||||
const selection = selector.selectModel('task-1', { agentType: 'explore' });
|
||||
expect(selection.model).toBe('haiku');
|
||||
expect(selection.usedUserDefault).toBe(false);
|
||||
});
|
||||
|
||||
it('should include optimizer recommendation when different', () => {
|
||||
selector.updateConfig({ defaultModel: 'opus' });
|
||||
const selection = selector.selectModel('task-1', {
|
||||
agentType: 'explore',
|
||||
recommendedModel: 'haiku',
|
||||
});
|
||||
expect(selection.model).toBe('opus');
|
||||
expect(selection.optimizerRecommendation).toBe('haiku');
|
||||
});
|
||||
|
||||
it('should select session mode for high token tasks', () => {
|
||||
const mode = selector.selectExecutionMode({ estimatedTokens: 60000 });
|
||||
expect(mode.mode).toBe('session');
|
||||
});
|
||||
|
||||
it('should select task-tool mode for low token explore tasks', () => {
|
||||
const mode = selector.selectExecutionMode({
|
||||
estimatedTokens: 10000,
|
||||
agentType: 'explore',
|
||||
});
|
||||
expect(mode.mode).toBe('task-tool');
|
||||
});
|
||||
|
||||
it('should return correct cost multipliers', () => {
|
||||
expect(selector.getModelCostMultiplier('opus')).toBe(5.0);
|
||||
expect(selector.getModelCostMultiplier('sonnet')).toBe(1.0);
|
||||
expect(selector.getModelCostMultiplier('haiku')).toBe(0.04);
|
||||
});
|
||||
});
|
||||
|
||||
describe('GroupScheduler', () => {
|
||||
let scheduler: GroupScheduler;
|
||||
|
||||
beforeEach(() => {
|
||||
resetGroupScheduler();
|
||||
scheduler = new GroupScheduler();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
resetGroupScheduler();
|
||||
});
|
||||
|
||||
it('should build schedule from plan items', () => {
|
||||
const items: PlanItem[] = [
|
||||
{ id: 't1', title: 'Task 1', description: 'Desc 1', parallelGroup: 0 },
|
||||
{ id: 't2', title: 'Task 2', description: 'Desc 2', parallelGroup: 0 },
|
||||
{ id: 't3', title: 'Task 3', description: 'Desc 3', parallelGroup: 1 },
|
||||
];
|
||||
|
||||
const schedule = scheduler.buildSchedule(items);
|
||||
|
||||
expect(schedule.groups).toHaveLength(2);
|
||||
expect(schedule.groups[0].tasks).toHaveLength(2);
|
||||
expect(schedule.groups[1].tasks).toHaveLength(1);
|
||||
expect(schedule.totalTasks).toBe(3);
|
||||
});
|
||||
|
||||
it('should order groups by number', () => {
|
||||
const items: PlanItem[] = [
|
||||
{ id: 't1', title: 'Task 1', description: 'Desc 1', parallelGroup: 2 },
|
||||
{ id: 't2', title: 'Task 2', description: 'Desc 2', parallelGroup: 0 },
|
||||
{ id: 't3', title: 'Task 3', description: 'Desc 3', parallelGroup: 1 },
|
||||
];
|
||||
|
||||
const schedule = scheduler.buildSchedule(items);
|
||||
|
||||
expect(schedule.groups[0].groupNumber).toBe(0);
|
||||
expect(schedule.groups[1].groupNumber).toBe(1);
|
||||
expect(schedule.groups[2].groupNumber).toBe(2);
|
||||
});
|
||||
|
||||
it('should track group dependencies', () => {
|
||||
const items: PlanItem[] = [
|
||||
{ id: 't1', title: 'Task 1', description: 'Desc 1', parallelGroup: 0 },
|
||||
{ id: 't2', title: 'Task 2', description: 'Desc 2', parallelGroup: 1, dependencies: ['t1'] },
|
||||
];
|
||||
|
||||
const schedule = scheduler.buildSchedule(items);
|
||||
|
||||
expect(schedule.groups[1].dependsOnGroups).toContain(0);
|
||||
});
|
||||
|
||||
it('should return first group as ready initially', () => {
|
||||
const items: PlanItem[] = [
|
||||
{ id: 't1', title: 'Task 1', description: 'Desc 1', parallelGroup: 0 },
|
||||
{ id: 't2', title: 'Task 2', description: 'Desc 2', parallelGroup: 1 },
|
||||
];
|
||||
|
||||
scheduler.buildSchedule(items);
|
||||
const nextGroup = scheduler.getNextReadyGroup();
|
||||
|
||||
expect(nextGroup).not.toBeNull();
|
||||
expect(nextGroup!.groupNumber).toBe(0);
|
||||
});
|
||||
|
||||
it('should update task status', () => {
|
||||
const items: PlanItem[] = [
|
||||
{ id: 't1', title: 'Task 1', description: 'Desc 1', parallelGroup: 0 },
|
||||
];
|
||||
|
||||
scheduler.buildSchedule(items);
|
||||
scheduler.startGroup(0);
|
||||
scheduler.updateTaskStatus('t1', 'completed');
|
||||
|
||||
const stats = scheduler.getStats();
|
||||
expect(stats.completedTasks).toBe(1);
|
||||
});
|
||||
|
||||
it('should determine execution mode based on task characteristics', () => {
|
||||
const items: PlanItem[] = [
|
||||
{
|
||||
id: 't1',
|
||||
title: 'High token task',
|
||||
description: 'Desc',
|
||||
parallelGroup: 0,
|
||||
estimatedTokens: 60000,
|
||||
},
|
||||
];
|
||||
|
||||
const schedule = scheduler.buildSchedule(items);
|
||||
expect(schedule.groups[0].executionMode).toBe('session');
|
||||
});
|
||||
});
|
||||
|
||||
describe('ExecutionBridge', () => {
|
||||
let bridge: ExecutionBridge;
|
||||
|
||||
beforeEach(() => {
|
||||
resetExecutionBridge();
|
||||
resetGroupScheduler();
|
||||
resetModelSelector();
|
||||
bridge = new ExecutionBridge();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
bridge.reset();
|
||||
resetExecutionBridge();
|
||||
resetGroupScheduler();
|
||||
resetModelSelector();
|
||||
});
|
||||
|
||||
it('should have idle status initially', () => {
|
||||
expect(bridge.status).toBe('idle');
|
||||
});
|
||||
|
||||
it('should load a plan', () => {
|
||||
const items: PlanItem[] = [
|
||||
{ id: 't1', title: 'Task 1', description: 'Desc 1', parallelGroup: 0 },
|
||||
{ id: 't2', title: 'Task 2', description: 'Desc 2', parallelGroup: 1 },
|
||||
];
|
||||
|
||||
const schedule = bridge.loadPlan(items);
|
||||
|
||||
expect(schedule).not.toBeNull();
|
||||
expect(schedule.totalTasks).toBe(2);
|
||||
expect(bridge.status).toBe('idle');
|
||||
});
|
||||
|
||||
it('should report progress', () => {
|
||||
const items: PlanItem[] = [
|
||||
{ id: 't1', title: 'Task 1', description: 'Desc 1' },
|
||||
];
|
||||
|
||||
bridge.loadPlan(items);
|
||||
const progress = bridge.getProgress();
|
||||
|
||||
expect(progress.totalTasks).toBe(1);
|
||||
expect(progress.completedTasks).toBe(0);
|
||||
expect(progress.status).toBe('idle');
|
||||
});
|
||||
|
||||
it('should update model config', () => {
|
||||
bridge.updateModelConfig({ defaultModel: 'opus' });
|
||||
const config = bridge.getModelConfig();
|
||||
expect(config.defaultModel).toBe('opus');
|
||||
});
|
||||
|
||||
it('should reset properly', () => {
|
||||
const items: PlanItem[] = [
|
||||
{ id: 't1', title: 'Task 1', description: 'Desc 1' },
|
||||
];
|
||||
|
||||
bridge.loadPlan(items);
|
||||
bridge.reset();
|
||||
|
||||
expect(bridge.status).toBe('idle');
|
||||
expect(bridge.getProgress().totalTasks).toBe(0);
|
||||
});
|
||||
|
||||
it('should throw when starting without a spawner', async () => {
|
||||
const items: PlanItem[] = [
|
||||
{ id: 't1', title: 'Task 1', description: 'Desc 1' },
|
||||
];
|
||||
|
||||
bridge.loadPlan(items);
|
||||
|
||||
await expect(bridge.start()).rejects.toThrow('No agent spawner configured');
|
||||
});
|
||||
|
||||
it('should track execution history', () => {
|
||||
const history = bridge.getHistory();
|
||||
expect(Array.isArray(history)).toBe(true);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user