diff --git a/docs/reference/cli.md b/docs/reference/cli.md index 0e42102b2..5267afc0b 100644 --- a/docs/reference/cli.md +++ b/docs/reference/cli.md @@ -323,7 +323,13 @@ schema bindings, and metadata projections remain provenance for inspection; they are not resume authorization. Smithers decides which finished rows can be reused and which newly rendered or unfinished tasks run. Ultrafuzz does not rewrite historical artifacts or automatically reset, replay, timetravel, or -fork completed work. +fork completed work. When the run's `smithers/resolved-config.json` parses as +the current resolved-config schema, agent adapters in the continued workflow +read the run's `smithers/execution-config.toml` (launch gave them a copy of the +same file); for a run whose resolved config does not parse, resume sets no +`ULTRAFUZZ_CONFIG_PATH`. If resume cannot prune stale task-worktree +registrations, it reports a `WORKFLOW_WORKTREE_REPAIR_FAILED` warning and +continues. `resume --refresh-controller` first renders the currently installed Ultrafuzz controller and stock adapters beside the historical source, then delegates to @@ -458,13 +464,18 @@ an available agent-written report. `pause` requests a graceful stop: no new tasks are scheduled, in-flight tasks finish, and the run settles in the resumable `paused` state. `resume` reports -`submitted: false` instead of launching a duplicate continuation when the linked -workflow is still in an active state (running, in-progress, started, queued, -retrying, or waiting). `resume --reset-node` retries one failed workflow node and -its dependents in the same linked run; the applied reset is recorded so retrying -the command after a failed continuation resumes the already-reset run instead of -repeating the reset. `fork` may start from a checkpoint frame and may reset one -workflow node before starting the fork. +`submitted: false` (text output `Run already active`) instead of launching a +duplicate continuation when the linked workflow is still active (its Smithers +run state is `running`, `recovering`, or one of the `waiting-*` states), and +leaves the run's recorded state and workflow deadline unchanged. A run still +finishing its in-flight tasks after `pause` is still active; resume it again +once `status` reports `paused`. Smithers also reports a run as `running` for up +to 30 seconds after its controller process exits (its heartbeat window), so +resume such a run again after that. `resume --reset-node` retries one failed +workflow node and its dependents in the same linked run; the applied reset is +recorded so retrying the command after a failed continuation resumes the +already-reset run instead of repeating the reset. `fork` may start from a +checkpoint frame and may reset one workflow node before starting the fork. Every command in this section takes an Ultrafuzz run ID and resolves the linked workflow run from existing product evidence; none of them require the diff --git a/docs/reference/configuration.md b/docs/reference/configuration.md index 1095c7155..1f5f87da9 100644 --- a/docs/reference/configuration.md +++ b/docs/reference/configuration.md @@ -139,7 +139,8 @@ resume. `workflow_deadline_seconds` is not a guaranteed wall-clock limit. Ultrafuzz records `workflow_deadline_at` in run state when the run is created, and again -from each `resume`, `replay`, or `fork`, but nothing enforces it on a timer: an +from each `replay`, `fork`, or `resume` that starts a controller (not one that +finds the run still active), but nothing enforces it on a timer: an unattended run keeps executing, and incurring provider cost, past its deadline. The deadline is checked only when a command synchronizes the run: `ultrafuzz status` (including each `--watch` poll), `inspect`, `why`, and `stats`, plus diff --git a/docs/schemas.md b/docs/schemas.md index c8e8caedc..6aa919644 100644 --- a/docs/schemas.md +++ b/docs/schemas.md @@ -84,7 +84,13 @@ a working-tree rebuild cannot change an active run. Ambient Node loader/search variables are removed and both ESM and CommonJS module resolution must stay inside that snapshot; document reads are unaffected. A path lookup alone is not a preflight. Ordinary resume now delegates continuation to Smithers instead of -using the historical launcher or closure as an authorization gate. A current +using the historical launcher or closure as an authorization gate. When resume +cannot re-verify the launcher, it reports a `WORKFLOW_TRUSTED_CLI_UNVERIFIED` +warning and, if `/trusted-bin/ultrafuzz` exists, keeps it first on `PATH` +rather than letting tasks reach another `ultrafuzz`. That launcher still +verifies its closure before every dispatch, so if its metadata or closure is +damaged, or the Node binary it names is gone, each task's validator preflight +fails; `resume --refresh-controller` does not repair such a launcher. A current controller refresh publishes a new controller path without rewriting the historical closure. diff --git a/packages/cli/src/commands/resume.ts b/packages/cli/src/commands/resume.ts index 2f3bb79da..44e701df8 100644 --- a/packages/cli/src/commands/resume.ts +++ b/packages/cli/src/commands/resume.ts @@ -5,6 +5,7 @@ import { cliEntrypoint, cliIo, commandFromRuntime, + diagnosticsText, emitCommandResult, globalFlags, projectRoot @@ -37,11 +38,14 @@ export default class Resume extends Command { resetNode: flags["reset-node"], env: cliIo().env }); - emitCommandResult( - this, - "resume", - commandFromRuntime("resume", result, (value) => `Submitted ${value.action}: ${value.workflow_run_id}\n`), - flags.json === true + const commandResult = commandFromRuntime("resume", result, (value) => + value.submitted + ? `Submitted ${value.action}: ${value.workflow_run_id}\n` + : `Run already active: ${value.workflow_run_id}; no new controller was started. If a pause is still draining, resume again once status reports paused; if its controller process just exited, resume again after 30 seconds.\n` ); + if (result.ok && result.diagnostics.length > 0) { + commandResult.text = `${commandResult.text ?? ""}${diagnosticsText(result.diagnostics)}`; + } + emitCommandResult(this, "resume", commandResult, flags.json); } } diff --git a/packages/cli/test/cli.test.ts b/packages/cli/test/cli.test.ts index 18676ca3a..750066e23 100644 --- a/packages/cli/test/cli.test.ts +++ b/packages/cli/test/cli.test.ts @@ -1918,6 +1918,26 @@ test("status surfaces a terminal product and live workflow lifecycle divergence" ); }); +test("resume of an already-active run says no controller was started instead of claiming a submission", async () => { + const project = tempProject(); + const env = fakeSmithersEnv(project); + assert.equal((await cli(project, ["init", "--json"], env)).code, 0); + writeSmallTopology(project); + const runId = "resume-already-active"; + const run = await cli(project, ["run", "--run-id", runId, "--json"], env); + assert.equal(run.code, 0, run.stderr); + + // The fake runner still reports the run as running, so resume only attaches. + const resumed = await cli(project, ["resume", runId], env); + + assert.equal(resumed.code, 0, `${resumed.stderr}\n${resumed.stdout}`); + assert.match( + resumed.stdout, + /^Run already active: ultrafuzz-resume-already-active; no new controller was started\./mu + ); + assert.doesNotMatch(resumed.stdout, /Submitted/u); +}); + test("status --watch stops immediately on a degraded verdict even while product state is nonterminal", async () => { const project = tempProject(); const env = fakeSmithersEnv(project); diff --git a/packages/runtime/src/start-run.ts b/packages/runtime/src/start-run.ts index ba32e842b..62b3a79e6 100644 --- a/packages/runtime/src/start-run.ts +++ b/packages/runtime/src/start-run.ts @@ -54,6 +54,7 @@ import { prepareTrustedCliEnvironment, runTrustedJsonValidatorPreflight, TRUSTED_CLI_ENVIRONMENT_VARIABLES, + ULTRAFUZZ_TRUSTED_BIN_ENV, type TrustedCliEnvironment } from "./trusted-cli.js"; import { hasRuntimeErrors, runtimeFailure, runtimeResult } from "./utils.js"; @@ -478,6 +479,7 @@ export async function resumeRun(input: WorkflowLifecycleInput) { async function submitSmithersContinuation(input: WorkflowLifecycleInput) { let releaseLifecycleLock: (() => Promise) | undefined; + const diagnostics: RuntimeDiagnostic[] = []; try { const projectRoot = path.resolve(input.projectRoot); const runsRoot = await runsRootForProject(projectRoot); @@ -526,6 +528,10 @@ async function submitSmithersContinuation(input: WorkflowLifecycleInput) { const smithersRoot = safeResolveInside(layout.root, "smithers", "Smithers evidence"); const tasksPath = safeResolveInside(smithersRoot, "tasks.json", "workflow task manifest"); const configPath = safeResolveInside(smithersRoot, "resolved-config.json", "workflow config"); + // Agent adapters parse ULTRAFUZZ_CONFIG_PATH as TOML; given the JSON above + // they find no agent tables and fall back to default auth. Launch writes + // the same config as TOML beside it and hands adapters a copy of that file. + const agentConfigPath = safeResolveInside(smithersRoot, "execution-config.toml", "workflow agent config"); let taskDocument: SmithersTaskManifestDocument | undefined; let config: ResolvedConfig | undefined; if (fs.existsSync(tasksPath)) { @@ -571,7 +577,7 @@ async function submitSmithersContinuation(input: WorkflowLifecycleInput) { ...forgeGuard.env, ULTRAFUZZ_ARTIFACTS_MODULE: import.meta.resolve("@ultrafuzz/artifacts"), ULTRAFUZZ_RUNTIME_MODULE: import.meta.resolve("@ultrafuzz/runtime"), - ...(config === undefined ? {} : { ULTRAFUZZ_CONFIG_PATH: configPath }), + ...(config === undefined ? {} : { ULTRAFUZZ_CONFIG_PATH: agentConfigPath }), ULTRAFUZZ_WORKFLOW_PERSISTED_PATH: workflowPath }; let trustedCli: TrustedCliEnvironment = { @@ -593,9 +599,27 @@ async function submitSmithersContinuation(input: WorkflowLifecycleInput) { }); if (prepared.active) runTrustedJsonValidatorPreflight({ layout, trusted: prepared }); trustedCli = prepared; - } catch { + } catch (error) { // Historical validator identity is task setup provenance, not authority - // to prevent Smithers from continuing the workflow. + // to prevent Smithers from continuing the workflow. Keep the run-owned + // launcher first on PATH anyway: it re-verifies its closure on every + // call, while dropping it lets tasks run whatever `ultrafuzz` is on PATH. + const launcher = path.join( + layout.root, + "trusted-bin", + process.platform === "win32" ? "ultrafuzz.cmd" : "ultrafuzz" + ); + const launcherKept = fs.existsSync(launcher); + if (launcherKept) trustedCli.env[ULTRAFUZZ_TRUSTED_BIN_ENV] = path.dirname(launcher); + diagnostics.push( + resumeWarning( + "WORKFLOW_TRUSTED_CLI_UNVERIFIED", + launcherKept + ? `resume could not re-verify the run's trusted Ultrafuzz CLI (tasks still call ${launcher}, and their preflight-json-validator step fails while that launcher cannot verify itself)` + : `resume could not re-verify the run's trusted Ultrafuzz CLI (${launcher} does not exist, so tasks call whatever \`ultrafuzz\` is on PATH)`, + error + ) + ); } } const agentRefs = tasks.flatMap((task) => task.agentChain.map((profile) => profile.agentRef)); @@ -611,7 +635,15 @@ async function submitSmithersContinuation(input: WorkflowLifecycleInput) { assertCurrentCloudAgentCredentialEnvironment(config, tasks, lifecycleEnvironment); } if (typeof metadata.source_revision === "string") { - repairPrunableRunWorktreeRegistrations({ projectRoot, runRoot: layout.root, runId }); + try { + repairPrunableRunWorktreeRegistrations({ projectRoot, runRoot: layout.root, runId }); + } catch (error) { + // Pruning is cleanup: a stale registration it leaves behind surfaces + // when Smithers recreates that task's worktree, so do not stop here. + diagnostics.push( + resumeWarning("WORKFLOW_WORKTREE_REPAIR_FAILED", "resume could not prune stale task worktrees", error) + ); + } } const result = await runSmithersLifecycleCommand({ action: "resume", @@ -668,20 +700,27 @@ async function submitSmithersContinuation(input: WorkflowLifecycleInput) { trustedCli.environmentVariableNames ) }); - recordNativeContinuationState({ - layout, - config, - requestedConcurrency: input.maxConcurrency, - alreadyRunning: result.alreadyRunning ?? false - }); - return runtimeResult(true, { - run_id: runId, - workflow_run_id: smithersRunId, - action: "resume" as const, - submitted: !result.alreadyRunning - }); + // An attach to a run Smithers still reports active started no controller, + // so it must not re-record status, lease or deadline; the resume that + // starts the next controller does. + if (result.alreadyRunning !== true) { + recordNativeContinuationState({ layout, config, requestedConcurrency: input.maxConcurrency }); + } + return runtimeResult( + true, + { + run_id: runId, + workflow_run_id: smithersRunId, + action: "resume" as const, + submitted: !result.alreadyRunning + }, + diagnostics + ); } catch (error) { - return runtimeFailure([smithersDiagnostic(error, "WORKFLOW_LIFECYCLE_FAILED")]); + return runtimeFailure([ + smithersDiagnostic(error, "WORKFLOW_LIFECYCLE_FAILED"), + ...diagnostics + ]); } finally { await releaseLifecycleLock?.(); } @@ -739,7 +778,6 @@ function recordNativeContinuationState(input: { layout: RunLayout; config: ResolvedConfig | undefined; requestedConcurrency: number | undefined; - alreadyRunning: boolean; }): void { try { const submittedAt = new Date().toISOString(); @@ -749,19 +787,17 @@ function recordNativeContinuationState(input: { state.status = "running"; state.started_at ??= submittedAt; delete state.finished_at; - if (!input.alreadyRunning) { - const leaseDurationMs = - (input.config?.run.controllerLeaseSeconds ?? Math.max(1, state.controller_lease.duration_ms / 1_000)) * 1_000; - state.controller_lease = { - ...state.controller_lease, - status: "active", - duration_ms: leaseDurationMs, - renewed_at: submittedAt, - expires_at: new Date(submittedAtMs + leaseDurationMs).toISOString() - }; - state.concurrency.requested_concurrency = - input.requestedConcurrency ?? input.config?.run.maxParallelAgents ?? state.concurrency.requested_concurrency; - } + const leaseDurationMs = + (input.config?.run.controllerLeaseSeconds ?? Math.max(1, state.controller_lease.duration_ms / 1_000)) * 1_000; + state.controller_lease = { + ...state.controller_lease, + status: "active", + duration_ms: leaseDurationMs, + renewed_at: submittedAt, + expires_at: new Date(submittedAtMs + leaseDurationMs).toISOString() + }; + state.concurrency.requested_concurrency = + input.requestedConcurrency ?? input.config?.run.maxParallelAgents ?? state.concurrency.requested_concurrency; if (input.config !== undefined) { state.workflow_deadline_at = new Date( submittedAtMs + input.config.run.workflowDeadlineSeconds * 1_000 @@ -779,6 +815,11 @@ function recordNativeContinuationState(input: { } } +function resumeWarning(code: string, context: string, error: unknown): RuntimeDiagnostic { + const diagnostic = smithersDiagnostic(error, code); + return { ...diagnostic, message: `${context}: ${diagnostic.message}`, severity: "warning", source: "runtime" }; +} + export async function replayRun(input: WorkflowLifecycleInput) { return submitLifecycleAction(input, "replay"); } diff --git a/packages/runtime/test/runtime.test.ts b/packages/runtime/test/runtime.test.ts index 357f3466d..65db202c1 100644 --- a/packages/runtime/test/runtime.test.ts +++ b/packages/runtime/test/runtime.test.ts @@ -26,8 +26,10 @@ import { artifactContractDefinition, artifactContractSchemaBinding, artifactSchemaBundleDigest, + artifactSchemaDirectory, artifactSchemaRegistry, artifactSchemaRegistryFromDirectory, + artifactValidatorSmokeFixturePath, ARTIFACT_VALIDATOR_SMOKE_FIXTURE_SHA256, GOAL_PLAN_JSON_SCHEMA_ID, THREAT_MODEL_JSON_SCHEMA_ID, @@ -36,6 +38,7 @@ import { goalPlanJsonSchema, layoutForRunRoot, manifestDigest, + parseJsonValidatorPreflightSuccessEnvelope, readPlannedGraphDocument, readRunState, promptArtifactAuthorityPathSelectorId, @@ -2503,6 +2506,16 @@ function controllerRefreshTerminalEnv( }); } +/** Make a fake lifecycle runner append `$` to `logPath` on every `up`. */ +function logFakeRunnerUpVariable(env: Record, name: string, logPath: string): void { + const shim = env.SMITHERS_BIN; + assert.ok(shim); + fs.writeFileSync( + shim, + fs.readFileSync(shim, "utf8").replace(" up)\n", ` up)\n printf '%s\\n' "$${name}" >> ${shellQuote(logPath)}\n`) + ); +} + function workflowEvents( workflowRunId: string, events: Array<{ @@ -15383,7 +15396,7 @@ test("native resume delegates the persisted workflow after mutable project sourc assert.equal(resumed.ok, true, JSON.stringify(resumed.diagnostics)); const consumed = fs.readFileSync(snapshotBytesLog, "utf8"); assert.equal(consumed.includes(`workflow=${mutableWorkflow}\n`), true); - assert.match(consumed, /^config=.*\/smithers\/resolved-config\.json$/mu); + assert.match(consumed, /^config=.*\/smithers\/execution-config\.toml$/mu); assert.match(consumed, /^agent=.*\/\.smithers\/agents\/codex\.ts$/mu); assert.match(consumed, /HostileReplacement/u); assert.match(consumed, /export const hostile/u); @@ -25050,6 +25063,16 @@ test("native continuation does not use historical trusted CLI identity as an aut assert.equal(refreshed.ok, true, JSON.stringify(refreshed.diagnostics)); assert.equal(refreshed.value?.submitted, true); assert.equal(fs.readFileSync(trustedMetadataPath, "utf8"), "{}\n"); + // The failed re-verification is reported instead of swallowed (#1143). + for (const resumed of [ordinary, refreshed]) { + assert.equal( + resumed.diagnostics.some( + (diagnostic) => diagnostic.code === "WORKFLOW_TRUSTED_CLI_UNVERIFIED" && diagnostic.severity === "warning" + ), + true, + JSON.stringify(resumed.diagnostics) + ); + } assert.equal( fs .readFileSync(env.SMITHERS_FAKE_LOG!, "utf8") @@ -25060,6 +25083,88 @@ test("native continuation does not use historical trusted CLI identity as an aut ); }); +test("a resume that cannot re-verify the trusted CLI leaves tasks on the run's own working launcher", async () => { + const project = tempProject(); + initProject({ projectRoot: project, force: true }); + writeSmallTopology(project); + const runId = "resume-keeps-trusted-launcher"; + const env = controllerRefreshTerminalEnv(project, runId); + const upPathLog = path.join(project, "fake-smithers-up-path.log"); + logFakeRunnerUpVariable(env, "PATH", upPathLog); + const launched = await startRun({ projectRoot: project, runId, env }); + assert.equal(launched.ok, true, JSON.stringify(launched.diagnostics)); + assert.ok(launched.value); + fs.writeFileSync(upPathLog, "", "utf8"); + + // Without a CLI entrypoint resume cannot re-verify the launcher, although + // the launcher and its closure are intact (#1143). + const resumed = await runtimeResumeRun({ projectRoot: project, runId, env }); + assert.equal(resumed.ok, true, JSON.stringify(resumed.diagnostics)); + assert.equal(resumed.value?.submitted, true); + + // Each task's validator preflight runs `ultrafuzz` from the runner's PATH. + // It must still reach the run's launcher, which verifies its closure and + // dispatches to the recorded (fake) CLI, never another `ultrafuzz`. + const runnerPath = fs.readFileSync(upPathLog, "utf8").trim(); + const resolved = runnerPath + .split(path.delimiter) + .map((entry) => path.join(entry, "ultrafuzz")) + .find((candidate) => fs.existsSync(candidate)); + assert.equal(resolved, path.join(launched.value.run_root, "trusted-bin", "ultrafuzz")); + const findings = artifactSchemaRegistry().find((entry) => entry.filename === "findings.schema.json"); + assert.ok(findings); + const stdout = execFileSync( + resolved, + [ + "json", + "validate", + "--schema", + path.join(artifactSchemaDirectory(), findings.filename), + "--file", + artifactValidatorSmokeFixturePath(), + "--json" + ], + { encoding: "utf8", env: { ...process.env, PATH: runnerPath } } + ); + parseJsonValidatorPreflightSuccessEnvelope(Buffer.from(stdout, "utf8")); + const warning = resumed.diagnostics.find((diagnostic) => diagnostic.code === "WORKFLOW_TRUSTED_CLI_UNVERIFIED"); + assert.equal(warning?.severity, "warning", JSON.stringify(resumed.diagnostics)); +}); + +test("native continuation hands generated agents the run's TOML config, so CodexAgent keeps API-key auth", async () => { + const project = tempProject(); + initProject({ projectRoot: project, force: true }); + writeSmallTopology(project); + const runId = "continuation-agent-config"; + const env = controllerRefreshTerminalEnv(project, runId); + const configPathLog = path.join(project, "fake-smithers-up-config-path.log"); + logFakeRunnerUpVariable(env, "ULTRAFUZZ_CONFIG_PATH", configPathLog); + const launched = await startRun({ projectRoot: project, runId, env }); + assert.equal(launched.ok, true, JSON.stringify(launched.diagnostics)); + fs.writeFileSync(configPathLog, "", "utf8"); + const resumed = await resumeRun({ projectRoot: project, runId, env }); + assert.equal(resumed.ok, true, JSON.stringify(resumed.diagnostics)); + assert.equal(resumed.value?.submitted, true); + + // Build the stock adapter from the config path the resumed runner received. + // The init config selects `auth = "api-key"` for CodexAgent; an adapter that + // cannot read it falls back to subscription auth and clears the key. + const { createCodexAgent } = await loadGeneratedCodexAgent(project); + const previous = { config: process.env.ULTRAFUZZ_CONFIG_PATH, key: process.env.OPENAI_API_KEY }; + process.env.ULTRAFUZZ_CONFIG_PATH = fs.readFileSync(configPathLog, "utf8").trim(); + process.env.OPENAI_API_KEY = "continuation-codex-key"; + try { + const agent = createCodexAgent() as { opts: { env: Record } }; + assert.equal(agent.opts.env.CODEX_API_KEY, "continuation-codex-key"); + assert.equal(agent.opts.env.OPENAI_API_KEY, "continuation-codex-key"); + } finally { + if (previous.config === undefined) delete process.env.ULTRAFUZZ_CONFIG_PATH; + else process.env.ULTRAFUZZ_CONFIG_PATH = previous.config; + if (previous.key === undefined) delete process.env.OPENAI_API_KEY; + else process.env.OPENAI_API_KEY = previous.key; + } +}); + test("controller refresh authenticates newly required sealed runner patches and rejects source drift", async () => { const project = tempProject(); initProject({ projectRoot: project, force: true }); @@ -26419,6 +26524,11 @@ test("ordinary resume checks active-run ownership before detached preflight", as }); const run = await startRun({ projectRoot: project, runId: "active-lifecycle-run", env }); assert.equal(run.ok, true, JSON.stringify(run.diagnostics)); + // An attach starts no controller, so the run's state (status, lease and + // workflow deadline) must stay exactly as its live owner left it (#1153). + assert.ok(run.value); + const statePath = path.join(run.value.run_root, "state.json"); + const stateBefore = fs.readFileSync(statePath, "utf8"); fs.writeFileSync(env.SMITHERS_FAKE_LOG!, "", "utf8"); // A duplicate `up --resume --detach` renders the workflow before Smithers // checks ownership. Keep that path fatal so this regression proves active @@ -26448,6 +26558,53 @@ test("ordinary resume checks active-run ownership before detached preflight", as const forcedCommands = fs.readFileSync(env.SMITHERS_FAKE_LOG!, "utf8"); assert.match(forcedCommands, /inspect ultrafuzz-active-lifecycle-run --format json --full-output/u); assert.doesNotMatch(forcedCommands, /^up /mu); + assert.equal(fs.readFileSync(statePath, "utf8"), stateBefore); +}); + +test("resume continues with a warning when stale task-worktree cleanup fails", async () => { + const project = tempProject(); + initProject({ projectRoot: project, force: true }); + writeSmallTopology(project); + const runId = "worktree-cleanup-failure"; + const env = controllerRefreshTerminalEnv(project, runId); + const commandLog = env.SMITHERS_FAKE_LOG; + assert.ok(commandLog); + const launched = await startRun({ projectRoot: project, runId, env }); + assert.equal(launched.ok, true, JSON.stringify(launched.diagnostics)); + assert.ok(launched.value); + const runRoot = launched.value.run_root; + // Cleanup runs only for runs launched from a Git revision. + const metadataPath = path.join(runRoot, "run.json"); + const metadata = JSON.parse(fs.readFileSync(metadataPath, "utf8")) as Record; + fs.writeFileSync(metadataPath, `${JSON.stringify({ ...metadata, source_revision: "0".repeat(40) }, null, 2)}\n`); + // Git lists a prunable registration owned by this run, then cannot remove it. + const previousPath = process.env.PATH ?? ""; + const gitBin = temporaryRoot("ufz-failing-git-"); + fs.writeFileSync( + path.join(gitBin, "git"), + [ + "#!/bin/sh", + 'case "$*" in', + ` *"worktree list --porcelain"*) printf 'worktree %s\\nbranch refs/heads/ultrafuzz/%s/stale\\nprunable\\n\\n' ${shellQuote(path.join(runRoot, "workspaces", "stale"))} ${runId} ;;`, + ' *"worktree remove"*) echo "fatal: synthetic removal failure" >&2; exit 1 ;;', + ` *) PATH=${shellQuote(previousPath)} exec git "$@" ;;`, + "esac", + "" + ].join("\n"), + { mode: 0o755 } + ); + process.env.PATH = [gitBin, previousPath].join(path.delimiter); + const resumed = await resumeRun({ projectRoot: project, runId, env }).finally(() => { + process.env.PATH = previousPath; + }); + + assert.equal(resumed.ok, true, JSON.stringify(resumed.diagnostics)); + assert.equal(resumed.value?.submitted, true); + const warning = resumed.diagnostics.find((diagnostic) => diagnostic.code === "WORKFLOW_WORKTREE_REPAIR_FAILED"); + assert.ok(warning, JSON.stringify(resumed.diagnostics)); + assert.equal(warning.severity, "warning"); + assert.match(warning.message, /synthetic removal failure/u); + assert.match(fs.readFileSync(commandLog, "utf8"), /^up .*--resume ultrafuzz-worktree-cleanup-failure/mu); }); test("resume derives reset identities from the canonical nodes of a failed workflow", async () => { diff --git a/packages/runtime/test/smithers-preparation-race.integration.test.ts b/packages/runtime/test/smithers-preparation-race.integration.test.ts index 204af09e9..7bed7acff 100644 --- a/packages/runtime/test/smithers-preparation-race.integration.test.ts +++ b/packages/runtime/test/smithers-preparation-race.integration.test.ts @@ -1,6 +1,6 @@ import assert from "node:assert/strict"; import { temporaryRoot } from "./temporary-root.js"; -import { execFileSync } from "node:child_process"; +import { execFileSync, spawnSync } from "node:child_process"; import fs from "node:fs"; import path from "node:path"; import test from "node:test"; @@ -226,6 +226,105 @@ process.stdout.write(JSON.stringify(result));` } }); +// #1153 asked Ultrafuzz for its own controller-generation fence around pause +// and continuation handoff. The pinned engine already provides the guarantee: +// a graceful pause lets in-flight tasks finish instead of aborting them, a +// second controller is refused while the owner is alive (`--force` does not +// take ownership), and a resume after the park runs only the remaining work. +// Pin those facts so an engine bump that regresses them fails here. +test("graceful pause drains in-flight tasks and refuses a second controller until the run parks", async () => { + const root = temporaryRoot("ultrafuzz-smithers-pause-handoff-"); + const workflowDir = path.join(root, ".smithers", "workflows"); + const workflowPath = path.join(workflowDir, "pause-handoff.tsx"); + const traceLog = path.join(root, "trace.log"); + const releasePath = path.join(root, "release-in-flight"); + const runId = `pause-handoff-${process.pid}-${Date.now()}`; + const runner = (args: string[]) => + spawnSync(smithersBinary(), args, { + cwd: root, + encoding: "utf8", + env: { ...process.env, SMITHERS_POST_FAILURE: "0" } + }); + const trace = (): string[] => + fs.existsSync(traceLog) ? fs.readFileSync(traceLog, "utf8").trim().split("\n").filter(Boolean) : []; + const pidOf = (lines: readonly string[], event: string) => + lines.find((line) => line.endsWith(` ${event}`))?.split(" ")[0]; + + try { + fs.mkdirSync(workflowDir, { recursive: true }); + initFixtureRepository(root); + const smithersPackageRoot = fs.realpathSync(path.join(runtimePackageRoot(), "node_modules", "smthrs")); + fs.symlinkSync(path.dirname(smithersPackageRoot), path.join(root, ".smithers", "node_modules"), "dir"); + fs.writeFileSync(workflowPath, pauseHandoffWorkflowSource({ traceLog, releasePath }), "utf8"); + + const launched = runner([ + "up", + workflowPath, + "--detach", + "--run-id", + runId, + "--root", + root, + "--input", + "{}", + "--format", + "json" + ]); + assert.equal(launched.status, 0, launched.stderr); + await waitUntil(() => trace().filter((line) => line.endsWith(" start")).length === 2, 60_000, "a and b start"); + + const pause = runner(["pause", runId, "--format", "json"]); + assert.match(pause.stdout, /"pause-requested"/u, pause.stderr); + // Both in-flight tasks are held open, so the owner is still draining: a + // replacement controller and a node reset are refused despite `--force`. + for (const takeover of [ + ["up", workflowPath, "--resume", runId, "--run-id", runId, "--force", "--detach", "--format", "json"], + ["timetravel", workflowPath, "--run-id", runId, "--node-id", "a", "--no-vcs", "--force", "--format", "json"] + ]) { + const refused = runner(takeover); + assert.notEqual(refused.status, 0, takeover.join(" ")); + assert.match(refused.stdout + refused.stderr, /RUN_OWNER_ALIVE/u, takeover.join(" ")); + } + // The draining run still reports an active state, so `ultrafuzz resume` + // only attaches to it instead of starting a controller. + const draining = JSON.parse(runner(["inspect", runId, "--format", "json", "--full-output"]).stdout) as { + data?: { runState?: { state?: string } }; + }; + assert.equal(draining.data?.runState?.state, "running"); + assert.equal(trace().length, 2, "no task ended or started while the pause drained"); + + fs.writeFileSync(releasePath, "", "utf8"); + await waitForStatus(root, runId, "paused", 60_000); + const parked = trace(); + const ownerPid = pidOf(parked, "a start"); + assert.equal(pidOf(parked, "a end"), ownerPid, "the draining owner finished a"); + assert.equal(pidOf(parked, "b end"), ownerPid, "the draining owner finished b"); + assert.equal(pidOf(parked, "c start"), undefined, "the pause stopped new scheduling"); + + const resumed = runner(["up", workflowPath, "--resume", runId, "--run-id", runId, "--detach", "--format", "json"]); + assert.equal(resumed.status, 0, resumed.stderr); + await waitForSuccessfulCompletion(root, runId, 60_000); + const finished = trace(); + for (const event of ["a start", "b start", "c start"]) { + const runs = finished.filter((line) => line.endsWith(` ${event}`)).length; + assert.equal(runs, 1, `${event} ran ${runs} times: ${JSON.stringify(finished)}`); + } + assert.notEqual(pidOf(finished, "c start"), ownerPid, "only the replacement controller ran c"); + } finally { + // Release any task still held, then give every engine that ran a task + // time to exit. Deleting the root first would delete the release marker + // too, leaving a detached engine polling until its 120 s hold expires and + // writing its logs back under the deleted root. + fs.writeFileSync(releasePath, "", "utf8"); + await waitUntil( + () => !trace().some((line) => processIsAlive(Number(line.split(" ")[0]))), + 30_000, + "the detached engines exit" + ).catch(() => undefined); + fs.rmSync(root, { recursive: true, force: true }); + } +}); + function initFixtureRepository(root: string): void { execGit(root, ["init", "--quiet", "--initial-branch=main"]); execGit(root, ["config", "user.name", "Ultrafuzz Synthetic Test"]); @@ -368,6 +467,45 @@ export default smithers((ctx) => ( `; } +function pauseHandoffWorkflowSource(input: { traceLog: string; releasePath: string }): string { + return `/** @jsxImportSource smthrs */ +import fs from "node:fs"; +import { createSmithers } from "smthrs"; +import { z } from "zod/v4"; + +const traceLog = ${JSON.stringify(input.traceLog)}; +const releasePath = ${JSON.stringify(input.releasePath)}; +const trace = (event) => fs.appendFileSync(traceLog, process.pid + " " + event + "\\n", "utf8"); +const { Workflow, Task, Parallel, Sequence, smithers, outputs } = createSmithers({ + input: z.object({}), + step: z.object({ done: z.literal(true) }) +}); +// Hold each in-flight task until the test releases it, bounded so a failed +// test cannot leave the detached engine polling forever. +const held = (id) => async () => { + trace(id + " start"); + const deadline = Date.now() + 120000; + while (!fs.existsSync(releasePath) && Date.now() < deadline) { + await new Promise((resolve) => setTimeout(resolve, 50)); + } + trace(id + " end"); + return { done: true }; +}; + +export default smithers(() => ( + + + + {held("a")} + {held("b")} + + {() => (trace("c start"), { done: true })} + + +)); +`; +} + async function waitForFile(filePath: string, timeoutMs: number): Promise { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { @@ -377,7 +515,29 @@ async function waitForFile(filePath: string, timeoutMs: number): Promise { throw new Error(`timed out waiting for ${filePath}`); } +function processIsAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch (error) { + return (error as NodeJS.ErrnoException).code === "EPERM"; + } +} + +async function waitUntil(condition: () => boolean, timeoutMs: number, label: string): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (condition()) return; + await new Promise((resolve) => setTimeout(resolve, 100)); + } + throw new Error(`timed out waiting until ${label}`); +} + async function waitForSuccessfulCompletion(root: string, runId: string, timeoutMs: number): Promise { + await waitForStatus(root, runId, "finished", timeoutMs); +} + +async function waitForStatus(root: string, runId: string, expected: string, timeoutMs: number): Promise { const deadline = Date.now() + timeoutMs; let status = "unknown"; while (Date.now() < deadline) { @@ -393,13 +553,13 @@ async function waitForSuccessfulCompletion(root: string, runId: string, timeoutM await new Promise((resolve) => setTimeout(resolve, 100)); continue; } - if (status === "finished") return; + if (status === expected) return; if (["failed", "cancelled", "canceled"].includes(status)) { throw new Error(`synthetic Smithers workflow ended with status ${status}`); } await new Promise((resolve) => setTimeout(resolve, 100)); } - throw new Error(`synthetic Smithers workflow did not finish; final status ${status}`); + throw new Error(`synthetic Smithers workflow did not reach ${expected}; final status ${status}`); } function smithersNodeOutput(root: string, runId: string, nodeId: string): Buffer {