import { ReviewExecutionError, } from "./types.ts"; import type { ExtensionContext } from "@earendil-works/pi-coding-agent"; import { buildClassifierTranscript, parseDecision, type ModelDecision, } from "../policy.ts"; import type { BoundaryRequest, BoundaryReview, BoundaryReviewerContext, } from "../broker/index.ts"; import type { CompletionMessage, Config, ReviewAttemptObservation, ReviewAttemptStatus, ReviewErrorClass, ReviewExecutionSummary, ReviewResult, ReviewerMeta, ReviewerRuntime, ReviewerTelemetryEvent, } from "./types.ts"; import { FORMAT_RETRY_INSTRUCTION } from "./consts.ts"; import { applyReviewerInputBudget, preflightPart, reviewPreflight, sharedReviewContext, textFromAssistant, } from "./input.ts"; import { abortableDelay, abortableOperation, classifyProviderFailure, incrementError, isFormatError, isRetryableError, modelCall, normalizedStopReason, observedUsage, parseErrorClass, resolveApiKeyAndHeaders, resolveReviewerMeta, retryDelayMs, reviewerSessionId, type ProviderAttemptMetadata, type ReviewerResolver, } from "./provider.ts"; export async function complete( ctx: ExtensionContext, config: Config, request: BoundaryRequest, reviewerContext?: BoundaryReviewerContext, resolve: ReviewerResolver = resolveReviewerMeta, observe?: (event: ReviewerTelemetryEvent) => void, ): Promise { const started = Date.now(); const selectedTranscript = buildClassifierTranscript( ctx.sessionManager.buildContextEntries(), config, { ...request, trustedRetryOriginalRequestId: reviewerContext?.userOverride?.originalRequestId, }, ); const transcript = applyReviewerInputBudget( request, selectedTranscript, reviewerContext, config.maxReviewerInputTokens, ); const sharedContext = sharedReviewContext( request, transcript, reviewerContext, ); const preflight = reviewPreflight( request, transcript, reviewerContext, sharedContext, config.maxReviewerInputTokens, ); // The outer controller is cancelled only by the session, so it stays // aborted for the rest of the review. Each attempt adds its own controller // for the per-attempt `timeoutMs` budget. const controller = new AbortController(); const onSessionAbort = () => controller.abort(); if (ctx.signal?.aborted) controller.abort(); else ctx.signal?.addEventListener("abort", onSessionAbort, { once: true }); const attempts: ReviewAttemptObservation[] = []; const errorCounts: ReviewExecutionSummary["errorCounts"] = {}; const summary = (): ReviewExecutionSummary => ({ attempts, errorCounts, durationMs: Date.now() - started, transcript, preflight, }); try { // A budget preflight failure is a sizing estimate, not a safety verdict. // When a human explicitly authorized this exact retry, their decision // must not be vetoed by an estimator: proceed with the truncated // evidence and let the reviewer see the override. The failureCode stays // on the transcript for observability. if (transcript.failureCode && !reviewerContext?.userOverride) { incrementError(errorCounts, transcript.failureCode); throw new ReviewExecutionError( transcript.failureCode, summary(), ); } if (controller.signal.aborted) { incrementError(errorCounts, "abort"); throw new ReviewExecutionError("abort", summary()); } let meta: ReviewerMeta; try { meta = await abortableOperation(resolve(ctx, config), controller.signal); } catch { const errorClass = controller.signal.aborted ? "abort" : "model_resolution"; incrementError(errorCounts, errorClass); throw new ReviewExecutionError(errorClass, summary()); } let lastError: unknown; let lastErrorClass: ReviewErrorClass = "unknown"; const retryErrors: ReviewErrorClass[] = []; const maxAttempts = config.retries + 1; const formatRetryFitsBudget = preflight.total.estimatedTokens + preflightPart(`\n\n${FORMAT_RETRY_INSTRUCTION}`).estimatedTokens <= config.maxReviewerInputTokens; let formatRetry = false; for (let attempt = 1; attempt <= maxAttempts; attempt++) { // Every attempt owns a fresh `timeoutMs` budget that covers its // authentication resolution and its model call; a retry starts over // with the full budget. `retries` is the only bound on how many // attempts a review may spend. The outer controller still cancels every // attempt when the session aborts. const attemptController = new AbortController(); let attemptTimedOut = false; const attemptTimeout = setTimeout(() => { attemptTimedOut = true; attemptController.abort(); }, config.timeoutMs); const onOuterAbort = () => attemptController.abort(); if (controller.signal.aborted) attemptController.abort(); else controller.signal.addEventListener("abort", onOuterAbort, { once: true }); // A session abort outranks the per-attempt timeout, so a cancelled // review is reported as an abort rather than a retryable timeout. const cancelledClass = (): "abort" | "timeout" | undefined => controller.signal.aborted ? "abort" : attemptTimedOut ? "timeout" : undefined; const attemptStarted = Date.now(); let message: CompletionMessage | undefined; let status: ReviewAttemptStatus = "transport_failure"; let errorClass: ReviewErrorClass = "unknown"; let decision: ModelDecision | undefined; let runtime: ReviewerRuntime | undefined; const providerMetadata: ProviderAttemptMetadata = {}; try { const auth = await abortableOperation( resolveApiKeyAndHeaders(ctx, meta.model), attemptController.signal, ); if (cancelledClass()) { throw new Error("review attempt cancelled during authentication"); } runtime = { ...meta, auth, sessionId: reviewerSessionId(ctx, config, meta.model, auth), }; message = await abortableOperation( modelCall( runtime, config, attemptController, sharedContext, config.maxTokens, config.timeoutMs, formatRetry, providerMetadata, ), attemptController.signal, ); // Fail closed when the attempt was cancelled while the provider was // still resolving: a provider that ignores its abort signal must not // be able to authorize a decision delivered after the attempt timed // out or the session aborted. if (cancelledClass()) { throw new Error("review attempt cancelled before a decision"); } if (message.stopReason !== "stop") { if (message.stopReason === "aborted") { status = "abort"; errorClass = "abort"; } else { status = message.stopReason === "error" ? "transport_failure" : "non_stop"; errorClass = message.stopReason === "length" ? "output_limit" : message.stopReason === "error" ? classifyProviderFailure(message, undefined, providerMetadata) : "provider_stop"; } throw new Error("reviewer returned a non-stop response"); } const text = textFromAssistant(message); if (!text) { status = "format_error"; errorClass = "empty_output"; throw new Error("reviewer returned empty output"); } try { decision = parseDecision(text); } catch (error) { status = "format_error"; errorClass = parseErrorClass(error); throw error; } status = "success"; errorClass = "none"; } catch (error) { lastError = error; const cancelled = cancelledClass(); if (cancelled) { status = cancelled; errorClass = cancelled; } else if (!runtime) { errorClass = "authentication"; } else if (errorClass === "unknown") { errorClass = classifyProviderFailure( message, error, providerMetadata, ); } } finally { clearTimeout(attemptTimeout); controller.signal.removeEventListener("abort", onOuterAbort); } const usage = observedUsage(message); const delayMs = retryDelayMs(errorClass, providerMetadata); const willRetry = !decision && isRetryableError(errorClass) && (!isFormatError(errorClass) || formatRetryFitsBudget) && attempt < maxAttempts && !controller.signal.aborted && // A `Retry-After` beyond the cap is reported as an infinite delay; // never retry on it. Number.isFinite(delayMs); const observation: ReviewAttemptObservation = { attempt: attempts.length + 1, model: message?.responseModel || `${meta.model.provider}/${meta.model.id}`, status, errorClass, stopReason: normalizedStopReason(message?.stopReason), durationMs: Date.now() - attemptStarted, willRetry, usageAvailability: usage.availability, usage: usage.usage, }; attempts.push(observation); observe?.({ type: "review_attempt", requestId: request.id, surface: request.surface, ...observation, }); if (decision) { return { decision, attempts: attempts.length, retryErrors, durationMs: Date.now() - started, transcript, summary: summary(), }; } lastErrorClass = errorClass; retryErrors.push(errorClass); incrementError(errorCounts, errorClass); if (controller.signal.aborted) break; if (!willRetry) break; formatRetry = isFormatError(errorClass); try { await abortableDelay(delayMs, controller.signal); } catch { lastErrorClass = "abort"; incrementError(errorCounts, lastErrorClass); break; } } void lastError; throw new ReviewExecutionError(lastErrorClass, summary()); } finally { ctx.signal?.removeEventListener("abort", onSessionAbort); } } export function modelDecisionToBoundaryReview( decision: ModelDecision, ): BoundaryReview { return { outcome: decision.outcome, riskLevel: decision.risk_level, userAuthorization: decision.user_authorization, rationale: decision.rationale, }; } export function currentTurnScope(ctx: ExtensionContext): string { const userMessages = ctx.sessionManager .buildContextEntries() .filter((entry) => { if (!entry || typeof entry !== "object" || Array.isArray(entry)) { return false; } // SAFETY: session entries are runtime-shaped objects; inspect message only after this structural boundary. const message = (entry as unknown as Record).message; return ( Boolean(message) && typeof message === "object" && !Array.isArray(message) && (message as Record).role === "user" ); }).length; return `${ctx.sessionManager.getSessionId()}:${userMessages}`; } export function denialLabel( denial: { request: BoundaryRequest; review: BoundaryReview; }, index: number, ): string { const target = denial.request.resolvedPath ?? denial.request.path ?? denial.request.destination ?? denial.request.command ?? denial.request.toolName ?? denial.request.operation; const compact = String(target).replace(/\s+/g, " ").slice(0, 90); return `${index + 1}. ${denial.request.surface}: ${compact} — ${denial.review.rationale.slice(0, 70)}`; }