/** * Single generic auto-resume loop for #416 (owner scope = single LLM call). * Call-local retries, at most the effective autoResumeLimit times (injected once * per call by the caller, #422 — never re-read from disk inside the loop), * in-place (same runId/session). * Unifies presentation: intermediate attempts use dummyIo, only final Terminal is presented. * * Owner 2026-08-23: a dispatch that exits by throwing used to bypass the entire * retry mechanism (the throw escaped the while-loop before the count check ever * ran). Every exception is now retained whole, in place, and the ordinary retry * path continues: same budget, same call-local count semantics. No failure-type * classification — every thrown value is treated identically. */ import { constants as fsConstants } from "node:fs"; import { randomUUID } from "node:crypto"; import { open } from "node:fs/promises"; import { join } from "node:path"; import type { DurablePrincipal, DurablePrincipalAuthority, SessionCustomEntryAppender, } from "../host-contracts.ts"; import { AUTO_RESUME_LIMIT, describeErrorIdentity, acquireRunWriterLease, markRunTerminal, RunWriterLeaseHeldError, type RunWriterLease, } from "./run-lifecycle.ts"; import { parseAutoResumeLimit } from "./config.ts"; import { processCancelSignalName } from "./process-cancel.ts"; import { isLawfulTypedTerminalOutcome, formatTerminalResult, type TerminalArtifactRef, type TerminalResult, type TerminalRoleName } from "./terminal.ts"; import { attachRecordedSubmissions, ensureRealArtifactsDirectory, presentFailureTerminal, retainPackageFault, } from "./settlement.ts"; export { ensureRealArtifactsDirectory }; import type { CliIo } from "./cli-io.ts"; import { serializeThrownValue } from "../serialize-thrown-value.ts"; import { errorText } from "../unknown-value.ts"; const dummyIo: CliIo = { stdout: () => {}, stderr: () => {} }; /** * Persist run-state after a host-turn result, outside the retried dispatch try. * The host outcome already decided the terminal. This write only seals the run. */ export async function persistReturnedRunState( admitted: { runDirectory: string }, ): Promise { await markRunTerminal(admitted.runDirectory); } export async function presentTerminal(terminal: TerminalResult, io: CliIo, runDirectory: string, held?: readonly string[]): Promise { try { if (held !== undefined) { for (const value of held) io.stdout(value); } else if (terminal.roleOutcome.kind === "failure" || terminal.roleOutcome.kind === "no_receipt") { presentFailureTerminal(terminal, io); } else { io.stdout(formatTerminalResult(terminal)); } } catch (error) { // Presentation consumes a settled report; it never reopens dispatch or // replaces the host's cause, identity or open details. await retainPackageFault({ runDirectory, diagnostic: `terminal presentation failed beside host terminal: ${describeErrorIdentity(error)}`, error, stderr: (text) => io.stderr(text), }); } } /** * Best-effort finalization of the durable run state before an exception-path * synthetic failure terminal is returned (#426 review: analyst-ledger classifies * running runs as live — an exhausted invocation must not remain live * indefinitely). Finalization failure must not mask the real cause. */ async function finalizeExceptionRunBestEffort(runDirectory: string, io: CliIo): Promise { try { await markRunTerminal(runDirectory); } catch (error) { io.stderr( `run terminal-state finalization failed (best-effort continue): ${describeErrorIdentity(error)}\n`, ); } } export type AutoResumeDispatchResult = { exitCode: number; terminal?: TerminalResult; /** * Pre-turn settlement under the writer lease (station child exhausted). * Loop presents this result and must not redispatch (#840 父子不层叠). */ skipAutoResume?: true; /** * Host turn actually started. Absent on beforeDispatch / pre-turn settlement * so the loop retries the initial payload instead of a session resume. */ turnDispatched?: true; /** * Dispatch skipped run-state persist so this loop seam owns it. * Absent when dispatch already persisted or never produced a terminal. */ needsPersist?: true; }; /** * Thrown by dispatchPostAdmissionTurn when a failure happens inside its own * settlement authority (presentControlledFailure) after the host turn * genuinely started — never to fabricate a replacement terminal (ADR 0080: * one settlement disposition owner — settlement stays exactly presentControlledFailure * / settleFailureTerminalResult, or — once the retry budget is exhausted — * this loop's own dispatchExceptionFailureTerminal). Its only job is to carry * the "turn already started" fact across the throw boundary so this loop * selects a resume payload on the next attempt instead of replaying the * initial one (#840 r9 判词 class 1 boundary — 覆盖 executeTurn 已启动后至 * 返回带 turnDispatched 结果前的全部异常). `cause` is the true failure; * retention below serializes it whole via the standard Error.cause chain. */ export class TurnDispatchedFailure extends Error { override readonly name = "TurnDispatchedFailure"; constructor(cause: unknown) { super( `settlement failed after the host turn genuinely started: ${ errorText(cause) }`, { cause }, ); } } /** Session custom-entry type carrying the pointer to one dispatch error file. */ export const DISPATCH_ERROR_RETENTION_ENTRY_TYPE = "ak_run_dispatch_error_retention" as const; /** Cycle- and bigint-safe JSON replacer so serialization itself cannot drop data. */ function jsonSafeReplacer(): (key: string, value: unknown) => unknown { const seen = new WeakSet(); return (_key: string, value: unknown): unknown => { if (typeof value === "bigint") return `${value}n`; if (typeof value === "object" && value !== null) { if (seen.has(value)) return "[circular]"; seen.add(value); } return value; }; } /** * Retain one throwing dispatch attempt's complete exception as an independent * per-attempt file under the run's artifacts directory, then leave an * addressable pointer in the session principal (custom entry). Exclusive-create * open (O_EXCL) with a per-attempt unique name enforces 史必追加 (#419): a later * attempt can never overwrite an earlier attempt's file. */ /** * Hardened create-once JSON write shared by every durable artifact this loop * retains directly (dispatch-error dumps, lawful-persist-failure error/ * evidence records): O_EXCL (fail loud on a colliding name, never overwrite) * + O_NOFOLLOW where the platform provides it (a planted symlink is never * followed) — mirrors settlement.ts's own hardened artifact writers. */ export async function writeHardenedArtifactFile( artifactsDir: string, namePrefix: string, payload: Record, ): Promise { const filePath = join(artifactsDir, `${namePrefix}-${randomUUID()}.json`); const body = `${JSON.stringify(payload, jsonSafeReplacer(), 2)}\n`; const noFollowFlag = typeof fsConstants.O_NOFOLLOW === "number" ? fsConstants.O_NOFOLLOW : 0; const handle = await open( filePath, fsConstants.O_WRONLY | fsConstants.O_CREAT | fsConstants.O_EXCL | noFollowFlag, 0o600, ); try { await handle.writeFile(body, "utf8"); } finally { await handle.close(); } return filePath; } async function retainDispatchError( admitted: { runDirectory: string; principal: DurablePrincipal }, principalAuthority: DurablePrincipalAuthority, sessionAppender: SessionCustomEntryAppender, attempt: number, error: unknown, ): Promise<{ file: string; pointerError?: unknown }> { const artifactsDir = await ensureRealArtifactsDirectory(admitted.runDirectory); // Whole-object dump: everything the thrown value carries, nothing picked. const filePath = await writeHardenedArtifactFile(artifactsDir, `dispatch-error-attempt-${attempt}`, { version: 1, attempt, recordedAt: new Date().toISOString(), error: serializeThrownValue(error), }); // Addressable pointer in the dossier (卷宗): Pi session custom-entry codec // (appendPiSessionCustomEntry). Lease still owned here with run-writer. let pointerLease: RunWriterLease; try { pointerLease = await acquireRunWriterLease(admitted.runDirectory); } catch (error) { if (error instanceof RunWriterLeaseHeldError) return { file: filePath }; throw error; } // Pointer-stage failure (#426 fix_now #5) is separated from the file write: // once the error file is durably on disk, a failed session append must not // reject through here and orphan it — the file path is still handed back. let pointerError: unknown; try { const timestamp = new Date().toISOString(); await sessionAppender( principalAuthority, admitted.principal, DISPATCH_ERROR_RETENTION_ENTRY_TYPE, { version: 1, attempt, file: filePath, recordedAt: timestamp }, ); } catch (error) { pointerError = error; } finally { await pointerLease.release(); } return pointerError === undefined ? { file: filePath } : { file: filePath, pointerError }; } /** * Unwrap TurnDispatchedFailure before this loop's own final presentation * (#840 r9 判词 class 1 — the auto-resume.ts final presentation boundary). * The wrapper only carries the "turn already genuinely started" signal * across the throw boundary so the loop above selects a resume payload; the * typed decisiveFacts and diagnostic text below must name the real cause it * wraps (接住可以,洗白不行 — 未识别异常不得冒用具体标签,真因必须落痕), not * the internal signal's own identity. */ function unwrapTurnDispatchedFailure(error: unknown): unknown { let current = error; while (current instanceof TurnDispatchedFailure) { current = current.cause; } return current; } async function attachDispatchExceptionTerminal( admitted: { readonly runDirectory: string; readonly runId: string; readonly projectRoot: string; }, terminal: TerminalResult, io: CliIo, ): Promise { try { return await attachRecordedSubmissions( { projectRoot: admitted.projectRoot, runId: admitted.runId, runDirectory: admitted.runDirectory, }, terminal, ); } catch (error) { // Existing diagnostic seam: attach true cause on stderr, original Terminal stays. io.stderr( `dispatch exception ledger attach failed (best-effort continue): ${describeErrorIdentity(error)}\n`, ); return terminal; } } /** * Typed failure terminal for a retry path that ended with only exceptions: * loud, non-lawful, carrying the last true cause and the pointers to the * full per-attempt error files. Never rethrows the raw exception at callers. */ function dispatchExceptionFailureTerminal(input: { role: TerminalRoleName; runId: string; causeError: unknown; errorFiles: readonly string[]; autoResumeAttempts: number; endReason: string; /** True only when every attempt threw; otherwise describe just the final attempt. */ everyAttemptThrew: boolean; }): TerminalResult { const causeError = unwrapTurnDispatchedFailure(input.causeError); // #426 review: this terminal fires whenever the FINAL dispatch throws, not // only when every attempt threw — do not misrepresent a mixed retry history. const history = input.everyAttemptThrew ? "dispatch threw an exception on every attempt" : "the final dispatch threw an exception"; const diagnostic = `${history} (${input.endReason}; resumes used ${input.autoResumeAttempts}); last cause: ${describeErrorIdentity(causeError)}`; // #881: no fabricated cause class — original error identity + error-file pointers carry the fact. const decisiveFacts: Record = { diagnostic, resumesUsed: input.autoResumeAttempts, dispatchErrorFiles: [...input.errorFiles], }; if (input.errorFiles.length > 0) { decisiveFacts.lastDispatchErrorFile = input.errorFiles[input.errorFiles.length - 1]; } const candidate = causeError as { name?: unknown; code?: unknown }; if (typeof candidate?.name === "string") decisiveFacts.errorName = candidate.name; if (typeof candidate?.code === "string" || typeof candidate?.code === "number") { decisiveFacts.errorCode = candidate.code; } const artifacts: TerminalArtifactRef[] = input.errorFiles.map((path) => ({ kind: "error", path, })); return { roleOutcome: { kind: "failure", role: input.role, diagnostic, decisiveFacts, }, navigator: { disposition: "no-advice" }, artifacts, runId: input.runId, autoResumeCount: input.autoResumeAttempts, }; } export async function runWithAutoResumeLoop< T extends AutoResumeDispatchResult, TPayload = unknown, >(options: { admitted: { runDirectory: string; /** Identity for the loop-owned typed failure terminal (dispatch-exception exhaustion). */ role: TerminalRoleName; runId: string; principal: DurablePrincipal; /** Required: exception terminals attach the ledger via this existing admitted fact. */ projectRoot: string; }; principalAuthority: DurablePrincipalAuthority; io: CliIo; sessionAppender: SessionCustomEntryAppender; /** * Effective ceiling (#422), resolved by the caller before the loop; never re-read * per round. undefined = package default (AUTO_RESUME_LIMIT). Domain-validated at * this single entry point (#422): NaN/negative/fractional/Infinity reject loudly * before the first dispatch instead of silently bypassing the ceiling comparison. */ autoResumeLimit?: number | undefined; /** * #855: process-cancel signal. Once aborted by a catchable process signal, * the loop must not dispatch another turn or spawn another child. */ signal?: AbortSignal; buildInitialPayload: () => TPayload; buildResumePayload: () => TPayload; /** * #987: the initial turn retains the existing writer lease/liveness guard; * actual resume attempts pass no lease so a live holder cannot pre-block the * host CLI's own continuation contract. */ dispatch: ( payload: TPayload, lease: RunWriterLease | undefined, isFirst: boolean, attemptIo: CliIo, ) => Promise; }): Promise { // #422 single-point resolution + domain validation. NaN would bypass every // `attempts >= limit` comparison (always false) — reject here, before any dispatch. const limit = options.autoResumeLimit ?? AUTO_RESUME_LIMIT; parseAutoResumeLimit(limit); let autoResumeAttempts = 0; let isFirst = true; let currentPayload = options.buildInitialPayload(); let dispatchOrdinal = 0; let lastThrownError: unknown; let everyAttemptThrew = true; const retainedErrorFiles: string[] = []; const failedAttempts: Array<{ attempt: number; diagnostic: string; decisiveFacts?: Readonly>; errorFile?: string; }> = []; let stoppedResult: T | undefined; let endReason: string | undefined; while (true) { let result: T | undefined; // Set only when the caught throw is a TurnDispatchedFailure (#840 r9 判词 // class 1): the host turn genuinely started this attempt even though // dispatch produced no result — the next payload must still be a resume, // not a replay of the initial one. let turnStartedBeforeThrow = false; try { // #987: guard the initial new turn with the existing lease, but do not // pre-block a real host CLI resume on a package-side live holder. const lease = isFirst ? await acquireRunWriterLease(options.admitted.runDirectory) : undefined; result = await options.dispatch(currentPayload, lease, isFirst, dummyIo); } catch (error) { // Owner 2026-08-23: 「出了异常,就原地记录错误信息,然后重试。」 // Retain the whole exception in place (per-attempt full file + dossier // pointer); recording failure must not break the retry path (PR #418 // diagnostic-sink-isolation precedent). When a lease is held inside // dispatch, that path owns release in its own finally. lastThrownError = error; const failedAttempt: (typeof failedAttempts)[number] = { attempt: dispatchOrdinal, diagnostic: describeErrorIdentity(unwrapTurnDispatchedFailure(error)), }; failedAttempts.push(failedAttempt); turnStartedBeforeThrow = error instanceof TurnDispatchedFailure; const attempt = dispatchOrdinal; try { // Track the file as soon as it is durably written (#426 review): // a pointer-stage failure comes back separately (pointerError) and must // never orphan the retained file (#426 fix_now #5). const { file, pointerError } = await retainDispatchError( options.admitted, options.principalAuthority, options.sessionAppender, attempt, error, ); retainedErrorFiles.push(file); failedAttempt.errorFile = file; options.io.stderr( `dispatch attempt ${attempt} threw (${describeErrorIdentity(error)}); full error retained at ${file}\n`, ); if (pointerError !== undefined) { options.io.stderr( `dispatch error retention failed (best-effort continue): ${describeErrorIdentity(pointerError)}\n`, ); } } catch (retentionError) { options.io.stderr( `dispatch error retention failed (best-effort continue): ${describeErrorIdentity(retentionError)}\n`, ); } } dispatchOrdinal += 1; // #836 r12 class 3: run-state persist for a dispatch that deferred it // (needsPersist) is settled by the dispatch closure itself, before this // loop ever sees the result — through the single existing // presentControlledFailure / settleFailureTerminalResult authority (ADR // 0080), never a second hand-rolled classify/artifact/Terminal here. This // loop only ever sees the already-resolved outcome: a lawful/failure // Terminal (with skipAutoResume set when persist failed after an // already-settled result) or the original non-lawful result unchanged. if (result !== undefined) { everyAttemptThrew = false; const terminal = (result as { terminal?: TerminalResult }).terminal; if (terminal !== undefined) { (terminal as { autoResumeCount?: number }).autoResumeCount = autoResumeAttempts; if (terminal.roleOutcome.kind === "failure") { failedAttempts.push({ attempt: dispatchOrdinal - 1, diagnostic: terminal.roleOutcome.diagnostic, decisiveFacts: terminal.roleOutcome.decisiveFacts, }); } } if (terminal !== undefined && isLawfulTypedTerminalOutcome(terminal.roleOutcome) || result.skipAutoResume === true || autoResumeAttempts >= limit || processCancelSignalName(options.signal) !== undefined) { stoppedResult = result; break; } } else { if (autoResumeAttempts >= limit) { endReason = "auto-resume budget exhausted"; break; } const cancelName = processCancelSignalName(options.signal); if (cancelName !== undefined) { endReason = `ak-role terminated by ${cancelName}`; break; } } autoResumeAttempts++; // Resume payload only after a host turn actually started — either a // returned result says so, or the attempt threw a TurnDispatchedFailure // after genuinely starting the turn (#840 r9 判词 class 1). Pre-turn // throws and beforeDispatch failures retry the initial payload (#840 / #416). if (result?.turnDispatched === true || turnStartedBeforeThrow) { currentPayload = options.buildResumePayload(); isFirst = false; } } const terminal = stoppedResult === undefined ? await attachDispatchExceptionTerminal( options.admitted, dispatchExceptionFailureTerminal({ role: options.admitted.role, runId: options.admitted.runId, causeError: lastThrownError, errorFiles: retainedErrorFiles, autoResumeAttempts, endReason: endReason!, everyAttemptThrew, }), options.io, ) : stoppedResult.terminal; if (stoppedResult === undefined) { await finalizeExceptionRunBestEffort(options.admitted.runDirectory, options.io); } const complete: TerminalResult | undefined = terminal !== undefined && failedAttempts.length > 0 ? { ...terminal, roleOutcome: { ...terminal.roleOutcome, decisiveFacts: { ...terminal.roleOutcome.decisiveFacts, failedAttempts }, } } as TerminalResult : terminal; if (complete !== undefined) { await presentTerminal(complete, options.io, options.admitted.runDirectory); } return stoppedResult === undefined ? { exitCode: 1, terminal: complete } as T : { ...stoppedResult, ...(complete === undefined ? {} : { terminal: complete }) } as T; }