diff --git a/docs/docs/features/planning.md b/docs/docs/features/planning.md index 21ea1a29a..f533d29cd 100644 --- a/docs/docs/features/planning.md +++ b/docs/docs/features/planning.md @@ -41,6 +41,10 @@ Planner Studio can include several kinds of context before generation: The goal is not to flood the model with every file. The goal is to make the proposed work easy to inspect before it runs, so reviewers can tell whether the agent saw enough relevant context. Repository summaries and indexing improve this step; see [Repository Knowledge](./repository-knowledge.md). +## How The Plan Is Written + +The planning agent writes the plan to files instead of returning it in its reply: one file per issue, checked by a validator it runs itself and fixes until the plan is complete. ProPR validates the files again before saving the plan. A plan too long for one model message therefore arrives whole, and a plan with an incomplete issue fails with a clear error instead of being saved. Set `PROPR_PLAN_GENERATION_MODE=response` to parse the plan from the reply instead (see the [configuration reference](../operations/configuration-reference.md)). + ## Review Before Running Plans stay in draft until you finalize them. Before running anything, you can: diff --git a/docs/docs/operations/configuration-reference.md b/docs/docs/operations/configuration-reference.md index f55238809..10b70e849 100644 --- a/docs/docs/operations/configuration-reference.md +++ b/docs/docs/operations/configuration-reference.md @@ -95,6 +95,8 @@ Unified image selection, per-agent credential paths, and execution limits. Codin | `CODEX_STREAM_IDLE_TIMEOUT_MS` | `1800000` (30 minutes) | Maximum quiet period on a Codex response stream before reconnecting. This is separate from the whole-task `CODEX_TIMEOUT_MS`. | Optional tuning. | | `CODEX_STREAM_MAX_RETRIES` | `5` | Number of Codex response-stream reconnect attempts. Zero disables retries. | Optional tuning. | | `CONTEXT_ANALYSIS_TIMEOUT_MS` | `3600000` (60 minutes) | Timeout for planner keyword extraction and semantic relevance scoring calls. | Optional. | +| `PROPR_PLAN_GENERATION_MODE` | `file` | `file`: the planning agent writes one JSON file per issue in a scratch workspace and runs the plan validator until it passes; ProPR re-validates before saving. `response`: the plan is parsed from the agent's reply. If the file workspace or agent cannot be started at all, ProPR falls back to `response` for that run; an invalid plan fails instead. | Optional. | +| `PROPR_PLAN_WORKSPACE_ROOT` | `/tmp/git-processor/plan-workspaces` | Scratch workspaces for plan agents. Must be the same path for the API and the Docker daemon, like the worktree root. | Custom worktree layouts only. | | `ANTIGRAVITY_TIMEOUT_MS` | `86400000` (24 hours) | Antigravity task run timeout. | Optional. | | `OPENCODE_TIMEOUT_MS` | `86400000` (24 hours) | OpenCode task run timeout. | Optional. | | `VIBE_MAX_TURNS` | `1000` | Maximum agent turns per Vibe run. | Optional. | diff --git a/packages/core/src/services/syntheticRoutingService.ts b/packages/core/src/services/syntheticRoutingService.ts index f41a203a7..7874a36e5 100644 --- a/packages/core/src/services/syntheticRoutingService.ts +++ b/packages/core/src/services/syntheticRoutingService.ts @@ -163,15 +163,20 @@ export class SyntheticRoutingSession { } } - async executeTask(options: AgentTaskOptions): Promise { + /** prepareWorkspace is awaited before every physical invocation, including failover. */ + async executeTask(options: AgentTaskOptions, prepareWorkspace?: () => Promise): Promise { this.constrain(estimateTaskRequiredTokens(options)); for (;;) { const selection = await this.select(); this.executionAttemptCount += 1; + // Preparation failures belong to the caller, not to a physical agent. + // Do not retry them or invoke an agent with an earlier attempt's workspace. + const worktreePath = prepareWorkspace ? await prepareWorkspace() : options.worktreePath; const attemptHistoryId = await this.service.recordAttempt(selection, options.taskId); try { const result = await selection.physicalAgent.executeTask({ ...options, + worktreePath, model: selection.physicalModel, isRetry: selection.attemptNumber > 1 || options.isRetry, retryReason: selection.attemptNumber > 1 diff --git a/packages/core/src/services/taskPlanning/llmCalling.ts b/packages/core/src/services/taskPlanning/llmCalling.ts index 9e7f22022..c5d7ddc20 100644 --- a/packages/core/src/services/taskPlanning/llmCalling.ts +++ b/packages/core/src/services/taskPlanning/llmCalling.ts @@ -13,6 +13,7 @@ import { import { enforceGranularity } from './granularity.js'; import { runPlanFileAgent } from './planFileAgent.js'; import { extractWholeJsonArray, incompletePlanItems, PLAN_FILE, PLAN_ORIGINAL_FILE } from './planValidation.js'; +import { resolvePlanGenerationMode, tryGeneratePlanWithFiles } from './planFileGeneration.js'; import type { Plan } from '../../claude/prompts/plannerPrompts.js'; import type { CallLLMOptions, CallLLMForPlanResult } from './types.js'; @@ -96,7 +97,20 @@ export async function callLLMForPlan(opts: CallLLMOptions): Promise { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return null; + throw new PlanningFailedError(`Could not inspect plan output ${file}: ${(error as Error).message}`); + }); +} + +/** Check parents as well as leaf files, against the workspace resolved before execution. */ +async function checkWorkspaceDirectory(directory: string): Promise { + const stats = await workspaceEntry(directory); + if (!stats) return false; + if (!stats.isDirectory() || await realpath(directory) !== directory) { + throw new PlanningFailedError(`Plan output directory ${directory} must be a real directory within the workspace, without symlinks.`); + } + return true; +} + +/** Read through a checked file handle so replacing the leaf cannot redirect a read. */ +async function readWorkspaceFile(file: string): Promise { + if (!await checkWorkspaceDirectory(path.dirname(file))) return null; + const stats = await workspaceEntry(file); + if (!stats) return null; + if (!stats.isFile()) throw new PlanningFailedError(`Plan output ${file} must be a regular file, without symlinks.`); + const handle = await open(file, constants.O_RDONLY | constants.O_NOFOLLOW | constants.O_NONBLOCK); + try { + const opened = await handle.stat(); + if (!opened.isFile()) throw new PlanningFailedError(`Plan output ${file} must be a regular file.`); + // Read at most the limit plus one byte, including if the file grows after stat. + if (opened.size > MAX_WORKSPACE_FILE_BYTES) throw new PlanningFailedError(`Plan output ${file} exceeds the ${MAX_WORKSPACE_FILE_BYTES}-byte limit.`); + const buffer = Buffer.alloc(MAX_WORKSPACE_FILE_BYTES + 1); + let length = 0; + while (length < buffer.length) { + const { bytesRead } = await handle.read(buffer, length, buffer.length - length, length); + if (bytesRead === 0) break; + length += bytesRead; + } + if (length > MAX_WORKSPACE_FILE_BYTES) throw new PlanningFailedError(`Plan output ${file} exceeds the ${MAX_WORKSPACE_FILE_BYTES}-byte limit.`); + return buffer.toString('utf8', 0, length); + } finally { + await handle.close(); + } +} + +async function readTaskFiles(directory: string): Promise> { + try { + if (!await checkWorkspaceDirectory(directory)) return {}; + const names = (await readdir(directory)).filter(name => name.endsWith('.json')).sort(); + if (names.length > MAX_TASK_FILES) throw new PlanningFailedError(`Plan output has ${names.length} task files; the limit is ${MAX_TASK_FILES}. No tasks were accepted.`); + const files: Record = {}; + for (const name of names) { + const content = await readWorkspaceFile(path.join(directory, name)); + if (content === null) throw new PlanningFailedError(`Could not read task file ${PLAN_TASKS_DIR}/${name}. No tasks were accepted.`); + files[name] = content; + } + return files; + } catch (error) { + if (error instanceof PlanningFailedError) throw error; + throw new PlanningFailedError(`Could not read plan task files: ${(error as Error).message}`); + } +} + +/** Output need not be valid, readable, or even complete to rule out a retry. */ +async function hasPlanOutput(workspace: string, taskFiles: boolean): Promise { + if (await workspaceEntry(path.join(workspace, PLAN_FILE))) return true; + if (!taskFiles) return false; + const directory = path.join(workspace, PLAN_TASKS_DIR); + if (!await checkWorkspaceDirectory(directory)) return false; + return (await readdir(directory)).length > 0; +} + export interface PlanFileAgentOptions { /** Workspace name prefix and log label. */ purpose: 'generation' | 'repair'; @@ -42,6 +131,11 @@ export interface PlanFileAgentOptions { files?: Record; /** Content the result must preserve (repair); validated against it. */ original?: string; + /** + * The agent writes one task object per file under `tasks/` instead of a + * single plan.json, so no single tool call has to hold the whole plan. + */ + taskFiles?: boolean; /** `agent:model`, as configured for planning. */ model: string; draftId: string; @@ -53,11 +147,16 @@ export interface PlanFileAgentOptions { routingSession?: SyntheticRoutingSession; } +// The registry, model aliases and log helpers open database and queue +// connections when first loaded; loading them on use keeps planning modules +// (and their tests) free of those side effects until an agent actually runs. async function resolveAgent(model: string, routingSession?: SyntheticRoutingSession) { + const { resolveModelAlias } = await import('../../config/modelAliases.js'); const separator = model.indexOf(':'); const alias = separator === -1 ? null : model.slice(0, separator); const modelName = separator === -1 ? model : model.slice(separator + 1); if (routingSession) return { runner: routingSession, model: resolveModelAlias(modelName) }; + const { AgentRegistry } = await import('../../agents/AgentRegistry.js'); const registry = AgentRegistry.getInstance(); await registry.ensureInitialized(); const agent = alias ? registry.getAgentByAlias(alias) : registry.getDefaultAgent(); @@ -65,21 +164,45 @@ async function resolveAgent(model: string, routingSession?: SyntheticRoutingSess return { runner: agent, model: resolveModelAlias(modelName) }; } -export async function runPlanFileAgent(options: PlanFileAgentOptions): Promise { - const { purpose, prompt, files = {}, original, model, draftId, repository, githubToken, executionType, correlationId, metadata, routingSession } = options; - const log = correlationId ? logger.withCorrelation(correlationId) : logger; +/** Track every attempt until cleanup, including one whose preparation fails. */ +async function preparePlanWorkspace(options: PlanFileAgentOptions, workspaces: string[]): Promise { const root = planWorkspaceRoot(); - await mkdir(root, { recursive: true }); - const workspace = await mkdtemp(path.join(root, `${purpose}-`)); + const workspace = await mkdir(root, { recursive: true }) + .then(() => mkdtemp(path.join(root, `${options.purpose}-`))) + .then(directory => realpath(directory)) + .catch(error => { throw new PlanFileAgentUnavailableError(`Could not create a plan workspace under ${root}: ${(error as Error).message}`); }); + workspaces.push(workspace); try { await writeFile(path.join(workspace, PLAN_VALIDATOR_FILE), PLAN_VALIDATOR_SCRIPT); - for (const [name, content] of Object.entries(files)) await writeFile(path.join(workspace, name), content); - // Some agent CLIs refuse to run outside a git repository. - await execFileAsync('git', ['init', '-q'], { cwd: workspace }).catch(() => undefined); + for (const [name, content] of Object.entries(options.files || {})) await writeFile(path.join(workspace, name), content); + if (options.taskFiles) await mkdir(path.join(workspace, PLAN_TASKS_DIR)); + } catch (error) { + throw new PlanFileAgentUnavailableError(`Could not prepare the plan workspace: ${(error as Error).message}`); + } + // Some agent CLIs refuse to run outside a git repository. + await execFileAsync('git', ['init', '-q'], { cwd: workspace }).catch(() => undefined); + return workspace; +} + +export async function runPlanFileAgent(options: PlanFileAgentOptions): Promise { + const { purpose, prompt, original, taskFiles = false, model, draftId, repository, githubToken, executionType, correlationId, metadata, routingSession } = options; + const log = correlationId ? logger.withCorrelation(correlationId) : logger; + const workspaces: string[] = []; + try { + let workspace = await preparePlanWorkspace(options, workspaces); + let firstAttempt = true; + const prepareAttemptWorkspace = async () => { + if (!firstAttempt) workspace = await preparePlanWorkspace(options, workspaces); + firstAttempt = false; + return workspace; + }; - const { runner, model: resolvedModel } = await resolveAgent(model, routingSession); + const { runner, model: resolvedModel } = await resolveAgent(model, routingSession).catch(error => { + throw new PlanFileAgentUnavailableError((error as Error).message); + }); + const { buildAnalysisWorkRef, withTaskLogAttribution } = await import('../../utils/llmLogger.js'); const [repoOwner = 'unknown', repoName = 'unknown'] = repository.split('/'); - log.info({ purpose, model, workspace }, 'Running plan file agent'); + log.info({ purpose, model, workspace, taskFiles }, 'Running plan file agent'); const result = await runner.executeTask({ worktreePath: workspace, issueRef: { number: 0, repoOwner, repoName }, @@ -87,28 +210,59 @@ export async function runPlanFileAgent(options: PlanFileAgentOptions): Promise

