FN-8909: isolate orphaned planning segment recovery failures

Keep planning-time recovery active when archived or soft-deleted tasks retain timing anchors.

- Restrict orphaned planning segment candidates to live, non-archived tasks.
- Isolate per-task and sweep failures so poisoned rows do not abort healthy finalization.
- Add PostgreSQL and unit regressions for archived, deleted, and racing rows.
- Add a patch changeset for the recovery fix.

Files changed:
 .../fn-8909-orphaned-planning-segment-sweep.md     |   7 ++
 ...phaned-planning-segment-poisoned-row.pg.test.ts |  91 ++++++++++++++++++
 packages/engine/src/__tests__/self-healing.test.ts |  75 +++++++++++++++
 packages/engine/src/self-healing.ts                | 103 +++++++++++++--------
 4 files changed, 237 insertions(+), 39 deletions(-)

Fusion-Task-Id: FN-8909

Fusion-Task-Lineage: 247867ea-90ce-4768-b21c-2441788170f3

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-09 13:51:52 -07:00
parent e72d9da1f7
commit c1c41ac60a
4 changed files with 237 additions and 39 deletions

View File

@@ -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.

View File

@@ -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<string>(),
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();
});
});

View File

@@ -8836,6 +8836,8 @@ describe("SelfHealingManager", () => {
vi.setSystemTime(new Date("2026-01-01T00:00:01.000Z")); vi.setSystemTime(new Date("2026-01-01T00:00:01.000Z"));
expect(await recovery.finalizeOrphanedPlanningSegments()).toBe(1); 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(updateTaskAtomic).toHaveBeenCalledOnce();
expect(task).toMatchObject({ cumulativePlanningMs: 1050, planningStartedAt: null }); expect(task).toMatchObject({ cumulativePlanningMs: 1050, planningStartedAt: null });
expect(recoveryStore.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ expect(recoveryStore.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({
@@ -8871,6 +8873,79 @@ describe("SelfHealingManager", () => {
recovery.stop(); 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<Task> | 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<string>() });
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<Task> | 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<string>() });
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<Task>) => 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<string>() });
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", () => { describe("recoverOrphanedPlanningTasks", () => {

View File

@@ -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 * A planning anchor is safe because triage ownership and graph Plan Review are
* exclusive. Recovery finalizes only when neither in-process owner is live; * exclusive. Recovery finalizes only when neither in-process owner is live;
* the atomic null-check makes restart and repeated maintenance idempotent. * 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<number> { async finalizeOrphanedPlanningSegments(): Promise<number> {
const planningIds = this.options.getPlanningTaskIds?.() ?? new Set<string>();
const tasks = await this.store.listTasks({});
let finalized = 0; let finalized = 0;
for (const task of tasks) { try {
if (!task.planningStartedAt || planningIds.has(task.id) || this.options.hasActivePlanningWorkflowSession?.(task.id)) continue; const planningIds = this.options.getPlanningTaskIds?.() ?? new Set<string>();
let applied = false; const tasks = await this.store.listTasks({ slim: true, includeArchived: false });
const endMs = Date.now(); for (const task of tasks) {
if (typeof this.store.updateTaskAtomic === "function") { try {
await this.store.updateTaskAtomic(task.id, (live) => { // Defense-in-depth for legacy/custom stores; archive snapshots require
if (!live.planningStartedAt || planningIds.has(live.id) || this.options.hasActivePlanningWorkflowSession?.(live.id)) return null; // includeArchived: false above because they deliberately omit deletedAt.
const patch = finalizePlanningSegment(live, endMs); if (
applied = patch.planningStartedAt === null; task.deletedAt ||
return patch; !task.planningStartedAt ||
}); planningIds.has(task.id) ||
} else { this.options.hasActivePlanningWorkflowSession?.(task.id)
const live = await this.store.getTask(task.id); ) continue;
if (live?.planningStartedAt && !planningIds.has(live.id) && !this.options.hasActivePlanningWorkflowSession?.(live.id)) { let applied = false;
const patch = finalizePlanningSegment(live, endMs); const endMs = Date.now();
if (patch.planningStartedAt === null) { await this.store.updateTask(task.id, patch); applied = true; } 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) { if (finalized === 0) {
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?.({ await this.store.recordRunAuditEvent?.({
taskId: task.id,
agentId: "self-healing", agentId: "self-healing",
runId: generateSyntheticRunId("orphaned-planning-segment", task.id), runId: generateSyntheticRunId("orphaned-planning-segment", "global"),
domain: "database", domain: "database",
mutationType: "task:reconcile-orphaned-planning-segment", mutationType: "task:reconcile-orphaned-planning-segment-no-action",
target: task.id, target: "planning-segments",
metadata: { taskId: task.id, finalizedCount: 1, reason: "no-live-planning-owner" }, 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<number> { async recoverOrphanedPlanningTasks(): Promise<number> {