diff --git a/docs/reference/artifacts-reports.md b/docs/reference/artifacts-reports.md index 656258f40..486521a4c 100644 --- a/docs/reference/artifacts-reports.md +++ b/docs/reference/artifacts-reports.md @@ -227,7 +227,16 @@ status changes afterwards. Finished, failed, timed-out, and cancelled attempts are recorded; a cancellation has outcome and category `canceled` and the Smithers cancellation reason as its message. A terminal event with no started attempt in the same Smithers activation, or stamped before that attempt -started, is skipped. A finished attempt whose node then fails, for example +started, is skipped. An attempt that was still running when its controller +stopped, for example because the controller was killed, has no terminal event: +when the run resumes, Smithers marks it cancelled without emitting one. Like a +failure, it is recorded unless it stopped before Smithers selected a rung: as +`canceled`, with the message `abandoned: the controller stopped during this +attempt, and the resumed run cancelled it`, and it counts toward the node's +`retry_count`. Its terminal is the next event of the same task, normally the +replacement attempt's start, because one resume can abandon several attempts; +it is not recorded before the task has such an event. A finished attempt whose +node then fails, for example because the verifier or artifact gates reject its output, is recorded as failed with category `invalid-output` for findings validation and `artifact-validation` otherwise. diff --git a/docs/reference/cli.md b/docs/reference/cli.md index 59097b75a..83f79ad0b 100644 --- a/docs/reference/cli.md +++ b/docs/reference/cli.md @@ -406,7 +406,11 @@ completed and currently elapsed execution time, token components, estimated spend, model, and attempt count. JSON output uses `ultrafuzz.stats.v1` inside the normal CLI envelope and additionally exposes retries, executed/reused counts, outcomes, failure categories, completeness, unattributed usage, and -cumulative run accounting. +cumulative run accounting. Run state records a task that Smithers cancelled, +for example by `ultrafuzz cancel`, as failed. When Smithers cancelled every +failed task of a node, `stats` reports the node with status `canceled` and +counts it under `canceled` rather than `failed`; `status` likewise counts +cancelled tasks apart from failures. The closed stats v1 field `pricing_complete` continues to mean complete cost coverage: it is the inverse of run accounting `partial_pricing`. Accounting v4's own `pricing_complete` field independently reports whether each event had diff --git a/docs/reference/development.md b/docs/reference/development.md index 4a7e68158..ce168448e 100644 --- a/docs/reference/development.md +++ b/docs/reference/development.md @@ -74,11 +74,13 @@ pnpm --filter @ultrafuzz/modal test `report`, and `events` as separate CLI processes, and the generated workflow runs on the pinned Smithers engine under Bun, with a stub `codex` executable in place of the model. It SIGKILLs the detached controller while one node is -running, resumes the run, and checks that it succeeds with a verified report and -that no finished task started again. It needs Linux, Bun, Git, and access to -the npm registry, because `run` and `resume` install the pinned engine from npm -as they do for any campaign. On SIGINT or SIGTERM, the test kills the detached -campaign and deletes its fixture, about 1 GB, before it exits. +running, resumes the run, and checks that it succeeds with a verified report, +that no finished task started again, and that `status` and `stats` count the +same agent attempts, including the one the kill interrupted. It needs Linux, +Bun, Git, and access to the npm registry, because `run` and `resume` install the +pinned engine from npm as they do for any campaign. On SIGINT or SIGTERM, the +test kills the detached campaign and deletes its fixture, about 1 GB, before it +exits. Package-local `typecheck` and `test` scripts may build direct workspace dependencies first because package exports point at `dist/**`. diff --git a/packages/cli/schema/cli-result.schema.json b/packages/cli/schema/cli-result.schema.json index 440909e09..538897b16 100644 --- a/packages/cli/schema/cli-result.schema.json +++ b/packages/cli/schema/cli-result.schema.json @@ -1561,6 +1561,7 @@ "running", "succeeded", "failed", + "canceled", "skipped", "timed-out", "reused-from-prior-run", @@ -1612,6 +1613,7 @@ "running", "succeeded", "failed", + "canceled", "skipped", "timed-out", "reused-from-prior-run", @@ -1625,6 +1627,7 @@ "running": { "$ref": "#/$defs/nonNegativeInteger" }, "succeeded": { "$ref": "#/$defs/nonNegativeInteger" }, "failed": { "$ref": "#/$defs/nonNegativeInteger" }, + "canceled": { "$ref": "#/$defs/nonNegativeInteger" }, "skipped": { "$ref": "#/$defs/nonNegativeInteger" }, "timed-out": { "$ref": "#/$defs/nonNegativeInteger" }, "reused-from-prior-run": { "$ref": "#/$defs/nonNegativeInteger" }, diff --git a/packages/cli/src/run-statistics.ts b/packages/cli/src/run-statistics.ts index cecfe02ee..974dd8de0 100644 --- a/packages/cli/src/run-statistics.ts +++ b/packages/cli/src/run-statistics.ts @@ -58,7 +58,7 @@ export interface TokenStatistics { models: string[]; } -export type NodeStatisticsStatus = NodeStatus | "unknown"; +export type NodeStatisticsStatus = NodeStatus | "canceled" | "unknown"; export interface NodeStatistics { node_id: string; @@ -484,6 +484,17 @@ function taskWorkflowAgentId(nodeState: NodeState | undefined): string | undefin return "agent_task_id" in provenance.workflow ? provenance.workflow.agent_task_id : undefined; } +/** Whether Smithers ended every failed task of the node cancelled. */ +function workflowCancelled(states: readonly NodeState[]): boolean { + const failedTaskStates = states.flatMap((nodeState) => { + const provenance = nodeState.provenance; + if (nodeState.status !== "failed" || provenance === undefined || !("workflow" in provenance)) return []; + const workflow = provenance.workflow; + return workflow !== undefined && "state" in workflow ? [workflow.state] : []; + }); + return failedTaskStates.length > 0 && failedTaskStates.every((state) => state === "cancelled"); +} + function usageNodeAliases( descriptors: readonly NodeDescriptor[], diagnostics: RuntimeDiagnostic[] @@ -548,7 +559,10 @@ function nodeStatistics( const canonicalState = descriptor.states.find((nodeState) => nodeState.node_id === descriptor.nodeId) ?? descriptor.states[0]; - const status = canonicalState?.status ?? aggregateNodeStatus(descriptor.states); + const stateStatus = canonicalState?.status ?? aggregateNodeStatus(descriptor.states); + // Run state records a task that Smithers cancelled as failed; `status` does + // not count it as a failure (#1087). + const status = stateStatus === "failed" && workflowCancelled(descriptor.states) ? "canceled" : stateStatus; const strategyStates = descriptor.states.filter((nodeState) => nodeState.node_id !== descriptor.nodeId); const timedStates = strategyStates.length > 0 ? strategyStates : descriptor.states; const currentElapsedMs = timedStates.reduce((total, nodeState) => { @@ -803,9 +817,20 @@ function aggregateNodeStatus(states: readonly NodeState[]): NodeStatisticsStatus } function aggregateAttemptOutcome(attempts: readonly NodeAttemptLedgerEntry[]): NodeAttemptOutcome | "mixed" | null { - const latestByStrategy = new Map(); - for (const attempt of attempts) latestByStrategy.set(attempt.strategy_attempt_id, attempt.outcome); - const outcomes = [...new Set(latestByStrategy.values())]; + const latestByStrategy = new Map(); + for (const attempt of attempts) { + const previous = latestByStrategy.get(attempt.strategy_attempt_id); + // Rows are appended as they are recorded, so a version that records more + // attempts can append an earlier attempt of a workflow run after a later one. + if ( + previous?.workflow_run_id === attempt.workflow_run_id && + previous.source_event_sequence > attempt.source_event_sequence + ) { + continue; + } + latestByStrategy.set(attempt.strategy_attempt_id, attempt); + } + const outcomes = [...new Set([...latestByStrategy.values()].map((attempt) => attempt.outcome))]; if (outcomes.length === 0) return null; return outcomes.length === 1 ? outcomes[0]! : "mixed"; } @@ -818,6 +843,7 @@ function emptyStatusCounts(): Record { running: 0, succeeded: 0, failed: 0, + canceled: 0, skipped: 0, "timed-out": 0, "reused-from-prior-run": 0, diff --git a/packages/cli/test/cli-contracts.test.ts b/packages/cli/test/cli-contracts.test.ts index 7a7bcbf36..03c268dc3 100644 --- a/packages/cli/test/cli-contracts.test.ts +++ b/packages/cli/test/cli-contracts.test.ts @@ -234,6 +234,7 @@ test("stats has an exact command discriminator and a fully closed result shape", running: 0, succeeded: 0, failed: 0, + canceled: 0, skipped: 0, "timed-out": 0, "reused-from-prior-run": 0, diff --git a/packages/cli/test/e2e/campaign-resume.test.ts b/packages/cli/test/e2e/campaign-resume.test.ts index 06d9b6a43..3b5765028 100644 --- a/packages/cli/test/e2e/campaign-resume.test.ts +++ b/packages/cli/test/e2e/campaign-resume.test.ts @@ -204,6 +204,7 @@ interface StatsValue { status: string; attempt_count: number | null; executed_attempt_count: number | null; + failure_categories: string[] | null; }>; totals: { node_count: number; status_counts: Record }; } @@ -498,21 +499,16 @@ test( assert.equal(stats.totals.node_count, AGENT_NODES.length); assert.equal(stats.totals.status_counts.succeeded, AGENT_NODES.length); assert.equal(stats.nodes.find((node) => node.node_id === "project-discovery")?.executed_attempt_count, 1); - // status counts every agent invocation, including the one the kill interrupted. + // status and stats count every agent invocation, including the one the kill interrupted, which + // stats records as canceled. const statusAttempts = health.model_mix.reduce((total, entry) => total + entry.attempts, 0); assert.equal(statusAttempts, starts.length); - mark("checks done"); - - await t.test( - "stats counts the agent attempt the controller crash interrupted", - { - todo: "Smithers emits no terminal event for the attempt it abandons when the resumed run starts, so attempts.jsonl never records it (#1187)" - }, - () => { - const statsAttempts = stats.nodes.reduce((total, node) => total + (node.attempt_count ?? 0), 0); - assert.equal(statsAttempts, statusAttempts); - } + assert.equal( + stats.nodes.reduce((total, node) => total + (node.attempt_count ?? 0), 0), + statusAttempts ); + assert.deepEqual(stats.nodes.find((node) => node.node_id === INTERRUPTED_NODE)?.failure_categories, ["canceled"]); + mark("checks done"); } finally { process.removeListener("SIGINT", interrupted); process.removeListener("SIGTERM", interrupted); diff --git a/packages/cli/test/run-statistics.test.ts b/packages/cli/test/run-statistics.test.ts index 94ac46c4b..3cc4dc1e4 100644 --- a/packages/cli/test/run-statistics.test.ts +++ b/packages/cli/test/run-statistics.test.ts @@ -12,6 +12,7 @@ import { manifestDigest, type NodeAttemptLedgerEntry, type NodeState, + type NodeWorkflowProvenance, type PlannedGraphDocument, type RunAccountingSummary, type RunMetadataDocument, @@ -20,6 +21,8 @@ import { } from "@ultrafuzz/artifacts"; import { assertRunMetadataAccountingUsageAuthority } from "@ultrafuzz/runtime"; +import { envelope } from "../src/command-shared.js"; +import { buildStatisticsCommandResult } from "../src/commands/stats.js"; import { deriveRunStatistics, type StatisticsEvidence } from "../src/run-statistics.js"; const RUN_ID = "stats-unit"; @@ -40,7 +43,12 @@ function attempt( nodeId: string, strategyAttemptId: string, sourceEventSequence: number, - options: { runId?: string; startedAt?: string; finishedAt?: string; outcome?: "succeeded" | "failed" } = {} + options: { + runId?: string; + startedAt?: string; + finishedAt?: string; + outcome?: "succeeded" | "failed" | "canceled"; + } = {} ): NodeAttemptLedgerEntry { const outcome = options.outcome ?? "succeeded"; return createNodeAttemptLedgerEntry( @@ -59,7 +67,7 @@ function attempt( outcome, inputManifestDigest: manifestDigest("input"), outputManifestDigest: outcome === "succeeded" ? manifestDigest("output") : null, - ...(outcome === "failed" ? { failureCategory: "executor-error" as const } : {}) + ...(outcome === "succeeded" ? {} : { failureCategory: outcome === "failed" ? "executor-error" : "canceled" }) } ); } @@ -168,6 +176,7 @@ test("stats derives closed per-node timing, usage, cost, and status totals", () running: 0, succeeded: 1, failed: 0, + canceled: 0, skipped: 0, "timed-out": 0, "reused-from-prior-run": 0, @@ -192,6 +201,69 @@ test("stats derives closed per-node timing, usage, cost, and status totals", () assert.deepEqual(derived.diagnostics, []); }); +test("stats counts a failed node whose task Smithers cancelled as canceled", () => { + type WorkflowState = "cancelled" | "failed"; + const failed = (nodeId: string, logicalNodeId: string, workflow: NodeWorkflowProvenance): NodeState => ({ + ...terminalNodeState(nodeId, logicalNodeId), + status: "failed", + outputs: [], + provenance: { workflow } + }); + const task = (nodeId: string, logicalNodeId: string, state: WorkflowState): NodeState => + failed(nodeId, logicalNodeId, { + run_id: WORKFLOW_RUN_ID, + task_id: `node:${nodeId}`, + agent_task_id: `node:${nodeId}`, + verifier_task_id: `verify:${nodeId}`, + state, + attempt: 1 + }); + const status = (states: NodeState[], graph?: PlannedGraphDocument) => { + // Smithers cancelled the tasks before it selected an agent, so no attempt was recorded. + const { value } = deriveRunStatistics( + evidence({ + ...(graph === undefined ? {} : { graph, usage: [] }), + state: runState(states, "canceled"), + attempts: [] + }), + Date.parse(FINISHED_AT) + ); + // The JSON envelope is validated against the closed CLI result schema. + envelope("stats", buildStatisticsCommandResult(value, [])); + return [value.nodes[0]?.status, value.totals.status_counts.failed, value.totals.status_counts.canceled]; + }; + + // Run state records a cancelled task as failed; `status` counts it apart from failures (#1087). + assert.deepEqual(status([task("node", "node", "cancelled")]), ["canceled", 0, 1]); + assert.deepEqual(status([task("node", "node", "failed")]), ["failed", 1, 0]); + // A fan-out node's canonical state aggregates its strategy tasks and has no Smithers state of its own. + const fanOut = (...states: WorkflowState[]) => + status( + [ + failed("fan", "fan", { run_id: WORKFLOW_RUN_ID, aggregate_attempt_statuses: ["failed", "failed"] }), + ...states.map((state, index) => task(`fan__model_${index}__attempt_0`, "fan", state)) + ], + graphDocument( + "fan", + states.map((_, index) => `node:fan__model_${index}__attempt_0`), + ["gpt-test", "gpt-other"] + ) + ); + assert.deepEqual(fanOut("cancelled", "cancelled"), ["canceled", 0, 1]); + assert.deepEqual(fanOut("cancelled", "failed"), ["failed", 1, 0]); +}); + +test("stats reports a node's latest attempt outcome in event order, not ledger row order", () => { + const outcome = (attempts: NodeAttemptLedgerEntry[]) => + deriveRunStatistics(evidence({ attempts }), Date.parse(FINISHED_AT)).value.nodes[0]?.outcome; + const abandoned = attempt("node", "node", 1, { outcome: "canceled" }); + const replacement = attempt("node", "node", 2); + + assert.equal(outcome([abandoned, replacement]), "succeeded"); + // A run synchronized by an earlier version records the abandoned attempt after its replacement. + assert.equal(outcome([replacement, abandoned]), "succeeded"); +}); + test("stats counts only the latest cumulative usage snapshot for each attempt", () => { const snapshots = [ usage(1, "node:node", { diff --git a/packages/cli/test/stats-command.test.ts b/packages/cli/test/stats-command.test.ts index d8152122b..1174adfd0 100644 --- a/packages/cli/test/stats-command.test.ts +++ b/packages/cli/test/stats-command.test.ts @@ -22,6 +22,7 @@ const statistics: RunStatisticsValue = { running: 0, succeeded: 0, failed: 0, + canceled: 0, skipped: 0, "timed-out": 0, "reused-from-prior-run": 0, diff --git a/packages/runtime/src/workflow-sync.ts b/packages/runtime/src/workflow-sync.ts index 30384cdac..f9002656c 100644 --- a/packages/runtime/src/workflow-sync.ts +++ b/packages/runtime/src/workflow-sync.ts @@ -179,6 +179,19 @@ interface TerminalWorkflowAttempt { superseded: boolean; } +type WorkflowAttemptStart = Pick< + TerminalWorkflowAttempt, + "retry" | "iteration" | "nodeId" | "startedSequence" | "startedAt" +>; + +type TerminalAttemptOutcome = Pick; + +const ABANDONED_ATTEMPT_OUTCOME: TerminalAttemptOutcome = { + outcome: "canceled", + failureCategory: "canceled", + failureMessage: "abandoned: the controller stopped during this attempt, and the resumed run cancelled it" +}; + type SmithersNodeAttemptAuthorities = ReadonlyMap; interface AccountingSummary { @@ -5012,33 +5025,56 @@ function nodeAttemptAgentProvenance(selection: SmithersAttemptAgentSelection): N * recordable after a reset (timetravel / retry-task) restarts the attempt * numbering: the earlier occurrence is only marked `superseded`, because * Smithers upserts the attempt row for the replacement (#1099). Unpairable - * events are skipped, never fatal. A start without a terminal was abandoned - * (Smithers cancels in-progress rows at the next RunStarted without an event), - * and a terminal without a live start in the same activation, such as a - * NodeCancelled for an attempt that already ended or never started, is not an - * occurrence (#1139). Nor is a terminal stamped before its start, which no - * ledger row can hold: Smithers stamps a run cancellation's NodeCancelled with - * an instant it takes before its transaction, so a start that commits meanwhile - * can carry a later timestamp. + * events are skipped, never fatal. A terminal without a live start in the same + * activation, such as a NodeCancelled for an attempt that already ended or + * never started, is not an occurrence (#1139). Nor is a terminal stamped before + * its start, which no ledger row can hold: Smithers stamps a run cancellation's + * NodeCancelled with an instant it takes before its transaction, so a start + * that commits meanwhile can carry a later timestamp. + * + * A start still open at the next RunStarted was abandoned by the activation + * that stopped, for example when its controller was killed. Smithers marks that + * attempt cancelled when the next activation starts, but emits no event for + * it. The occurrence is canceled and ends at the next event of the same task, + * normally its replacement's NodeStarted: one RunStarted can abandon several + * attempts, and each occurrence needs its own terminal event. */ function terminalWorkflowAttempts(events: readonly WorkflowEvent[]): TerminalWorkflowAttempt[] { - const active = new Map< - string, - Pick - >(); + const active = new Map(); + // Starts that an earlier activation left open, keyed by task and iteration. + const abandoned = new Map(); const latestByIdentity = new Map(); const attempts: TerminalWorkflowAttempt[] = []; + const end = (started: WorkflowAttemptStart, event: WorkflowEvent, terminal: TerminalAttemptOutcome): void => { + if (event.timestampMs < Date.parse(started.startedAt)) return; + const occurrence: TerminalWorkflowAttempt = { + ...started, + finishedSequence: event.sourceEventSequence, + finishedAt: new Date(event.timestampMs).toISOString(), + ...terminal, + superseded: false + }; + latestByIdentity.set(JSON.stringify([started.nodeId, started.iteration, started.retry]), occurrence); + attempts.push(occurrence); + }; for (const event of events) { if (event.type === "RunStarted") { + for (const started of active.values()) { + abandoned.set(JSON.stringify([started.nodeId, started.iteration]), started); + } active.clear(); continue; } if (!["NodeStarted", "NodeFinished", "NodeFailed", "NodeCancelled"].includes(event.type)) continue; const { nodeId, iteration, attempt: retry } = event.payload; + if (typeof nodeId !== "string" || !Number.isSafeInteger(iteration)) continue; + const task = JSON.stringify([nodeId, iteration]); + const lost = abandoned.get(task); + abandoned.delete(task); + if (lost !== undefined) end(lost, event, ABANDONED_ATTEMPT_OUTCOME); // A NodeCancelled for a node with no live attempt carries `attempt: null`. - if (typeof nodeId !== "string" || !Number.isSafeInteger(iteration) || !Number.isSafeInteger(retry)) continue; + if (!Number.isSafeInteger(retry)) continue; const identity = JSON.stringify([nodeId, iteration, retry]); - const timestamp = new Date(event.timestampMs).toISOString(); if (event.type === "NodeStarted") { const previous = latestByIdentity.get(identity); if (previous !== undefined) previous.superseded = true; @@ -5047,25 +5083,15 @@ function terminalWorkflowAttempts(events: readonly WorkflowEvent[]): TerminalWor iteration: iteration as number, nodeId, startedSequence: event.sourceEventSequence, - startedAt: timestamp + startedAt: new Date(event.timestampMs).toISOString() }); continue; } const started = active.get(identity); const terminal = terminalOutcomeForEvent(event); - if (started === undefined || terminal === undefined || event.timestampMs < Date.parse(started.startedAt)) continue; + if (started === undefined || terminal === undefined) continue; active.delete(identity); - const occurrence: TerminalWorkflowAttempt = { - ...started, - finishedSequence: event.sourceEventSequence, - finishedAt: timestamp, - outcome: terminal.outcome, - ...(terminal.failureCategory === undefined ? {} : { failureCategory: terminal.failureCategory }), - ...(terminal.failureMessage === undefined ? {} : { failureMessage: terminal.failureMessage }), - superseded: false - }; - latestByIdentity.set(identity, occurrence); - attempts.push(occurrence); + end(started, event, terminal); } return attempts; } @@ -5079,13 +5105,7 @@ function recordedTerminalAttemptSequences( ); } -function terminalOutcomeForEvent(event: WorkflowEvent): - | { - outcome: NodeAttemptOutcome; - failureCategory?: NodeAttemptFailureCategory; - failureMessage?: string; - } - | undefined { +function terminalOutcomeForEvent(event: WorkflowEvent): TerminalAttemptOutcome | undefined { switch (event.type) { case "NodeFinished": return { outcome: "succeeded" }; diff --git a/packages/runtime/test/runtime.test.ts b/packages/runtime/test/runtime.test.ts index c3b96850c..f6b4da1ca 100644 --- a/packages/runtime/test/runtime.test.ts +++ b/packages/runtime/test/runtime.test.ts @@ -99,6 +99,7 @@ import { pauseRun, prepareControllerGeneration, readLinkedWorkflowEvidence, + readReportPublicationStatus, replayRun as runtimeReplayRun, resumeRun as runtimeResumeRun, startRun as runtimeStartRun, @@ -24191,7 +24192,10 @@ test("syncRun tolerates duplicate active starts for one attempt identity", async assert.equal(fs.readFileSync(path.join(run.value.run_root, "attempts.jsonl"), "utf8"), ""); }); -test("syncRun abandons an unterminated occurrence at a later run activation boundary", async () => { +const ABANDONED_ATTEMPT_MESSAGE = + "abandoned: the controller stopped during this attempt, and the resumed run cancelled it"; + +test("syncRun records an attempt abandoned at a run activation boundary even when a reset reuses its number", async () => { const project = tempProject(); initProject({ projectRoot: project, force: true }); writeSmallTopology(project); @@ -24214,11 +24218,182 @@ test("syncRun abandons an unterminated occurrence at a later run activation boun }); const run = await startRun({ projectRoot: project, runId, env }); assert.equal(run.ok, true, JSON.stringify(run.diagnostics)); + assert.ok(run.value); const sync = await syncRun({ projectRoot: project, runId, env }); assert.equal(sync.ok, true, JSON.stringify(sync.diagnostics)); assert.equal(sync.value?.status, "running"); + // The replacement reuses attempt number 1, as after `resume --reset-node`, so + // Smithers' attempt row describes it and the abandoned attempt has no agent. + assert.deepEqual( + attemptLedgerRows(run.value.run_root).map((entry) => [ + entry.started_event_sequence, + entry.source_event_sequence, + entry.outcome, + entry.failure_category, + entry.failure_message, + entry.agent + ]), + [[1, 3, "canceled", "canceled", ABANDONED_ATTEMPT_MESSAGE, undefined]] + ); +}); + +test("syncRun records each attempt that a killed controller abandoned as canceled", async () => { + const project = tempProject(); + initProject({ projectRoot: project, force: true }); + writeOptionalSpecialistTopology(project); + const runId = "sync-crash-abandoned-attempts"; + const workflowRunId = `ultrafuzz-${runId}`; + const tasks = ["direct-strategy", "optional-specialist"]; + // Resume marks each attempt the dead controller left in progress cancelled, + // emits no event for it, and restarts the task with the next attempt number. + const nodeDetails = Object.fromEntries( + tasks.map((task) => { + const nodeId = `node:${task}`; + const meta = { agentChainIndex: 0, agentId: `ultrafuzz-agent:${task}:0:default`, agentModel: "gpt-5.5" }; + const attempts = [ + { nodeId, attempt: 1, state: "cancelled", meta }, + { nodeId, attempt: 2, state: "in-progress", meta } + ]; + return [nodeId, { node: { nodeId, lastAttempt: 2 }, attempts }]; + }) + ); + const env = fakeLifecycleSmithersEnv(project, { + inspect: workflowInspect({ + workflowRunId, + status: "running", + state: "running", + steps: tasks.map((task) => ({ id: `node:${task}`, state: "in-progress" as const, attempt: 2 })) + }), + events: workflowEvents(workflowRunId, [ + { type: "RunStarted" }, + { type: "NodeStarted", nodeId: "node:direct-strategy", attempt: 1 }, + { type: "NodeStarted", nodeId: "node:optional-specialist", attempt: 1 }, + { type: "RunStarted" }, + { type: "NodeStarted", nodeId: "node:optional-specialist", attempt: 2 }, + { type: "NodeStarted", nodeId: "node:direct-strategy", attempt: 2 } + ]), + nodeDetails + }); + const run = await startRun({ projectRoot: project, runId, env }); + assert.equal(run.ok, true, JSON.stringify(run.diagnostics)); + assert.ok(run.value); + + for (let observation = 0; observation < 2; observation += 1) { + const sync = await syncRun({ projectRoot: project, runId, env }); + assert.equal(sync.ok, true, JSON.stringify(sync.diagnostics)); + assert.equal(sync.value?.status, "running"); + assert.deepEqual(sync.diagnostics, []); + } + + // One RunStarted abandoned both attempts, so each ends at its own task's + // replacement start; a repeated synchronization appends nothing. + assert.deepEqual( + attemptLedgerRows(run.value.run_root) + .map((entry) => [ + entry.node_id, + entry.attempt, + entry.started_event_sequence, + entry.source_event_sequence, + entry.outcome, + entry.failure_message, + (entry.reuse as { status?: string }).status, + (entry.agent as { profile_id?: string } | undefined)?.profile_id + ]) + .sort((left, right) => String(left[0]).localeCompare(String(right[0]))), + [ + ["direct-strategy", 1, 1, 5, "canceled", ABANDONED_ATTEMPT_MESSAGE, "executed", "default"], + ["optional-specialist", 1, 2, 4, "canceled", ABANDONED_ATTEMPT_MESSAGE, "executed", "default"] + ] + ); + const nodes = readRunState(layoutForRunRoot(run.value.run_root, runId)).nodes; + assert.deepEqual( + tasks.map((task) => nodes[task]?.retry_count), + [1, 1] + ); +}); + +test("syncRun records an abandoned attempt of an already reported run and publishes the report again", async () => { + const project = tempProject(); + initProject({ projectRoot: project, force: true }); + writeSingleFinalReportTopology(project); + const runId = "sync-reported-abandoned-attempt"; + const workflowRunId = `ultrafuzz-${runId}`; + const nodeId = "node:final-report"; + const meta = { agentChainIndex: 0, agentId: "ultrafuzz-agent:final-report:0:default", agentModel: "gpt-5.5" }; + const base = Date.parse("2026-07-03T00:00:00.000Z"); + const events = [ + { type: "RunStarted" }, + { type: "NodeStarted", nodeId, attempt: 1 }, + { type: "RunStarted" }, + { type: "NodeStarted", nodeId, attempt: 2 }, + { type: "NodeFinished", nodeId, attempt: 2 }, + { type: "NodeFinished", nodeId: "verify:final-report", attempt: 1 }, + { type: "RunFinished" } + ].map((event, sequence) => ({ ...event, sequence, timestampMs: base + sequence * 100 })); + const lifecycle = (observed: typeof events) => + fakeLifecycleSmithersEnv(project, { + inspect: workflowInspect({ + workflowRunId, + steps: [ + { id: nodeId, state: "finished", attempt: 2 }, + { id: "verify:final-report", state: "finished", attempt: 1 } + ] + }), + events: workflowEvents(workflowRunId, observed), + nodeDetails: { + [nodeId]: { + node: { nodeId, lastAttempt: 2 }, + attempts: [ + { nodeId, attempt: 1, state: "cancelled", meta }, + { nodeId, attempt: 2, state: "finished", meta } + ] + } + } + }); + // An earlier version finalized and reported the run without recording the abandoned attempt. + const earlier = lifecycle(events.filter((event) => event.sequence !== 1)); + const run = await startRun({ projectRoot: project, runId, env: earlier }); + assert.equal(run.ok, true, JSON.stringify(run.diagnostics)); + assert.ok(run.value); + const runRoot = run.value.run_root; + writeEmptyFinalReportArtifactSet(runRoot, runId); + const observe = async (env: Record) => { + const sync = await syncRun({ projectRoot: project, runId, env }); + assert.equal(sync.ok, true, JSON.stringify(sync.diagnostics)); + assert.equal(sync.value?.status, "succeeded", JSON.stringify(sync.diagnostics)); + const state = readRunState(layoutForRunRoot(runRoot, runId)); + const report = readReportPublicationStatus(runRoot, state); + return { + ledger: attemptLedgerRows(runRoot).map((entry) => [ + entry.attempt, + entry.started_event_sequence, + entry.source_event_sequence, + entry.outcome + ]), + retryCount: state.nodes["final-report"]?.retry_count, + // What `status` reports, and the verified report itself. + report: [report.status, report.verification, loadCurrentFinalReportSnapshot(runRoot).artifacts.source] + }; + }; + + assert.deepEqual(await observe(earlier), { + ledger: [[2, 3, 4, "succeeded"]], + retryCount: 0, + report: ["available", "verified", "verified-runtime-report"] + }); + // The abandoned attempt is appended after its replacement and raises the + // immutable node's retry count, which changes the state the published report + // describes; the same synchronization publishes the report again. + assert.deepEqual(await observe(lifecycle(events)), { + ledger: [ + [2, 3, 4, "succeeded"], + [1, 1, 3, "canceled"] + ], + retryCount: 1, + report: ["available", "verified", "verified-runtime-report"] + }); }); test("syncRun records external wait reasons from workflow events", async () => {