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:
7
.changeset/fn-8909-orphaned-planning-segment-sweep.md
Normal file
7
.changeset/fn-8909-orphaned-planning-segment-sweep.md
Normal 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.
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
@@ -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<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", () => {
|
||||
|
||||
@@ -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<number> {
|
||||
const planningIds = this.options.getPlanningTaskIds?.() ?? new Set<string>();
|
||||
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<string>();
|
||||
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<number> {
|
||||
|
||||
Reference in New Issue
Block a user