diff --git a/.changeset/orphaned-pending-rewrite-to-failed.md b/.changeset/orphaned-pending-rewrite-to-failed.md new file mode 100644 index 0000000000..84cd68a61a --- /dev/null +++ b/.changeset/orphaned-pending-rewrite-to-failed.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Orphaned in-flight review steps are now marked failed for re-review instead of silently skipped at merge. +category: fix +dev: `resolveOrphanedPendingStepResults` rewrites orphans to `status:"failed"` (never deletes — deletion satisfied the merge gate and skipped review); the sweep also runs in periodic maintenance, skips `in-progress` rows, re-reads before writing, and the audit event is registered in `DatabaseMutationType` with metadata `{taskId, column, orphanedCount, resultCount}`. diff --git a/AGENTS.md b/AGENTS.md index 62c707cb03..ea2ab1af64 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -283,7 +283,7 @@ Scoped exception (FN-5819): shared-branch-group members (`branchContext.assignme - FN-8004: `agent:heartbeat-move-skipped-soft-delete` records a heartbeat move that races a soft-deleted task without parking the durable agent. Metadata remains ids/timestamps/source only (`agentId`, optional `taskId`/`deletedAt`, `moveAttemptedAt`, optional `source`); it never stores error prose. - FN-8141: the executor's `fn_task_done(outcome="blocked", reason=..., blockedBy?=[...])` honest-blocked exit emits `task:execution-blocked-parked` when an executor parks a genuinely-impossible task `failed` (`error = "BLOCKED: "`) instead of laundering it to `done` by skipping steps. It bypasses the completion/verdict/bulk-completion gates (blocked is not a completion claim), leaves steps in their true statuses, preserves worktree/branch, records `blockedBy` as real `task.dependencies` edges so the task requeues behind the blocker, and does NOT hand off to review — the parked row is honored by the executor's `status === "failed"` post-loop branch and is not auto-recovered into in-review by `recoverStrandedCompletedTodoTasks` (steps are not all done/skipped and `task.error` is set). Metadata stays ids/outcomes-only (`taskId`, `blockedBy` ids, `hasReason` boolean — never the reason prose). - FN-8305: durable symbol-lock operations emit `symbol-lock:acquired`, `symbol-lock:acquire-conflict`, `symbol-lock:renewed`, `symbol-lock:released`, `symbol-lock:reconcile-stale`, and deduplicated `symbol-lock:reconcile-stale-no-action`. Metadata is ids/counts/outcomes-only; normalized opaque symbol keys are permitted IDs, while raw symbol prose is not. -- FN-8492: the self-healing STARTUP sweep `reconcile-orphaned-pending-step-results` emits `task:reconcile-orphaned-pending-step-results` when it clears `pending` workflow-step results with no live session behind them (canonical liveness triple: `activeSessionRegistry` path, `executingTaskLock`, `isTaskActive`). It runs right after legacy adoption and before the in-review recovery steps, because an orphaned `pending` result reads as "work in flight" to the merge gate and rides the in-review stall escalator to a `failed` deadlock park. Metadata is ids/counts-only (`taskId`, `column`, `clearedCount`, `remainingCount`); user pauses and live sessions are never disturbed. +- FN-8492: the self-healing sweep `reconcile-orphaned-pending-step-results` (startup, right after legacy adoption, plus periodic maintenance) emits `task:reconcile-orphaned-pending-step-results` when it REWRITES `pending` workflow-step results with no live session behind them to `failed` (canonical liveness triple: `activeSessionRegistry` path, `executingTaskLock`, `isTaskActive`). It must never DELETE an orphaned entry — the merge gate blocks on pending/failed results, not on an enabled step with no result, so deletion silently satisfies the gate and the task merges with its review skipped; the `failed` rewrite keeps the gate closed and hands re-run/bypass to the failed-pre-merge-steps recovery and FN-7720 operator-bypass paths. `in-progress` rows are always skipped (executor-owned; resume is deferred at startup), the row is re-read immediately before the write, and user pauses are never disturbed. Metadata is ids/counts-only (`taskId`, `column`, `orphanedCount`, `resultCount`). - FN-8356: self-healing emits `task:reconcile-stale-duplicate-decision` when it clears a triage-marker duplicate-decision pause against a missing, deleted, done, or archived canonical. Metadata is ids/outcomes-only (`taskId`, `canonicalId`, `canonicalColumn`, `canonicalDeleted`, `priorPausedReason`); active canonical decisions and user pauses remain untouched. - U9b (R10/KTD-8): the self-healing STARTUP recovery step `adopt-legacy-task-rows` emits `task:reconcile-legacy-adoption` when it adopts a pre-cutover row through the KTD-8 adoption table (clearing a legacy `task.status` whose writer the cutover deleted so the graph re-enters at its owning node, and/or landing the one-time `reviewLevel` -> `enabledWorkflowSteps` preset backfill), and `task:reconcile-legacy-adoption-unmappable` when an UNKNOWN status parks the row `paused` for a human with its status deliberately left in place. Metadata is ids/counts/outcomes-only (`taskId`, `action`, `priorStatus`, `column`, `backfilledStepCount`, `reason`), where `reason` is a fixed adoption-table note and never row prose. Adoption runs FIRST in startup recovery (every later step reasons about `task.status`), stamps `task.legacyAdoptedAt` only on rows it actually mutates (so upgrade does not mass-write every `done` row), and never touches a user pause or a `preserve` gate. `planLegacyAdoption` in `packages/core/src/legacy-adoption.ts` is the single shared decision used by both this sweep and the store-open reconcile so the two cannot drift. - U10 (R9): the pre-graph cutover machinery is DELETED and stays deleted, ratcheted by `packages/engine/src/__tests__/legacy-tombstones.test.ts`. Gone: `workflow-cutover.ts`, `workflow-authoritative-driver.ts`, `workflow-parity-observer.ts`, the `graphCompletionInterceptors` re-entry map, triage's out-of-graph `runPlanReviewBeforeExecution` gate, and the in-session `fn_review_step` tool with its RETHINK git-reset/session-rewind, per-step conversation checkpoints, deferred reviewer provider-error channel, and review-level prompt scaffolding. Plan/code/browser review are owned EXCLUSIVELY by workflow-graph nodes — do not re-introduce a second review authority inside the implementation session; that duplicate-Plan-Review race is what the cutover removed. The tombstone test strips comments before searching, so the FNXC notes that explain each deletion are expected to remain in source while the code must not. diff --git a/packages/core/src/__tests__/legacy-adoption.test.ts b/packages/core/src/__tests__/legacy-adoption.test.ts index b70197931e..aa0e224e98 100644 --- a/packages/core/src/__tests__/legacy-adoption.test.ts +++ b/packages/core/src/__tests__/legacy-adoption.test.ts @@ -290,25 +290,40 @@ FNXC:LegacyAdoption 2026-07-19-04:40 (U9b / KTD-8): Orphaned pending step results. A pre-cutover crash leaves a `pending` result with no live session and the graph waits on it forever; a LEASED one is real work in flight. */ -describe("resolveOrphanedPendingStepResults (U9b)", () => { - it("clears pending results with no live session and preserves live ones", () => { - const results = [ +describe("resolveOrphanedPendingStepResults (U9b, FN-8492 rewrite-to-failed)", () => { + it("marks dead-session pending results failed and preserves live/completed ones", () => { + const input = [ { stepIndex: 0, status: "done" }, { stepIndex: 1, status: "pending" }, // orphaned { stepIndex: 2, status: "pending" }, // live — leased { stepIndex: 3, status: "failed" }, ]; - const { cleared, clearedCount } = resolveOrphanedPendingStepResults( - results, + const { results, orphanedCount } = resolveOrphanedPendingStepResults( + input, (r) => r.stepIndex === 2, + { output: "orphan-note", completedAt: "2026-07-22T23:00:00.000Z" }, ); - expect(clearedCount).toBe(1); - expect(cleared.map((r) => r.stepIndex)).toEqual([0, 2, 3]); + expect(orphanedCount).toBe(1); + expect(results.map((r) => r.status)).toEqual(["done", "failed", "pending", "failed"]); + expect(results[1]).toMatchObject({ output: "orphan-note", completedAt: "2026-07-22T23:00:00.000Z" }); + }); + + /* + FNXC:OrphanedPendingSteps 2026-07-22-16:35 (FN-8492 review follow-up): + Deletion is a severity inversion: the merge gate blocks on pending/failed results, not + on an enabled step with NO result, so deleting a dead review's pending entry silently + satisfied the gate and the task merged with its review skipped. Rewrite, never delete. + */ + it("NEVER deletes an orphaned entry — deletion silently satisfied the merge gate", () => { + const input = [{ stepIndex: 0, status: "pending" }]; + const { results } = resolveOrphanedPendingStepResults(input, () => false); + expect(results).toHaveLength(1); + expect(results[0]?.status).toBe("failed"); }); it("is a no-op on empty/absent results", () => { - expect(resolveOrphanedPendingStepResults([], () => false).clearedCount).toBe(0); - expect(resolveOrphanedPendingStepResults(null, () => false).clearedCount).toBe(0); + expect(resolveOrphanedPendingStepResults([], () => false).orphanedCount).toBe(0); + expect(resolveOrphanedPendingStepResults(null, () => false).orphanedCount).toBe(0); }); }); diff --git a/packages/core/src/legacy-adoption.ts b/packages/core/src/legacy-adoption.ts index c88cdb05a4..7af3463e82 100644 --- a/packages/core/src/legacy-adoption.ts +++ b/packages/core/src/legacy-adoption.ts @@ -57,6 +57,8 @@ export interface LegacyAdoptionAction { export const LEGACY_STATUS_ADOPTION: Readonly> = { // ── Triage plan-review statuses whose writers U3 deleted → graph re-entry ── "planning": { kind: "resume-graph", note: "re-enter planning node" }, + "plan-review-unavailable": { kind: "resume-graph", note: "plan-review retry (leased)" }, + // ── Live graph signals with post-cutover writers — preserve, never clear ─── /* FNXC:LegacyAdoption 2026-07-22-15:55 (FN-8498 incident): `needs-replan` is NOT legacy — post-U3 it is written live by the graph's own plan-replan @@ -67,7 +69,6 @@ export const LEGACY_STATUS_ADOPTION: Readonly( results: readonly T[] | null | undefined, isLive: (result: T) => boolean, -): { cleared: T[]; clearedCount: number } { - if (!results || results.length === 0) return { cleared: [], clearedCount: 0 }; - let clearedCount = 0; - const cleared = results.filter((result) => { - if (result.status !== "pending") return true; - if (isLive(result)) return true; - clearedCount++; - return false; + orphanMark?: { output?: string; completedAt?: string }, +): { results: T[]; orphanedCount: number } { + if (!results || results.length === 0) return { results: [], orphanedCount: 0 }; + let orphanedCount = 0; + const next = results.map((result) => { + if (result.status !== "pending") return result; + if (isLive(result)) return result; + orphanedCount++; + return { ...result, status: "failed", ...(orphanMark ?? {}) } as T; }); - return { cleared, clearedCount }; + return { results: next, orphanedCount }; } diff --git a/packages/engine/src/__tests__/self-healing-orphaned-pending-step-results.test.ts b/packages/engine/src/__tests__/self-healing-orphaned-pending-step-results.test.ts index c075c1ef3e..cd89b0fff6 100644 --- a/packages/engine/src/__tests__/self-healing-orphaned-pending-step-results.test.ts +++ b/packages/engine/src/__tests__/self-healing-orphaned-pending-step-results.test.ts @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { EventEmitter } from "node:events"; import type { Settings, Task, TaskStore, WorkflowStepResult } from "@fusion/core"; @@ -14,14 +14,21 @@ vi.mock("../run-audit.js", async (importOriginal) => { }); import { SelfHealingManager } from "../self-healing.js"; +import { activeSessionRegistry, executingTaskLock } from "../active-session-registry.js"; /* FNXC:OrphanedPendingSteps 2026-07-22-16:20 (FN-8492 incident): An engine restart killed an in-flight pre-merge Code Review session, leaving its `pending` workflowStepResult with no live session behind it. The merge gate read that as "incomplete pre-merge workflow steps" and after 3 identical 30-minute stalls the deadlock -disposer parked the task `failed`. These tests pin the startup sweep that clears such -orphans — and the liveness veto that keeps it from eating a genuinely live session. +disposer parked the task `failed`. These tests pin the sweep that recovers such orphans — +and the liveness veto that keeps it from eating a genuinely live session. + +FNXC:OrphanedPendingSteps 2026-07-22-16:35 (review follow-up): +The sweep REWRITES orphans to status:"failed" — it must never delete them. Deleting a +pending review entry silently satisfied the merge gate (an enabled step with no result +does not block) and FN-8492 merged with Code Review skipped. The rewrite keeps the gate +closed and routes re-run/bypass through the failed-pre-merge-steps paths. */ function stepResult(overrides: Partial = {}): WorkflowStepResult { @@ -55,7 +62,13 @@ function storeFor(tasks: Task[]): TaskStore & EventEmitter { const tasksById = new Map(tasks.map((entry) => [entry.id, entry])); return Object.assign(new EventEmitter(), { getSettings: vi.fn(async () => ({ globalPause: false, enginePaused: false } as Settings)), - listTasks: vi.fn(async () => [...tasksById.values()]), + // Honors limit/offset so the >500-row pagination path is actually exercised. + listTasks: vi.fn(async (options?: { limit?: number; offset?: number }) => { + const all = [...tasksById.values()]; + const offset = options?.offset ?? 0; + const limit = options?.limit ?? all.length; + return all.slice(offset, offset + limit); + }), getTask: vi.fn(async (id: string) => tasksById.get(id)), updateTask: vi.fn(async (id: string, patch: Partial) => { const next = { ...tasksById.get(id)!, ...patch } as Task; @@ -65,10 +78,14 @@ function storeFor(tasks: Task[]): TaskStore & EventEmitter { }) as unknown as TaskStore & EventEmitter; } -describe("FN-8492: reconcile orphaned pending step results at startup", () => { +describe("FN-8492: reconcile orphaned pending step results", () => { beforeEach(() => vi.clearAllMocks()); + afterEach(() => { + executingTaskLock._clearForTest(); + for (const path of ["/wt/registry-live"]) activeSessionRegistry.unregisterPath(path); + }); - it("clears a dead-session pending result, keeps completed ones, and audits ids/counts only", async () => { + it("rewrites a dead-session pending result to failed (never deletes) and audits ids/counts only", async () => { const stranded = task("FN-1", { workflowStepResults: [ stepResult({ status: "passed", verdict: "APPROVE" }), @@ -80,47 +97,91 @@ describe("FN-8492: reconcile orphaned pending step results at startup", () => { expect(await manager.reconcileOrphanedPendingStepResults()).toBe(1); const recovered = await store.getTask("FN-1"); - expect(recovered?.workflowStepResults).toHaveLength(1); + // Rewrite-to-failed: same length, gate stays closed via the failed entry. + expect(recovered?.workflowStepResults).toHaveLength(2); expect(recovered?.workflowStepResults?.[0]?.status).toBe("passed"); + expect(recovered?.workflowStepResults?.[1]?.status).toBe("failed"); + expect(recovered?.workflowStepResults?.[1]?.completedAt).toBeTruthy(); expect(recordRunAuditEventMock).toHaveBeenCalledTimes(1); expect(recordRunAuditEventMock).toHaveBeenCalledWith(expect.objectContaining({ type: "task:reconcile-orphaned-pending-step-results", target: "FN-1", - metadata: expect.objectContaining({ taskId: "FN-1", clearedCount: 1, remainingCount: 1 }), + metadata: expect.objectContaining({ taskId: "FN-1", orphanedCount: 1, resultCount: 2 }), })); }); - it("never clears a pending result while the task session is live (executor resumed it)", async () => { - const live = task("FN-LIVE", { - workflowStepResults: [stepResult({ status: "pending" })], - }); - const store = storeFor([live]); + it("vetoes on every leg of the liveness triple: isTaskActive, registry path, executing lock", async () => { + const viaCallback = task("FN-CB", { workflowStepResults: [stepResult({ status: "pending" })] }); + const viaRegistry = task("FN-REG", { workflowStepResults: [stepResult({ status: "pending" })] }); + const viaLock = task("FN-LOCK", { workflowStepResults: [stepResult({ status: "pending" })] }); + activeSessionRegistry.registerPath("/wt/registry-live", { taskId: "FN-REG", kind: "workflow-step", ownerKey: "test" }); + expect(executingTaskLock.tryClaim("FN-LOCK")).toBe(true); + const store = storeFor([viaCallback, viaRegistry, viaLock]); const manager = new SelfHealingManager(store, { rootDir: "/repo", - isTaskActive: (id: string) => id === "FN-LIVE", + isTaskActive: (id: string) => id === "FN-CB", }); expect(await manager.reconcileOrphanedPendingStepResults()).toBe(0); - expect((await store.getTask("FN-LIVE"))?.workflowStepResults).toHaveLength(1); + for (const id of ["FN-CB", "FN-REG", "FN-LOCK"]) { + expect((await store.getTask(id))?.workflowStepResults?.[0]?.status).toBe("pending"); + } expect(recordRunAuditEventMock).not.toHaveBeenCalled(); }); - it("leaves user-paused tasks and tasks with no pending results untouched", async () => { + it("skips user-paused and in-progress rows, and tasks with no pending results", async () => { const userPaused = task("FN-PAUSED", { userPaused: true, paused: true, workflowStepResults: [stepResult({ status: "pending" })], }); + // Executor-owned: resumeOrphaned re-attaches its session on a deferred timer, so + // startup liveness is unprovable — the sweep must never judge in-progress rows. + const inProgress = task("FN-INPROG", { + column: "in-progress", + workflowStepResults: [stepResult({ status: "pending" })], + }); const complete = task("FN-DONE-STEPS", { workflowStepResults: [stepResult({ status: "passed" }), stepResult({ status: "failed" })], }); const noResults = task("FN-NONE"); - const store = storeFor([userPaused, complete, noResults]); + const store = storeFor([userPaused, inProgress, complete, noResults]); const manager = new SelfHealingManager(store, { rootDir: "/repo" }); expect(await manager.reconcileOrphanedPendingStepResults()).toBe(0); - expect((await store.getTask("FN-PAUSED"))?.workflowStepResults).toHaveLength(1); + expect((await store.getTask("FN-PAUSED"))?.workflowStepResults?.[0]?.status).toBe("pending"); + expect((await store.getTask("FN-INPROG"))?.workflowStepResults?.[0]?.status).toBe("pending"); expect((await store.getTask("FN-DONE-STEPS"))?.workflowStepResults).toHaveLength(2); expect(recordRunAuditEventMock).not.toHaveBeenCalled(); }); + + it("paginates past 500 rows and recovers orphans on every page", async () => { + const many = Array.from({ length: 502 }, (_, i) => + task(`FN-P${i}`, { workflowStepResults: [stepResult({ status: "pending" })] })); + const store = storeFor(many); + const manager = new SelfHealingManager(store, { rootDir: "/repo" }); + + expect(await manager.reconcileOrphanedPendingStepResults()).toBe(502); + expect((await store.getTask("FN-P501"))?.workflowStepResults?.[0]?.status).toBe("failed"); + // Two pages of 500 + the short page signalling the end. + expect((store.listTasks as ReturnType).mock.calls.length).toBeGreaterThanOrEqual(2); + }); + + it("isolates a per-task updateTask failure: the other orphan is still recovered and counted", async () => { + const failing = task("FN-FAILS", { workflowStepResults: [stepResult({ status: "pending" })] }); + const healthy = task("FN-OK", { workflowStepResults: [stepResult({ status: "pending" })] }); + const store = storeFor([failing, healthy]); + const passthrough = (store.updateTask as ReturnType).getMockImplementation()! as + (id: string, patch: Partial) => Promise; + (store.updateTask as ReturnType).mockImplementation(async (id: string, patch: Partial) => { + if (id === "FN-FAILS") throw new Error("write refused"); + return passthrough(id, patch); + }); + const manager = new SelfHealingManager(store, { rootDir: "/repo" }); + + expect(await manager.reconcileOrphanedPendingStepResults()).toBe(1); + expect((await store.getTask("FN-OK"))?.workflowStepResults?.[0]?.status).toBe("failed"); + expect(recordRunAuditEventMock).toHaveBeenCalledTimes(1); + expect(recordRunAuditEventMock).toHaveBeenCalledWith(expect.objectContaining({ target: "FN-OK" })); + }); }); diff --git a/packages/engine/src/run-audit.ts b/packages/engine/src/run-audit.ts index 18bac47709..84a1ca3749 100644 --- a/packages/engine/src/run-audit.ts +++ b/packages/engine/src/run-audit.ts @@ -571,6 +571,14 @@ export type DatabaseMutationType = * deliberately left in place so the operator can see what it carried. Same metadata shape. */ | "task:reconcile-legacy-adoption-unmappable" + /* + FNXC:OrphanedPendingSteps 2026-07-22-16:35 (FN-8492): + Startup/periodic rewrite of orphaned `pending` workflow-step results (no live session + behind them) to `failed`, so the merge gate stays closed and failed-pre-merge-steps + recovery owns the re-run. Metadata ids/counts-only: + { taskId, column, orphanedCount, resultCount }. + */ + | "task:reconcile-orphaned-pending-step-results" /** * FNXC:MergeQueue 2026-07-15-10:05: * Wedged single-flight merge reclaim. Metadata ids/outcomes-only: diff --git a/packages/engine/src/self-healing.ts b/packages/engine/src/self-healing.ts index 683efaa0c2..37e994b124 100644 --- a/packages/engine/src/self-healing.ts +++ b/packages/engine/src/self-healing.ts @@ -2698,6 +2698,11 @@ export class SelfHealingManager { { name: "reconcile-done-task-integrity", fn: () => this.reconcileDoneTaskIntegrity() }, { name: "reconcile-stale-merger-status", fn: () => this.reconcileStaleMergerStatus() }, { name: "reconcile-stale-duplicate-decision", fn: () => this.reconcileStaleDuplicateDecisionPause() }, + // FNXC:OrphanedPendingSteps 2026-07-22-16:35 (FN-8492 review follow-up): also + // steady-state — a step session can die without an engine restart, and startup-only + // cadence left that case riding the 3×30-min stall escalator to a deadlock park. + // Live sessions register their worktree path, so the liveness veto holds here. + { name: "reconcile-orphaned-pending-step-results", fn: () => this.reconcileOrphanedPendingStepResults() }, { name: "recover-mergeable-review", fn: () => this.recoverMergeableReviewTasks() }, // FNXC:Workspace 2026-06-22-09:30 (Phase D U1) — workspace-mode reconcilers. { name: "reconcile-workspace-partial-lands", fn: () => this.reconcileWorkspacePartialLands() }, @@ -6673,18 +6678,30 @@ export class SelfHealingManager { /* FNXC:OrphanedPendingSteps 2026-07-22-16:20 (FN-8492 incident): - Startup consumer of `resolveOrphanedPendingStepResults` — the U9b helper shipped with NO - caller (the same gap U9 left for the adoption table), so an engine restart that killed an + Consumer of `resolveOrphanedPendingStepResults` — the U9b helper shipped with NO caller + (the same gap U9 left for the adoption table), so an engine restart that killed an in-flight pre-merge step session left its `pending` workflowStepResult behind forever. The merge gate read it as "incomplete pre-merge workflow steps", surfaced an identical stall every 30 minutes, and after 3 stalls the deadlock disposer parked the task `failed` (FN-8492: Code Review died in the 21:29 restart, task parked two hours later). - Runs at STARTUP, right after legacy adoption: sessions do not survive a restart, so a - `pending` result with no live session behind it is orphaned by construction. Liveness - uses the canonical triple (activeSessionRegistry path, executingTaskLock, isTaskActive) - because runStartupRecovery runs AFTER the executor resumes orphaned sessions — a resumed - session re-registers its path and must veto the clear. User pauses are never disturbed. + FNXC:OrphanedPendingSteps 2026-07-22-16:35 (review follow-up, same day): + Orphans are REWRITTEN to status:"failed" (never deleted) — deleting a pending review + entry silently satisfied the merge gate and FN-8492 merged with Code Review skipped; the + failed rewrite keeps the gate closed and hands re-run/bypass to the existing + failed-pre-merge-steps recovery and FN-7720 operator-bypass paths. + + Runs at startup (right after legacy adoption) AND in periodic maintenance: a step + session can die without an engine restart, and waiting for the next boot leaves the + task riding the 3×30-min stall escalator to a deadlock park. Liveness uses the + canonical triple (activeSessionRegistry path, executingTaskLock, isTaskActive) — live + workflow-step sessions register their worktree path (kind "workflow-step") for their + whole run, so a live review always vetoes. `in-progress` rows are skipped entirely: + they are executor-owned and `resumeOrphaned()` re-attaches their sessions on a + DEFERRED timer (getResumeOrphanDelayMs, ~30s), so at startup their liveness cannot be + proven yet and must not be judged here. The row is re-read immediately before the + write so the whole-array update cannot clobber a lease written after the page + snapshot. User pauses are never disturbed. */ async reconcileOrphanedPendingStepResults(): Promise { try { @@ -6692,24 +6709,47 @@ export class SelfHealingManager { let offset = 0; let recovered = 0; + const isSessionLive = (taskId: string): boolean => { + const livePaths = activeSessionRegistry.pathsForTask(taskId); + return livePaths.some((path) => activeSessionRegistry.isPathActive(path)) + || executingTaskLock.has(taskId) + || this.options.isTaskActive?.(taskId) === true; + }; + for (;;) { const tasks = await this.store.listTasks({ slim: true, includeArchived: false, limit: pageSize, offset }); for (const task of tasks) { // An operator park is authoritative; this sweep must not reach through it. if (task.userPaused === true) continue; - const results = task.workflowStepResults; - if (!results?.some((result) => result.status === "pending")) continue; + // Executor-owned rows: resume is deferred at startup, liveness unprovable here. + if (task.column === "in-progress") continue; + if (!task.workflowStepResults?.some((result) => result.status === "pending")) continue; + if (isSessionLive(task.id)) continue; - const livePaths = activeSessionRegistry.pathsForTask(task.id); - const hasActiveRegisteredPath = livePaths.some((path) => activeSessionRegistry.isPathActive(path)); - const sessionLive = hasActiveRegisteredPath || executingTaskLock.has(task.id) - || this.options.isTaskActive?.(task.id) === true; - - const { cleared, clearedCount } = resolveOrphanedPendingStepResults(results, () => sessionLive); - if (clearedCount === 0) continue; + // Re-read the live row before mutating: the page snapshot can be stale against + // a merger/planner that wrote a fresh pending lease after the page was fetched. + const fresh = await this.store.getTask(task.id); + if (!fresh || fresh.userPaused === true || fresh.column === "in-progress") continue; + const { results, orphanedCount } = resolveOrphanedPendingStepResults( + fresh.workflowStepResults, + () => isSessionLive(task.id), + { + output: "Step session did not survive an engine restart or crash; marked failed by self-healing (FN-8492).", + completedAt: new Date().toISOString(), + }, + ); + if (orphanedCount === 0) continue; try { - await this.store.updateTask(task.id, { workflowStepResults: cleared }); + await this.store.updateTask(task.id, { workflowStepResults: results }); + // Counted on the successful mutation; a failed audit emit below must not + // understate how many tasks were actually recovered. + recovered += 1; + } catch (error) { + log.warn(`reconcileOrphanedPendingStepResults: failed for ${task.id}: ${error instanceof Error ? error.message : String(error)}`); + continue; + } + try { await createRunAuditor(this.store, { runId: generateSyntheticRunId("reconcile-orphaned-pending-steps", task.id), agentId: "self-healing", @@ -6717,25 +6757,24 @@ export class SelfHealingManager { taskLineageId: task.lineageId, phase: "reconcile-orphaned-pending-step-results", }).database({ - type: "task:reconcile-orphaned-pending-step-results" as DatabaseMutationType, + type: "task:reconcile-orphaned-pending-step-results", target: task.id, // ids/counts/outcomes only — never step output or reviewer prose. metadata: { taskId: task.id, column: task.column, - clearedCount, - remainingCount: cleared.length, + orphanedCount, + resultCount: results.length, }, }); - recovered += 1; } catch (error) { - log.warn(`reconcileOrphanedPendingStepResults: failed for ${task.id}: ${error instanceof Error ? error.message : String(error)}`); + log.warn(`reconcileOrphanedPendingStepResults: audit emit failed for ${task.id}: ${error instanceof Error ? error.message : String(error)}`); } } if (tasks.length < pageSize) break; offset += tasks.length; } - if (recovered > 0) log.log(`Cleared orphaned pending step results on ${recovered} task(s)`); + if (recovered > 0) log.log(`Marked orphaned pending step results failed on ${recovered} task(s)`); return recovered; } catch (error) { log.error(`reconcileOrphanedPendingStepResults failed: ${error instanceof Error ? error.message : String(error)}`);