diff --git a/core/run-audit.ts b/core/run-audit.ts index a8903a31..cff4869b 100644 --- a/core/run-audit.ts +++ b/core/run-audit.ts @@ -33,14 +33,71 @@ export type AuditOutcome = | { kind: 'commit'; inScope: string[]; outOfScope: string[]; unchanged: boolean; needsAmendment: boolean }; const MAX_CHANGES = 10_000; +/** Every path in the report together; keeps the saved result (scope lists included) far below its 1 MiB limit. */ +const MAX_PATH_BYTES = 480 * 1024; +const KINDS = new Set(['add', 'modify', 'delete', 'rename', 'mode']); +/** + * Whether a path part names Git's metadata directory in any spelling Git itself refuses (read-cache.c, verify_path): + * any case; on NTFS with trailing dots or spaces and as the 8.3 short name `git~1`; on HFS with ignorable code points. + */ +export function isDotGit(part: string): boolean { + // NTFS ends a name at a stream separator (`:`), then drops trailing dots and spaces. Callers split on `\` as well + // as `/`, since NTFS reads a backslash as a directory separator. + const name = part.split(':')[0]!; + const plain = name.replace(/[\u200c-\u200f\u202a-\u202e\u206a-\u206f\ufeff]/g, '').toLowerCase().replace(/[. ]+$/, ''); + return plain === '.git' || plain === 'git~1'; +} +/** Agent-controlled text in a finding is quoted (AGENTS.md), and each list is cut short, so a reason stays readable. */ +const q = (text: string) => JSON.stringify(text.length > 300 ? `${text.slice(0, 300)}…` : text); +const list = (items: readonly string[]) => items.slice(0, 5).map(q).join(', ') + (items.length > 5 ? ` and ${items.length - 5} more` : ''); +const TYPES = new Set(['file', 'symlink', 'gitlink', 'directory', 'other']); + +/** + * Every field the audit reads, checked before it reads any (AGENTS.md: partial records fail closed). A missing boolean + * or entry type must never read as "clean". Returns the first problem, or null. + */ +function malformed(manifest: ChangeManifest): string | null { + if (!manifest || typeof manifest !== 'object') return 'The change report is missing.'; + if (!Array.isArray(manifest.changes)) return 'The change report has no change list.'; + for (const field of ['agentCommits', 'linkTargetChanges', 'nestedGitlinkContent'] as const) + if (!Array.isArray(manifest[field]) || manifest[field].some(entry => typeof entry !== 'string')) return `The change report has no ${field} list.`; + if (typeof manifest.metadataChanged !== 'boolean') return 'The change report does not say whether Git metadata changed.'; + let bytes = 0; + for (const change of manifest.changes) { + if (!change || typeof change !== 'object' || typeof change.path !== 'string' || !change.path) return 'A change has no path.'; + // Only a rename has an old path; on any other kind it would be staged as if it were part of the change. + if (change.kind === 'rename' ? typeof change.oldPath !== 'string' || !change.oldPath : change.oldPath !== undefined) return `The change at ${q(change.path)} has an invalid old path.`; + if (!KINDS.has(change.kind)) return `The change at ${q(change.path)} has an unknown kind.`; + if (typeof change.underGit !== 'boolean') return `The change at ${q(change.path)} does not say whether it is under .git.`; + // add has only a new entry, delete only an old one, every other kind both. + const needsOld = change.kind !== 'add', needsNew = change.kind !== 'delete'; + if ((needsOld && !TYPES.has(change.oldType as string)) || (!needsOld && change.oldType !== undefined)) return `The change at ${q(change.path)} has an invalid old entry type.`; + if ((needsNew && !TYPES.has(change.newType as string)) || (!needsNew && change.newType !== undefined)) return `The change at ${q(change.path)} has an invalid new entry type.`; + if (change.newType === 'symlink' && typeof change.newLinkTarget !== 'string') return `The link at ${q(change.path)} has no target.`; + if (change.linkTargetTraversesLink !== undefined && typeof change.linkTargetTraversesLink !== 'boolean') return `The link at ${q(change.path)} has an invalid traversal flag.`; + // Both paths can be saved (an undeclared rename source is a scope finding), measured as the JSON that stores them. + for (const path of [change.path, ...(change.oldPath ? [change.oldPath] : [])]) bytes += Buffer.byteLength(JSON.stringify(path)) + 1; + } + if (bytes > MAX_PATH_BYTES) return 'The change report is too large to audit.'; + return null; +} /** A stored link target must stay inside the repo, outside `.git`, without an absolute path. */ function unsafeLinkTarget(linkPath: string, target: string): string | null { if (!target || target.includes('\0')) return 'empty or invalid target'; if (target.startsWith('/')) return 'absolute target'; - const resolved = posix.normalize(posix.join(posix.dirname(linkPath), target)); + // A backslash or a drive prefix would be a separator or an absolute path on NTFS, where the audit's other checks + // (which use POSIX paths) could not see an escape: link text in this repository is POSIX-only. + if (target.includes('\\') || /^[A-Za-z]:/.test(target)) return 'target uses a Windows path form'; + // Checked as written too: a `.git` part that a later `..` cancels (`.git/../src`) still names the metadata on the way. + if (target.split(/[/\\]/).some(isDotGit)) return 'target enters .git'; + // A trailing slash names the same directory: `./` and `a/../` are the root, like `.`. + const resolved = posix.normalize(posix.join(posix.dirname(linkPath), target)).replace(/\/+$/, '') || '.'; if (resolved === '..' || resolved.startsWith('../')) return 'target leaves the repository'; - if (resolved === '.git' || resolved.startsWith('.git/')) return 'target enters .git'; + // The repository root contains .git: a link to it reaches the metadata through one more path part. + if (resolved === '.') return 'target is the repository root'; + // As for paths: Git's metadata is `.git` in any case, at any depth. + if (resolved.split(/[/\\]/).some(isDotGit)) return 'target enters .git'; return null; } @@ -50,31 +107,61 @@ function unsafeLinkTarget(linkPath: string, target: string): string | null { */ export function auditRun(item: PlanItem, manifest: ChangeManifest, pathKey: (path: string) => string): AuditOutcome { const violations: string[] = []; - if (!Array.isArray(manifest.changes) || manifest.changes.length > MAX_CHANGES) return { kind: 'violation', violations: ['The change report is missing or too large to audit.'] }; + if (!Array.isArray(manifest?.changes) || manifest.changes.length > MAX_CHANGES) return { kind: 'violation', violations: ['The change report is missing or too large to audit.'] }; + const problem = malformed(manifest); + if (problem) return { kind: 'violation', violations: [problem] }; if (manifest.metadataChanged) violations.push('The agent changed Git metadata under .git.'); - for (const path of manifest.linkTargetChanges) violations.push(`A declared symlink target changed: ${path}.`); - for (const path of manifest.nestedGitlinkContent) violations.push(`Content appeared under a gitlink: ${path}.`); + // Agents never commit: the metadata volume is read-only to them, so any agent commit is a violation, never undone (#66). + if (manifest.agentCommits.length) violations.push(`The agent made its own commits: ${list(manifest.agentCommits)}.`); + if (manifest.linkTargetChanges.length) violations.push(`A declared symlink target changed: ${list(manifest.linkTargetChanges)}.`); + if (manifest.nestedGitlinkContent.length) violations.push(`Content appeared under a gitlink: ${list(manifest.nestedGitlinkContent)}.`); const declared = new Set(item.files.flatMap(file => [file.path, ...(file.renamed_from ? [file.renamed_from] : [])]).map(pathKey)); + // Each path appears once, under the trusted path identity: two entries for one path contradict each other. + const seen = new Map(); + for (const change of manifest.changes) { + // A case-only rename's two sides are one path under a folding identity; count it once for this change. + const keys = new Set([change.path, ...(change.oldPath ? [change.oldPath] : [])].map(pathKey)); + for (const key of keys) { + const entries = [...(seen.get(key) ?? []), change]; + seen.set(key, entries); + if (entries.length === 1) continue; + // The only second entry allowed is a case-only rename that Git reports as one delete and one add (the file also + // changed a lot): exactly two entries, one delete and one add, with different spellings. + const [a, b] = entries as [ManifestChange, ManifestChange]; + const splitRename = entries.length === 2 && a.path !== b.path && [a.kind, b.kind].sort().join() === 'add,delete'; + if (!splitRename) return { kind: 'violation', violations: [`The change report lists ${q(key)} more than once.`] }; + } + } for (const change of manifest.changes) { const paths = [change.path, ...(change.oldPath ? [change.oldPath] : [])]; - if (change.underGit || paths.some(path => path === '.git' || path.startsWith('.git/'))) { violations.push(`The agent changed ${change.path} under .git.`); continue; } - if (paths.some(path => path.startsWith('/') || posix.normalize(path).startsWith('../') || path.includes('\0'))) - { violations.push(`Invalid path in the change report: ${change.path}.`); continue; } - if (change.oldType === 'gitlink' || change.newType === 'gitlink') { violations.push(`Plan items cannot change gitlinks: ${change.path}.`); continue; } + // A path must be in canonical form: another spelling (./, a/../, //, a trailing /) could reach .git or hide a match. + if (paths.some(path => path !== posix.normalize(path) || path.endsWith('/') || path.startsWith('./'))) + { violations.push(`Path not in canonical form in the change report: ${q(change.path)}.`); continue; } + // Git refuses a .git part in any spelling it treats as .git, at any depth, so the audit does too. + if (change.underGit || paths.some(path => path.split(/[/\\]/).some(isDotGit))) { violations.push(`The agent changed ${q(change.path)} under .git.`); continue; } + if (paths.some(path => path.startsWith('/') || posix.normalize(path).startsWith('../') || ['.', '..'].includes(posix.normalize(path)) || path.includes('\0'))) + { violations.push(`Invalid path in the change report: ${q(change.path)}.`); continue; } + if (change.oldType === 'gitlink' || change.newType === 'gitlink') { violations.push(`Plan items cannot change gitlinks: ${q(change.path)}.`); continue; } + // Only a declared pre-existing link may change at all: deleting it, turning it into a file, or renaming it from an + // undeclared path is a link change too. + if (change.oldType === 'symlink' && !declared.has(pathKey(change.oldPath ?? change.path))) + { violations.push(`A pre-existing symlink was changed at an undeclared path: ${q(change.oldPath ?? change.path)}.`); continue; } if (change.newType === 'symlink') { - if (change.oldType !== 'symlink') { violations.push(`New symlink or file-to-symlink conversion: ${change.path}.`); continue; } - if (!declared.has(pathKey(change.path))) { violations.push(`A pre-existing symlink changed at an undeclared path: ${change.path}.`); continue; } + if (change.oldType !== 'symlink') { violations.push(`New symlink or file-to-symlink conversion: ${q(change.path)}.`); continue; } + if (!declared.has(pathKey(change.path))) { violations.push(`A pre-existing symlink changed at an undeclared path: ${q(change.path)}.`); continue; } const unsafe = unsafeLinkTarget(change.path, change.newLinkTarget ?? ''); - if (unsafe) { violations.push(`Unsafe symlink target at ${change.path}: ${unsafe}.`); continue; } - if (change.linkTargetTraversesLink !== false) { violations.push(`The symlink target at ${change.path} traverses another link, or was not checked.`); continue; } + if (unsafe) { violations.push(`Unsafe symlink target at ${q(change.path)}: ${unsafe}.`); continue; } + if (change.linkTargetTraversesLink !== false) { violations.push(`The symlink target at ${q(change.path)} traverses another link, or was not checked.`); continue; } } - for (const type of [change.oldType, change.newType]) if (type === 'directory' || type === 'other') violations.push(`Unexpected ${type} entry: ${change.path}.`); + for (const type of [change.oldType, change.newType]) if (type === 'directory' || type === 'other') violations.push(`Unexpected ${type} entry: ${q(change.path)}.`); } if (violations.length) return { kind: 'violation', violations }; const inScope: string[] = [], outOfScope: string[] = []; for (const change of manifest.changes) { const paths = [change.path, ...(change.oldPath ? [change.oldPath] : [])]; - (paths.every(path => declared.has(pathKey(path))) ? inScope : outOfScope).push(change.path); + const undeclared = paths.filter(path => !declared.has(pathKey(path))); + // Each undeclared path is the scope finding itself, a rename's source included: never only the declared other side. + if (undeclared.length) outOfScope.push(...undeclared); else inScope.push(change.path); } return { kind: 'commit', inScope, outOfScope, unchanged: manifest.changes.length === 0, needsAmendment: outOfScope.length > 0 }; } diff --git a/docs/implementation/runner-lifecycle.md b/docs/implementation/runner-lifecycle.md index 1cc62884..265e0971 100644 --- a/docs/implementation/runner-lifecycle.md +++ b/docs/implementation/runner-lifecycle.md @@ -292,7 +292,7 @@ This runs before the coordinator opens. ## HTTP and UI contract -**Status reads.** `GET /api/runner` returns only the task and attempt rows plus `stateVersion`, `retryable`, `unresolved` and `stopRequested`. `stopRequested` comes from the in-memory job: null, or `{ attemptId, reason, saved }`, where `saved` is false while the first-reason write has failed. The UI shows "Stopping (not saved yet)" only from this field. It does not change `stateVersion`, so user actions still compare against the durable state version. `unresolved` is computed by the server from the in-memory marker: null, or `{ attemptId, reason: "result-not-saved" | "start-not-saved" }`. The UI shows "Needs restart: result could not be saved" only from this field, because the durable row alone may still look active. It does not rebuild Git history or the full review. +**Status reads.** `GET /api/runner` returns only the task and attempt rows plus `stateVersion`, `retryable`, `unresolved` and `stopRequested`. `stopRequested` comes from the in-memory job: null, or `{ attemptId, reason, saved }`, where `saved` is false while the first-reason write has failed. The UI shows "Stopping (not saved yet)" only from this field. It does not change `stateVersion`, so user actions still compare against the durable state version. `unresolved` is computed by the server from the in-memory marker: null, or `{ attemptId, reason: "result-not-saved" | "start-not-saved" | "preparation-not-removed" | "storage-not-removed" }`: the terminal write or the start could not be saved; the host-side preparation files could not be removed (before the terminal write when D never ran, after it once D settled); or the task storage could not be removed after a saved terminal write. The UI shows "Needs restart" with that cause only from this field, because the durable row alone may still look active. It does not rebuild Git history or the full review. **User actions.** Cancel, retry and "run again" requests send `attemptId`, `expectedStateVersion` and an `actionId` idempotency key (see "Feedback-event contract"). A replayed `actionId` returns the saved outcome. The server takes the plan identity from its trusted configuration, never from the request, and every `Store` call is scoped by that identity. A mismatch returns HTTP 409 with the current state. The UI then shows that state and keeps any draft. diff --git a/runner/coordinator.ts b/runner/coordinator.ts index fd3aea1e..e35c149f 100644 --- a/runner/coordinator.ts +++ b/runner/coordinator.ts @@ -1,6 +1,6 @@ import { identityKey, type PlanIdentity } from '../core/identity.ts'; import { captureInvocation, type InvocationHandle, type InvocationInput, type InvocationResult, type StopReason, type TaskClone, type UnreleasedResource } from '../agents/contract.ts'; -import type { AttemptRecord, Store } from './store.ts'; +import type { AttemptRecord, LedgerEntry, Store } from './store.ts'; import { ATTEMPT_PHASES, GuardRefusal, ShuttingDownError, WRITABLE_KINDS, bounded, sameContext, type AttemptKind, type Classification, type FirstReason, type ShutdownCapability, settleWith } from './lifecycle.ts'; /** What F's host-side preparation hands to D's start call. */ @@ -8,6 +8,20 @@ export interface PreparedAttempt { readonly clone: TaskClone; readonly vendor: 'claude' | 'codex'; readonly approvedArgv: readonly (readonly string[])[]; + /** Opaque data the deps keep for their own finish/release steps (for example the task workspace). */ + readonly private?: unknown; +} +/** The ledger record saved with `completed` in the same transaction (runner-lifecycle.md, publication step 3). */ +export interface HistoryRecord { readonly base: string; readonly head: string; readonly entries: readonly LedgerEntry[] } +/** A finish step's failure with its own actionable diagnostic (for example a safety violation). */ +export class FinishFailure extends Error {} +/** + * Preparation failed after it allocated task storage. `allocated` lets the coordinator remove that storage after the + * terminal write, as on every other path (runner-lifecycle.md, "Task storage is never removed before the terminal write"). + */ +export class PreparationFailure extends Error { + readonly allocated: PreparedAttempt; + constructor(cause: unknown, allocated: PreparedAttempt) { super(cause instanceof Error ? cause.message : String(cause), { cause }); this.allocated = allocated; } } export interface RunnerDeps { /** @@ -29,6 +43,14 @@ export interface RunnerDeps { start(input: InvocationInput, prepared: PreparedAttempt): InvocationHandle; /** Validate a clean result; throw with an actionable reason if it is invalid. Returns the value to persist. */ validate(attempt: AttemptRecord, result: InvocationResult): unknown; + /** + * Optional asynchronous replacement for validate, used by writable attempts: audit, make the runner commit inside + * task storage, and return the value plus the ledger record. Nothing is written to the Store here; the record is + * saved with `completed` in one transaction. Throw FinishFailure with an actionable diagnostic to fail the attempt. + */ + finish?(attempt: AttemptRecord, result: InvocationResult, prepared: PreparedAttempt, signal: AbortSignal): Promise<{ value: unknown; history?: HistoryRecord }>; + /** Optional: remove task storage after the terminal write and before the slot is freed. A failure keeps the slot under a marker. */ + release?(attempt: AttemptRecord, prepared: PreparedAttempt): Promise; now?(): number; } export interface SlotLimits { readonly writable: number; readonly readOnly: number } @@ -42,10 +64,10 @@ export interface RunnerStatus { unresolved: { attemptId: string; reason: UnresolvedReason } | null; } /** - * Why a task's slot stays held until restart: the terminal write failed, pending -> running failed, or the host-side - * preparation files could not be removed. + * Why a task's slot stays held until restart: the terminal write failed, pending -> running failed, the host-side + * preparation files could not be removed, or task storage could not be removed after a saved terminal write. */ -export type UnresolvedReason = 'result-not-saved' | 'start-not-saved' | 'preparation-not-removed'; +export type UnresolvedReason = 'result-not-saved' | 'start-not-saved' | 'preparation-not-removed' | 'storage-not-removed'; type Group = 'writable' | 'readOnly'; interface Job { identity: PlanIdentity; key: string; group: Group; attemptId: string; attempt?: AttemptRecord; @@ -65,10 +87,11 @@ const D_REASON: Record = { cancelled: 'cancelled', stal /** setTimeout accepts at most 2^31-1 ms; longer waits are re-armed. */ const MAX_TIMER = 2_147_483_647; const PREPARATION_TIMEOUT = 'Timed out while preparing.'; -const NEEDS_RESTART: Record = { +export const NEEDS_RESTART: Readonly> = { 'result-not-saved': 'Needs restart: the last result could not be saved.', 'start-not-saved': 'Needs restart: the start of the last attempt could not be saved.', 'preparation-not-removed': 'Needs restart: the last attempt\'s preparation files could not be removed.', + 'storage-not-removed': 'Needs restart: the last attempt\'s task storage could not be removed.', }; const FOREIGN_RESULT = 'The agent returned a result for a different attempt; it was not saved.'; const NOT_STARTED_UNRELEASED = 'Not started: an earlier agent\'s cleanup could not be confirmed. Restart codeboost to run it again.'; @@ -245,28 +268,31 @@ export class RunnerCoordinator { if (job.firstReason) return await this.#endBeforeLaunch(job, attempt, {}); let prepared: PreparedAttempt; try { prepared = await this.#deps.prepare(attempt, job.controller.signal); } - catch (error) { return await this.#endBeforeLaunch(job, attempt, this.#preparationDetail(job, error)); } - if (job.firstReason || job.preparationTimedOut) return await this.#endBeforeLaunch(job, attempt, this.#preparationDetail(job)); + catch (error) { + // Storage that preparation allocated before it failed is removed after the terminal write, like every other path. + return await this.#endBeforeLaunch(job, attempt, this.#preparationDetail(job, error), error instanceof PreparationFailure ? error.allocated : undefined); + } + if (job.firstReason || job.preparationTimedOut) return await this.#endBeforeLaunch(job, attempt, this.#preparationDetail(job), prepared); // Launch check: one synchronous turn, no await between the checks and D's start call. const now = this.#now(), row = this.#store.getAttempt(job.identity, attempt.id), task = this.#store.getTask(job.identity); if (row.firstReason && !job.firstReason) job.firstReason = row.firstReason; - if (row.state !== 'pending' || job.firstReason) return await this.#endBeforeLaunch(job, attempt, {}); + if (row.state !== 'pending' || job.firstReason) return await this.#endBeforeLaunch(job, attempt, {}, prepared); // A context change comes before both time checks, as in the settlement order and startup recovery. if (!sameContext(row.context, this.#store.currentContext(job.identity))) { // Recorded like any stale stop, so the row keeps it even if a cancel task lands during cleanup. this.#requestStop(job, 'stale'); - return await this.#endBeforeLaunch(job, attempt, {}); + return await this.#endBeforeLaunch(job, attempt, {}, prepared); } - if (task.budgetDeadline !== null && now >= task.budgetDeadline) { this.#requestStop(job, 'time-limit'); return await this.#endBeforeLaunch(job, attempt, {}); } - if (now >= attempt.deadline) { job.preparationTimedOut = true; return await this.#endBeforeLaunch(job, attempt, { detail: PREPARATION_TIMEOUT }); } + if (task.budgetDeadline !== null && now >= task.budgetDeadline) { this.#requestStop(job, 'time-limit'); return await this.#endBeforeLaunch(job, attempt, {}, prepared); } + if (now >= attempt.deadline) { job.preparationTimedOut = true; return await this.#endBeforeLaunch(job, attempt, { detail: PREPARATION_TIMEOUT }, prepared); } // Fail closed: once D reported resources it could not remove, no new invocation starts, even one already admitted. - if (this.#unreleased) return await this.#endBeforeLaunch(job, attempt, { detail: NOT_STARTED_UNRELEASED }); + if (this.#unreleased) return await this.#endBeforeLaunch(job, attempt, { detail: NOT_STARTED_UNRELEASED }, prepared); let handle: InvocationHandle; try { const input = captureInvocation({ clone: prepared.clone, phase: ATTEMPT_PHASES[attempt.kind], vendor: prepared.vendor, approvedArgv: prepared.approvedArgv, deadline: attempt.deadline, attemptId: attempt.id, runnerOwner: this.#deps.runnerOwner, context: attempt.context }, now); handle = this.#deps.start(input, prepared); - } catch (error) { return await this.#endBeforeLaunch(job, attempt, { detail: `Launch failed: ${message(error)}` }); } + } catch (error) { return await this.#endBeforeLaunch(job, attempt, { detail: `Launch failed: ${message(error)}` }, prepared); } job.handle = handle; let running: boolean | undefined; try { running = this.#write(() => this.#store.markRunning(job.identity, attempt.id)); } catch { running = undefined; } @@ -291,20 +317,29 @@ export class RunnerCoordinator { // Accept only the result of this exact invocation, as the question path does. Anything else is never validated // or saved: the attempt fails closed. if (result.attemptId !== attempt.id || !result.context || !sameContext(result.context, attempt.context)) { - this.#settle(job, { stopReason: 'capture-failure', exitCode: null, signal: null, valid: false, + const foreignSaved = this.#settle(job, { stopReason: 'capture-failure', exitCode: null, signal: null, valid: false, detail: job.firstReason === 'stale' ? job.staleCause : FOREIGN_RESULT }); job.decided = true; + // Task storage and host-side preparation files wait for the terminal write, as on every other path. + if (foreignSaved) await this.#release(job, attempt, prepared); + if (foreignSaved && !(await this.#removePreparation(job, attempt))) this.#holdForPreparation(job); return; } - let valid = false, value: unknown, detail = result.stderr ? bounded(result.stderr) : undefined; + // The agent's stderr is its own text: quote it (AGENTS.md), so it cannot forge a runner line in the diagnostic. + let valid = false, value: unknown, history: HistoryRecord | undefined, detail = result.stderr ? JSON.stringify(bounded(result.stderr)) : undefined; if (!job.firstReason && result.exitCode === 0 && !result.stopReason) { - try { value = this.#deps.validate(attempt, result); valid = true; } - catch (error) { detail = `Invalid output: ${message(error)}`; } + try { + if (this.#deps.finish) { const done = await this.#deps.finish(attempt, result, prepared, job.controller.signal); value = done.value; history = done.history; } + else value = this.#deps.validate(attempt, result); + valid = true; + } catch (error) { detail = error instanceof FinishFailure ? bounded(error.message) : `Invalid output: ${message(error)}`; } } // A stale stop keeps its own cause; the agent's stderr is not a reason the attempt went stale. if (job.firstReason === 'stale') detail = job.staleCause; - const saved = this.#settle(job, { stopReason: result.stopReason, exitCode: result.exitCode, signal: result.signal, valid, result: value, detail }); + const saved = this.#settle(job, { stopReason: result.stopReason, exitCode: result.exitCode, signal: result.signal, valid, result: value, detail, history }); job.decided = true; + // Task storage goes after the terminal write too; a failed removal holds the slot under a marker. + if (saved) await this.#release(job, attempt, prepared); // Host-side preparation files go after the terminal write, so a failed write leaves them for startup recovery. if (saved && !(await this.#removePreparation(job, attempt))) this.#holdForPreparation(job); } catch (error) { @@ -324,14 +359,17 @@ export class RunnerCoordinator { #preparationDetail(job: Job, error?: unknown): { detail?: string } { // D never ran, so there is no D stop reason; without one the Store keeps this text instead of "Timed out.". if (job.preparationTimedOut && !job.firstReason) return { detail: PREPARATION_TIMEOUT }; - return error === undefined || job.firstReason ? {} : { detail: `Preparation failed: ${message(error)}` }; + // Preparation errors can name repository paths an agent chose (an earlier item's files): quote them (AGENTS.md). + return error === undefined || job.firstReason ? {} : { detail: `Preparation failed: ${JSON.stringify(message(error))}` }; } /** Ending without a handle: host-side cleanup, then the terminal write from the first reason. */ - async #endBeforeLaunch(job: Job, attempt: AttemptRecord, s: { detail?: string }): Promise { + async #endBeforeLaunch(job: Job, attempt: AttemptRecord, s: { detail?: string }, prepared?: PreparedAttempt): Promise { // Stops that land while preparation finishes are taken into account; once the job is ending, the outcome is fixed. job.decided = true; const removed = await this.#removePreparation(job, attempt); - this.#settle(job, { exitCode: null, signal: null, valid: false, detail: job.firstReason === 'stale' ? job.staleCause : s.detail }); + const saved = this.#settle(job, { exitCode: null, signal: null, valid: false, detail: job.firstReason === 'stale' ? job.staleCause : s.detail }); + // Task storage (if preparation allocated it) waits for the terminal write, like every other path. + if (saved && prepared) await this.#release(job, attempt, prepared); if (!removed) this.#holdForPreparation(job); } /** @@ -341,7 +379,7 @@ export class RunnerCoordinator { async #removePreparation(job: Job, attempt: AttemptRecord): Promise { try { await this.#deps.cleanupPreparation(attempt); return true; } catch (error) { - console.error(`Runner job ${job.attemptId} could not remove its preparation files: ${message(error)}`); + console.error(`Runner job ${job.attemptId} could not remove its preparation files: ${JSON.stringify(message(error))}`); return false; } } @@ -349,7 +387,16 @@ export class RunnerCoordinator { #holdForPreparation(job: Job): void { if (!this.#markers.has(job.key)) this.#markers.set(job.key, { group: job.group, attemptId: job.attemptId, reason: 'preparation-not-removed' }); } - #settle(job: Job, s: { stopReason?: StopReason; exitCode: number | null; signal: string | null; valid: boolean; result?: unknown; detail?: string }): Classification | undefined { + /** After the terminal write: remove task storage, then the slot is freed. A failure keeps the slot under a marker. */ + async #release(job: Job, attempt: AttemptRecord, prepared: PreparedAttempt): Promise { + if (!this.#deps.release) return; + try { await this.#deps.release(attempt, prepared); } + catch (error) { + console.error(`Runner job ${job.attemptId} could not remove its task storage: ${JSON.stringify(message(error))}`); + if (!this.#markers.has(job.key)) this.#markers.set(job.key, { group: job.group, attemptId: job.attemptId, reason: 'storage-not-removed' }); + } + } + #settle(job: Job, s: { stopReason?: StopReason; exitCode: number | null; signal: string | null; valid: boolean; result?: unknown; detail?: string; history?: HistoryRecord }): Classification | undefined { try { return this.#write(() => this.#store.settleAttempt(job.identity, job.attemptId, { ...s, firstReason: job.firstReason })); } catch { @@ -361,7 +408,7 @@ export class RunnerCoordinator { #unexpected(job: Job, error: unknown): void { // Fail closed: an unexpected error keeps the slot held until restart. this.#markers.set(job.key, { group: job.group, attemptId: job.attemptId, reason: 'result-not-saved' }); - console.error(`Runner job ${job.attemptId} failed unexpectedly: ${message(error)}`); + console.error(`Runner job ${job.attemptId} failed unexpectedly: ${JSON.stringify(message(error))}`); } } const message = (error: unknown) => bounded(error instanceof Error ? error.message : String(error)); diff --git a/runner/execution.ts b/runner/execution.ts new file mode 100644 index 00000000..41a9bfcc --- /dev/null +++ b/runner/execution.ts @@ -0,0 +1,335 @@ +import { identityKey, type PlanIdentity } from '../core/identity.ts'; +import type { PlanContext } from '../core/plan.ts'; +import type { InvocationContext, InvocationHandle, InvocationInput, TaskClone } from '../agents/contract.ts'; +import { prepareExecution } from '../core/execution-prompt.ts'; +import { auditRun, type ChangeManifest } from '../core/run-audit.ts'; +import { FinishFailure, NEEDS_RESTART, PreparationFailure, type PreparedAttempt, type RunnerCoordinator, type RunnerDeps } from './coordinator.ts'; +import type { AttemptRecord, Store } from './store.ts'; +import { CLOSED_STATUSES, GuardRefusal, ShuttingDownError, bounded, sameContext, settleWith, type ShutdownCapability } from './lifecycle.ts'; + +/** + * F2b: per-item execution and the runner's commit step (design, "How codeboost runs a plan"; plan-format.md, "After + * each run"). The workspace operations are D's (#66); until they exist this module is exercised with a fake. + */ +export interface WorkspaceRef { readonly clone: TaskClone; readonly storage: unknown } +export interface TaskWorkspace { + /** A fresh task filesystem from the recorded trusted head; never a reset of a used one. Abortable; settles only when its work stopped. */ + materialize(attempt: AttemptRecord, head: string, signal: AbortSignal): Promise; + /** + * No-follow snapshot, taken before launch, of every declared path that is a symlink in this workspace. F passes every + * path the item declares: an earlier item may have renamed or added links, so only D sees the actual entries. + */ + snapshotDeclaredLinks(workspace: WorkspaceRef, paths: readonly string[], signal: AbortSignal): Promise; + /** The change manifest after the agent settled, plus a digest the commit step must match. */ + inspectChanges(workspace: WorkspaceRef, input: { baseHead: string; linkSnapshot: unknown }, signal: AbortSignal): Promise; + /** + * Commit exactly `paths` (both sides of every rename) on top of `baseHead` with hooks off; refuse if the tree no longer + * matches `digest`. Agent commits are never undone: a manifest with any is a safety violation and never gets here (#66). + */ + commit(workspace: WorkspaceRef, input: { baseHead: string; paths: readonly string[]; message: string; trailers: Readonly>; digest: string }, signal: AbortSignal): Promise; + release(workspace: WorkspaceRef): Promise; +} +/** D's start call for an execute/fix phase with this prompt; returns at once (see #51). */ +export type AgentLauncher = (input: InvocationInput, prompt: string, workspace: WorkspaceRef) => InvocationHandle; +/** Trusted runner-side sources for a task. Issue text and lessons are untrusted data inside the prompt. */ +export interface ExecutionSources { + planContext(identity: PlanIdentity): PlanContext; + issue(identity: PlanIdentity): { number: number; title: string; body: string; comments: readonly string[] }; + lessons(identity: PlanIdentity): readonly string[]; + vendor(identity: PlanIdentity): 'claude' | 'codex'; +} +/** Prefix of the diagnostic for an audit safety violation. For people only: the executor never reads it back. */ +export const SAFETY_VIOLATION = 'Safety violation:'; +/** + * Safety violations the runner's own audit found, by attempt ID. executionDeps records them and ItemExecutor takes + * them, so agent output (stderr) can never be mistaken for one, and a violation survives a later stale or stop outcome. + */ +export class SafetyFindings { + #found = new Map(); + record(attemptId: string, reason: string): void { this.#found.set(attemptId, reason); } + /** A finding stays owed until its task has been moved to needs human (or closed); only then is it settled. */ + get(attemptId: string): string | undefined { return this.#found.get(attemptId); } + settle(attemptId: string): void { this.#found.delete(attemptId); } +} +export interface ExecutionResult { head: string; unchanged: boolean; inScope: string[]; outOfScope: string[] } +interface Private { workspace: WorkspaceRef; prompt: string; baseHead: string; linkSnapshot: unknown } + +/** + * RunnerDeps for execute attempts: fresh workspace, prompt, agent, then audit and the runner's own commit. + * `runnerOwner` is the database's runner token (`Store.runnerOwnerToken`); `workspace` must allocate task storage under it. + */ +export function executionDeps(store: Store, workspace: TaskWorkspace, launch: AgentLauncher, sources: ExecutionSources, runnerOwner: string, + findings: SafetyFindings): RunnerDeps { + const identityOf = (attempt: AttemptRecord): PlanIdentity => findIdentity(store, attempt); + return { + runnerOwner, + async prepare(attempt, signal) { + if (attempt.kind !== 'execute' || !attempt.item) throw new Error('Execution deps run execute attempts for one plan item.'); + const identity = identityOf(attempt), plan = store.getPlan(identity, attempt.context.planRevision); + const item = plan.items.find(entry => entry.id === attempt.item)!; + const context = sources.planContext(identity); + const baseHead = store.getSnapshot(identity, attempt.context.snapshotId).head; + const request = prepareExecution({ identity, attemptId: attempt.id, mode: 'execute', plan, itemId: item.id, + issue: sources.issue(identity), approvedLessons: sources.lessons(identity), allowedCommands: context.allowedCommands }); + const vendor = sources.vendor(identity); + const declaredPaths = [...new Set(item.files.flatMap(file => [file.path, ...(file.renamed_from ? [file.renamed_from] : [])]))]; + const ws = await workspace.materialize(attempt, baseHead, signal); + // From here task storage exists: a failure hands it to the coordinator, which removes it after the terminal write. + const data: Private = { workspace: ws, prompt: request.prompt, baseHead, linkSnapshot: undefined }; + const prepared = { clone: ws.clone, vendor, approvedArgv: request.approvedArgv, private: data }; + try { data.linkSnapshot = await workspace.snapshotDeclaredLinks(ws, declaredPaths, signal); } + catch (error) { throw new PreparationFailure(error, prepared); } + return prepared; + }, + async cleanupPreparation() { /* host-side files belong to D's materialize; task storage waits for release */ }, + start(input, prepared) { const data = prepared.private as Private; return launch(input, data.prompt, data.workspace); }, + validate() { throw new Error('Execute attempts publish through finish().'); }, + async finish(attempt, _result, prepared, signal) { + const data = prepared.private as Private, identity = identityOf(attempt); + const plan = store.getPlan(identity, attempt.context.planRevision), item = plan.items.find(entry => entry.id === attempt.item)!; + const violation = (reason: string): never => { + const text = bounded(`${SAFETY_VIOLATION} ${reason}`); + findings.record(attempt.id, text); + throw new FinishFailure(text); + }; + let manifest: ChangeManifest & { digest: string }; + try { manifest = await workspace.inspectChanges(data.workspace, { baseHead: data.baseHead, linkSnapshot: data.linkSnapshot }, signal); } + catch (error) { + // Only the stop's own abort error is the stop. Any other refusal is a finding, even if a stop is also pending. + if (signal.aborted && (error === signal.reason || (error instanceof Error && error.name === 'AbortError'))) throw error; + // Contract (Publishing step 2): an inspection that refuses sends the task to needs human. + return violation(`The change inspection refused: ${JSON.stringify(error instanceof Error ? error.message : String(error))}`); + } + // The commit step refuses a tree that no longer matches this digest; without one that guard has nothing to check. + if (typeof manifest?.digest !== 'string' || !manifest.digest) return violation('The change report has no digest.'); + let outcome: ReturnType; + try { outcome = auditRun(item, manifest, sources.planContext(identity).pathKey); } + catch (error) { return violation(`The change report could not be audited: ${JSON.stringify(error instanceof Error ? error.message : String(error))}`); } + if (outcome.kind === 'violation') return violation(outcome.violations.join(' ')); + if (outcome.unchanged) return { value: { head: data.baseHead, unchanged: true, inScope: [], outOfScope: [] } satisfies ExecutionResult }; + // Last check before the commit, after the last await: a stop, shutdown or context change makes nothing. + if (signal.aborted) throw signal.reason; + if (!sameContext(attempt.context, store.currentContext(identity))) throw new FinishFailure('The plan, snapshot or assignment changed during the audit; nothing was committed.'); + // Every change is in or out of scope here; a rename stages both its old and its new path. + const paths = [...new Set(manifest.changes.flatMap(entry => [entry.path, ...(entry.oldPath ? [entry.oldPath] : [])]))]; + let head: string; + try { head = await workspace.commit(data.workspace, { + baseHead: data.baseHead, paths, digest: manifest.digest, + // The title is plan text: on one line with no control or bidi/format characters, it cannot open a trailer block + // that forges Plan-Item or Plan-Revision, put terminal escapes into git log, or reorder how git log shows it. + message: `${item.id}: ${item.title.replace(/[\u0000-\u001f\u007f-\u009f\u200b-\u200f\u2028-\u202e\u2060-\u206f\ufeff]+/g, ' ').replace(/ {2,}/g, ' ').trim()}`, trailers: { 'Plan-Item': item.id, 'Plan-Revision': `r${plan.revision}` }, + }, signal); } + catch (error) { + // D's refusal text can name agent-chosen paths: quote it (AGENTS.md). A stop records its first reason before it + // aborts, so a stopped commit still ends as that stop, whatever this text says. + throw new FinishFailure(`The runner commit was refused: ${JSON.stringify(error instanceof Error ? error.message : String(error))}`); + } + // The ID goes into the ledger inside the terminal write; a malformed one must fail the attempt, not that write. + if (typeof head !== 'string' || !/^(?:[0-9a-f]{40}|[0-9a-f]{64})$/.test(head)) throw new FinishFailure('The workspace returned an invalid commit ID; nothing was published.'); + if (head === data.baseHead) throw new FinishFailure('The workspace made no new commit for a changed item; nothing was published.'); + const snapshot = store.getSnapshot(identity, attempt.context.snapshotId); + return { + value: { head, unchanged: false, inScope: outcome.inScope, outOfScope: outcome.outOfScope } satisfies ExecutionResult, + history: { base: snapshot.base, head, entries: [{ sha: head, owner: item.id, origin: 'owned', sourceSha: null }] }, + }; + }, + async release(_attempt, prepared) { + const data = prepared.private as Private; + await workspace.release(data.workspace); + }, + }; +} +/** The plan identity that owns an attempt; attempts are stored per plan key. */ +function findIdentity(store: Store, attempt: AttemptRecord): PlanIdentity { + const key = store.attemptOwner(attempt.id); + if (!key) throw new Error('Unknown attempt.'); + const [repositoryId, taskId, planId] = JSON.parse(key) as string[]; + return { repositoryId: repositoryId!, taskId: taskId!, planId: planId! }; +} + +/** + * `completed` lists the items this run finished before it ended. A thrown error (storage or a bug) carries no list; + * the finished items are still recorded durably, as completed attempts and ledger entries. + */ +export type ExecutionOutcome = + | { kind: 'executed'; items: string[]; unchanged: string[] } + | { kind: 'needs amendment'; item: string; outOfScope: string[]; checkpointId: string; completed: string[] } + | { kind: 'needs human'; item: string; reason: string; completed: string[] } + | { kind: 'stopped'; item: string; state: string; reason: string | null; completed: string[] }; +const refusal = (error: unknown) => error instanceof GuardRefusal || error instanceof ShuttingDownError; +/** A write outside settlement: no shutdown capability. */ +const direct = (fn: () => T): T => fn(); +/** + * Statuses that wait for a person; leaving one needs its own user action (runner-lifecycle.md), so a safety finding is + * owed there instead. Review statuses are not gates: a finding moves them to needs human, so the task cannot be merged. + */ +const HUMAN_GATES: readonly string[] = ['needs amendment', 'needs approval', 'possibly already fixed']; + +/** + * Runs a task's plan items in order, one execute attempt each. Stops at the first item that does not complete cleanly: + * out-of-scope files pause the task in needs amendment with a checkpoint; a safety violation moves it to needs human. + * The run stops too if the plan gets a new revision while it runs: every item runs against the revision it started on. + */ +export class ItemExecutor { + #store: Store; #runner: RunnerCoordinator; #sources: ExecutionSources; #findings: SafetyFindings; + #deadlineMs: number; + /** F2's status changes after an attempt settles are settlement writes: they still land after the shutdown gate closes. */ + #write: (fn: () => T) => T; + constructor(store: Store, runner: RunnerCoordinator, sources: ExecutionSources, findings: SafetyFindings, + options: { deadlineMs?: number; capability?: ShutdownCapability } = {}) { + this.#store = store; this.#runner = runner; this.#sources = sources; this.#findings = findings; + this.#deadlineMs = options.deadlineMs ?? 10 * 60_000; this.#write = settleWith(options.capability); + } + /** + * Tasks with a runTask in progress here, from its start to its return: a second one waits for none of its steps. + * The guard is per instance, so the server must keep one ItemExecutor per Store (as it keeps one coordinator). + */ + #inFlight = new Set(); + async runTask(identity: PlanIdentity, options: { fromItem?: string } = {}): Promise { + const key = identityKey(identity); + // Held for the whole run, so a second run can never pay this run's pause or finding, even in the microtasks between + // the coordinator dropping the job and this run resuming. + if (this.#inFlight.has(key)) return { kind: 'stopped', item: options.fromItem ?? this.#store.getPlan(identity).items[0]!.id, state: 'not started', + reason: 'An earlier run of this task is still finishing; start it again when that run has ended.', completed: [] }; + this.#inFlight.add(key); + try { return await this.#runTask(identity, options); } + finally { this.#inFlight.delete(key); } + } + async #runTask(identity: PlanIdentity, options: { fromItem?: string }): Promise { + const plan = this.#store.getPlan(identity); + const start = options.fromItem ? plan.items.findIndex(item => item.id === options.fromItem) : 0; + if (start < 0) throw new Error('Unknown plan item.'); + const done: string[] = [], unchanged: string[] = []; + const stopped = (item: string, state: string, reason: string | null): ExecutionOutcome => ({ kind: 'stopped', item, state, reason, completed: [...done] }); + // Shutdown began (admission is closed): pay nothing owed and start nothing; the next run after restart does. + if (this.#runner.closing) return stopped(options.fromItem ?? plan.items[start]!.id, 'not started', 'The review server is shutting down.'); + // An earlier run of this task that is still finishing (its storage release) settles its own findings and pause. + if (this.#runner.isActive(identity)) + return stopped(options.fromItem ?? plan.items[start]!.id, 'not started', 'An earlier run of this task is still finishing; start it again when that run has ended.'); + // A safety finding not yet acted on (a failed write, a human gate at the time) goes to needs human first. + for (const earlier of this.#store.getAttempts(identity)) { + const finding = this.#findings.get(earlier.id); + if (finding) return this.#escalate(identity, earlier, finding, stopped, [], true); + } + // A scope finding whose pause was never recorded (a failed write, the write gate, a crash) pauses now, before any item. + const owed = this.#unpausedScopeFinding(identity); + if (owed) return this.#pause(identity, owed.row, owed.result, stopped, [], true); + // Continuing after a scope pause (plan-format.md: reconcile the executed prefix with the audited head, validate the + // remaining items from that checkpoint) is not built yet (#88), so a paused task runs no further items: fail closed. + const checkpoint = this.#store.latestCheckpoint(identity); + if (checkpoint) + return stopped(options.fromItem ?? plan.items[start]!.id, 'not started', + `${checkpoint.item} changed files outside its plan item. Continuing after a scope pause is not supported yet (#88), so this task runs no further items.`); + /** Where the next item must start: the context the previous item left, or the current one for the first item. */ + let expected: InvocationContext | null = null; + for (const item of plan.items.slice(start)) { + // Admission reads the context in this same turn, so it cannot notice a change saved during an earlier item. + const current = this.#store.currentContext(identity); + // After an item, the task is still running unless someone changed its status meanwhile: then the run stops. + if (expected && this.#store.getTask(identity).status !== 'running') + return stopped(item.id, 'not started', `The task's status changed to ${this.#store.getTask(identity).status} during the run; ${item.id} was not started.`); + if (this.#store.getPlan(identity).revision !== plan.revision) + return stopped(item.id, 'not started', `The plan changed to a new revision during the run; review it before running ${item.id}.`); + if (expected && !sameContext(current, expected)) + return stopped(item.id, 'not started', `The task's snapshot or assignment changed during the run; review it before running ${item.id}.`); + let attempt: AttemptRecord; + try { + attempt = this.#runner.start(identity, { + expectedStateVersion: this.#store.getTask(identity).stateVersion, kind: 'execute', item: item.id, + expectedContext: current, deadline: Date.now() + this.#deadlineMs, + }); + } catch (error) { + if (!refusal(error)) throw error; + return stopped(item.id, 'not started', error instanceof Error ? error.message : String(error)); + } + await this.#runner.settled(identity); + const row = this.#store.getAttempt(identity, attempt.id); + // Only the runner's own audit records a finding; it wins over any later stale or stop outcome. + const violation = this.#findings.get(attempt.id); + if (violation) return this.#escalate(identity, row, violation, stopped, done); + if (row.state !== 'completed') { + // Still pending or running: the terminal write failed and the slot is held until restart. + const unresolved = this.#runner.status(identity).unresolved; + return stopped(item.id, row.state, row.state === 'pending' || row.state === 'running' + ? (unresolved ? NEEDS_RESTART[unresolved.reason] : `Needs restart: the outcome of ${item.id} could not be saved.`) : row.diagnostic); + } + const result = row.result as ExecutionResult; + done.push(item.id); + if (result.unchanged) unchanged.push(item.id); + if (result.outOfScope.length) return this.#pause(identity, row, result, stopped, done); + const snapshotId = result.unchanged ? row.context.snapshotId : this.#store.snapshotWithHead(identity, result.head); + if (!snapshotId) throw new Error(`The snapshot of ${item.id}'s commit is missing.`); + // The item's own commit (recordHistory) raised the context generation by exactly one; any other change is not ours. + expected = { ...row.context, snapshotId, stateVersion: row.context.stateVersion + (result.unchanged ? 0 : 1) }; + } + return { kind: 'executed', items: done, unchanged }; + } + /** + * A safety finding sends the task to needs human (plan-format.md, "After each run"). From running, queued or a review + * status it moves now. A human gate is kept, because leaving it needs its own user action (runner-lifecycle.md), and + * the finding stays owed: the task's next run escalates it before anything else. A closed task needs nothing. + * The finding is settled only once acted on, so a failed write leaves it owed too. + */ + #escalate(identity: PlanIdentity, row: AttemptRecord, violation: string, + stopped: (item: string, state: string, reason: string | null) => ExecutionOutcome, done: string[], owed = false): ExecutionOutcome { + const item = row.item!, task = this.#store.getTask(identity); + // The terminal write failed: the attempt still counts as active, so the move waits for restart; keep the finding owed. + if (row.state === 'pending' || row.state === 'running') { + const unresolved = this.#runner.status(identity).unresolved; + return stopped(item, owed ? 'not started' : row.state, `${violation} ${unresolved ? NEEDS_RESTART[unresolved.reason] : 'Needs restart: the attempt\'s outcome could not be saved.'}`); + } + if (CLOSED_STATUSES.includes(task.status)) { + this.#findings.settle(row.id); + return stopped(item, owed ? 'not started' : row.state, `${violation} The task is ${task.status}, so it was not moved to needs human.`); + } + // Already where the finding sends it: nothing is owed. + if (task.status === 'needs human') { + this.#findings.settle(row.id); + return { kind: 'needs human', item, reason: violation, completed: [...done] }; + } + if (HUMAN_GATES.includes(task.status)) + return stopped(item, owed ? 'not started' : row.state, `${violation} The task is ${task.status}; it moves to needs human when it next runs.`); + // Settling this run's own attempt writes through the shutdown capability; paying an owed finding at the start of a + // new run is that run's decision, so it does not, and the closed write gate refuses it like any other. + try { (owed ? direct : this.#write)(() => this.#store.transitionTask(identity, task.stateVersion, 'needs human')); } + catch (error) { + if (!(error instanceof GuardRefusal) && !(owed && error instanceof ShuttingDownError)) throw error; + return stopped(item, owed ? 'not started' : row.state, `${violation} The task could not be moved to needs human yet: ${(error as Error).message}`); + } + this.#findings.settle(row.id); + return { kind: 'needs human', item, reason: violation, completed: [...done] }; + } + /** The earliest completed execute attempt whose out-of-scope files have no checkpoint yet (its pause was lost). */ + #unpausedScopeFinding(identity: PlanIdentity): { row: AttemptRecord; result: ExecutionResult } | null { + // Every completed execute attempt, not only the latest: a later clean one must not hide an earlier owed pause. + for (const row of this.#store.getAttempts(identity)) { + const result = row.result as ExecutionResult | undefined; + if (row.kind !== 'execute' || row.state !== 'completed' || !result?.outOfScope?.length) continue; + if (!this.#store.checkpointAtHead(identity, result.head)) return { row, result }; + } + return null; + } + /** + * The scope pause. The checkpoint names the revision the item ran against and the snapshot its own commit created, so + * a revision or HEAD observation saved since cannot erase the finding. Only a refused pause (the task is closed or no + * longer running) returns stopped; anything else is thrown, and the next run pauses first. + */ + #pause(identity: PlanIdentity, row: AttemptRecord, result: ExecutionResult, + stopped: (item: string, state: string, reason: string | null) => ExecutionOutcome, done: string[], owed = false): ExecutionOutcome { + const item = row.item!; + const snapshotId = this.#store.snapshotWithHead(identity, result.head); + if (!snapshotId) throw new Error(`The snapshot of ${item}'s commit is missing.`); + const items = this.#store.getPlan(identity, row.context.planRevision).items; + let checkpointId: string; + try { + checkpointId = (owed ? direct : this.#write)(() => this.#store.pauseForAmendment(identity, { revision: row.context.planRevision, snapshotId }, { + item, baseEntries: this.#sources.planContext(identity).baseEntries, + completedItems: items.slice(0, items.findIndex(entry => entry.id === item) + 1).map(entry => entry.id), outOfScopePaths: result.outOfScope, + }, { owed })).id; + } catch (error) { + if (!(error instanceof GuardRefusal) && !(owed && error instanceof ShuttingDownError)) throw error; + return stopped(item, owed ? 'not started' : row.state, `${item} changed files outside its plan item, but the task could not pause for amendment: ${(error as Error).message}`); + } + return { kind: 'needs amendment', item, outOfScope: result.outOfScope, checkpointId, completed: [...done] }; + } +} diff --git a/runner/store.ts b/runner/store.ts index d370f338..e9594f4a 100644 --- a/runner/store.ts +++ b/runner/store.ts @@ -503,6 +503,53 @@ export class Store { this.#run('INSERT INTO checkpoints VALUES (?,?,?)', key, checkpoint.id, encode(checkpoint)); return checkpoint; }); } + /** + * F2's scope pause. `ranAt` is where the item actually ran: its plan revision and the snapshot its own commit + * created, not whatever is current, so a revision or HEAD observation saved since cannot erase the scope finding + * (plan-format.md, "After each run"). The checkpoint and the move to needs amendment commit together, so a refused + * status change (a closed task, an active attempt or merge) records no checkpoint either. + */ + pauseForAmendment(identity: PlanIdentity, ranAt: ReviewState, evidence: Omit, options: { owed?: boolean } = {}): Checkpoint { + const key = identityKey(identity); + return this.#transaction(() => { + if (!this.#get('SELECT 1 FROM snapshots WHERE key=? AND id=?', key, ranAt.snapshotId)) throw new Error('Unknown snapshot.'); + const ids = this.getPlan(identity, ranAt.revision).items.map(item => item.id); + if (!ids.includes(evidence.item) || evidence.completedItems.at(-1) !== evidence.item || new Set(evidence.completedItems).size !== evidence.completedItems.length || evidence.completedItems.some((item, i) => item !== ids[i])) + throw new Error('Checkpoint must describe the executed plan prefix.'); + // Only the executor's own task pauses. In the run that found it, that is a running task, or a review status someone + // set since (a merge must not go past the finding); a queued status set since is kept, and the next run pays the + // pause from there. A human gate is kept too. A pause owed from an earlier run is paid from queued as well. + const status = this.#task(key).status; + const pausable = status === 'running' || status === 'in review' || status === 'approved but merge blocked' || (options.owed === true && status === 'queued'); + if (!pausable) throw new GuardRefusal(`The task is ${status}, so it was not paused for amendment.`); + this.transitionTask(identity, this.#task(key).state_version as number, 'needs amendment'); + const checkpoint = { ...evidence, revision: ranAt.revision, snapshotId: ranAt.snapshotId, id: randomUUID() }; + this.#run('INSERT INTO checkpoints VALUES (?,?,?)', key, checkpoint.id, encode(checkpoint)); + return checkpoint; + }); + } + /** + * The scope checkpoint recorded for a runner commit, found by the commit's head (not by the latest snapshot, which a + * later snapshot with the same head would shadow), or null if its pause was never recorded. + */ + checkpointAtHead(identity: PlanIdentity, head: string): Checkpoint | null { + for (const row of this.#db.prepare('SELECT data FROM checkpoints WHERE key=? ORDER BY rowid DESC').all(identityKey(identity))) { + const checkpoint = decode(row.data); + if (this.getSnapshot(identity, checkpoint.snapshotId).head === head) return checkpoint; + } + return null; + } + /** The most recent scope checkpoint of this plan, or null. */ + latestCheckpoint(identity: PlanIdentity): Checkpoint | null { + const row = this.#get('SELECT data FROM checkpoints WHERE key=? ORDER BY rowid DESC LIMIT 1', identityKey(identity)); + return row ? decode(row.data) : null; + } + /** The latest snapshot of this plan whose head is `head` (the one a runner commit created), or null. */ + snapshotWithHead(identity: PlanIdentity, head: string): string | null { + for (const row of this.#db.prepare('SELECT id, data FROM snapshots WHERE key=? ORDER BY rowid DESC').all(identityKey(identity))) + if (decode(row.data).head === head) return row.id as string; + return null; + } getCheckpoint(identity: PlanIdentity, id: string): Checkpoint { const row = this.#get('SELECT data FROM checkpoints WHERE key=? AND id=?', identityKey(identity), id); if (!row) throw new Error('Unknown checkpoint.'); return decode(row.data); @@ -761,6 +808,8 @@ export class Store { */ settleAttempt(identity: PlanIdentity, id: string, settlement: Omit & { signal?: string | null; result?: unknown; diagnosticRef?: string | null; + /** Writable attempts: the runner's commit, recorded with `completed` in this same transaction. */ + history?: { base: string; head: string; entries: readonly LedgerEntry[] }; }): Classification { if (settlement.firstReason !== null && !FIRST_REASONS.includes(settlement.firstReason)) throw new GuardRefusal('Unknown stop reason.'); const key = identityKey(identity); @@ -783,6 +832,11 @@ export class Store { this.#run(`UPDATE attempts SET state=?, first_reason=?, stop_reason=?, exit_code=?, signal=?, result=?, diagnostic=?, diagnostic_ref=?, settled_at=? WHERE id=?`, outcome.state, firstReason, settlement.stopReason ?? null, settlement.exitCode, settlement.signal ?? null, result, outcome.reason, settlement.diagnosticRef ?? null, new Date().toISOString(), id); + // The guards above ran first; recording history now advances the context without invalidating this attempt. + if (outcome.state === 'completed' && settlement.history) { + const context = decode(row.context); + this.recordHistory(identity, { revision: context.planRevision, snapshotId: context.snapshotId }, settlement.history.base, settlement.history.head, settlement.history.entries); + } // A pending cancel task wins over everything, including the time limit. if (task.cancel_requested !== null && !this.#closed(task.status)) this.#closeTask(key, 'cancelled', task.cancel_requested as string); else { diff --git a/test/run-audit.test.ts b/test/run-audit.test.ts index f02c8f45..813b70c0 100644 --- a/test/run-audit.test.ts +++ b/test/run-audit.test.ts @@ -18,14 +18,160 @@ describe('post-run audit', () => { .toEqual({ kind: 'commit', inScope: ['src/retry.ts', 'docs/New.md'], outOfScope: [], unchanged: false, needsAmendment: false }); expect(auditRun(item, manifest([file('src/retry.ts'), file('src/extra.ts', { kind: 'add', oldType: undefined })]), exact)) .toMatchObject({ kind: 'commit', outOfScope: ['src/extra.ts'], needsAmendment: true }); + // Both sides of a rename are findings when neither is declared. + expect(auditRun(item, manifest([file('x.ts', { kind: 'rename', oldPath: 'y.ts' })]), exact)).toMatchObject({ kind: 'commit', outOfScope: ['x.ts', 'y.ts'] }); + // The undeclared source is the finding, not the declared destination. expect(auditRun(item, manifest([file('docs/New.md', { kind: 'rename', oldPath: 'docs/unlisted.md' })]), exact)) - .toMatchObject({ kind: 'commit', outOfScope: ['docs/New.md'] }); + .toMatchObject({ kind: 'commit', inScope: [], outOfScope: ['docs/unlisted.md'] }); }); it('reports planned-but-unchanged, and uses the trusted path identity', () => { - expect(auditRun(item, manifest([], { agentCommits: ['abc'] }), exact)).toMatchObject({ kind: 'commit', unchanged: true }); + expect(auditRun(item, manifest([]), exact)).toMatchObject({ kind: 'commit', unchanged: true }); expect(auditRun(item, manifest([file('SRC/Retry.ts')]), folded)).toMatchObject({ inScope: ['SRC/Retry.ts'] }); expect(auditRun(item, manifest([file('SRC/Retry.ts')]), exact)).toMatchObject({ outOfScope: ['SRC/Retry.ts'] }); }); + it('treats any agent commit as a safety violation, even with no file changes (#66: never undone)', () => { + expect(auditRun(item, manifest([], { agentCommits: ['abc'] }), exact)).toEqual({ kind: 'violation', violations: ['The agent made its own commits: "abc".'] }); + for (const field of ['agentCommits', 'linkTargetChanges', 'nestedGitlinkContent'] as const) + expect(auditRun(item, manifest([file('src/retry.ts')], { [field]: undefined as unknown as string[] }), exact)).toEqual({ kind: 'violation', violations: [`The change report has no ${field} list.`] }); + }); + it('fails closed on a partial record: every field the audit reads must be present and well-formed', () => { + const cases: [string, ChangeManifest, RegExp][] = [ + ['no metadata flag', manifest([file('src/retry.ts')], { metadataChanged: undefined as unknown as boolean }), /whether Git metadata changed/], + ['no underGit', manifest([{ path: 'src/retry.ts', kind: 'modify', oldType: 'file', newType: 'file' } as ManifestChange]), /under \.git/], + ['add without a new type', manifest([{ path: 'src/new.ts', kind: 'add', underGit: false } as ManifestChange]), /invalid new entry type/], + ['modify without an old type', manifest([file('src/retry.ts', { oldType: undefined })]), /invalid old entry type/], + ['unknown kind', manifest([file('src/retry.ts', { kind: 'chmod' as ManifestChange['kind'] })]), /unknown kind/], + ['rename without an old path', manifest([file('docs/New.md', { kind: 'rename' })]), /invalid old path/], + ['old path on a modify', manifest([file('src/retry.ts', { oldPath: 'src/other.ts' })]), /invalid old path/], + ['non-string old path', manifest([file('docs/New.md', { kind: 'rename', oldPath: 7 as unknown as string })]), /invalid old path/], + ['delete with a new type', manifest([file('src/retry.ts', { kind: 'delete' })]), /invalid new entry type/], + ['add with an old type', manifest([file('src/new.ts', { kind: 'add' })]), /invalid old entry type/], + ['non-boolean traversal flag', manifest([file('link', { oldType: 'symlink', newType: 'symlink', newLinkTarget: 'x', linkTargetTraversesLink: 'no' as unknown as boolean })]), /traversal flag/], + ['non-string list entry', manifest([file('src/retry.ts')], { nestedGitlinkContent: [{ path: 'm' }] as unknown as string[] }), /no nestedGitlinkContent list/], + ['link without a target', manifest([file('link', { oldType: 'symlink', newType: 'symlink' })]), /has no target/], + ['too many path bytes', manifest(Array.from({ length: 600 }, (_, i) => file(`${'x'.repeat(1000)}${i}`))), /too large/], + ]; + for (const [label, report, reason] of cases) { + const outcome = auditRun(item, report, exact); + expect(outcome.kind, label).toBe('violation'); + expect((outcome as { violations: string[] }).violations.join(' '), label).toMatch(reason); + } + }); + it('refuses a report that lists one path twice', () => { + expect(auditRun(item, manifest([file('src/retry.ts'), file('src/retry.ts', { kind: 'delete', newType: undefined })]), exact)) + .toEqual({ kind: 'violation', violations: ['The change report lists "src/retry.ts" more than once.'] }); + expect(auditRun(item, manifest([file('docs/New.md', { kind: 'rename', oldPath: 'docs/old.md' }), file('docs/old.md')]), exact)).toMatchObject({ kind: 'violation' }); + }); + it('refuses a declared link retargeted into .git in any case or at any depth', () => { + for (const target of ['.GIT/config', '.Git', 'vendor/.git/hooks', 'sub/.GiT']) + expect(auditRun(item, manifest([file('link', { oldType: 'symlink', newType: 'symlink', newLinkTarget: target, linkTargetTraversesLink: false })]), exact), target) + .toEqual({ kind: 'violation', violations: ['Unsafe symlink target at "link": target enters .git.'] }); + }); + it('refuses removing, replacing or renaming a pre-existing symlink from an undeclared path, and allows it at a declared one', () => { + for (const change of [file('lnk', { kind: 'delete', oldType: 'symlink', newType: undefined }), file('lnk', { oldType: 'symlink', newType: 'file' })]) + expect(auditRun(item, manifest([change]), exact)).toEqual({ kind: 'violation', violations: ['A pre-existing symlink was changed at an undeclared path: "lnk".'] }); + expect(auditRun(item, manifest([file('link', { kind: 'rename', oldPath: 'lnk', oldType: 'symlink', newType: 'symlink', newLinkTarget: 'src/retry.ts', linkTargetTraversesLink: false })]), exact)) + .toEqual({ kind: 'violation', violations: ['A pre-existing symlink was changed at an undeclared path: "lnk".'] }); + expect(auditRun(item, manifest([file('link', { kind: 'delete', oldType: 'symlink', newType: undefined })]), exact)).toMatchObject({ kind: 'commit', inScope: ['link'] }); + }); + it('refuses every spelling Git treats as .git, in paths and link targets', () => { + for (const path of ['.git /config', '.git./hooks', 'GIT~1/config', 'sub/.g\u200cit/hooks/post-checkout', '.GIT\ufeff']) + expect(auditRun(item, manifest([file(path)]), exact), path).toMatchObject({ kind: 'violation' }); + for (const target of ['.git.', 'GIT~1', '.g\u200dit/config']) + expect(auditRun(item, manifest([file('link', { oldType: 'symlink', newType: 'symlink', newLinkTarget: target, linkTargetTraversesLink: false })]), exact), target) + .toEqual({ kind: 'violation', violations: ['Unsafe symlink target at "link": target enters .git.'] }); + expect(auditRun(item, manifest([file('src/.github-notes.md', { kind: 'add', oldType: undefined })]), exact)).toMatchObject({ kind: 'commit' }); + }); + it('ends a name where NTFS does (a stream separator or a backslash) before comparing it with .git', () => { + for (const path of ['.git:x', '.git::$INDEX_ALLOCATION/config', 'git~1:s', '.git\\config', 'GIT~1\\hooks', 'a/.g\u200eit', 'a/.g\u202ait', 'a/.g\u206bit', 'a/.gi\u200ft']) + expect(auditRun(item, manifest([file(path, { kind: 'add', oldType: undefined })]), exact), path).toMatchObject({ kind: 'violation' }); + }); + it('reads a backslash as a directory separator when looking for .git', () => { + for (const path of ['x\\.git\\hooks\\post-checkout', 'a/b\\.GIT']) + expect(auditRun(item, manifest([file(path, { kind: 'add', oldType: undefined })]), exact), path).toMatchObject({ kind: 'violation' }); + }); + it('refuses bad link targets: a Windows path form, .git cancelled by .., empty, or with a NUL', () => { + const link = (target: string) => manifest([file('link', { oldType: 'symlink', newType: 'symlink', newLinkTarget: target, linkTargetTraversesLink: false })]); + for (const target of ['sub\\.git\\config', '..\\..\\outside', 'C:\\Windows', 'c:outside']) + expect(auditRun(item, link(target), exact), target).toEqual({ kind: 'violation', violations: ['Unsafe symlink target at "link": target uses a Windows path form.'] }); + expect(auditRun(item, link('.git/../src/retry.ts'), exact)).toEqual({ kind: 'violation', violations: ['Unsafe symlink target at "link": target enters .git.'] }); + for (const target of ['', 'src/re\0try.ts']) + expect(auditRun(item, link(target), exact), JSON.stringify(target)).toEqual({ kind: 'violation', violations: ['Unsafe symlink target at "link": empty or invalid target.'] }); + }); + it('refuses a NUL in a path, and a directory, other or gitlink old entry', () => { + expect(auditRun(item, manifest([file('src/re\0try.ts')]), exact)).toMatchObject({ kind: 'violation' }); + expect(auditRun(item, manifest([file('build', { kind: 'delete', oldType: 'directory', newType: undefined })]), exact)) + .toEqual({ kind: 'violation', violations: ['Unexpected directory entry: "build".'] }); + expect(auditRun(item, manifest([file('fifo', { kind: 'delete', oldType: 'other', newType: undefined })]), exact)) + .toEqual({ kind: 'violation', violations: ['Unexpected other entry: "fifo".'] }); + expect(auditRun(item, manifest([file('vendor/lib', { kind: 'delete', oldType: 'gitlink', newType: undefined })]), exact)) + .toEqual({ kind: 'violation', violations: ['Plan items cannot change gitlinks: "vendor/lib".'] }); + }); + it('covers the whole HFS ignorable ranges, and cuts a quoted path at 300 characters', () => { + for (const mark of ['\u200c', '\u200f', '\u202a', '\u202e', '\u206a', '\u206f', '\ufeff']) + expect(auditRun(item, manifest([file(`a/.gi${mark}t`, { kind: 'add', oldType: undefined })]), exact), JSON.stringify(mark)).toMatchObject({ kind: 'violation' }); + const long = `${'x'.repeat(400)}/.git`; + const outcome = auditRun(item, manifest([file(long, { kind: 'add', oldType: undefined })]), exact) as { violations: string[] }; + expect(outcome.violations[0]).toBe(`The agent changed ${JSON.stringify(`${long.slice(0, 300)}…`)} under .git.`); + }); + it('trusts D\'s underGit flag on its own', () => { + expect(auditRun(item, manifest([file('src/retry.ts', { underGit: true })]), exact)).toEqual({ kind: 'violation', violations: ['The agent changed "src/retry.ts" under .git.'] }); + }); + it('refuses a declared link retargeted to the directory that holds the repository', () => { + expect(auditRun(item, manifest([file('link', { oldType: 'symlink', newType: 'symlink', newLinkTarget: '..', linkTargetTraversesLink: false })]), exact)) + .toEqual({ kind: 'violation', violations: ['Unsafe symlink target at "link": target leaves the repository.'] }); + }); + it('refuses a declared link retargeted to the repository root', () => { + for (const target of ['.', './', 'a/..']) + expect(auditRun(item, manifest([file('link', { oldType: 'symlink', newType: 'symlink', newLinkTarget: target, linkTargetTraversesLink: false })]), exact), target) + .toEqual({ kind: 'violation', violations: ['Unsafe symlink target at "link": target is the repository root.'] }); + }); + it('refuses a declared link renamed to an undeclared path, a directory entry, and an absolute path', () => { + expect(auditRun(item, manifest([file('lnk2', { kind: 'rename', oldPath: 'link', oldType: 'symlink', newType: 'symlink', newLinkTarget: 'src/retry.ts', linkTargetTraversesLink: false })]), exact)) + .toEqual({ kind: 'violation', violations: ['A pre-existing symlink changed at an undeclared path: "lnk2".'] }); + expect(auditRun(item, manifest([file('build', { kind: 'add', oldType: undefined, newType: 'directory' })]), exact)) + .toEqual({ kind: 'violation', violations: ['Unexpected directory entry: "build".'] }); + expect(auditRun(item, manifest([file('/etc/passwd', { kind: 'add', oldType: undefined })]), exact)) + .toEqual({ kind: 'violation', violations: ['Invalid path in the change report: "/etc/passwd".'] }); + }); + it('counts a case-only rename once under a case-folding identity, whether reported as a rename or as a delete and an add', () => { + const renamed: PlanItem = { ...item, files: [{ path: 'README.md', kind: 'rename', renamed_from: 'Readme.md', change: 'x' }] }; + expect(auditRun(renamed, manifest([file('README.md', { kind: 'rename', oldPath: 'Readme.md' })]), folded)).toMatchObject({ kind: 'commit', inScope: ['README.md'] }); + const split = manifest([file('Readme.md', { kind: 'delete', newType: undefined }), file('README.md', { kind: 'add', oldType: undefined })]); + expect(auditRun(renamed, split, folded)).toMatchObject({ kind: 'commit', inScope: ['Readme.md', 'README.md'], outOfScope: [] }); + // Undeclared, the same pair is a scope finding, not a safety violation; two adds of one folded path are still refused. + expect(auditRun(item, manifest([file('notes', { kind: 'delete', newType: undefined }), file('NOTES', { kind: 'add', oldType: undefined })]), folded)).toMatchObject({ kind: 'commit', outOfScope: ['notes', 'NOTES'] }); + expect(auditRun(item, manifest([file('notes', { kind: 'add', oldType: undefined }), file('NOTES', { kind: 'add', oldType: undefined })]), folded)).toMatchObject({ kind: 'violation' }); + // Only one delete and one add, with different spellings and no rename entry, make a split rename. + const add = (path: string) => file(path, { kind: 'add', oldType: undefined }), del = (path: string) => file(path, { kind: 'delete', newType: undefined }); + for (const [label, changes] of [ + ['add, delete, add', [add('src/aB.ts'), del('src/Ab.ts'), add('src/ab.ts')]], + ['delete, add, delete, add', [del('src/Ab.ts'), add('src/aB.ts'), del('src/AB.ts'), add('src/ab.ts')]], + ['delete and add of one spelling', [del('src/ab.ts'), add('src/ab.ts')]], + ['a rename and an add', [file('src/ab.ts', { kind: 'rename', oldPath: 'src/x.ts' }), add('src/AB.ts')]], + ] as [string, ManifestChange[]][]) + expect(auditRun(item, manifest(changes), folded), label).toMatchObject({ kind: 'violation' }); + }); + it('counts two spellings of one path under a case-folding identity as a duplicate', () => { + expect(auditRun(item, manifest([file('SRC/Retry.ts'), file('src/retry.ts')]), folded)).toMatchObject({ kind: 'violation' }); + expect(auditRun(item, manifest([file('SRC/Retry.ts'), file('src/retry.ts')]), exact)).toMatchObject({ kind: 'commit' }); + }); + it('quotes agent-controlled paths in findings and cuts long lists short', () => { + const outcome = auditRun(item, manifest([file('src/retry.ts')], { nestedGitlinkContent: ['a', 'b', 'c', 'd', 'e', 'f', 'g"; rm -rf /'] }), exact); + expect(outcome).toEqual({ kind: 'violation', violations: ['Content appeared under a gitlink: "a", "b", "c", "d", "e" and 2 more.'] }); + }); + it('refuses a change path of . or .., any non-canonical spelling, and .git in any case at any depth', () => { + for (const path of ['..', '.', 'a/../..', 'x/../.git/hooks/pre-commit', './.git/config', './a.ts', 'a//b', 'dir/', 'sub/.git/config', '.GIT/config', 'src/.Git']) + expect(auditRun(item, manifest([file(path)]), exact), path).toMatchObject({ kind: 'violation' }); + expect(auditRun(item, manifest([file('docs/New.md', { kind: 'rename', oldPath: 'x/../.git/config' })]), exact)).toMatchObject({ kind: 'violation' }); + }); + it('limits the saved paths as JSON, a rename\'s old path included (an undeclared one is saved as the finding)', () => { + // 100 paths of 4,000 control characters are under the raw byte cap but about 2.4 MB as JSON. + expect(auditRun(item, manifest(Array.from({ length: 100 }, (_, i) => file(`${'\u0001'.repeat(4000)}${i}`, { kind: 'add', oldType: undefined }))), exact)) + .toEqual({ kind: 'violation', violations: ['The change report is too large to audit.'] }); + expect(auditRun(item, manifest([file('docs/New.md', { kind: 'rename', oldPath: `docs/${'o'.repeat(600_000)}.md` })]), exact)) + .toEqual({ kind: 'violation', violations: ['The change report is too large to audit.'] }); + }); it('stops on every safety violation before any scope decision', () => { const cases: [string, ChangeManifest][] = [ ['metadata', manifest([file('src/retry.ts')], { metadataChanged: true })], diff --git a/test/runner-coordinator.test.ts b/test/runner-coordinator.test.ts index 337f5c43..cf9d9431 100644 --- a/test/runner-coordinator.test.ts +++ b/test/runner-coordinator.test.ts @@ -180,7 +180,7 @@ describe('stops and settlement', () => { await runner.settled(A); expect(launches).toHaveLength(0); expect(cleaned()).toBe(1); - expect(store.getAttempt(A, attempt.id)).toMatchObject({ state: 'failed', firstReason: null, diagnostic: 'Preparation failed: clone failed' }); + expect(store.getAttempt(A, attempt.id)).toMatchObject({ state: 'failed', firstReason: null, diagnostic: 'Preparation failed: "clone failed"' }); expect(() => runner.start(B, request(store, B))).not.toThrow(); }); it('fails with the launch error when D start throws and no stop is recorded', async () => { @@ -508,8 +508,8 @@ describe('copilot review', () => { it.each([ ['another attempt', (input: InvocationInput) => ({ attemptId: randomUUID() })], ['another context', (input: InvocationInput) => ({ context: { ...input.context, stateVersion: input.context.stateVersion + 1 } })], - ] as const)('never validates or saves a result for %s', async (_label, foreign) => { - const { store, runner, launches, preparations, deps } = setup(); + ] as const)('never validates or saves a result for %s, and still removes its preparation files', async (_label, foreign) => { + const { store, runner, launches, preparations, deps, cleaned } = setup(); const validate = vi.spyOn(deps, 'validate'); const attempt = runner.start(A, request(store, A)); await until(() => preparations.length === 1, 'preparation'); preparations[0]!.resolve(); @@ -519,8 +519,20 @@ describe('copilot review', () => { expect(validate).not.toHaveBeenCalled(); expect(store.getAttempt(A, attempt.id)).toMatchObject({ state: 'failed', result: null, diagnostic: 'The agent returned a result for a different attempt; it was not saved.' }); + expect(cleaned()).toBe(1); expect(() => runner.start(B, request(store, B))).not.toThrow(); }); + it('keeps preparation files for startup recovery when a foreign result\'s terminal write fails', async () => { + const { store, runner, launches, preparations, cleaned } = setup(); + vi.spyOn(store, 'settleAttempt').mockImplementation(() => { throw Object.assign(new Error('disk full'), { code: 'ERR_SQLITE_ERROR' }); }); + runner.start(A, request(store, A)); + await until(() => preparations.length === 1, 'preparation'); preparations[0]!.resolve(); + await until(() => launches.length === 1, 'launch'); + launches[0]!.settle({ attemptId: randomUUID() }); + await runner.settled(A); + expect(cleaned()).toBe(0); + expect(runner.status(A).unresolved).toMatchObject({ reason: 'result-not-saved' }); + }); it('holds the slot when preparation files cannot be removed after D settles', async () => { const { store, runner, launches, preparations, deps } = setup(); vi.spyOn(console, 'error').mockImplementation(() => undefined); @@ -544,7 +556,7 @@ describe('copilot review', () => { preparations[0]!.reject(new Error('clone failed')); await runner.settled(A); expect(launches).toHaveLength(0); - expect(store.getAttempt(A, attempt.id)).toMatchObject({ state: 'failed', diagnostic: 'Preparation failed: clone failed' }); + expect(store.getAttempt(A, attempt.id)).toMatchObject({ state: 'failed', diagnostic: 'Preparation failed: "clone failed"' }); expect(runner.status(A).unresolved).toEqual({ attemptId: attempt.id, reason: 'preparation-not-removed' }); expect(() => runner.start(B, request(store, B))).toThrow(/No free runner slot/); }); diff --git a/test/runner-execution.test.ts b/test/runner-execution.test.ts new file mode 100644 index 00000000..c223e168 --- /dev/null +++ b/test/runner-execution.test.ts @@ -0,0 +1,910 @@ +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { DatabaseSync } from 'node:sqlite'; +import { join } from 'node:path'; +import { randomUUID } from 'node:crypto'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { Store } from '../runner/store.ts'; +import { RunnerCoordinator } from '../runner/coordinator.ts'; +import { ItemExecutor, SAFETY_VIOLATION, SafetyFindings, executionDeps, type ExecutionOutcome, type ExecutionSources, type TaskWorkspace, type WorkspaceRef } from '../runner/execution.ts'; +import { MAX_REASON, ShuttingDownError, type ShutdownCapability } from '../runner/lifecycle.ts'; +import type { ChangeManifest, ManifestChange } from '../core/run-audit.ts'; +import type { InvocationResult } from '../agents/contract.ts'; +import type { Plan, PlanContext } from '../core/plan.ts'; + +const oid = (n: number) => n.toString(16).padStart(40, '0'); +const RUNNER_OWNER = '0123456789abcdef0123456789abcdef'; +const identity = { repositoryId: 'repo', taskId: 'task', planId: 'plan' }; +const plan: Plan = { schema_version: 1, issue: 1, revision: 1, summary: 'Two items', questions: [], items: [ + { id: 'P1', title: 'First', intent: 'Change a', files: [{ path: 'a.ts', kind: 'edit', renamed_from: null, change: 'x' }], acceptance: [{ type: 'cmd', text: 'npm test' }], depends_on: [] }, + { id: 'P2', title: 'Second', intent: 'Change b', files: [{ path: 'b.ts', kind: 'edit', renamed_from: null, change: 'y' }], acceptance: [{ type: 'check', text: 'b reads well' }], depends_on: ['P1'] }, +] }; +const context: PlanContext = { identity, issue: 1, baseEntries: [{ path: 'a.ts', kind: 'file' }, { path: 'b.ts', kind: 'file' }], pathKey: p => p, allowedCommands: [['npm', 'test']] }; +const dirs: string[] = [], cleanups: (() => Promise | void)[] = []; +afterEach(async () => { for (const c of cleanups.splice(0).reverse()) await c(); for (const d of dirs.splice(0)) rmSync(d, { recursive: true, force: true }); }); +const change = (path: string, over: Partial = {}): ManifestChange => ({ path, kind: 'modify', oldType: 'file', newType: 'file', underGit: false, ...over }); +const manifest = (changes: ManifestChange[], over: Partial = {}): ChangeManifest & { digest: string } => + ({ changes, agentCommits: [], metadataChanged: false, linkTargetChanges: [], nestedGitlinkContent: [], digest: `digest-${changes.length}`, ...over }); + +function setup(options: { manifests?: Record; exit?: Record>; + commit?: (item: string) => Promise; release?: () => Promise; startError?: Error; + inspect?: (item: string, signal: AbortSignal) => Promise; snapshotError?: Error; capability?: (store: Store) => ShutdownCapability; settleError?: boolean; + plan?: Plan; commitHead?: string; pathKeyError?: Error; materializeError?: Error } = {}) { + const dir = mkdtempSync(join(tmpdir(), 'codeboost-exec-')); dirs.push(dir); + const path = join(dir, 'state.sqlite'), store = new Store(path); + store.createPlan(JSON.stringify(options.plan ?? plan), 'json', context, oid(1), oid(2)); + store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); + const log: string[] = [], commits: { item: string; baseHead: string; paths: readonly string[]; trailers: Record; digest: string; message: string }[] = []; + let next = 100; + const itemOf = (ws: WorkspaceRef) => (ws.storage as { item: string }).item; + const workspace: TaskWorkspace = { + async materialize(attempt, head) { log.push(`materialize ${attempt.item} @${head.slice(-3)}`); if (options.materializeError) throw options.materializeError; return { clone: { id: `c-${attempt.id}`, taskId: 'task', directory: '/tmp/x', head }, storage: { item: attempt.item, attemptId: attempt.id } }; }, + async snapshotDeclaredLinks(ws, paths) { log.push(`snapshot ${itemOf(ws)} [${paths.join(',')}]`); if (options.snapshotError) throw options.snapshotError; return { item: itemOf(ws) }; }, + async inspectChanges(ws, input, signal) { + log.push(`inspect ${itemOf(ws)} @${input.baseHead.slice(-3)}`); await options.inspect?.(itemOf(ws), signal); + return options.manifests?.[itemOf(ws)] ?? manifest([change(itemOf(ws) === 'P1' ? 'a.ts' : 'b.ts')]); + }, + async commit(ws, input) { + await options.commit?.(itemOf(ws)); + const head = options.commitHead ?? oid(next++); commits.push({ item: itemOf(ws), baseHead: input.baseHead, paths: input.paths, trailers: { ...input.trailers }, digest: input.digest, message: input.message }); + log.push(`commit ${itemOf(ws)} -> ${head.slice(-3)}`); return head; + }, + async release(ws) { + const attemptId = (ws.storage as { attemptId: string }).attemptId; + log.push(`release ${itemOf(ws)} after ${store.getAttempt(identity, attemptId).state}`); + await options.release?.(); + }, + }; + const auditContext: PlanContext = options.pathKeyError ? { ...context, pathKey: () => { throw options.pathKeyError; } } : context; + const sources: ExecutionSources = { planContext: () => auditContext, issue: () => ({ number: 1, title: 'Issue', body: 'Please fix', comments: [] }), lessons: () => [], vendor: () => 'claude' }; + const prompts: string[] = [], argv: (readonly (readonly string[])[])[] = [], owners: string[] = []; + const findings = new SafetyFindings(), capability = options.capability?.(store); + if (options.settleError) store.settleAttempt = () => { throw Object.assign(new Error('disk full'), { code: 'ERR_SQLITE_ERROR' }); }; + const deps = executionDeps(store, workspace, (input, prompt, ws) => { + if (options.startError) throw options.startError; + log.push(`start ${itemOf(ws)}`); prompts.push(prompt); argv.push(input.approvedArgv); owners.push(input.runnerOwner); + return { attemptId: input.attemptId, settled: Promise.resolve({ attemptId: input.attemptId, context: input.context, exitCode: 0, signal: null, stdout: 'done', stderr: '', ...options.exit?.[itemOf(ws)] }), cancel: () => undefined }; + }, sources, RUNNER_OWNER, findings); + const runner = new RunnerCoordinator(store, deps, undefined, capability); + cleanups.push(async () => { await runner.close(); store.close(); }); + return { store, path, workspace, findings, runner, executor: new ItemExecutor(store, runner, sources, findings, { capability }), log, commits, prompts, argv, owners }; +} + +describe('item execution', () => { + it('runs items in order, commits each with trailers, and records owned ledger entries', async () => { + const { store, executor, log, commits, prompts, argv, owners } = setup(); + expect(await executor.runTask(identity)).toEqual({ kind: 'executed', items: ['P1', 'P2'], unchanged: [] }); + expect(commits.map(c => [c.item, c.baseHead.slice(-3), c.trailers, c.paths, c.message])).toEqual([ + ['P1', '002', { 'Plan-Item': 'P1', 'Plan-Revision': 'r1' }, ['a.ts'], 'P1: First'], + ['P2', '064', { 'Plan-Item': 'P2', 'Plan-Revision': 'r1' }, ['b.ts'], 'P2: Second']]); + expect(store.getLedger(identity)).toEqual(expect.arrayContaining([ + { sha: oid(100), owner: 'P1', origin: 'owned', sourceSha: null }, { sha: oid(101), owner: 'P2', origin: 'owned', sourceSha: null }])); + expect(store.getSnapshot(identity).head).toBe(oid(101)); + expect(log).toEqual([ + 'materialize P1 @002', 'snapshot P1 [a.ts]', 'start P1', 'inspect P1 @002', 'commit P1 -> 064', 'release P1 after completed', + 'materialize P2 @064', 'snapshot P2 [b.ts]', 'start P2', 'inspect P2 @064', 'commit P2 -> 065', 'release P2 after completed']); + expect(prompts[0]).toContain(''); + expect(argv[0]).toEqual([['npm', 'test']]); + expect(owners[0]).toBe(RUNNER_OWNER); + expect(argv[1]).toEqual([]); + }); + it('reports a planned-but-unchanged item without committing', async () => { + const { executor, commits } = setup({ manifests: { P1: manifest([]) } }); + expect(await executor.runTask(identity)).toEqual({ kind: 'executed', items: ['P1', 'P2'], unchanged: ['P1'] }); + expect(commits.map(c => c.item)).toEqual(['P2']); + }); + it('commits out-of-scope files with the item, records a checkpoint, and pauses in needs amendment', async () => { + const { store, executor, commits, log } = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) } }); + const outcome = await executor.runTask(identity); + expect(outcome).toMatchObject({ kind: 'needs amendment', item: 'P1', outOfScope: ['extra.ts'] }); + expect(commits[0]!.paths).toEqual(['a.ts', 'extra.ts']); + expect(store.getCheckpoint(identity, (outcome as { checkpointId: string }).checkpointId)).toMatchObject({ item: 'P1', completedItems: ['P1'], outOfScopePaths: ['extra.ts'] }); + expect(store.getTask(identity).status).toBe('needs amendment'); + expect(log.some(line => line.includes('P2'))).toBe(false); + }); + it('stops a safety violation before any commit, fails the attempt, and moves the task to needs human', async () => { + const { store, executor, commits, log } = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) } }); + const outcome = await executor.runTask(identity); + expect(outcome).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect((outcome as { reason: string }).reason.startsWith(SAFETY_VIOLATION)).toBe(true); + expect(commits).toEqual([]); + expect(store.getTask(identity).status).toBe('needs human'); + expect(store.getLedger(identity).some(entry => entry.owner === 'P1')).toBe(false); + expect(log).toContain('release P1 after failed'); + }); + it('stops on an agent failure without inspecting or committing', async () => { + const { store, executor, log } = setup({ exit: { P1: { exitCode: 1, stderr: 'agent crashed' } } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'failed', reason: '"agent crashed"' }); + expect(log.some(line => line.startsWith('inspect'))).toBe(false); + expect(store.getSnapshot(identity).head).toBe(oid(2)); + }); + it('records no ledger entry when the commit is refused', async () => { + const { store, executor } = setup({ commit: async () => { throw new Error('work tree changed after the audit'); } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'failed', reason: 'The runner commit was refused: "work tree changed after the audit"' }); + expect(store.getLedger(identity)).toEqual([]); + }); + it('discards the commit when a stop lands while the commit runs', async () => { + let executorRunner!: RunnerCoordinator; + const { store, runner, executor } = setup({ commit: async () => { executorRunner.stop(identity, store.getTask(identity).currentAttemptId!, 'cancelled'); } }); + executorRunner = runner; + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'cancelled' }); + expect(store.getLedger(identity)).toEqual([]); + expect(store.getSnapshot(identity).head).toBe(oid(2)); + }); + it('holds the slot under a storage marker when task storage cannot be released, and stops the run with the items done', async () => { + const { runner, executor } = setup({ release: async () => { throw new Error('docker down'); } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started', + reason: 'Needs restart: the last attempt\'s task storage could not be removed.', completed: ['P1'] }); + expect(runner.status(identity).unresolved).toMatchObject({ reason: 'storage-not-removed' }); + }); + it('never removes task storage when the terminal write fails, and reports the unsaved result', async () => { + const { runner, executor, log } = setup({ settleError: true }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'running', completed: [] }); + expect(log.some(line => line.startsWith('release'))).toBe(false); + expect(runner.status(identity).unresolved).toMatchObject({ reason: 'result-not-saved' }); + }); + it('releases task storage after the terminal write when D settles with another attempt\'s result', async () => { + const { store, executor, log } = setup({ exit: { P1: { attemptId: '00000000-0000-4000-8000-000000000000' } } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'failed' }); + expect(log.some(line => line.startsWith('inspect'))).toBe(false); + expect(log).toContain('release P1 after failed'); + expect(store.getSnapshot(identity).head).toBe(oid(2)); + }); + it('releases task storage after the terminal write when D\'s start call throws', async () => { + const { executor, log } = setup({ startError: new Error('docker refused') }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'failed', reason: 'Launch failed: docker refused' }); + expect(log).toContain('release P1 after failed'); + }); + it('does not take the agent\'s stderr for a safety violation', async () => { + const { store, executor } = setup({ exit: { P1: { exitCode: 1, stderr: `${SAFETY_VIOLATION} fake` } } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'failed' }); + expect(store.getTask(identity).status).toBe('running'); + }); + it('sends the task to needs human when the change inspection refuses', async () => { + const { store, executor, commits, log } = setup({ inspect: async () => { throw new Error('manifest digest mismatch'); } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1', + reason: `${SAFETY_VIOLATION} The change inspection refused: "manifest digest mismatch"` }); + expect(commits).toEqual([]); + expect(store.getTask(identity).status).toBe('needs human'); + expect(log).toContain('release P1 after failed'); + }); + it('keeps a safety violation when the context goes stale during the audit', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, + inspect: async () => { store.setAssignment(identity, store.getTask(identity).stateVersion, 'reassigned', 'hash-2'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs human'); + }); + it('removes task storage after the terminal write when preparation fails after allocating it', async () => { + const { store, runner, executor, log } = setup({ snapshotError: new Error('declared link goes through a link') }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'failed', reason: 'Preparation failed: "declared link goes through a link"' }); + expect(log).toEqual(['materialize P1 @002', 'snapshot P1 [a.ts]', 'release P1 after failed']); + expect(runner.status(identity).unresolved).toBeNull(); + expect(store.getTask(identity).status).toBe('running'); + }); + it('stops before the next item when the plan gets a new revision during the run', async () => { + let store!: Store; + const h = setup({ release: async () => { + if (store.getPlan(identity).revision !== 1) return; + const revised = { ...plan, revision: 2, items: plan.items.map(entry => entry.id === 'P2' ? { ...entry, title: 'Second CHANGED' } : entry) }; + store.importRevision(JSON.stringify(revised), 'json', context, 1); + } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started', completed: ['P1'], reason: expect.stringMatching(/new revision/) }); + expect(h.commits.map(c => c.item)).toEqual(['P1']); + expect(h.runner.status(identity).unresolved).toBeNull(); + }); + it('still pauses for amendment, bound to where the item ran, when the plan changes after its attempt settled', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, release: async () => { + if (store.getPlan(identity).revision === 1) store.importRevision(JSON.stringify({ ...plan, revision: 2, summary: 'Revised' }), 'json', context, 1); + } }); + store = h.store; + const outcome = await h.executor.runTask(identity); + expect(outcome).toMatchObject({ kind: 'needs amendment', item: 'P1', outOfScope: ['extra.ts'], completed: ['P1'] }); + // The revision really changed during release (a failed import there would be absorbed as a storage failure). + expect(store.getPlan(identity).revision).toBe(2); + expect(h.runner.status(identity).unresolved).toBeNull(); + expect(store.getCheckpoint(identity, (outcome as { checkpointId: string }).checkpointId)).toMatchObject({ + revision: 1, snapshotId: store.snapshotWithHead(identity, oid(100)) }); + expect(store.getTask(identity).status).toBe('needs amendment'); + // The finding is not skipped by resuming at the next item. + expect(await h.executor.runTask(identity, { fromItem: 'P2' })).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started' }); + expect(h.commits.map(c => c.item)).toEqual(['P1']); + }); + it('returns a stopped outcome, with no checkpoint, when the task closed before it could pause for amendment', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, release: async () => { + store.cancelTask(identity, store.getTask(identity).stateVersion, randomUUID()); + } }); + store = h.store; + const outcome = await h.executor.runTask(identity); + expect(outcome).toMatchObject({ kind: 'stopped', item: 'P1', completed: ['P1'] }); + expect(store.getTask(identity).status).toBe('cancelled'); + // The refused pause recorded no checkpoint either (one transaction). + const db = new DatabaseSync(h.path); + try { expect(db.prepare('SELECT COUNT(*) AS n FROM checkpoints').get()).toEqual({ n: 0 }); } finally { db.close(); } + }); + it('pauses for amendment through the capability after the shutdown write gate closed', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, + capability: s => s.shutdownCapability(), release: async () => { store.closeWrites(); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs amendment', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs amendment'); + }); + it('stages both paths of a rename in the runner commit', async () => { + const { executor, commits } = setup({ manifests: { P1: manifest([change('a.ts', { kind: 'rename', oldPath: 'old.ts' })]) } }); + await executor.runTask(identity); + expect(commits[0]!.paths).toEqual(['a.ts', 'old.ts']); + }); + it('treats an inspection aborted by a stop as that stop, not a finding', async () => { + let runner!: RunnerCoordinator, store!: Store; + const h = setup({ inspect: async (_item, signal) => { + runner.stop(identity, store.getTask(identity).currentAttemptId!, 'cancelled'); + throw signal.reason; + } }); + runner = h.runner; store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'cancelled' }); + expect(store.getTask(identity).status).toBe('running'); + expect(h.commits).toEqual([]); + }); + it('keeps a real inspection refusal as a finding even when a stop is pending', async () => { + let runner!: RunnerCoordinator, store!: Store; + const h = setup({ inspect: async () => { + runner.stop(identity, store.getTask(identity).currentAttemptId!, 'cancelled'); + throw new Error('metadata digest changed'); + } }); + runner = h.runner; store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1', + reason: `${SAFETY_VIOLATION} The change inspection refused: "metadata digest changed"` }); + expect(store.getTask(identity).status).toBe('needs human'); + }); + it('makes no commit when a stop lands during an inspection that ignores the abort', async () => { + let runner!: RunnerCoordinator, store!: Store; + const h = setup({ inspect: async () => { runner.stop(identity, store.getTask(identity).currentAttemptId!, 'cancelled'); } }); + runner = h.runner; store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'cancelled' }); + expect(h.commits).toEqual([]); + }); + it('sends a malformed change report to needs human', async () => { + const { store, executor, commits } = setup({ manifests: { P1: manifest([change('a.ts')], { linkTargetChanges: undefined as unknown as string[] }) } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs human'); + expect(commits).toEqual([]); + }); + it('makes no commit when the context changes during the audit', async () => { + let store!: Store; + const h = setup({ inspect: async () => { store.setAssignment(identity, store.getTask(identity).stateVersion, 'reassigned', 'hash-2'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'stale' }); + expect(h.commits).toEqual([]); + }); + it('pauses first on the next run when a scope pause was never recorded, and never runs past it', async () => { + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) } }); + const pause = h.store.pauseForAmendment.bind(h.store); + let fail = true; + h.store.pauseForAmendment = (...args) => { if (fail) { fail = false; throw Object.assign(new Error('disk full'), { code: 'ERR_SQLITE_ERROR' }); } return pause(...args); }; + await expect(h.executor.runTask(identity)).rejects.toThrow(/disk full/); + expect(h.store.getTask(identity).status).toBe('running'); + expect(await h.executor.runTask(identity, { fromItem: 'P2' })).toMatchObject({ kind: 'needs amendment', item: 'P1', outOfScope: ['extra.ts'] }); + expect(h.store.getTask(identity).status).toBe('needs amendment'); + expect(h.commits.map(c => c.item)).toEqual(['P1']); + }); + it('does not take a closed write gate for a refused pause when it has no capability', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, release: async () => { store.closeWrites(); } }); + store = h.store; + await expect(h.executor.runTask(identity)).rejects.toBeInstanceOf(ShuttingDownError); + }); + it('stops before the next item when the assignment changes during the run', async () => { + let store!: Store; + const h = setup({ release: async () => { + if (store.getTask(identity).currentAttemptId && store.getAttempts(identity).length === 1) + store.setAssignment(identity, store.getTask(identity).stateVersion, 'reassigned', 'hash-2'); + } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started', completed: ['P1'] }); + expect(h.commits.map(c => c.item)).toEqual(['P1']); + expect(h.runner.status(identity).unresolved).toBeNull(); + }); + it('stops before the next item, and binds a pause to the item\'s own commit, when HEAD is observed during the run', async () => { + let store!: Store; + const observe = () => { const snapshot = store.getSnapshot(identity); store.recordHistory(identity, { revision: 1, snapshotId: snapshot.id }, snapshot.base, oid(999), []); }; + const clean = setup({ release: async () => { if (store.getAttempts(identity).length === 1) observe(); } }); + store = clean.store; + expect(await clean.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started', reason: expect.stringMatching(/snapshot or assignment changed/) }); + expect(clean.runner.status(identity).unresolved).toBeNull(); + const scoped = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, release: async () => observe() }); + store = scoped.store; + const outcome = await scoped.executor.runTask(identity); + expect(outcome).toMatchObject({ kind: 'needs amendment', item: 'P1' }); + const checkpoint = store.getCheckpoint(identity, (outcome as { checkpointId: string }).checkpointId); + expect(checkpoint.snapshotId).toBe(store.snapshotWithHead(identity, oid(100))); + expect(checkpoint.snapshotId).not.toBe(store.getSnapshot(identity).id); + }); + it('fails the attempt, instead of breaking the terminal write, when the workspace returns an invalid commit ID', async () => { + const { runner, executor, log } = setup({ commitHead: 'HEAD' }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'failed', reason: 'The workspace returned an invalid commit ID; nothing was published.' }); + expect(runner.status(identity).unresolved).toBeNull(); + expect(log).toContain('release P1 after failed'); + }); + it('keeps a human gate set during release, and escalates a safety violation owed from it when the task next runs', async () => { + let store!: Store; + const toApproval = async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs approval'); }; + const scoped = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, release: toApproval }); + store = scoped.store; + expect(await scoped.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs approval'); + const unsafe = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, release: toApproval }); + store = unsafe.store; + expect(await unsafe.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', reason: expect.stringMatching(/needs approval; it moves to needs human when it next runs/) }); + expect(store.getTask(identity).status).toBe('needs approval'); + // Still at the gate on the next run: the finding stays owed, reported as not started. + expect(await unsafe.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'not started', reason: expect.stringMatching(/needs approval; it moves to needs human/) }); + // A person releases the gate; the owed finding goes to needs human before any item runs again. + store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); + expect(await unsafe.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs human'); + expect(store.getAttempts(identity)).toHaveLength(1); + }); + it('treats an audit that throws as a safety violation', async () => { + const { store, executor, commits } = setup({ pathKeyError: new Error('Non-ASCII case-insensitive paths require an adapter.') }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1', + reason: `${SAFETY_VIOLATION} The change report could not be audited: "Non-ASCII case-insensitive paths require an adapter."` }); + expect(store.getTask(identity).status).toBe('needs human'); + expect(commits).toEqual([]); + }); + it('snapshots both sides of a declared rename before launch', async () => { + const renamed: Plan = { ...plan, items: [{ ...plan.items[0]!, files: [{ path: 'c.ts', kind: 'rename', renamed_from: 'a.ts', change: 'move' }] }, plan.items[1]!] }; + const h = setup({ plan: renamed, manifests: { P1: manifest([change('c.ts', { kind: 'rename', oldPath: 'a.ts' })]) } }); + await h.executor.runTask(identity); + expect(h.log).toContain('snapshot P1 [c.ts,a.ts]'); + }); + it('pays a pause owed from an earlier run when the task is queued again, instead of getting stuck', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs approval'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1' }); + store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); + expect(await h.executor.runTask(identity, { fromItem: 'P2' })).toMatchObject({ kind: 'needs amendment', item: 'P1', outOfScope: ['extra.ts'] }); + expect(store.getTask(identity).status).toBe('needs amendment'); + expect(h.commits.map(c => c.item)).toEqual(['P1']); + }); + it('runs no further items once a task has a scope checkpoint, even after an approved continuation (not supported yet)', async () => { + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) } }); + const paused = await h.executor.runTask(identity); + expect(paused).toMatchObject({ kind: 'needs amendment', item: 'P1' }); + const store = h.store; + store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); + const refused = { kind: 'stopped', state: 'not started', reason: expect.stringMatching(/Continuing after a scope pause is not supported yet/) }; + expect(await h.executor.runTask(identity, { fromItem: 'P2' })).toMatchObject(refused); + const amended = { ...plan, revision: 2, items: [{ ...plan.items[0]!, files: [...plan.items[0]!.files, { path: 'extra.ts', kind: 'add', renamed_from: null, change: 'z' }] }, plan.items[1]!] }; + store.importRevision(JSON.stringify(amended), 'json', context, 1); + store.approveContinuation(identity, (paused as { checkpointId: string }).checkpointId, { revision: 2, snapshotId: store.getSnapshot(identity).id }); + expect(await h.executor.runTask(identity, { fromItem: 'P2' })).toMatchObject(refused); + // A later snapshot with the same head does not make the recorded pause owed again. + const snapshot = store.getSnapshot(identity); + store.recordHistory(identity, { revision: 2, snapshotId: snapshot.id }, snapshot.base, snapshot.head, []); + expect(await h.executor.runTask(identity, { fromItem: 'P2' })).toMatchObject(refused); + expect(store.getTask(identity).status).toBe('queued'); + expect(h.commits.map(c => c.item)).toEqual(['P1']); + }); + it('fails the attempt when the workspace makes no new commit for a changed item', async () => { + const { store, executor } = setup({ commitHead: oid(2) }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'failed', reason: 'The workspace made no new commit for a changed item; nothing was published.' }); + expect(store.getLedger(identity).some(entry => entry.owner === 'P1')).toBe(false); + }); + it('names the real cause when the start of an attempt could not be saved', async () => { + const h = setup(); + h.store.markRunning = () => { throw Object.assign(new Error('disk full'), { code: 'ERR_SQLITE_ERROR' }); }; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'pending', reason: 'Needs restart: the start of the last attempt could not be saved.' }); + }); + it('stops before the next item when only the referenced code changes during the run', async () => { + let store!: Store; + const h = setup({ release: async () => { + if (store.getAttempts(identity).length !== 1) return; + store.setAssignment(identity, store.getTask(identity).stateVersion, store.currentContext(identity).assignmentId, 'new-code-hash'); + } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started', completed: ['P1'], reason: expect.stringMatching(/snapshot or assignment changed/) }); + expect(h.runner.status(identity).unresolved).toBeNull(); + }); + it('treats an AbortError from the inspection as the stop', async () => { + let runner!: RunnerCoordinator, store!: Store; + const h = setup({ inspect: async () => { + runner.stop(identity, store.getTask(identity).currentAttemptId!, 'cancelled'); + throw Object.assign(new Error('The operation was aborted'), { name: 'AbortError' }); + } }); + runner = h.runner; store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'cancelled' }); + expect(store.getTask(identity).status).toBe('running'); + }); + it('keeps a queued status set during release instead of pausing, but escalates a safety violation so the item is not re-run', async () => { + let store!: Store; + const toQueued = async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); }; + const scoped = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, release: toQueued }); + store = scoped.store; + expect(await scoped.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1' }); + expect(store.getTask(identity).status).toBe('queued'); + const unsafe = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, release: toQueued }); + store = unsafe.store; + expect(await unsafe.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs human'); + expect(await unsafe.executor.runTask(identity)).toMatchObject({ kind: 'stopped', state: 'not started' }); + expect(store.getAttempts(identity)).toHaveLength(1); + }); + it('keeps task storage when a foreign result\'s terminal write fails', async () => { + const { runner, executor, log } = setup({ settleError: true, exit: { P1: { attemptId: '00000000-0000-4000-8000-000000000000' } } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'running' }); + expect(log.some(line => line.startsWith('release'))).toBe(false); + expect(runner.status(identity).unresolved).toMatchObject({ reason: 'result-not-saved' }); + }); + it('leaves a safety violation\'s task alone when it was cancelled during release', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, + release: async () => { store.cancelTask(identity, store.getTask(identity).stateVersion, randomUUID()); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', reason: expect.stringMatching(/is cancelled, so it was not moved/) }); + expect(store.getTask(identity).status).toBe('cancelled'); + }); + it('removes task storage after the terminal write when a stop lands during preparation', async () => { + let runner!: RunnerCoordinator, store!: Store; + const h = setup(); + runner = h.runner; store = h.store; + const snapshot = h.workspace.snapshotDeclaredLinks.bind(h.workspace); + h.workspace.snapshotDeclaredLinks = async (ws, paths, signal) => { runner.stop(identity, store.getTask(identity).currentAttemptId!, 'cancelled'); return snapshot(ws, paths, signal); }; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'cancelled' }); + expect(h.log).toEqual(['materialize P1 @002', 'snapshot P1 [a.ts]', 'release P1 after cancelled']); + }); + it('removes task storage after the terminal write when the context goes stale before launch', async () => { + let store!: Store; + const h = setup(); + store = h.store; + const snapshot = h.workspace.snapshotDeclaredLinks.bind(h.workspace); + h.workspace.snapshotDeclaredLinks = async (ws, paths, signal) => { + store.setAssignment(identity, store.getTask(identity).stateVersion, 'reassigned', 'hash-2'); return snapshot(ws, paths, signal); + }; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'stale' }); + expect(h.log).toEqual(['materialize P1 @002', 'snapshot P1 [a.ts]', 'release P1 after stale']); + }); + it('refuses a scope pause whose executed prefix does not match the plan at the item\'s revision', () => { + const { store } = setup(); + const snapshotId = store.getSnapshot(identity).id; + expect(() => store.pauseForAmendment(identity, { revision: 1, snapshotId }, { item: 'P2', baseEntries: [], completedItems: ['P2'], outOfScopePaths: ['x'] })) + .toThrow(/executed plan prefix/); + expect(() => store.pauseForAmendment(identity, { revision: 1, snapshotId }, { item: 'P1', baseEntries: [], completedItems: ['P1', 'P2'], outOfScopePaths: ['x'] })) + .toThrow(/executed plan prefix/); + }); + it('bounds a finding whose text comes from the workspace', async () => { + const { executor } = setup({ inspect: async () => { throw new Error('x'.repeat(10_000)); } }); + const outcome = await executor.runTask(identity) as { kind: string; reason: string }; + expect(outcome.kind).toBe('needs human'); + expect(outcome.reason.length).toBeLessThanOrEqual(MAX_REASON); + }); + it('removes task storage after the terminal write when a cancel task lands on the row during preparation', async () => { + let store!: Store; + const h = setup(); + store = h.store; + const snapshot = h.workspace.snapshotDeclaredLinks.bind(h.workspace); + h.workspace.snapshotDeclaredLinks = async (ws, paths, signal) => { + store.cancelTask(identity, store.getTask(identity).stateVersion, randomUUID()); return snapshot(ws, paths, signal); + }; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'cancelled' }); + expect(h.log).toEqual(['materialize P1 @002', 'snapshot P1 [a.ts]', 'release P1 after cancelled']); + expect(store.getTask(identity).status).toBe('cancelled'); + }); + it('keeps a safety finding owed when moving to needs human fails, and escalates it before the next run launches anything', async () => { + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) } }); + const transition = h.store.transitionTask.bind(h.store); + let fail = true; + h.store.transitionTask = (...args) => { if (fail && args[2] === 'needs human') { fail = false; throw Object.assign(new Error('disk full'), { code: 'ERR_SQLITE_ERROR' }); } return transition(...args); }; + await expect(h.executor.runTask(identity)).rejects.toThrow(/disk full/); + expect(h.store.getTask(identity).status).toBe('running'); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(h.store.getAttempts(identity)).toHaveLength(1); + }); + it('stops before the next item when someone changes the status during release', async () => { + let store!: Store; + const h = setup({ release: async () => { if (store.getAttempts(identity).length === 1) store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started', reason: expect.stringMatching(/changed to queued/) }); + expect(h.commits.map(c => c.item)).toEqual(['P1']); + }); + it('refuses an owed pause from a human gate, and validates the prefix against the item\'s own revision', () => { + const { store } = setup(); + const snapshotId = store.getSnapshot(identity).id; + // Revision 2 drops P2; the prefix [P1, P2] is still right for revision 1, where the item ran. + store.importRevision(JSON.stringify({ ...plan, revision: 2, items: [plan.items[0]!] }), 'json', context, 1); + const evidence = { item: 'P2', baseEntries: [], completedItems: ['P1', 'P2'], outOfScopePaths: ['x'] }; + store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs approval'); + expect(() => store.pauseForAmendment(identity, { revision: 1, snapshotId }, evidence, { owed: true })).toThrow(/is needs approval/); + store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); + expect(store.pauseForAmendment(identity, { revision: 1, snapshotId }, evidence, { owed: true })).toMatchObject({ revision: 1, completedItems: ['P1', 'P2'] }); + }); + it('escalates a safety violation over a review status, so the task cannot be merged past it', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'in review'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs human'); + }); + it('settles a finding whose task is already needs human, so it is not escalated again later', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, + release: async () => { if (store.getAttempts(identity).length === 1) store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs human'); } }); + store = h.store; + let version = 0; + const releaseDone = h.workspace.release.bind(h.workspace); + h.workspace.release = async ws => { await releaseDone(ws); version = store.getTask(identity).stateVersion; }; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(h.findings.get(store.getAttempts(identity)[0]!.id)).toBeUndefined(); + // Settled without a needs human -> needs human write. + expect(store.getTask(identity).stateVersion).toBe(version); + }); + it('keeps a finding owed when a merge in progress refuses the escalation, and escalates it on the next run', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, release: async () => { + store.transitionTask(identity, store.getTask(identity).stateVersion, 'in review'); + const snapshot = store.getSnapshot(identity); + store.beginMergeAttempt(identity, { revision: 1, snapshotId: snapshot.id, reviewVersion: store.reviewVersion(identity) }, snapshot.head, null, 'direct'); + } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', reason: expect.stringMatching(/could not be moved to needs human yet: A merge is in progress/) }); + const attemptId = store.getAttempts(identity)[0]!.id; + expect(h.findings.get(attemptId)).toBeDefined(); + store.finishMergeAttempt(identity, store.getMergeAttempt(identity)!.id, { state: 'failed', reason: 'GitHub refused.' }); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(h.findings.get(attemptId)).toBeUndefined(); + }); + it('escalates through the capability after the shutdown write gate closed', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, capability: s => s.shutdownCapability(), + release: async () => { store.closeWrites(); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs human'); + }); + it('settles an owed finding when the task was closed before its next run', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs approval'); } }); + store = h.store; + await h.executor.runTask(identity); + const attemptId = store.getAttempts(identity)[0]!.id; + expect(h.findings.get(attemptId)).toBeDefined(); + store.cancelTask(identity, store.getTask(identity).stateVersion, randomUUID()); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'not started', reason: expect.stringMatching(/is cancelled/) }); + expect(h.findings.get(attemptId)).toBeUndefined(); + }); + it('records completed and the ledger entry in one transaction: a failed history write leaves neither', async () => { + const h = setup(); + h.store.recordHistory = () => { throw Object.assign(new Error('disk full'), { code: 'ERR_SQLITE_ERROR' }); }; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'running' }); + expect(h.store.getLedger(identity)).toEqual([]); + expect(h.runner.status(identity).unresolved).toMatchObject({ reason: 'result-not-saved' }); + }); + it('pauses for amendment over a review status set during release, so a merge cannot go past the finding', async () => { + for (const status of ['in review', 'approved but merge blocked'] as const) { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, status); } }); + store = h.store; + expect(await h.executor.runTask(identity), status).toMatchObject({ kind: 'needs amendment', item: 'P1' }); + expect(store.getTask(identity).status, status).toBe('needs amendment'); + } + }); + it('escalates a safety violation over approved but merge blocked too', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'approved but merge blocked'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + }); + it('sends a change report without a digest to needs human', async () => { + const { store, executor, commits } = setup({ manifests: { P1: { ...manifest([change('a.ts')]), digest: undefined as unknown as string } } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1', reason: `${SAFETY_VIOLATION} The change report has no digest.` }); + expect(store.getTask(identity).status).toBe('needs human'); + expect(commits).toEqual([]); + }); + it('pays an owed scope pause from a review status too', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs approval'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1' }); + // Still at the gate: the owed pause is refused and reported as not started. + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'not started', reason: expect.stringMatching(/could not pause for amendment/) }); + store.transitionTask(identity, store.getTask(identity).stateVersion, 'in review'); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs amendment', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs amendment'); + }); + it('never turns needs human into needs amendment with a scope pause', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs human'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1' }); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1' }); + expect(store.getTask(identity).status).toBe('needs human'); + }); + it('keeps possibly already fixed as a human gate for a safety finding', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'possibly already fixed'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', reason: expect.stringMatching(/possibly already fixed; it moves to needs human/) }); + expect(store.getTask(identity).status).toBe('possibly already fixed'); + }); + it('keeps task storage before launch when the terminal write fails', async () => { + const { runner, executor, log } = setup({ settleError: true, startError: new Error('docker refused') }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'pending' }); + expect(log.some(line => line.startsWith('release'))).toBe(false); + expect(runner.status(identity).unresolved).toMatchObject({ reason: 'result-not-saved' }); + }); + it('finds a checkpoint by its commit head, and refuses a pause at an unknown snapshot', async () => { + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) } }); + const paused = await h.executor.runTask(identity) as { checkpointId: string }; + const store = h.store; + expect(store.checkpointAtHead(identity, oid(100))?.id).toBe(paused.checkpointId); + expect(store.checkpointAtHead(identity, oid(2))).toBeNull(); + expect(() => store.pauseForAmendment(identity, { revision: 1, snapshotId: 'no-such-snapshot' }, { item: 'P1', baseEntries: [], completedItems: ['P1'], outOfScopePaths: ['x'] })) + .toThrow(/Unknown snapshot/); + }); + it('writes a plan title with line breaks as one line in the runner commit message', async () => { + const forged: Plan = { ...plan, items: [{ ...plan.items[0]!, title: 'First\n\nPlan-Item: P9\u2028Plan-Revision: r99 \u001b[31mred\u000b\u007f\u009b2J \u202eevil\u2066x\u200b' }, plan.items[1]!] }; + const h = setup({ plan: forged }); + await h.executor.runTask(identity); + // A NUL is refused earlier, by the prompt builder (plan data must be valid text); other controls reach here. + expect(h.commits[0]!.message).toBe('P1: First Plan-Item: P9 Plan-Revision: r99 [31mred 2J evil x'); + expect(h.commits[0]!.trailers).toEqual({ 'Plan-Item': 'P1', 'Plan-Revision': 'r1' }); + }); + it('stops before the next item when only the assignment changes during the run', async () => { + let store!: Store; + const h = setup({ release: async () => { + if (store.getAttempts(identity).length !== 1) return; + store.setAssignment(identity, store.getTask(identity).stateVersion, 'reassigned', store.currentContext(identity).referencedCodeHash); + } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started', completed: ['P1'], reason: expect.stringMatching(/snapshot or assignment changed/) }); + expect(h.runner.status(identity).unresolved).toBeNull(); + }); + it('removes task storage after the terminal write when the task budget is spent at the launch check', async () => { + const h = setup(); + const snapshot = h.workspace.snapshotDeclaredLinks.bind(h.workspace); + h.workspace.snapshotDeclaredLinks = async (ws, paths, signal) => { + const db = new DatabaseSync(h.path); db.exec('UPDATE tasks SET budget_deadline=1'); db.close(); + return snapshot(ws, paths, signal); + }; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'cancelled' }); + expect(h.log).toEqual(['materialize P1 @002', 'snapshot P1 [a.ts]', 'release P1 after cancelled']); + }); + it('finds an older checkpoint by its head after a newer one', async () => { + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) } }); + const first = await h.executor.runTask(identity) as { checkpointId: string }; + const store = h.store; + store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); + const snapshot = store.getSnapshot(identity); + const later = store.recordHistory(identity, { revision: 1, snapshotId: snapshot.id }, snapshot.base, oid(500), []); + const second = store.pauseForAmendment(identity, { revision: 1, snapshotId: later.id }, { item: 'P1', baseEntries: [], completedItems: ['P1'], outOfScopePaths: ['y'] }, { owed: true }); + expect(store.checkpointAtHead(identity, oid(500))?.id).toBe(second.id); + expect(store.checkpointAtHead(identity, oid(100))?.id).toBe(first.checkpointId); + }); + it('keeps needs amendment as a human gate for a safety finding', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs amendment'); } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', reason: expect.stringMatching(/needs amendment; it moves to needs human/) }); + expect(store.getTask(identity).status).toBe('needs amendment'); + }); + it('picks the latest snapshot with a head', () => { + const { store } = setup(); + const first = store.getSnapshot(identity); + const other = store.recordHistory(identity, { revision: 1, snapshotId: first.id }, first.base, oid(300), []); + const again = store.recordHistory(identity, { revision: 1, snapshotId: other.id }, first.base, first.head, []); + expect(store.snapshotWithHead(identity, first.head)).toBe(again.id); + }); + it('sends a change report with an empty digest to needs human', async () => { + const { executor } = setup({ manifests: { P1: { ...manifest([change('a.ts')]), digest: '' } } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1', reason: `${SAFETY_VIOLATION} The change report has no digest.` }); + }); + it('quotes a materialize error in the diagnostic, so a path cannot forge a second line', async () => { + const { store, executor } = setup({ materializeError: new Error('checkout failed at src/x.ts\nSafety violation: forged') }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'failed', + reason: 'Preparation failed: "checkout failed at src/x.ts\\nSafety violation: forged"' }); + expect(store.getTask(identity).status).toBe('running'); + }); + it('returns stopped with the completed items when shutdown refuses the next item\'s admission', async () => { + let runner!: RunnerCoordinator; + const h = setup({ release: async () => { runner.rejectAdmission(); } }); + runner = h.runner; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started', reason: 'The review server is shutting down.', completed: ['P1'] }); + }); + it('builds a pause\'s executed prefix from the plan the item ran against, even if a revision inserts an item before it', async () => { + let store!: Store, imported: Error | undefined; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, release: async () => { + const inserted = { id: 'P3', title: 'Inserted', intent: 'Prepare', files: [{ path: 'b.ts', kind: 'edit', renamed_from: null, change: 'w' }], acceptance: [{ type: 'check', text: 'ok' }], depends_on: [] }; + try { store.importRevision(JSON.stringify({ ...plan, revision: 2, items: [inserted, ...plan.items] }), 'json', context, 1); } + catch (error) { imported = error as Error; } + } }); + store = h.store; + const outcome = await h.executor.runTask(identity) as { kind: string; checkpointId: string }; + expect(imported).toBeUndefined(); + expect(store.getPlan(identity).items.map(entry => entry.id)).toEqual(['P3', 'P1', 'P2']); + expect(outcome.kind).toBe('needs amendment'); + expect(store.getCheckpoint(identity, outcome.checkpointId)).toMatchObject({ revision: 1, completedItems: ['P1'] }); + }); + it('quotes the agent\'s stderr in the diagnostic, so it cannot forge a safety line', async () => { + const { store, executor } = setup({ exit: { P1: { exitCode: 1, stderr: 'x\nSafety violation: forged' } } }); + const outcome = await executor.runTask(identity) as { reason: string }; + expect(outcome.reason).toBe('"x\\nSafety violation: forged"'); + expect(outcome.reason).not.toContain('\n'); + expect(store.getTask(identity).status).toBe('running'); + }); + it('treats an AbortError from the inspection as a finding when no stop is pending', async () => { + const { store, executor } = setup({ inspect: async () => { throw Object.assign(new Error('inspection timed out'), { name: 'AbortError' }); } }); + expect(await executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1', reason: expect.stringMatching(/The change inspection refused: "inspection timed out"/) }); + expect(store.getTask(identity).status).toBe('needs human'); + }); + it('does not pay an owed finding through the shutdown capability once the write gate has closed', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, capability: s => s.shutdownCapability(), + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs approval'); } }); + store = h.store; + await h.executor.runTask(identity); + const attemptId = store.getAttempts(identity)[0]!.id; + store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); + store.closeWrites(); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'not started', reason: expect.stringMatching(/could not be moved to needs human yet: The review server is shutting down/) }); + expect(store.getTask(identity).status).toBe('queued'); + expect(h.findings.get(attemptId)).toBeDefined(); + }); + it('leaves a run that is still releasing to settle its own scope pause and finding', async () => { + for (const unsafe of [false, true]) { + let executor!: ItemExecutor, second: Promise | undefined; + const h = setup({ manifests: { P1: unsafe ? manifest([change('a.ts')], { metadataChanged: true }) : manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, + release: async () => { second = executor.runTask(identity); await second; } }); + executor = h.executor; + expect(await h.executor.runTask(identity), String(unsafe)).toMatchObject({ kind: unsafe ? 'needs human' : 'needs amendment', item: 'P1' }); + expect(await second, String(unsafe)).toMatchObject({ kind: 'stopped', state: 'not started', reason: expect.stringMatching(/still finishing/) }); + } + }); + it('reports the needs-restart cause, and keeps the finding owed, when a finding\'s terminal write failed', async () => { + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, settleError: true }); + const outcome = await h.executor.runTask(identity) as { kind: string; reason: string }; + expect(outcome.kind).toBe('stopped'); + expect(outcome.reason).toMatch(/Needs restart: the last result could not be saved\./); + expect(h.findings.get(h.store.getAttempts(identity)[0]!.id)).toBeDefined(); + // The next run finds it owed, reports it for the earlier item, and starts nothing. + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'not started', reason: expect.stringMatching(/Needs restart/) }); + }); + it('ends as the stop, not a refused commit, when a stop lands and the commit then rejects', async () => { + let runner!: RunnerCoordinator, store!: Store; + const h = setup({ commit: async () => { + runner.stop(identity, store.getTask(identity).currentAttemptId!, 'cancelled'); + throw new Error('commit aborted'); + } }); + runner = h.runner; store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', state: 'cancelled' }); + expect(store.getLedger(identity)).toEqual([]); + }); + it('does not pay an owed scope pause through the shutdown capability once the write gate has closed', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, capability: s => s.shutdownCapability(), + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs approval'); } }); + store = h.store; + await h.executor.runTask(identity); + store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); + store.closeWrites(); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P1', reason: expect.stringMatching(/could not pause for amendment: The review server is shutting down/) }); + expect(store.getTask(identity).status).toBe('queued'); + expect(store.latestCheckpoint(identity)).toBeNull(); + }); + it('records the snapshot\'s base, not the item\'s base head, in the ledger record', async () => { + const h = setup(); + await h.executor.runTask(identity); + const snapshot = h.store.getSnapshot(identity); + expect(snapshot.head).toBe(oid(101)); + expect(snapshot.base).toBe(oid(1)); + }); + it('quotes D\'s error text in the coordinator\'s log lines', async () => { + const lines: string[] = []; + const spy = vi.spyOn(console, 'error').mockImplementation((line: unknown) => { lines.push(String(line)); }); + try { + const h = setup({ release: async () => { throw new Error('rm failed\nRunner job forged: ok'); } }); + await h.executor.runTask(identity); + } finally { spy.mockRestore(); } + expect(lines.some(line => line.endsWith('could not remove its task storage: "rm failed\\nRunner job forged: ok"'))).toBe(true); + expect(lines.every(line => !line.includes('\n'))).toBe(true); + }); + it('finds an owed scope pause even when a later clean item completed after it', async () => { + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) } }); + const store = h.store; + // Record P1's out-of-scope result as if its pause had been lost, then let a clean P2 complete after it. + const pause = store.pauseForAmendment.bind(store); + store.pauseForAmendment = () => { throw Object.assign(new Error('disk full'), { code: 'ERR_SQLITE_ERROR' }); }; + await expect(h.executor.runTask(identity)).rejects.toThrow(/disk full/); + store.pauseForAmendment = pause; + const p2 = h.runner.start(identity, { expectedStateVersion: store.getTask(identity).stateVersion, kind: 'execute', item: 'P2', + expectedContext: store.currentContext(identity), deadline: Date.now() + 60_000 }); + await h.runner.settled(identity); + expect(store.getAttempt(identity, p2.id).state).toBe('completed'); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs amendment', item: 'P1', outOfScope: ['extra.ts'] }); + // Paid from the durable result: no item ran again. + expect(store.getAttempts(identity)).toHaveLength(2); + }); + it('stops before the next item when only the context generation changes during the run', async () => { + let store!: Store; + const h = setup({ release: async () => { + if (store.getAttempts(identity).length !== 1) return; + const current = store.currentContext(identity); + store.setAssignment(identity, store.getTask(identity).stateVersion, current.assignmentId, current.referencedCodeHash); + } }); + store = h.store; + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', item: 'P2', state: 'not started', completed: ['P1'], reason: expect.stringMatching(/snapshot or assignment changed/) }); + expect(h.runner.status(identity).unresolved).toBeNull(); + }); + it('pays nothing owed once shutdown began', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, + release: async () => { store.transitionTask(identity, store.getTask(identity).stateVersion, 'needs approval'); } }); + store = h.store; + await h.executor.runTask(identity); + store.transitionTask(identity, store.getTask(identity).stateVersion, 'queued'); + h.runner.rejectAdmission(); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'stopped', state: 'not started', reason: 'The review server is shutting down.' }); + expect(store.getTask(identity).status).toBe('queued'); + expect(h.findings.get(store.getAttempts(identity)[0]!.id)).toBeDefined(); + }); + it('never lets a second run pay the first run\'s pause, whatever microtask it starts in', async () => { + for (let steps = 0; steps <= 24; steps++) { + let executor!: ItemExecutor, second: Promise | undefined; + const h = setup({ manifests: { P1: manifest([change('a.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) }, release: async () => { + void (async () => { for (let i = 0; i < steps; i++) await null; second = executor.runTask(identity); })(); + } }); + executor = h.executor; + expect(await h.executor.runTask(identity), `after ${steps} steps`).toMatchObject({ kind: 'needs amendment', item: 'P1' }); + for (let i = 0; i < 50 && !second; i++) await null; + expect((await second)?.kind, `after ${steps} steps`).toBe('stopped'); + } + }); + it('does not start while a job started outside the executor is still releasing its storage', async () => { + let executor!: ItemExecutor, during: Promise | undefined; + const h = setup({ release: async () => { during = executor.runTask(identity); await during; } }); + executor = h.executor; + h.runner.start(identity, { expectedStateVersion: h.store.getTask(identity).stateVersion, kind: 'execute', item: 'P1', + expectedContext: h.store.currentContext(identity), deadline: Date.now() + 60_000 }); + await h.runner.settled(identity); + expect(await during).toMatchObject({ kind: 'stopped', state: 'not started', reason: expect.stringMatching(/still finishing/) }); + expect(h.store.getAttempts(identity)).toHaveLength(1); + }); + it('pauses at a later item with the whole executed prefix', async () => { + const h = setup({ manifests: { P2: manifest([change('b.ts'), change('extra.ts', { kind: 'add', oldType: undefined })]) } }); + const outcome = await h.executor.runTask(identity) as { kind: string; item: string; checkpointId: string; completed: string[] }; + expect(outcome).toMatchObject({ kind: 'needs amendment', item: 'P2', completed: ['P1', 'P2'] }); + expect(h.store.getCheckpoint(identity, outcome.checkpointId)).toMatchObject({ item: 'P2', completedItems: ['P1', 'P2'], outOfScopePaths: ['extra.ts'] }); + }); + it('escalates an older owed finding even when a later attempt completed cleanly', async () => { + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) } }); + const store = h.store; + // Both attempts start outside the executor, so nothing acts on P1's finding yet. + h.runner.start(identity, { expectedStateVersion: store.getTask(identity).stateVersion, kind: 'execute', item: 'P1', expectedContext: store.currentContext(identity), deadline: Date.now() + 60_000 }); + await h.runner.settled(identity); + h.runner.start(identity, { expectedStateVersion: store.getTask(identity).stateVersion, kind: 'execute', item: 'P2', expectedContext: store.currentContext(identity), deadline: Date.now() + 60_000 }); + await h.runner.settled(identity); + expect(store.getAttempts(identity).map(row => row.state)).toEqual(['failed', 'completed']); + expect(await h.executor.runTask(identity)).toMatchObject({ kind: 'needs human', item: 'P1' }); + expect(store.getAttempts(identity)).toHaveLength(2); + }); + it('throws, rather than reporting stopped, when this run\'s own escalation meets a closed write gate without a capability', async () => { + let store!: Store; + const h = setup({ manifests: { P1: manifest([change('a.ts')], { metadataChanged: true }) }, release: async () => { store.closeWrites(); } }); + store = h.store; + await expect(h.executor.runTask(identity)).rejects.toBeInstanceOf(ShuttingDownError); + expect(h.findings.get(store.getAttempts(identity)[0]!.id)).toBeDefined(); + }); +});