{ + // Preserve terminal failures before inspecting output: an empty workspace + // does not authorize retrying explicit cancellation or other terminal errors. + // Usage limits also keep their identity for upstream requeueing. + if ((error as Error)?.name === 'UsageLimitError' || isNonRetryableSyntheticFailure(error)) throw error; + let producedOutput: boolean; + try { + // A later empty attempt must not erase evidence that an earlier one ran. + producedOutput = false; + for (const attemptWorkspace of workspaces) { + if (await hasPlanOutput(attemptWorkspace, taskFiles)) { + producedOutput = true; + break; + } + } + } catch (inspectionError) { + // Absence of output must be established before allowing response fallback. + throw new PlanningFailedError(`Plan ${purpose} agent failed: ${(error as Error).message}. Could not establish absence of plan output: ${(inspectionError as Error).message}`); + } + if (producedOutput) throw new PlanningFailedError(`Plan ${purpose} agent failed after producing output: ${(error as Error).message}`); + throw new PlanFileAgentUnavailableError(`Plan ${purpose} agent could not run: ${(error as Error).message}`); }); - const planText = await readFile(path.join(workspace, PLAN_FILE), 'utf8').catch(() => null); const agentFailure = result.success ? '' : ` The agent reported: ${(result.error || 'execution failed').slice(0, 300)}`; - if (planText === null) throw new PlanningFailedError(`Plan ${purpose} produced no ${PLAN_FILE}.${agentFailure}`); + if (purpose === 'generation' && !result.success) throw new PlanningFailedError(`Plan generation did not finish.${agentFailure}`); // Never trust the workspace copy of the validator or of the original. - const report = await validatePlanText(planText, original); + let planText: string | null; + let report; + if (taskFiles) { + const written = await readTaskFiles(path.join(workspace, PLAN_TASKS_DIR)); + if (Object.keys(written).length === 0) throw new PlanningFailedError(`Plan ${purpose} wrote no task files to ${PLAN_TASKS_DIR}/.${agentFailure}`); + ({ planText, ...report } = await validatePlanTaskFiles(written)); + } else { + planText = await readWorkspaceFile(path.join(workspace, PLAN_FILE)); + if (planText === null) throw new PlanningFailedError(`Plan ${purpose} produced no ${PLAN_FILE}.${agentFailure}`); + report = await validatePlanText(planText, original); + } if (!report.valid) { log.warn({ purpose, model, errors: report.errors.slice(0, 10) }, 'Plan file agent result failed validation'); throw new PlanningFailedError(`Plan ${purpose} did not produce a valid plan: ${report.errors.slice(0, 3).join('; ')}.${agentFailure}`); } log.info({ purpose, model, taskCount: report.taskCount }, 'Plan file agent produced a valid plan'); - return JSON.parse(planText) as PlanItem[]; + return JSON.parse(planText!) as PlanItem[]; } finally { - await rm(workspace, { recursive: true, force: true }).catch(error => { - log.warn({ workspace, error: (error as Error).message }, 'Failed to remove plan workspace'); - }); + for (const workspace of workspaces) { + await rm(workspace, { recursive: true, force: true }).catch(error => { + log.warn({ workspace, error: (error as Error).message }, 'Failed to remove plan workspace'); + }); + } } } diff --git a/packages/core/src/services/taskPlanning/planFileGeneration.ts b/packages/core/src/services/taskPlanning/planFileGeneration.ts new file mode 100644 index 000000000..628bcbf8f --- /dev/null +++ b/packages/core/src/services/taskPlanning/planFileGeneration.ts @@ -0,0 +1,85 @@ +/** + * File-based plan generation: the planning agent writes each task to its own + * file in a scratch workspace and runs the plan validator until it passes. + * + * The plan then no longer depends on the agent's final chat message, which is + * cut at the provider's per-message output limit for very large plans, and a + * syntax mistake is fixed by the agent in place instead of by a separate + * repair model re-emitting the whole plan. + */ + +import type { Plan } from '../../claude/prompts/plannerPrompts.js'; +import logger from '../../utils/logger.js'; +import type { SyntheticRoutingSession } from '../syntheticRoutingService.js'; +import { PlanFileAgentUnavailableError, runPlanFileAgent } from './planFileAgent.js'; +import { PLAN_TASKS_DIR, PLAN_VALIDATOR_FILE } from './planValidation.js'; + +export type PlanGenerationMode = 'file' | 'response'; + +/** + * `PROPR_PLAN_GENERATION_MODE=response` restores the previous behaviour, where + * the plan is parsed from the model's reply. Anything else selects files. + */ +export function resolvePlanGenerationMode(env: NodeJS.ProcessEnv = process.env): PlanGenerationMode { + return env.PROPR_PLAN_GENERATION_MODE?.trim().toLowerCase() === 'response' ? 'response' : 'file'; +} + +/** The planner prompt plus the file contract, which replaces its reply format. */ +export function buildPlanFilePrompt(fullContext: string): string { + return `${fullContext} + +--- +## How to deliver the plan (this replaces the output format above) + +Do not put the plan in your reply. Write it to files in the current directory: + +1. Write each task as its own file in plan order: \`${PLAN_TASKS_DIR}/001.json\`, \`${PLAN_TASKS_DIR}/002.json\`, and so on. Each file holds ONE JSON object with the string fields "title", "body" and "implementation", written exactly as the instructions above describe. Write one task per tool call; never put the whole plan into a single file or message. +2. Run \`node ${PLAN_VALIDATOR_FILE} --tasks ${PLAN_TASKS_DIR}\`. It assembles the task files into plan.json and prints a JSON report. +3. If the report lists errors, fix the files it names and run it again. Repeat until it exits with status 0 and reports "valid": true. +4. Do not edit ${PLAN_VALIDATOR_FILE} and do not write any other files. The workspace holds only these files; everything you need from the repository is in this prompt. +5. When the validator passes, reply with one short line, for example "Plan written: 4 tasks."`; +} + +export interface FilePlanGenerationOptions { + draftId: string; + fullContext: string; + model: string; + repository: string; + githubToken: string; + correlationId?: string; + metadata?: Record; + routingSession?: SyntheticRoutingSession; +} + +/** + * Generates the plan through files when that mode is selected. Returns null + * when the response mode should be used instead: either it is selected, or + * the agent workspace could not be set up or run at all. A plan that the + * agent wrote but that fails validation is an error, not a fallback: running + * the whole generation again as a reply would double its cost and would hit + * the same model limits. + */ +export async function tryGeneratePlanWithFiles(options: FilePlanGenerationOptions): Promise { + if (resolvePlanGenerationMode() !== 'file') return null; + const { draftId, fullContext, model, repository, githubToken, correlationId, metadata, routingSession } = options; + const log = correlationId ? logger.withCorrelation(correlationId) : logger; + try { + return await runPlanFileAgent({ + purpose: 'generation', + prompt: buildPlanFilePrompt(fullContext), + taskFiles: true, + model, + draftId, + repository, + githubToken, + executionType: 'plan-generation', + correlationId, + metadata: { ...metadata, planGenerationMode: 'file' }, + routingSession, + }); + } catch (error) { + if (!(error instanceof PlanFileAgentUnavailableError)) throw error; + log.warn({ model, error: error.message }, 'File-based plan generation is unavailable, generating from the model reply instead'); + return null; + } +} diff --git a/packages/core/src/services/taskPlanning/planValidation.ts b/packages/core/src/services/taskPlanning/planValidation.ts index 763791b17..3c5b3c552 100644 --- a/packages/core/src/services/taskPlanning/planValidation.ts +++ b/packages/core/src/services/taskPlanning/planValidation.ts @@ -8,7 +8,7 @@ */ import { execFile } from 'node:child_process'; -import { mkdtemp, rm, writeFile } from 'node:fs/promises'; +import { mkdir, mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import path from 'node:path'; import { promisify } from 'node:util'; @@ -18,6 +18,8 @@ const execFileAsync = promisify(execFile); export const PLAN_FILE = 'plan.json'; export const PLAN_ORIGINAL_FILE = 'original.txt'; export const PLAN_VALIDATOR_FILE = 'validate-plan.mjs'; +/** Incremental contract: one task object per file, assembled in file-name order. */ +export const PLAN_TASKS_DIR = 'tasks'; /** * Positions (1-based) of plan items that are not complete issues: a model @@ -100,25 +102,51 @@ function findJsonArrayEnd(text: string, start: number): number | null { /** * Standalone validator. `node validate-plan.mjs [plan.json] [original.txt]` - * prints a JSON report and exits 0 only when the plan is valid. With an + * prints a JSON report and exits 0 only when the plan is valid. With + * `--tasks

` it first assembles `/*.json`, one task object per file + * in file-name order, into plan.json, so a plan too large for one message can + * be written task by task. With an * original present it also proves that the plan keeps the original's content: * every task and field matches its original occurrence, in task order. An * ambiguous original fails closed rather than authorizing content loss. */ export const PLAN_VALIDATOR_SCRIPT = String.raw`#!/usr/bin/env node // Validates a ProPR plan. Do not edit this file: ProPR re-runs its own copy. -import { existsSync, readFileSync } from 'node:fs'; +import { existsSync, readdirSync, readFileSync, writeFileSync } from 'node:fs'; -const planPath = process.argv[2] || 'plan.json'; -const originalPath = process.argv[3] || 'original.txt'; +const args = process.argv.slice(2); +const tasksFlag = args.indexOf('--tasks'); +const tasksDir = tasksFlag === -1 ? null : (args.splice(tasksFlag, 2)[1] || 'tasks'); +const planPath = args[0] || 'plan.json'; +const originalPath = args[1] || 'original.txt'; const REQUIRED = ['title', 'body', 'implementation']; const errors = []; +let taskFiles = []; let plan; -try { - plan = JSON.parse(readFileSync(planPath, 'utf8')); -} catch (error) { - errors.push(planPath + ' is not valid JSON: ' + error.message); +if (tasksDir !== null) { + taskFiles = existsSync(tasksDir) ? readdirSync(tasksDir).filter(name => name.endsWith('.json')).sort() : []; + if (taskFiles.length === 0) errors.push('no task files: write one task object per file as ' + tasksDir + '/001.json, ' + tasksDir + '/002.json, ...'); + const tasks = []; + for (const name of taskFiles) { + try { + const task = JSON.parse(readFileSync(tasksDir + '/' + name, 'utf8')); + if (Array.isArray(task)) errors.push(tasksDir + '/' + name + ' must hold one task object, not an array'); + else tasks.push(task); + } catch (error) { + errors.push(tasksDir + '/' + name + ' is not valid JSON: ' + error.message); + } + } + if (errors.length === 0) writeFileSync(planPath, JSON.stringify(tasks, null, 2) + '\n'); +} +const label = index => 'task ' + (index + 1) + (taskFiles[index] ? ' (' + tasksDir + '/' + taskFiles[index] + ')' : ''); + +if (errors.length === 0) { + try { + plan = JSON.parse(readFileSync(planPath, 'utf8')); + } catch (error) { + errors.push(planPath + ' is not valid JSON: ' + error.message); + } } if (plan !== undefined) { @@ -126,12 +154,12 @@ if (plan !== undefined) { else if (plan.length === 0) errors.push('the plan has no tasks'); else plan.forEach((task, index) => { if (!task || typeof task !== 'object' || Array.isArray(task)) { - errors.push('task ' + (index + 1) + ' is not an object'); + errors.push(label(index) + ' is not an object'); return; } for (const field of REQUIRED) { if (typeof task[field] !== 'string' || !task[field].trim()) { - errors.push('task ' + (index + 1) + ' has no non-empty "' + field + '" string'); + errors.push(label(index) + ' has no non-empty "' + field + '" string'); } } }); @@ -214,7 +242,7 @@ if (errors.length === 0 && existsSync(originalPath)) { } source.forEach((task, index) => { if (!same(task, plan[index])) { - errors.push('task ' + (index + 1) + ' does not match ' + originalPath + ': fix syntax only, keep every task and field in its original task'); + errors.push(label(index) + ' does not match ' + originalPath + ': fix syntax only, keep every task and field in its original task'); } }); } catch (error) { @@ -233,6 +261,17 @@ export interface PlanValidationReport { errors: string[]; } +async function runOwnValidator(directory: string, args: string[]): Promise { + await writeFile(path.join(directory, PLAN_VALIDATOR_FILE), PLAN_VALIDATOR_SCRIPT); + const { stdout } = await execFileAsync(process.execPath, [PLAN_VALIDATOR_FILE, ...args], { + cwd: directory, timeout: 30_000, maxBuffer: 1024 * 1024, + }).catch((error: { stdout?: string; message: string }) => { + if (typeof error.stdout === 'string' && error.stdout.trim()) return { stdout: error.stdout }; + throw error; + }); + return JSON.parse(stdout) as PlanValidationReport; +} + /** * Validates plan text with ProPR's own copy of the validator, in a private * directory, so a workspace's (possibly edited) validator is never trusted. @@ -240,16 +279,29 @@ export interface PlanValidationReport { export async function validatePlanText(planText: string, original?: string): Promise { const directory = await mkdtemp(path.join(tmpdir(), 'propr-plan-validate-')); try { - await writeFile(path.join(directory, PLAN_VALIDATOR_FILE), PLAN_VALIDATOR_SCRIPT); await writeFile(path.join(directory, PLAN_FILE), planText); if (original !== undefined) await writeFile(path.join(directory, PLAN_ORIGINAL_FILE), original); - const { stdout } = await execFileAsync(process.execPath, [PLAN_VALIDATOR_FILE, PLAN_FILE, PLAN_ORIGINAL_FILE], { - cwd: directory, timeout: 30_000, maxBuffer: 1024 * 1024, - }).catch((error: { stdout?: string; message: string }) => { - if (typeof error.stdout === 'string' && error.stdout.trim()) return { stdout: error.stdout }; - throw error; - }); - return JSON.parse(stdout) as PlanValidationReport; + return await runOwnValidator(directory, [PLAN_FILE, PLAN_ORIGINAL_FILE]); + } finally { + await rm(directory, { recursive: true, force: true }); + } +} + +/** + * Assembles and validates task files (name → content, one task object each) + * with ProPR's own validator. Returns the assembled plan text when valid. + */ +export async function validatePlanTaskFiles(taskFiles: Record): Promise { + const directory = await mkdtemp(path.join(tmpdir(), 'propr-plan-validate-')); + try { + await mkdir(path.join(directory, PLAN_TASKS_DIR)); + for (const [name, content] of Object.entries(taskFiles)) { + if (path.basename(name) !== name || !name.endsWith('.json')) throw new Error(`Invalid task file name: ${name}`); + await writeFile(path.join(directory, PLAN_TASKS_DIR, name), content); + } + const report = await runOwnValidator(directory, ['--tasks', PLAN_TASKS_DIR, PLAN_FILE]); + const planText = report.valid ? await readFile(path.join(directory, PLAN_FILE), 'utf8') : null; + return { ...report, planText }; } finally { await rm(directory, { recursive: true, force: true }); } diff --git a/test/planFileGeneration.test.ts b/test/planFileGeneration.test.ts new file mode 100644 index 000000000..fa163c4af --- /dev/null +++ b/test/planFileGeneration.test.ts @@ -0,0 +1,612 @@ +import assert from 'node:assert/strict'; +import { existsSync, mkdirSync, mkdtempSync, readdirSync, readFileSync, rmSync, symlinkSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; +import { spawnSync } from 'node:child_process'; +import { after, beforeEach, describe, mock, test } from 'node:test'; +import type { Agent, AnalyzeOptions } from '../packages/core/src/agents/types.js'; +import type { SyntheticRoutingSession as RoutingSession } from '../packages/core/src/services/syntheticRoutingService.js'; + +const workspaceRoot = mkdtempSync(path.join(tmpdir(), 'propr-plan-workspaces-')); +process.env.PROPR_PLAN_WORKSPACE_ROOT = workspaceRoot; +after(() => rmSync(workspaceRoot, { recursive: true, force: true })); + +const logger = { info: mock.fn(), warn: mock.fn(), error: mock.fn(), debug: mock.fn() }; +await mock.module('../packages/core/src/utils/logger.js', { + defaultExport: { ...logger, withCorrelation: mock.fn(() => logger) }, +}); + +class PlanningFailedError extends Error { + constructor(message: string) { + super(message); + this.name = 'PlanningFailedError'; + } +} +await mock.module('../packages/core/src/services/planning/index.js', { + namedExports: { + PlanningFailedError, + updateTraceForRun: mock.fn(async () => undefined), + validatePromptTokens: mock.fn(async (prompt: string) => ({ valid: true, tokenCount: Math.ceil(prompt.length / 4), source: 'tiktoken' as const })), + CLAUDE_CODE_OVERHEAD: 5_000, + getModelHardLimit: mock.fn(() => 200_000), + getRawInputCharLimit: mock.fn(() => null), + }, +}); +await mock.module('../packages/core/src/config/modelAliases.js', { + namedExports: { resolveModelAlias: (model: string) => model }, +}); +await mock.module('../packages/core/src/utils/llmLogger.js', { + namedExports: { + buildAnalysisWorkRef: (executionType: string, taskId: string, repository: string) => ({ workType: 'plan', planDraftId: taskId, workRepository: repository, executionType }), + withTaskLogAttribution: (metadata: Record, attribution: unknown) => ({ ...metadata, proprLogAttribution: attribution }), + }, +}); +await mock.module('../packages/core/src/utils/llmEstimation.js', { + namedExports: { estimateLlmDuration: mock.fn(async () => ({ estimatedDurationMs: 1, isHistoricalEstimate: false, sampleCount: 0, avgMsPerToken: 0 })) }, +}); +const replies: string[] = []; +const runLightweightLLMAnalysis = mock.fn(async (options: { prompt: string; routingSession?: RoutingSession }) => { + if (options.routingSession instanceof SyntheticRoutingSession) { + const result = await options.routingSession.analyze(options.prompt); + assert.equal(result.success, true); + return result.response; + } + const reply = replies.shift(); + if (reply === undefined) throw new Error('Unexpected reply-mode call'); + return reply; +}); +await mock.module('../packages/core/src/claude/claudeService.js', { namedExports: { runLightweightLLMAnalysis } }); + +// Exercise the real routing retry loop without opening database/queue connections. +await mock.module('../packages/core/src/db/connection.js', { namedExports: { db: {} } }); +await mock.module('../packages/core/src/config/configManager.js', { namedExports: { loadSyntheticAgents: async () => [] } }); +await mock.module('../packages/core/src/services/syntheticUsageSnapshotProvider.js', { + namedExports: { AliasSpecificAgentTankSnapshotProvider: class {} }, +}); +await mock.module('../packages/core/src/utils/tokenCalculation.js', { + namedExports: { estimateTokens: (text: string) => Math.ceil(text.length / 4) }, +}); +const { SyntheticRoutingSession, SyntheticRoutingService, SyntheticPoolExhaustedError } = await import('../packages/core/src/services/syntheticRoutingService.js'); + +const { PLAN_VALIDATOR_SCRIPT, validatePlanTaskFiles, validatePlanText } = await import('../packages/core/src/services/taskPlanning/planValidation.js'); +const { runPlanFileAgent, PlanFileAgentUnavailableError } = await import('../packages/core/src/services/taskPlanning/planFileAgent.js'); +const { buildPlanFilePrompt, resolvePlanGenerationMode, tryGeneratePlanWithFiles } = await import('../packages/core/src/services/taskPlanning/planFileGeneration.js'); +const { callLLMForPlan } = await import('../packages/core/src/services/taskPlanning/llmCalling.js'); + +const task = (title: string) => ({ title, body: `Why ${title} matters`, implementation: `~~~diff\n+ ${title}\n~~~` }); +const json = (value: unknown) => JSON.stringify(value, null, 2); + +type TaskOptions = { worktreePath: string; prompt: string; model?: string; taskId?: string; maxTurns?: number; metadata?: Record }; +/** A routing session whose agent writes files into the workspace it is given. */ +function fakeAgent(write: (workspace: string, options: TaskOptions) => void | Promise, result: Record = { success: true }) { + const calls: TaskOptions[] = []; + return { + calls, + session: { + fork: () => ({}), + executeTask: async (options: TaskOptions) => { + calls.push(options); + await write(options.worktreePath, options); + return { modifiedFiles: [], logs: '', modelUsed: options.model, executionTimeMs: 1, ...result }; + }, + }, + }; +} +/** Stub member selection only; physical invocation and failover use the production loop. */ +function routedAgents(...agents: ReturnType[]) { + let next = 0; + return new SyntheticRoutingSession({ + select: async () => { + const agent = agents[next++]; + if (!agent) throw new Error('routing pool exhausted'); + return { + physicalAgent: agent.session, physicalModel: 'opus', synthetic: true, + memberId: `member-${next}`, attemptNumber: next, + }; + }, + metadataFor: () => ({}), + recordAttempt: async () => undefined, + } as never, { requestedAgentAlias: 'pool', requestedModel: 'smart', requiredTokens: 0, callId: 'routing-test' }); +} +const writeTasks = (workspace: string, files: Record) => { + for (const [name, content] of Object.entries(files)) writeFileSync(path.join(workspace, 'tasks', name), content); +}; +const baseOptions = { model: 'claude:opus', draftId: 'draft-1', repository: 'acme/repo', githubToken: 'token', correlationId: 'plan-1' }; + +beforeEach(() => { + delete process.env.PROPR_PLAN_GENERATION_MODE; + replies.length = 0; + runLightweightLLMAnalysis.mock.resetCalls(); +}); + +describe('incremental task-file contract in the validator', () => { + test('assembles task files in name order into plan.json', async () => { + const report = await validatePlanTaskFiles({ '002.json': json(task('Second')), '001.json': json(task('First')) }); + assert.equal(report.valid, true, report.errors.join('; ')); + assert.equal(report.taskCount, 2); + assert.deepEqual(JSON.parse(report.planText!).map((item: { title: string }) => item.title), ['First', 'Second']); + }); + + test('names the file that needs fixing', async () => { + const report = await validatePlanTaskFiles({ '001.json': json(task('Fine')), '002.json': '{"title": "Broken",', '003.json': json([task('Array')]) }); + assert.equal(report.valid, false); + assert.equal(report.planText, null); + assert.match(report.errors.join('\n'), /tasks\/002\.json is not valid JSON/); + assert.match(report.errors.join('\n'), /tasks\/003\.json must hold one task object, not an array/); + const incomplete = await validatePlanTaskFiles({ '001.json': json({ title: 'Only a title' }) }); + assert.match(incomplete.errors[0], /^task 1 \(tasks\/001\.json\) has no non-empty "body" string$/); + }); + + test('reports a missing or empty tasks directory and keeps the plain plan.json form unchanged', async () => { + const directory = mkdtempSync(path.join(tmpdir(), 'propr-validator-')); + try { + writeFileSync(path.join(directory, 'validate-plan.mjs'), PLAN_VALIDATOR_SCRIPT); + const empty = spawnSync(process.execPath, ['validate-plan.mjs', '--tasks', 'tasks'], { cwd: directory, encoding: 'utf8' }); + assert.equal(empty.status, 1); + assert.match(JSON.parse(empty.stdout).errors[0], /no task files: write one task object per file as tasks\/001\.json/); + assert.equal(existsSync(path.join(directory, 'plan.json')), false, 'nothing is assembled from no tasks'); + } finally { + rmSync(directory, { recursive: true, force: true }); + } + const plain = await validatePlanText(json([{ title: 'Only a title' }])); + assert.equal(plain.errors[0], 'task 1 has no non-empty "body" string'); + }); +}); + +describe('plan file agent with task files', () => { + test('returns the plan the agent wrote task by task and cleans the workspace up', async () => { + const agent = fakeAgent(workspace => writeTasks(workspace, { '001.json': json(task('One')), '002.json': json(task('Two')) })); + const plan = await runPlanFileAgent({ ...baseOptions, purpose: 'generation', prompt: 'Plan it', taskFiles: true, executionType: 'plan-generation', routingSession: agent.session as never }); + assert.deepEqual(plan.map(item => item.title), ['One', 'Two']); + assert.equal(agent.calls[0].model, 'opus'); + assert.equal(agent.calls[0].maxTurns, 200, 'one turn per task file must fit, beyond the shipped CLAUDE_MAX_TURNS=10'); + assert.equal((agent.calls[0].metadata?.proprLogAttribution as { executionType: string }).executionType, 'plan-generation'); + assert.deepEqual(readdirSync(workspaceRoot), [], 'workspace removed'); + }); + + for (const failure of ['result', 'throw']) { + test(`isolates task files across routing retries after a failed ${failure}`, async () => { + const abandoned = fakeAgent(workspace => { + writeTasks(workspace, { + '001.json': json(task('Old one')), '002.json': json(task('Old two')), '003.json': json(task('Abandoned')), + }); + const validation = spawnSync(process.execPath, ['validate-plan.mjs', '--tasks', 'tasks'], { cwd: workspace, encoding: 'utf8' }); + assert.equal(validation.status, 0, validation.stderr); + writeFileSync(path.join(workspace, 'validate-plan.mjs'), 'tampered'); + if (failure === 'throw') throw new Error('transport failed'); + }, { success: false, error: 'transport failed' }); + const replacement = fakeAgent(async workspace => { + assert.notEqual(workspace, abandoned.calls[0].worktreePath); + assert.deepEqual(readdirSync(path.join(workspace, 'tasks')), []); + assert.equal(existsSync(path.join(workspace, 'plan.json')), false); + assert.equal(readFileSync(path.join(workspace, 'validate-plan.mjs'), 'utf8'), PLAN_VALIDATOR_SCRIPT); + await Promise.resolve(); + // Even a late write to the abandoned workspace cannot enter this plan. + writeTasks(abandoned.calls[0].worktreePath, { '004.json': json(task('Late abandoned task')) }); + writeTasks(workspace, { '001.json': json(task('Replacement one')), '002.json': json(task('Replacement two')) }); + const validation = spawnSync(process.execPath, ['validate-plan.mjs', '--tasks', 'tasks'], { cwd: workspace, encoding: 'utf8' }); + assert.equal(validation.status, 0, validation.stderr); + }); + const result = await callLLMForPlan({ + ...baseOptions, runId: 'run-retry', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession: routedAgents(abandoned, replacement), + }); + assert.deepEqual(result.plan.map(item => item.title), ['Replacement one', 'Replacement two']); + assert.equal(abandoned.calls.length, 1); + assert.equal(replacement.calls.length, 1); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 0); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + } + + for (const taskFiles of [true, false]) { + test(`rejects an empty successful retry instead of accepting abandoned ${taskFiles ? 'task files' : 'plan.json'}`, async () => { + const abandoned = fakeAgent(workspace => { + if (taskFiles) writeTasks(workspace, { '001.json': json(task('Abandoned')) }); + else writeFileSync(path.join(workspace, 'plan.json'), json([task('Abandoned')])); + }, { success: false, error: 'transport failed' }); + const replacement = fakeAgent(() => undefined); + await assert.rejects(runPlanFileAgent({ + ...baseOptions, purpose: 'generation', prompt: 'Plan it', taskFiles, executionType: 'plan-generation', + routingSession: routedAgents(abandoned, replacement), + }), taskFiles ? /wrote no task files/ : /produced no plan.json/); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + } + + test('retains earlier output evidence when routing exhausts after an empty retry', async () => { + const abandoned = fakeAgent(workspace => writeTasks(workspace, { '001.json': json(task('Abandoned')) }), { success: false }); + const replacement = fakeAgent(() => { throw new Error('could not run'); }); + await assert.rejects(tryGeneratePlanWithFiles({ + ...baseOptions, fullContext: 'Plan it', routingSession: routedAgents(abandoned, replacement), + }), (error: Error) => error instanceof PlanningFailedError && !(error instanceof PlanFileAgentUnavailableError) + && /failed after producing output/.test(error.message)); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + + test('does not invoke another agent if preparing its isolated workspace fails', async () => { + const files = { 'context.txt': 'original context' }; + const abandoned = fakeAgent(workspace => { + writeTasks(workspace, { '001.json': json(task('Abandoned')) }); + // Make preparation fail on the retry, after the first attempt produced output. + Object.assign(files, { 'missing/context.txt': 'cannot write here' }); + }, { success: false }); + const replacement = fakeAgent(() => undefined); + await assert.rejects(runPlanFileAgent({ + ...baseOptions, purpose: 'generation', prompt: 'Plan it', taskFiles: true, files, executionType: 'plan-generation', + routingSession: routedAgents(abandoned, replacement), + }), (error: Error) => error instanceof PlanningFailedError && !(error instanceof PlanFileAgentUnavailableError) + && /Could not prepare the plan workspace/.test(error.message)); + assert.equal(replacement.calls.length, 0); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + + test('re-validates with its own validator', async () => { + const agent = fakeAgent(workspace => { + writeFileSync(path.join(workspace, 'validate-plan.mjs'), 'console.log(JSON.stringify({ valid: true, taskCount: 1, errors: [] }))'); + writeTasks(workspace, { '001.json': json({ title: 'Stub' }) }); + }); + await assert.rejects( + runPlanFileAgent({ ...baseOptions, purpose: 'generation', prompt: 'Plan it', taskFiles: true, executionType: 'plan-generation', routingSession: agent.session as never }), + /task 1 \(tasks\/001\.json\) has no non-empty "body"/, + ); + }); + + test('accepts the file count limit and rejects overflow without returning a prefix', async () => { + for (const count of [200, 201]) { + const agent = fakeAgent(workspace => { + for (let index = 1; index <= count; index++) { + writeTasks(workspace, { [`${String(index).padStart(3, '0')}.json`]: json(task(`Task ${index}`)) }); + } + }); + const result = runPlanFileAgent({ ...baseOptions, purpose: 'generation', prompt: 'Plan it', taskFiles: true, executionType: 'plan-generation', routingSession: agent.session as never }); + if (count === 200) assert.equal((await result).length, count); + else await assert.rejects(result, /201 task files; the limit is 200/); + } + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + + const unacceptableFiles = { + oversized: (workspace: string) => writeTasks(workspace, { '002.json': json(task('Large')).padEnd(8 * 1024 * 1024 + 1, ' ') }), + directory: (workspace: string) => mkdirSync(path.join(workspace, 'tasks', '002.json')), + symlink: (workspace: string) => symlinkSync(path.join(workspace, 'tasks', '001.json'), path.join(workspace, 'tasks', '002.json')), + 'dangling symlink': (workspace: string) => symlinkSync(path.join(workspace, 'missing.json'), path.join(workspace, 'tasks', '002.json')), + }; + for (const [kind, write] of Object.entries(unacceptableFiles)) { + test(`rejects a task file of type ${kind} alongside a valid task`, async () => { + const agent = fakeAgent(workspace => { + writeTasks(workspace, { '001.json': json(task('Valid prefix')) }); + write(workspace); + }); + await assert.rejects( + tryGeneratePlanWithFiles({ ...baseOptions, fullContext: 'Plan it', routingSession: agent.session as never }), + (error: Error) => error instanceof PlanningFailedError && !(error instanceof PlanFileAgentUnavailableError) + && /002\.json/.test(error.message) && /exceeds|regular file/.test(error.message), + ); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + } + + for (const target of ['file', 'directory']) { + test(`rejects a task ${target} symlink to outside the workspace`, async () => { + const outside = mkdtempSync(path.join(tmpdir(), 'propr-outside-tasks-')); + writeFileSync(path.join(outside, '001.json'), json(task('Outside task'))); + try { + const agent = fakeAgent(workspace => { + if (target === 'directory') { + rmSync(path.join(workspace, 'tasks'), { recursive: true }); + symlinkSync(outside, path.join(workspace, 'tasks')); + } else { + writeTasks(workspace, { '001.json': json(task('Valid prefix')) }); + symlinkSync(path.join(outside, '001.json'), path.join(workspace, 'tasks', '002.json')); + } + }); + await assert.rejects( + tryGeneratePlanWithFiles({ ...baseOptions, fullContext: 'Plan it', routingSession: agent.session as never }), + /without symlinks/, + ); + assert.equal(existsSync(path.join(outside, '001.json')), true, 'cleanup leaves the outside directory untouched'); + assert.deepEqual(readdirSync(workspaceRoot), []); + } finally { + rmSync(outside, { recursive: true, force: true }); + } + }); + } + + test('rejects an empty task directory even when execution reports success', async () => { + const agent = fakeAgent(() => undefined); + await assert.rejects( + tryGeneratePlanWithFiles({ ...baseOptions, fullContext: 'Plan it', routingSession: agent.session as never }), + /wrote no task files/, + ); + }); + + test('rejects unsuccessful generation in the single-file mode too', async () => { + const agent = fakeAgent(workspace => writeFileSync(path.join(workspace, 'plan.json'), json([task('Prefix')])), { success: false }); + await assert.rejects( + runPlanFileAgent({ ...baseOptions, purpose: 'generation', prompt: 'Plan it', executionType: 'plan-generation', routingSession: agent.session as never }), + /Plan generation did not finish/, + ); + }); + + test('an agent that ran but wrote nothing is a planning failure, not an unavailable agent', async () => { + const agent = fakeAgent(() => undefined, { success: false, error: 'model refused' }); + await assert.rejects( + runPlanFileAgent({ ...baseOptions, purpose: 'generation', prompt: 'Plan it', taskFiles: true, executionType: 'plan-generation', routingSession: agent.session as never }), + (error: Error) => !(error instanceof PlanFileAgentUnavailableError) && /did not finish\. The agent reported: model refused/.test(error.message), + ); + }); + + test('an agent task that throws is unavailable, except for usage limits, which keep their type', async () => { + const throwing = (error: Error) => ({ executeTask: async () => { throw error; } }); + await assert.rejects( + runPlanFileAgent({ ...baseOptions, purpose: 'generation', prompt: 'x', taskFiles: true, executionType: 'plan-generation', routingSession: throwing(new Error('docker unavailable')) as never }), + (error: Error) => error instanceof PlanFileAgentUnavailableError && /docker unavailable/.test(error.message), + ); + const usageLimit = Object.assign(new Error('limit reached'), { name: 'UsageLimitError' }); + await assert.rejects( + runPlanFileAgent({ ...baseOptions, purpose: 'generation', prompt: 'x', taskFiles: true, executionType: 'plan-generation', routingSession: throwing(usageLimit) as never }), + (error: Error) => error === usageLimit, + ); + }); +}); + +describe('file-based plan generation', () => { + test('selects files unless the reply mode is chosen explicitly', () => { + assert.equal(resolvePlanGenerationMode({}), 'file'); + assert.equal(resolvePlanGenerationMode({ PROPR_PLAN_GENERATION_MODE: 'file' }), 'file'); + assert.equal(resolvePlanGenerationMode({ PROPR_PLAN_GENERATION_MODE: ' Response ' }), 'response'); + }); + + test('keeps the planner prompt and replaces its reply format with the file contract', () => { + const prompt = buildPlanFilePrompt('everything'); + assert.ok(prompt.startsWith('everything')); + assert.match(prompt, /tasks\/001\.json/); + assert.match(prompt, /node validate-plan\.mjs --tasks tasks/); + assert.match(prompt, /Do not put the plan in your reply/); + }); + + test('returns null in reply mode without running an agent', async () => { + process.env.PROPR_PLAN_GENERATION_MODE = 'response'; + const agent = fakeAgent(() => { throw new Error('must not run'); }); + assert.equal(await tryGeneratePlanWithFiles({ ...baseOptions, fullContext: 'x', routingSession: agent.session as never }), null); + assert.equal(agent.calls.length, 0); + }); + + test('callLLMForPlan uses the files and still enforces granularity', async () => { + const agent = fakeAgent((workspace, options) => { + assert.match(options.prompt, /^Generate a plan/); + writeTasks(workspace, { '001.json': json(task('A')), '002.json': json(task('B')) }); + }); + const result = await callLLMForPlan({ + ...baseOptions, runId: 'run-1', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'single', routingSession: agent.session as never, + }); + assert.equal(agent.calls.length, 1); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 0, 'the reply path is not used'); + assert.equal(result.plan.length, 1, 'single granularity merges the written tasks'); + assert.equal(result.enforcementMetadata.enforced, true); + }); + + test('rejects an unsuccessful execution even after it wrote a valid task', async () => { + const agent = fakeAgent(async workspace => { + writeTasks(workspace, { '001.json': json(task('Completed prefix')) }); + await Promise.resolve(); + }, { success: false, error: 'execution failed' }); + await assert.rejects(callLLMForPlan({ + ...baseOptions, runId: 'run-1', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession: agent.session as never, + }), /Plan generation did not finish\. The agent reported: execution failed/); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 0); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + + const partialOutputs = { + 'valid task': (workspace: string) => writeTasks(workspace, { '001.json': json(task('Prefix')) }), + 'invalid task': (workspace: string) => writeTasks(workspace, { '001.json': '{' }), + 'empty task': (workspace: string) => writeTasks(workspace, { '001.json': '' }), + 'temporary task': (workspace: string) => writeTasks(workspace, { '001.json.tmp': '{' }), + 'assembled plan': (workspace: string) => writeFileSync(path.join(workspace, 'plan.json'), json([task('Plan')])), + 'uninspectable task directory': (workspace: string) => { + rmSync(path.join(workspace, 'tasks'), { recursive: true }); + writeFileSync(path.join(workspace, 'tasks'), 'incomplete output'); + }, + }; + for (const [kind, write] of Object.entries(partialOutputs)) { + test(`does not fall back after execution throws with ${kind} output`, async () => { + const agent = fakeAgent(async workspace => { + write(workspace); + await Promise.resolve(); + throw new Error('execution crashed'); + }); + await assert.rejects(callLLMForPlan({ + ...baseOptions, runId: 'run-1', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession: agent.session as never, + }), (error: Error) => error instanceof PlanningFailedError && !(error instanceof PlanFileAgentUnavailableError) && /execution crashed/.test(error.message)); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 0); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + } + + test('preserves usage limits after producing output', async () => { + const usageLimit = Object.assign(new Error('limit reached'), { name: 'UsageLimitError' }); + const agent = fakeAgent(workspace => { + writeTasks(workspace, { '001.json': json(task('Prefix')) }); + throw usageLimit; + }); + await assert.rejects( + tryGeneratePlanWithFiles({ ...baseOptions, fullContext: 'Plan it', routingSession: agent.session as never }), + (error: Error) => error === usageLimit, + ); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + + const terminalFailures = { + ExecutionAbortedError: { name: 'ExecutionAbortedError' }, + IndexingCancelledError: { name: 'IndexingCancelledError' }, + SecurityException: { name: 'SecurityException' }, + ContextTokenLimitError: { name: 'ContextTokenLimitError' }, + 'invalid configuration': { code: 'INVALID_CONFIGURATION' }, + 'explicit cancellation reason': { terminationReason: 'user_cancelled' }, + 'explicit cancellation message': { message: 'Task was cancelled' }, + }; + for (const [kind, details] of Object.entries(terminalFailures)) { + test(`preserves ${kind} without routing retry or response fallback`, async () => { + const terminal = Object.assign(new Error('terminal failure'), details); + const agent = fakeAgent(async (_workspace, options) => { + assert.equal(options.taskId, baseOptions.draftId); + await Promise.resolve(); + throw terminal; + }); + const replacement = fakeAgent(() => undefined); + await assert.rejects(callLLMForPlan({ + ...baseOptions, runId: 'run-1', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession: routedAgents(agent, replacement), + }), (error: unknown) => error === terminal); + assert.equal(agent.calls.length, 1); + assert.equal(replacement.calls.length, 0); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 0); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + } + + for (const kind of ['no output', 'valid task', 'uninspectable task directory']) { + test(`preserves direct execution cancellation with ${kind}`, async () => { + const cancellation = Object.assign(new Error('stopped'), { name: 'ExecutionAbortedError' }); + const agent = fakeAgent(async workspace => { + if (kind !== 'no output') partialOutputs[kind as keyof typeof partialOutputs](workspace); + await Promise.resolve(); + throw cancellation; + }); + await assert.rejects(callLLMForPlan({ + ...baseOptions, runId: 'run-1', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession: agent.session as never, + }), (error: unknown) => error === cancellation); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 0); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + } + + test('preserves cancellation from the shared file repair runner', async () => { + process.env.PROPR_PLAN_GENERATION_MODE = 'response'; + replies.push('[{"title":"Broken","body":"Body" "implementation":"Fix"}]'); + const cancellation = Object.assign(new Error('stopped'), { name: 'ExecutionAbortedError' }); + const repair = fakeAgent(async workspace => { + assert.equal(existsSync(path.join(workspace, 'plan.json')), true); + await Promise.resolve(); + throw cancellation; + }); + await assert.rejects(callLLMForPlan({ + ...baseOptions, runId: 'run-1', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', repairRoutingSession: repair.session as never, + }), (error: unknown) => error === cancellation); + assert.equal(repair.calls.length, 1); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 1); + assert.deepEqual(readdirSync(workspaceRoot), []); + }); + + test('an unqualified transport abort still permits response fallback', async () => { + replies.push(json([task('From reply')])); + const agent = fakeAgent(() => { + throw Object.assign(new Error('transport interrupted'), { name: 'AbortError', code: 'ABORT_ERR' }); + }); + const result = await callLLMForPlan({ + ...baseOptions, runId: 'run-1', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession: agent.session as never, + }); + assert.deepEqual(result.plan.map(item => item.title), ['From reply']); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 1); + }); + + for (const failure of ['throw', 'result']) { + test(`response fallback starts a fresh routing call after all file members fail by ${failure}`, async () => { + const fileCalls: string[] = []; + const responseCalls: AnalyzeOptions[] = []; + const agents = ['first', 'second'].map(alias => ({ + config: { alias, enabled: true, supportedModels: ['opus'] }, + executeTask: async () => { + fileCalls.push(alias); + await Promise.resolve(); + if (failure === 'throw') throw new Error('temporary transport failure'); + return { success: false, error: 'temporary transport failure' }; + }, + analyze: async (_prompt: string, options: AnalyzeOptions) => { + assert.deepEqual(fileCalls, ['first', 'second'], 'file routing exhausted before response generation'); + assert.deepEqual(readdirSync(workspaceRoot), [], 'all empty file workspaces have been cleaned up'); + responseCalls.push(options); + // Response generation also retains its own normal routing failover. + if (alias === 'first') throw new Error('temporary response transport failure'); + return { success: true, response: json([task('From fresh route')]), modelUsed: 'opus', executionTimeMs: 1 }; + }, + })); + const router = new SyntheticRoutingService({ + loadSyntheticConfigs: async () => [{ + id: 'pool', alias: 'pool', enabled: true, defaultModel: 'smart', + models: [{ + id: 'smart', enabled: true, strategy: 'usage_based', + members: agents.map((agent, index) => ({ + id: agent.config.alias, directAgentAlias: agent.config.alias, model: 'opus', enabled: true, priority: 100 - index, + })), + }], + }], + getDirectAgent: alias => agents.find(agent => agent.config.alias === alias) as Agent | undefined, + }); + // Persistence is outside this regression; selection, retries and forks are real. + mock.method(router, 'recordAttempt', async () => null); + const routingSession = router.begin({ requestedAgentAlias: 'pool', requestedModel: 'smart', requiredTokens: 50_000 }); + const result = await callLLMForPlan({ + ...baseOptions, runId: 'run-fallback', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession, + }); + assert.deepEqual(result.plan.map(item => item.title), ['From fresh route']); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 1); + const responseSession = runLightweightLLMAnalysis.mock.calls[0].arguments[0].routingSession!; + assert.notEqual(responseSession, routingSession); + assert.notEqual(responseSession.callId, routingSession.callId); + assert.equal(responseSession.requiredTokens, routingSession.requiredTokens); + assert.equal(responseCalls.length, 2); + assert.deepEqual(responseCalls.map(options => { + const routing = options.metadata?.syntheticRouting as { callId: string; attemptNumber: number }; + return { callId: routing.callId, attemptNumber: routing.attemptNumber }; + }), [ + { callId: responseSession.callId, attemptNumber: 1 }, + { callId: responseSession.callId, attemptNumber: 2 }, + ]); + assert.deepEqual([...routingSession.attemptedMembers], ['first', 'second']); + await assert.rejects(routingSession.select(), SyntheticPoolExhaustedError); + }); + } + + test('explicit response mode retains the supplied routing session', async () => { + process.env.PROPR_PLAN_GENERATION_MODE = 'response'; + replies.push(json([task('From reply')])); + const agent = fakeAgent(() => { throw new Error('must not run'); }); + const fork = mock.method(agent.session, 'fork'); + const result = await callLLMForPlan({ + ...baseOptions, runId: 'run-response', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession: agent.session as never, + }); + assert.deepEqual(result.plan.map(item => item.title), ['From reply']); + assert.equal(agent.calls.length, 0); + assert.equal(fork.mock.callCount(), 0); + assert.equal(runLightweightLLMAnalysis.mock.calls[0].arguments[0].routingSession, agent.session); + }); + + test('falls back to the model reply only when the agent could not run', async () => { + replies.push(json([task('From reply')])); + const unavailable = { fork: () => ({}), executeTask: async () => { throw new Error('no container runtime'); } }; + const result = await callLLMForPlan({ + ...baseOptions, runId: 'run-1', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession: unavailable as never, + }); + assert.deepEqual(result.plan.map(item => item.title), ['From reply']); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 1); + + const invalid = fakeAgent(workspace => writeTasks(workspace, { '001.json': json({ title: 'Stub' }) })); + await assert.rejects(callLLMForPlan({ + ...baseOptions, runId: 'run-2', fullContext: 'Generate a plan', worktreePath: '/tmp/worktree', tokenLimit: 100_000, + repairModel: 'codex:gpt', granularity: 'balanced', routingSession: invalid.session as never, + }), /Plan generation did not produce a valid plan: task 1 \(tasks\/001\.json\) has no non-empty "body"/); + assert.equal(runLightweightLLMAnalysis.mock.callCount(), 1, 'an invalid plan is not regenerated as a reply'); + }); +}); diff --git a/test/planGenerationJsonRepair.test.ts b/test/planGenerationJsonRepair.test.ts index c7afd9ad2..f63d3e39f 100644 --- a/test/planGenerationJsonRepair.test.ts +++ b/test/planGenerationJsonRepair.test.ts @@ -1,6 +1,9 @@ import assert from 'node:assert/strict'; import { beforeEach, mock, test } from 'node:test'; +// These tests cover parsing the model reply; file-based generation has its own tests. +process.env.PROPR_PLAN_GENERATION_MODE = 'response'; + type AnalysisCall = { prompt: string; model: string; @@ -31,8 +34,9 @@ const runPlanFileAgent = mock.fn(async (options: RepairCall) => { if (repairResult instanceof Error) throw repairResult; return repairResult; }); +class PlanFileAgentUnavailableError extends Error {} await mock.module('../packages/core/src/services/taskPlanning/planFileAgent.js', { - namedExports: { runPlanFileAgent }, + namedExports: { runPlanFileAgent, PlanFileAgentUnavailableError }, }); const correlatedLogger = {