import { finalizeFindings } from "./findings.ts"; import { loadPrompt, type PromptName } from "./prompt-loader.ts"; import { parseReviewerOutput, parseSummary, parseTriage, parseValidatorOutput, } from "./results.ts"; import type { AgentOutcome, Finding, ModelSelection, PullRequestSnapshot, ReviewerOutput, ReviewResult, ValidatorOutput, } from "./types.ts"; export interface AgentExecution { prompt: string; task: string; selection: ModelSelection; signal?: AbortSignal; tools?: readonly string[]; attempt?: number; label?: string; } export type AgentExecutor = ( execution: AgentExecution, ) => Promise; export type ReviewStage = | "triage" | "summary" | "reviewers" | "validation" | "aggregation"; export type ReviewProgressEvent = | { type: "stage"; stage: ReviewStage; total?: number } | { type: "triage-complete"; review: boolean; reason: string } | { type: "triage-fallback"; reason: string } | { type: "summary-complete"; summary: string } | { type: "reviewer-coverage-retry"; reviewer: string; missing: number; unexpected: number; incomplete: boolean; } | { type: "reviewer-recovery-retry"; reviewer: string; attempt: number } | { type: "reviewer-complete"; completed: number; total: number } | { type: "reviewers-complete"; total: number; candidates: number } | { type: "validator-complete"; completed: number; total: number } | { type: "validation-complete"; total: number; accepted: number } | { type: "aggregation-complete"; findings: number } | { type: "failed"; stage: ReviewStage; message: string } | { type: "aborted"; stage: ReviewStage }; export interface OrchestratorRequest { snapshot: PullRequestSnapshot; manifestPath: string; selection: ModelSelection; execute: AgentExecutor; signal?: AbortSignal; onProgress?: (event: ReviewProgressEvent) => void | Promise; } const TRIAGE_CONTRACT = `{"schemaVersion":1,"review":true|false,"reason":"..."}`; const SUMMARY_CONTRACT = `{"schemaVersion":1,"summary":"..."}`; const REVIEWER_CONTRACT = `{"schemaVersion":1,"coverage":{"complete":boolean,"reviewedFiles":["path"]},"findings":[Finding]}. Finding is either {"kind":"code","id":"...","category":"instructions|bug|regression|permission|security|maintenance|goal|scope","title":"...","description":"...","evidence":"...","introducedByPr":"...","path":"...","line":1,"side":"LEFT|RIGHT","startLine":1 optional,"suggestion":"..." optional} or {"kind":"task-omission", same base fields, "requirement":"...", "path":"..." optional}.`; const VALIDATOR_CONTRACT = `{"schemaVersion":1,"findingId":"...","valid":boolean,"confidence":0-100 integer,"evidence":"...","rejectionReason":"..." optional,"finding":the complete corrected Finding}.`; const TRIAGE_FALLBACK_REASON = "Triage did not return structured output; continuing with review by default."; const REVIEWER_RECOVERY_MAX_ATTEMPTS = 2; class InvalidStageOutputError extends Error { constructor( readonly stage: PromptName, message: string, ) { super(message); this.name = "InvalidStageOutputError"; } } export interface ExecutionDiagnostics { stage: ReviewStage | PromptName; label?: string; requestedModel: string; thinking: string; attempt: number; piRetryAttempts: number; stopReason?: string; errorMessage?: string; partialText?: string; } /** * Raised when an isolated agent execution terminally fails (terminated, * aborted, no output, model mismatch). Carries the diagnostics needed to * produce an attributable failure report. */ export class AgentExecutionError extends Error { constructor( readonly diagnostics: ExecutionDiagnostics, message: string, ) { super(message); this.name = "AgentExecutionError"; } } function requestedProviderModel(selection: ModelSelection): string { return selection.model; } function resolvedProviderModel(resolved: { provider: string; id: string; }): string { return `${resolved.provider}/${resolved.id}`; } function validateResolvedModel( selection: ModelSelection, outcome: AgentOutcome, options: ExecutionOptions, ): void { if (!outcome.resolvedModel) return; const requested = requestedProviderModel(selection); const resolved = resolvedProviderModel(outcome.resolvedModel); if (requested !== resolved) { throw new AgentExecutionError( { stage: options.stage, label: options.label, requestedModel: requested, thinking: selection.thinking, attempt: options.attempt ?? 1, piRetryAttempts: outcome.piRetryAttempts, stopReason: outcome.stopReason, errorMessage: `model mismatch: requested ${requested}, resolved ${resolved}`, partialText: outcome.ok ? outcome.text : outcome.partialText, }, `model mismatch: requested ${requested}, resolved ${resolved} (context ${outcome.resolvedModel.contextWindow}, maxTokens ${outcome.resolvedModel.maxTokens})`, ); } } async function reportProgress( request: OrchestratorRequest, event: ReviewProgressEvent, ): Promise { try { await request.onProgress?.(event); } catch { // Progress is observational and must never affect review semantics. } } interface ExecutionOptions { stage: ReviewStage | PromptName; label?: string; attempt?: number; } async function runExecutor( request: OrchestratorRequest, prompt: string, task: string, options: ExecutionOptions, tools?: readonly string[], ): Promise { const outcome = await request.execute({ prompt, task, selection: request.selection, signal: request.signal, tools, attempt: options.attempt ?? 1, label: options.label, }); validateResolvedModel(request.selection, outcome, options); if (!outcome.ok) { throw new AgentExecutionError( { stage: options.stage, label: options.label, requestedModel: requestedProviderModel(request.selection), thinking: request.selection.thinking, attempt: options.attempt ?? 1, piRetryAttempts: outcome.piRetryAttempts, stopReason: outcome.stopReason, errorMessage: outcome.errorMessage, partialText: outcome.partialText, }, outcome.errorMessage ?? `Pi stopped: ${outcome.stopReason}`, ); } return outcome.text; } async function executeParsed( request: OrchestratorRequest, promptName: PromptName, task: string, contract: string, parse: (text: string) => T, options: ExecutionOptions = { stage: promptName }, ): Promise { const prompt = await loadPrompt(promptName); const output = await runExecutor(request, prompt, task, options); try { return parse(output); } catch (error) { const repairPrompt = await loadPrompt("repair-json"); const repaired = await runExecutor( request, repairPrompt, `Required JSON contract:\n\n${contract}\n\n\nInvalid prior response:\n\n${output}\n\n\nParser error: ${error instanceof Error ? error.message : "invalid output"}`, options, [], ); try { return parse(repaired); } catch (repairError) { const preview = repaired.replace(/\s+/g, " ").slice(0, 500); throw new InvalidStageOutputError( promptName, `${promptName} returned invalid JSON after repair: ${repairError instanceof Error ? repairError.message : "invalid output"}. Preview: ${preview || "(empty)"}`, ); } } } export function expectedChangedFiles(snapshot: PullRequestSnapshot): string[] { return [...new Set(snapshot.files.map((file) => file.path))].sort(); } interface CoverageMismatch { incomplete: boolean; missingFiles: string[]; unexpectedFiles: string[]; } function compareCoverage( output: ReviewerOutput, expectedFiles: string[], ): CoverageMismatch | undefined { const actual = [...new Set(output.coverage.reviewedFiles)].sort(); const expectedSet = new Set(expectedFiles); const actualSet = new Set(actual); const mismatch = { incomplete: !output.coverage.complete, missingFiles: expectedFiles.filter((file) => !actualSet.has(file)), unexpectedFiles: actual.filter((file) => !expectedSet.has(file)), }; return mismatch.incomplete || mismatch.missingFiles.length > 0 || mismatch.unexpectedFiles.length > 0 ? mismatch : undefined; } function coverageInstructions(expectedFiles: string[]): string { return `Expected changed files (copy these paths exactly into coverage.reviewedFiles after inspecting every file):\n\n${JSON.stringify(expectedFiles, null, 2)}\n\nInspect documentation, tests, changelogs, and AGENTS.md files as well as source files. If every listed file was inspected, return complete:true and the exact array. Otherwise return complete:false and only paths actually inspected. Do not add ./ prefixes, absolute paths, aliases, or previous rename paths.`; } function reviewerRepairContract(expectedFiles: string[]): string { return `${REVIEWER_CONTRACT}\n\n${coverageInstructions(expectedFiles)}\nDuring JSON repair, use complete:true with the exact expected array only when the invalid prior response explicitly says every expected changed file was reviewed. Otherwise use complete:false and include only supported paths.`; } function coverageFailure( reviewer: string, mismatch: CoverageMismatch, selection: ModelSelection, attempt: number, ): AgentExecutionError { const details: string[] = []; if (mismatch.incomplete) details.push("Reviewer reported complete:false."); if (mismatch.missingFiles.length > 0) details.push(`Missing files: ${mismatch.missingFiles.join(", ")}`); if (mismatch.unexpectedFiles.length > 0) details.push(`Unexpected files: ${mismatch.unexpectedFiles.join(", ")}`); const message = `${reviewer} coverage mismatch after retry. ${details.join(" ")}`; return new AgentExecutionError( { stage: "reviewers", label: reviewer, requestedModel: requestedProviderModel(selection), thinking: selection.thinking, attempt, piRetryAttempts: 0, stopReason: "coverage", errorMessage: message, }, message, ); } interface ReviewerFailure { label: string; diagnostics: ExecutionDiagnostics; message: string; } interface ReviewerBatchResult { reviewers: ReviewerOutput[]; failures: ReviewerFailure[]; completed: number; } function isRecoverableTermination(error: unknown): boolean { if (!(error instanceof AgentExecutionError)) return false; const { stopReason, errorMessage } = error.diagnostics; if (stopReason === "aborted") return false; if (errorMessage?.includes("model mismatch")) return false; if (errorMessage?.includes("coverage mismatch")) return false; return stopReason === "error" || stopReason === "terminated"; } /** * Run all reviewers in parallel with full isolation. A single reviewer's * transient termination is recovered once with the same model, thinking, * prompt, tools, manifest, and coverage contract. Other reviewers keep * running and their results are preserved. Partial output is never used; it * stays on the failure diagnostics for reporting only. */ async function runReviewers( request: OrchestratorRequest, specs: Array<{ name: PromptName; label: string }>, baseTask: string, summary: string, coverageTask: string, repairContract: string, expectedFiles: string[], ): Promise { let completed = 0; const settled = await Promise.allSettled( specs.map(async ({ name, label }) => { try { const reviewer = await runSingleReviewer( request, name, label, baseTask, summary, coverageTask, repairContract, expectedFiles, ); completed += 1; await reportProgress(request, { type: "reviewer-complete", completed, total: specs.length, }); return reviewer; } catch (error) { if ( error instanceof AgentExecutionError && isRecoverableTermination(error) && error.diagnostics.attempt < REVIEWER_RECOVERY_MAX_ATTEMPTS ) { await reportProgress(request, { type: "reviewer-recovery-retry", reviewer: label, attempt: error.diagnostics.attempt + 1, }); const recovered = await runSingleReviewer( request, name, label, baseTask, summary, coverageTask, repairContract, expectedFiles, error.diagnostics.attempt + 1, ); completed += 1; await reportProgress(request, { type: "reviewer-complete", completed, total: specs.length, }); return recovered; } throw error; } }), ); const reviewers: ReviewerOutput[] = []; const failures: ReviewerFailure[] = []; for (const result of settled) { if (result.status === "fulfilled") { reviewers.push(result.value); } else { const error = result.reason; if (error instanceof AgentExecutionError) { failures.push({ label: error.diagnostics.label ?? "reviewer", diagnostics: error.diagnostics, message: error.message, }); } else if (error instanceof InvalidStageOutputError) { failures.push({ label: "reviewer", diagnostics: { stage: error.stage, requestedModel: requestedProviderModel(request.selection), thinking: request.selection.thinking, attempt: 1, piRetryAttempts: 0, errorMessage: error.message, }, message: error.message, }); } else { const message = error instanceof Error ? error.message : "Unknown reviewer failure"; failures.push({ label: "reviewer", diagnostics: { stage: "reviewers", requestedModel: requestedProviderModel(request.selection), thinking: request.selection.thinking, attempt: 1, piRetryAttempts: 0, errorMessage: message, }, message, }); } } } return { reviewers, failures, completed: reviewers.length }; } async function runSingleReviewer( request: OrchestratorRequest, name: PromptName, label: string, baseTask: string, summary: string, coverageTask: string, repairContract: string, expectedFiles: string[], attempt = 1, ): Promise { const reviewerTask = `${baseTask}\n\nReviewer instance: ${label}\n\nPR summary: ${summary}\n\n${coverageTask}\n\nReturn JSON matching this contract:\n${REVIEWER_CONTRACT}`; let reviewer = await executeParsed( request, name, reviewerTask, repairContract, parseReviewerOutput, { stage: "reviewers", label, attempt }, ); let mismatch = compareCoverage(reviewer, expectedFiles); if (mismatch) { await reportProgress(request, { type: "reviewer-coverage-retry", reviewer: label, missing: mismatch.missingFiles.length, unexpected: mismatch.unexpectedFiles.length, incomplete: mismatch.incomplete, }); const retryTask = `${reviewerTask}\n\nCoverage retry required. Inspect any missing files, re-check the complete expected list, preserve still-valid findings, and return one complete replacement reviewer result.\n\nPrevious structured reviewer output:\n\n${JSON.stringify(reviewer)}\n\n\nMissing files:\n${JSON.stringify(mismatch.missingFiles)}\n\nUnexpected files:\n${JSON.stringify(mismatch.unexpectedFiles)}\n\nPrevious complete flag: ${reviewer.coverage.complete}`; reviewer = await executeParsed( request, name, retryTask, repairContract, parseReviewerOutput, { stage: "reviewers", label, attempt }, ); mismatch = compareCoverage(reviewer, expectedFiles); if (mismatch) throw coverageFailure(label, mismatch, request.selection, attempt); } return reviewer; } function formatReviewerFailures( failures: ReviewerFailure[], completed: number, total: number, ): string { const lines = failures.map((failure) => { const d = failure.diagnostics; const parts = [ `- reviewer: ${failure.label}`, `- model: ${d.requestedModel}`, `- thinking: ${d.thinking}`, `- attempt: ${d.attempt}`, `- stop reason: ${d.stopReason ?? "unknown"}`, `- Pi retry attempts: ${d.piRetryAttempts}`, `- error: ${d.errorMessage ?? failure.message}`, ]; return parts.join("\n"); }); return `Review incomplete during independent review\n${lines.join("\n")}\n- reviewers completed: ${completed}/${total}\n- failed reviewer output excluded from review results`; } async function mapLimited( values: T[], limit: number, mapper: (value: T) => Promise, ): Promise { const results: R[] = new Array(values.length); let next = 0; const workers = Array.from( { length: Math.min(limit, values.length) }, async () => { while (next < values.length) { const index = next++; results[index] = await mapper(values[index]); } }, ); await Promise.all(workers); return results; } function validatorPrompt(finding: Finding): PromptName { if (finding.category === "instructions") return "validate-instructions"; if (finding.category === "goal" || finding.category === "scope") return "validate-objective"; return "validate-bug"; } export async function runReview( request: OrchestratorRequest, ): Promise { let activeStage: ReviewStage = "triage"; try { const baseTask = `Snapshot manifest: ${request.manifestPath}. Read every referenced file needed for your role.`; await reportProgress(request, { type: "stage", stage: activeStage }); let triage: ReturnType; try { triage = await executeParsed( request, "triage", baseTask, TRIAGE_CONTRACT, parseTriage, ); await reportProgress(request, { type: "triage-complete", review: triage.review, reason: triage.reason, }); } catch (error) { if (!(error instanceof InvalidStageOutputError)) throw error; triage = { review: true, reason: TRIAGE_FALLBACK_REASON }; await reportProgress(request, { type: "triage-fallback", reason: TRIAGE_FALLBACK_REASON, }); } if (!triage.review) { return { status: "skipped", snapshot: request.snapshot, findings: [], reason: triage.reason, failures: [], }; } activeStage = "summary"; await reportProgress(request, { type: "stage", stage: activeStage }); const summary = await executeParsed( request, "summary", baseTask, SUMMARY_CONTRACT, parseSummary, ); await reportProgress(request, { type: "summary-complete", summary, }); const reviewerSpecs: Array<{ name: PromptName; label: string }> = [ { name: "instructions-review", label: "instructions-review #1" }, { name: "instructions-review", label: "instructions-review #2" }, { name: "bug-review", label: "bug-review #1" }, { name: "bug-review", label: "bug-review #2" }, { name: "objective-review", label: "objective-review #1" }, ]; const expectedFiles = expectedChangedFiles(request.snapshot); const coverageTask = coverageInstructions(expectedFiles); const repairContract = reviewerRepairContract(expectedFiles); activeStage = "reviewers"; await reportProgress(request, { type: "stage", stage: activeStage, total: reviewerSpecs.length, }); const reviewerResults = await runReviewers( request, reviewerSpecs, baseTask, summary, coverageTask, repairContract, expectedFiles, ); if (reviewerResults.failures.length > 0) { const completed = reviewerResults.completed; const message = formatReviewerFailures( reviewerResults.failures, completed, reviewerSpecs.length, ); await reportProgress(request, { type: "failed", stage: activeStage, message, }); return { status: "incomplete", snapshot: request.snapshot, findings: [], reason: "A required review stage failed.", failures: [message], }; } const reviewers = reviewerResults.reviewers; const candidates = reviewers.flatMap((reviewer) => reviewer.findings); await reportProgress(request, { type: "reviewers-complete", total: reviewerSpecs.length, candidates: candidates.length, }); activeStage = "validation"; await reportProgress(request, { type: "stage", stage: activeStage, total: candidates.length, }); let validatorCompleted = 0; const validations = await mapLimited(candidates, 4, async (finding) => { const validation = await executeParsed( request, validatorPrompt(finding), `${baseTask}\n\nCandidate finding:\n${JSON.stringify(finding)}\n\n${VALIDATOR_CONTRACT}`, VALIDATOR_CONTRACT, parseValidatorOutput, ); if ( validation.findingId !== finding.id || validation.finding.id !== finding.id ) { throw new Error(`Validator changed finding identity ${finding.id}`); } validatorCompleted += 1; await reportProgress(request, { type: "validator-complete", completed: validatorCompleted, total: candidates.length, }); return validation; }); const accepted = validations.filter( (validation) => validation.valid && validation.confidence >= 80 && validation.findingId === validation.finding.id, ).length; await reportProgress(request, { type: "validation-complete", total: candidates.length, accepted, }); activeStage = "aggregation"; await reportProgress(request, { type: "stage", stage: activeStage }); const findings = finalizeFindings(validations, request.snapshot.diff); await reportProgress(request, { type: "aggregation-complete", findings: findings.length, }); return { status: "complete", snapshot: request.snapshot, findings, failures: [], }; } catch (error) { if (request.signal?.aborted) { await reportProgress(request, { type: "aborted", stage: activeStage }); return { status: "aborted", snapshot: request.snapshot, findings: [], failures: [], }; } const message = error instanceof Error ? error.message : "Unknown review failure"; await reportProgress(request, { type: "failed", stage: activeStage, message, }); return { status: "incomplete", snapshot: request.snapshot, findings: [], reason: "A required review stage failed.", failures: [message], }; } } export type { ValidatorOutput };