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"));
|
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", () => {
|
||||||
|
|||||||
@@ -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> {
|
||||||
|
|||||||
Reference in New Issue
Block a user