diff --git a/.changeset/fn-8909-orphaned-planning-segment-sweep.md b/.changeset/fn-8909-orphaned-planning-segment-sweep.md new file mode 100644 index 0000000000..b2ecf4ffcb --- /dev/null +++ b/.changeset/fn-8909-orphaned-planning-segment-sweep.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Keep planning-time recovery running when archived tasks retain timing anchors. +category: fix +dev: Self-healing now enumerates live non-archived tasks and isolates per-task failures. diff --git a/packages/engine/src/__tests__/orphaned-planning-segment-poisoned-row.pg.test.ts b/packages/engine/src/__tests__/orphaned-planning-segment-poisoned-row.pg.test.ts new file mode 100644 index 0000000000..a933037d80 --- /dev/null +++ b/packages/engine/src/__tests__/orphaned-planning-segment-poisoned-row.pg.test.ts @@ -0,0 +1,91 @@ +/* +FNXC:TaskTiming 2026-08-09-20:34: +Archive cold-storage reconstruction retains planningStartedAt but omits deletedAt, +while archiving tombstones the live row. A mock listTasks fixture cannot prove +that cross-store behavior, so these cases use a real PostgreSQL TaskStore and a +throwaway database to keep the orphan-finalization invariant production-real. +*/ +import { afterAll, afterEach, beforeAll, beforeEach, expect, it } from "vitest"; +import "@fusion/core"; +import type { TaskStore } from "@fusion/core"; + +import { + createSharedPgTaskStoreTestHarness, + pgDescribe, + type SharedPgTaskStoreHarness, +} from "../../../core/src/__test-utils__/pg-test-harness.js"; +import { SelfHealingManager } from "../self-healing.js"; + +pgDescribe("orphaned planning segment poisoned rows", () => { + const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ + prefix: "fusion_orphaned_planning_segment", + }); + + beforeAll(h.beforeAll); + afterAll(h.afterAll); + beforeEach(async () => { await h.beforeEach(); }); + afterEach(async () => { await h.afterEach(); }); + + function createManager(store: TaskStore): SelfHealingManager { + return new SelfHealingManager(store, { + rootDir: "/tmp/test-project", + getPlanningTaskIds: () => new Set(), + hasActivePlanningWorkflowSession: () => false, + }); + } + + it("finalizes a live orphan without touching an archive snapshot that retains its planning anchor", async () => { + const store = h.store(); + const archived = await store.createTask({ description: "archived planning anchor" }); + const live = await store.createTask({ description: "live planning anchor" }); + const archivedStartedAt = "2026-08-08T16:40:09.736Z"; + const liveStartedAt = "2026-08-08T16:41:09.736Z"; + + await store.updateTask(archived.id, { planningStartedAt: archivedStartedAt, cumulativePlanningMs: 25 }); + await store.archiveTask(archived.id, { cleanup: false }); + await store.updateTask(live.id, { planningStartedAt: liveStartedAt, cumulativePlanningMs: 50 }); + + const manager = createManager(store); + await expect(manager.finalizeOrphanedPlanningSegments()).resolves.toBe(1); + + const finalizedLive = await store.getTask(live.id); + expect(finalizedLive!.planningStartedAt).toBeUndefined(); + expect(finalizedLive!.cumulativePlanningMs).toBeGreaterThan(50); + + const archivedEntries = await store.listArchivedTasks({ slim: true }); + expect(archivedEntries.tasks.find((task) => task.id === archived.id)).toMatchObject({ + planningStartedAt: archivedStartedAt, + cumulativePlanningMs: 25, + }); + + const audit = await store.getRunAuditEventsAsync(); + expect(audit.some((event) => + event.mutationType === "task:reconcile-orphaned-planning-segment" && event.taskId === live.id, + )).toBe(true); + expect(audit.some((event) => + event.mutationType === "task:reconcile-orphaned-planning-segment" && event.taskId === archived.id, + )).toBe(false); + + manager.stop(); + }); + + it("keeps plain soft-deleted rows out of the candidate list while finalizing a co-existing live orphan", async () => { + const store = h.store(); + const deleted = await store.createTask({ description: "deleted planning anchor" }); + const live = await store.createTask({ description: "live planning anchor" }); + + await store.updateTask(deleted.id, { planningStartedAt: "2026-08-08T16:40:09.736Z", cumulativePlanningMs: 25 }); + await store.deleteTask(deleted.id); + await store.updateTask(live.id, { planningStartedAt: "2026-08-08T16:41:09.736Z", cumulativePlanningMs: 50 }); + + const manager = createManager(store); + await expect(manager.finalizeOrphanedPlanningSegments()).resolves.toBe(1); + + const finalizedLive = await store.getTask(live.id); + expect(finalizedLive!.planningStartedAt).toBeUndefined(); + expect(finalizedLive!.cumulativePlanningMs).toBeGreaterThan(50); + expect((await store.listTasks({ slim: true, includeArchived: false })).map((task) => task.id)).not.toContain(deleted.id); + + manager.stop(); + }); +}); diff --git a/packages/engine/src/__tests__/self-healing.test.ts b/packages/engine/src/__tests__/self-healing.test.ts index 6c752e0eb9..e4c76b96af 100644 --- a/packages/engine/src/__tests__/self-healing.test.ts +++ b/packages/engine/src/__tests__/self-healing.test.ts @@ -8836,6 +8836,8 @@ describe("SelfHealingManager", () => { vi.setSystemTime(new Date("2026-01-01T00:00:01.000Z")); expect(await recovery.finalizeOrphanedPlanningSegments()).toBe(1); + // Call shape only: the PostgreSQL regression test proves this excludes archive snapshots. + expect(recoveryStore.listTasks).toHaveBeenCalledWith(expect.objectContaining({ slim: true, includeArchived: false })); expect(updateTaskAtomic).toHaveBeenCalledOnce(); expect(task).toMatchObject({ cumulativePlanningMs: 1050, planningStartedAt: null }); expect(recoveryStore.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ @@ -8871,6 +8873,79 @@ describe("SelfHealingManager", () => { recovery.stop(); }); + + it.each([ + ["before", ["FN-POISON", "FN-HEALTHY"]], + ["after", ["FN-HEALTHY", "FN-POISON"]], + ])("skips a deleted candidate listed %s a healthy orphan", async (_position, ids) => { + const poisoned = { id: "FN-POISON", deletedAt: "2026-08-08T16:42:53.336Z", planningStartedAt: "2026-01-01T00:00:00.000Z" } as Task; + const healthy = { id: "FN-HEALTHY", planningStartedAt: "2026-01-01T00:00:00.000Z", cumulativePlanningMs: 50 } as Task; + const tasks = ids.map((id) => id === poisoned.id ? poisoned : healthy); + const updateTaskAtomic = vi.fn(async (id: string, updater: (live: Task) => Partial | null) => { + const task = id === healthy.id ? healthy : poisoned; + const patch = updater(task); + if (patch) Object.assign(task, patch); + return patch; + }); + const recoveryStore = createMockStore({ listTasks: vi.fn().mockResolvedValue(tasks), updateTaskAtomic }); + const recovery = new SelfHealingManager(recoveryStore, { rootDir: "/tmp/test-project", getPlanningTaskIds: () => new Set() }); + vi.setSystemTime(new Date("2026-01-01T00:00:01.000Z")); + + await expect(recovery.finalizeOrphanedPlanningSegments()).resolves.toBe(1); + expect(updateTaskAtomic).toHaveBeenCalledTimes(1); + expect(updateTaskAtomic).toHaveBeenCalledWith(healthy.id, expect.any(Function)); + expect(healthy).toMatchObject({ planningStartedAt: null, cumulativePlanningMs: 1050 }); + + recovery.stop(); + }); + + it("contains a TOCTOU soft-delete failure and continues with other candidates", async () => { + const racing = { id: "FN-RACING", planningStartedAt: "2026-01-01T00:00:00.000Z", cumulativePlanningMs: 50 } as Task; + const healthy = { id: "FN-HEALTHY", planningStartedAt: "2026-01-01T00:00:00.000Z", cumulativePlanningMs: 50 } as Task; + const updateTaskAtomic = vi.fn(async (id: string, updater: (live: Task) => Partial | null) => { + if (id === racing.id) throw new Error("Task FN-RACING is soft-deleted (deletedAt=2026-08-08T16:42:53.336Z) and cannot be read or mutated"); + const patch = updater(healthy); + if (patch) Object.assign(healthy, patch); + return patch; + }); + const recoveryStore = createMockStore({ listTasks: vi.fn().mockResolvedValue([racing, healthy]), updateTaskAtomic }); + const recovery = new SelfHealingManager(recoveryStore, { rootDir: "/tmp/test-project", getPlanningTaskIds: () => new Set() }); + getSelfHealingLogger().warn.mockClear(); + vi.setSystemTime(new Date("2026-01-01T00:00:01.000Z")); + + await expect(recovery.finalizeOrphanedPlanningSegments()).resolves.toBe(1); + expect(healthy).toMatchObject({ planningStartedAt: null, cumulativePlanningMs: 1050 }); + expect(getSelfHealingLogger().warn).toHaveBeenCalledWith(expect.stringContaining("Failed to finalize orphaned planning segment for FN-RACING")); + + recovery.stop(); + }); + + it("contains fallback getTask failures and still finalizes healthy candidates around deleted rows", async () => { + const firstPoison = { id: "FN-POISON-1", deletedAt: "2026-08-08T16:42:53.336Z", planningStartedAt: "2026-01-01T00:00:00.000Z" } as Task; + const healthy = { id: "FN-HEALTHY", planningStartedAt: "2026-01-01T00:00:00.000Z", cumulativePlanningMs: 50 } as Task; + const secondPoison = { id: "FN-POISON-2", planningStartedAt: "2026-01-01T00:00:00.000Z" } as Task; + const getTask = vi.fn(async (id: string) => { + if (id === secondPoison.id) throw new Error("Task FN-POISON-2 is soft-deleted (deletedAt=2026-08-08T16:42:53.336Z) and cannot be read or mutated"); + return healthy; + }); + const updateTask = vi.fn(async (_id: string, patch: Partial) => Object.assign(healthy, patch)); + const recoveryStore = createMockStore({ + listTasks: vi.fn().mockResolvedValue([firstPoison, healthy, secondPoison]), + updateTaskAtomic: undefined, + getTask, + updateTask, + }); + const recovery = new SelfHealingManager(recoveryStore, { rootDir: "/tmp/test-project", getPlanningTaskIds: () => new Set() }); + getSelfHealingLogger().warn.mockClear(); + vi.setSystemTime(new Date("2026-01-01T00:00:01.000Z")); + + await expect(recovery.finalizeOrphanedPlanningSegments()).resolves.toBe(1); + expect(updateTask).toHaveBeenCalledWith(healthy.id, expect.objectContaining({ planningStartedAt: null, cumulativePlanningMs: 1050 })); + expect(getTask).not.toHaveBeenCalledWith(firstPoison.id); + expect(getSelfHealingLogger().warn).toHaveBeenCalledWith(expect.stringContaining("Failed to finalize orphaned planning segment for FN-POISON-2")); + + recovery.stop(); + }); }); describe("recoverOrphanedPlanningTasks", () => { diff --git a/packages/engine/src/self-healing.ts b/packages/engine/src/self-healing.ts index ba0ae75ed8..bdc9142129 100644 --- a/packages/engine/src/self-healing.ts +++ b/packages/engine/src/self-healing.ts @@ -14740,56 +14740,81 @@ const movedTask = await this.store.moveTask(task.id, completeLane); * A planning anchor is safe because triage ownership and graph Plan Review are * exclusive. Recovery finalizes only when neither in-process owner is live; * the atomic null-check makes restart and repeated maintenance idempotent. + * + * FNXC:TaskTiming 2026-08-09-20:34: + * listTasks merges archive cold-storage snapshots by default. Those snapshots + * retain planningStartedAt but carry no deletedAt, while their tombstoned live + * rows are refused by updateTaskAtomic and getTask. Enumerate live rows only; + * a deletedAt filter alone could not contain this archive path. Per-task and + * sweep guards ensure one race or poisoned row cannot disable finalization + * store-wide during startup or maintenance (Runfusion/Fusion#3386). */ async finalizeOrphanedPlanningSegments(): Promise { - const planningIds = this.options.getPlanningTaskIds?.() ?? new Set(); - const tasks = await this.store.listTasks({}); let finalized = 0; - for (const task of tasks) { - if (!task.planningStartedAt || planningIds.has(task.id) || this.options.hasActivePlanningWorkflowSession?.(task.id)) continue; - let applied = false; - const endMs = Date.now(); - if (typeof this.store.updateTaskAtomic === "function") { - await this.store.updateTaskAtomic(task.id, (live) => { - if (!live.planningStartedAt || planningIds.has(live.id) || this.options.hasActivePlanningWorkflowSession?.(live.id)) return null; - const patch = finalizePlanningSegment(live, endMs); - applied = patch.planningStartedAt === null; - return patch; - }); - } else { - const live = await this.store.getTask(task.id); - if (live?.planningStartedAt && !planningIds.has(live.id) && !this.options.hasActivePlanningWorkflowSession?.(live.id)) { - const patch = finalizePlanningSegment(live, endMs); - if (patch.planningStartedAt === null) { await this.store.updateTask(task.id, patch); applied = true; } + try { + const planningIds = this.options.getPlanningTaskIds?.() ?? new Set(); + const tasks = await this.store.listTasks({ slim: true, includeArchived: false }); + for (const task of tasks) { + try { + // Defense-in-depth for legacy/custom stores; archive snapshots require + // includeArchived: false above because they deliberately omit deletedAt. + if ( + task.deletedAt || + !task.planningStartedAt || + planningIds.has(task.id) || + this.options.hasActivePlanningWorkflowSession?.(task.id) + ) continue; + let applied = false; + const endMs = Date.now(); + if (typeof this.store.updateTaskAtomic === "function") { + await this.store.updateTaskAtomic(task.id, (live) => { + if (!live.planningStartedAt || planningIds.has(live.id) || this.options.hasActivePlanningWorkflowSession?.(live.id)) return null; + const patch = finalizePlanningSegment(live, endMs); + applied = patch.planningStartedAt === null; + return patch; + }); + } else { + const live = await this.store.getTask(task.id); + if (live?.planningStartedAt && !planningIds.has(live.id) && !this.options.hasActivePlanningWorkflowSession?.(live.id)) { + const patch = finalizePlanningSegment(live, endMs); + if (patch.planningStartedAt === null) { await this.store.updateTask(task.id, patch); applied = true; } + } + } + if (!applied) continue; + finalized++; + // FNXC:TaskTiming 2026-07-30-21:40: this recovery is operator-auditable + // without persisting duration prose; the atomically finalized task id + // and fixed no-live-owner reason are sufficient forensic evidence. + await this.store.recordRunAuditEvent?.({ + taskId: task.id, + agentId: "self-healing", + runId: generateSyntheticRunId("orphaned-planning-segment", task.id), + domain: "database", + mutationType: "task:reconcile-orphaned-planning-segment", + target: task.id, + metadata: { taskId: task.id, finalizedCount: 1, reason: "no-live-planning-owner" }, + }); + } catch (err: unknown) { + const errorMessage = err instanceof Error ? err.message : String(err); + log.warn(`Failed to finalize orphaned planning segment for ${task.id}: ${errorMessage}`); } } - if (applied) { - finalized++; - // FNXC:TaskTiming 2026-07-30-21:40: this recovery is operator-auditable - // without persisting duration prose; the atomically finalized task id - // and fixed no-live-owner reason are sufficient forensic evidence. + if (finalized === 0) { await this.store.recordRunAuditEvent?.({ - taskId: task.id, agentId: "self-healing", - runId: generateSyntheticRunId("orphaned-planning-segment", task.id), + runId: generateSyntheticRunId("orphaned-planning-segment", "global"), domain: "database", - mutationType: "task:reconcile-orphaned-planning-segment", - target: task.id, - metadata: { taskId: task.id, finalizedCount: 1, reason: "no-live-planning-owner" }, + mutationType: "task:reconcile-orphaned-planning-segment-no-action", + target: "planning-segments", + metadata: { finalizedCount: 0, reason: "no-eligible-orphan" }, }); } + return finalized; + } catch (err: unknown) { + const errorMessage = err instanceof Error ? err.message : String(err); + log.error(`Orphaned planning segment finalization failed: ${errorMessage}`); + return finalized; } - if (finalized === 0) { - await this.store.recordRunAuditEvent?.({ - agentId: "self-healing", - runId: generateSyntheticRunId("orphaned-planning-segment", "global"), - domain: "database", - mutationType: "task:reconcile-orphaned-planning-segment-no-action", - target: "planning-segments", - metadata: { finalizedCount: 0, reason: "no-eligible-orphan" }, - }); - } - return finalized; } async recoverOrphanedPlanningTasks(): Promise {