FN-8919: harden stale agent link recovery

Normalize missing task lookup errors so stale durable-agent links do not halt self-healing sweeps.

- Treat thrown not-found and deleted task lookups as stale links.
- Isolate transient task lookup failures to their affected agent.
- Add recovery regression coverage and a patch changeset.

Files changed:
 .changeset/fn-8919-agent-link-sweep-fail-open.md   |  7 ++
 .../self-healing-agent-link-drift.test.ts          | 58 +++++++++++++--
 .../self-healing-path-utils-task-miss.test.ts      | 51 +++++++++++++
 packages/engine/src/__tests__/self-healing.test.ts | 84 +++++++++++++++++++++-
 .../engine/src/healing/self-healing-path-utils.ts  | 34 ++++++++-
 packages/engine/src/self-healing.ts                | 57 +++++++++++----
 6 files changed, 272 insertions(+), 19 deletions(-)

Fusion-Task-Id: FN-8919

Fusion-Task-Lineage: 6dca2c1d-8ff9-47fb-b6c1-24f5af57b499

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-10 03:31:50 -07:00
parent 313eea1461
commit 08a3f2851b
6 changed files with 272 additions and 19 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Stale agent task links no longer stop self-healing from reconciling later agents.
category: fix
dev: Harden recoverAgentsRunningOnInactiveTasks and recoverDriftedAgentTaskLinks with isMissingTaskLookupError/readLinkedTaskOrUndefined for Runfusion/Fusion#3397.

View File

@@ -1,6 +1,6 @@
import { describe, expect, it, vi } from "vitest";
import { isEphemeralAgent, type Agent, type AgentStore, type Task } from "@fusion/core";
import { isEphemeralAgent, TaskDeletedError, TaskNotFoundError, type Agent, type AgentStore, type Task } from "@fusion/core";
import { SelfHealingManager } from "../self-healing.js";
@@ -9,9 +9,13 @@ function makeAgent(id: string, taskId: string, state: Agent["state"] = "active")
}
describe("FN-4296: self-healing agent link drift", () => {
function buildManager(agents: Agent[], tasks: Record<string, Task | null>, hasActiveAgentExecution?: (agentId: string) => boolean) {
function buildManager(agents: Agent[], tasks: Record<string, Task | null | Error>, hasActiveAgentExecution?: (agentId: string) => boolean) {
const store = {
getTask: vi.fn(async (taskId: string) => tasks[taskId] ?? null),
getTask: vi.fn(async (taskId: string) => {
const result = tasks[taskId] ?? null;
if (result instanceof Error) throw result;
return result;
}),
recordRunAuditEvent: vi.fn(async () => {}),
} as any;
@@ -47,9 +51,12 @@ describe("FN-4296: self-healing agent link drift", () => {
it("FN-4296: durable agent linked to archived task is cleared by sweep", async () => {
const agents = [makeAgent("agent-1", "FN-1")];
const { manager } = buildManager(agents, { "FN-1": { id: "FN-1", column: "archived" } as Task });
const { manager, store } = buildManager(agents, { "FN-1": { id: "FN-1", column: "archived" } as Task });
await manager.recoverDriftedAgentTaskLinks();
expect(agents[0].taskId).toBeUndefined();
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({
metadata: expect.objectContaining({ reason: "linked task in terminal column archived" }),
}));
manager.stop();
});
@@ -185,6 +192,49 @@ describe("FN-4296: self-healing agent link drift", () => {
manager.stop();
});
it("FN-8919: a throwing task miss clears itself and later stale links", async () => {
const agents = [makeAgent("agent-poison", "ERR-024"), makeAgent("agent-stale", "FN-2")];
const { manager, agentStore, store } = buildManager(agents, {
"ERR-024": new TaskNotFoundError("ERR-024"),
"FN-2": null,
});
await expect(manager.recoverDriftedAgentTaskLinks()).resolves.toBe(2);
expect(agentStore.syncExecutionTaskLink).toHaveBeenCalledWith("agent-poison", undefined);
expect(agentStore.syncExecutionTaskLink).toHaveBeenCalledWith("agent-stale", undefined);
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({
target: "agent-poison",
metadata: expect.objectContaining({ reason: "linked task missing" }),
}));
manager.stop();
});
it("FN-8919: TaskDeletedError is a missing task and transient errors isolate one agent", async () => {
const agents = [
makeAgent("agent-transient", "FN-connection"),
makeAgent("agent-deleted", "KB-1"),
makeAgent("agent-stale", "FN-2"),
];
const { manager, agentStore, store } = buildManager(agents, {
"FN-connection": new Error("connection terminated unexpectedly"),
"KB-1": new TaskDeletedError("KB-1", "2026-08-10T00:00:00.000Z"),
"FN-2": null,
});
await expect(manager.recoverDriftedAgentTaskLinks()).resolves.toBe(2);
expect(agents[0].taskId).toBe("FN-connection");
expect(agentStore.syncExecutionTaskLink).not.toHaveBeenCalledWith("agent-transient", undefined);
expect(agentStore.syncExecutionTaskLink).toHaveBeenCalledWith("agent-deleted", undefined);
expect(agentStore.syncExecutionTaskLink).toHaveBeenCalledWith("agent-stale", undefined);
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({
target: "agent-deleted",
metadata: expect.objectContaining({ reason: "linked task missing" }),
}));
manager.stop();
});
it("FN-4296: ephemeral agents are not touched", async () => {
const durable = makeAgent("agent-1", "FN-1");
const ephemeral = makeAgent("temp-worker", "FN-2");

View File

@@ -0,0 +1,51 @@
import { TaskDeletedError, TaskNotFoundError, type Task } from "@fusion/core";
import { describe, expect, it, vi } from "vitest";
import { isMissingTaskLookupError, readLinkedTaskOrUndefined } from "../healing/self-healing-path-utils.js";
describe("isMissingTaskLookupError", () => {
it("recognizes typed and structural task lookup misses across task-id prefixes", () => {
expect(isMissingTaskLookupError(new TaskNotFoundError("ERR-024"))).toBe(true);
expect(isMissingTaskLookupError({ name: "TaskNotFoundError", code: "TASK_NOT_FOUND" })).toBe(true);
expect(isMissingTaskLookupError(new TaskDeletedError("KB-1", "2026-08-10T00:00:00.000Z"))).toBe(true);
expect(isMissingTaskLookupError({ name: "TaskDeletedError" })).toBe(true);
expect(isMissingTaskLookupError(new Error("Task ERR-024 not found"))).toBe(true);
});
it("rejects unrelated, non-error values", () => {
expect(isMissingTaskLookupError(new Error("connection terminated unexpectedly"))).toBe(false);
expect(isMissingTaskLookupError(undefined)).toBe(false);
expect(isMissingTaskLookupError("Task ERR-024 not found")).toBe(false);
});
});
describe("readLinkedTaskOrUndefined", () => {
it("preserves resolved live and archive tasks", async () => {
const live = { id: "FN-1", column: "todo" } as Task;
const archived = { id: "FN-2", column: "archived" } as Task;
const getTask = vi.fn()
.mockResolvedValueOnce(live)
.mockResolvedValueOnce(archived);
const store = { getTask } as any;
await expect(readLinkedTaskOrUndefined(store, "FN-1")).resolves.toBe(live);
await expect(readLinkedTaskOrUndefined(store, "FN-2")).resolves.toBe(archived);
});
it("normalizes nullable and throwing task misses", async () => {
const getTask = vi.fn()
.mockResolvedValueOnce(null)
.mockRejectedValueOnce(new TaskNotFoundError("ERR-024"));
const store = { getTask } as any;
await expect(readLinkedTaskOrUndefined(store, "FN-1")).resolves.toBeUndefined();
await expect(readLinkedTaskOrUndefined(store, "ERR-024")).resolves.toBeUndefined();
});
it("rethrows transient lookup failures", async () => {
const error = new Error("connection terminated unexpectedly");
const store = { getTask: vi.fn().mockRejectedValue(error) } as any;
await expect(readLinkedTaskOrUndefined(store, "FN-1")).rejects.toBe(error);
});
});

View File

@@ -143,7 +143,7 @@ vi.mock("../merger.js", () => ({
import { SelfHealingManager, isBranchAheadOfBase, MAX_AUTO_MERGE_RETRIES } from "../self-healing.js";
import { HEARTBEAT_ERROR_RECOVERY_METADATA_KEY, HEARTBEAT_ERROR_RETRY_EXHAUSTED_PAUSE_REASON, HEARTBEAT_ERROR_UNRECOVERABLE_PAUSE_REASON, readHeartbeatErrorRetryCount } from "../agent-heartbeat.js";
import type { TaskStore, Settings, Task, AgentStore, Agent, NotificationProvider } from "@fusion/core";
import { TaskDeletedError, TaskNotFoundError, type TaskStore, type Settings, type Task, type AgentStore, type Agent, type NotificationProvider } from "@fusion/core";
import { EventEmitter } from "node:events";
import { execSync } from "node:child_process";
import { existsSync, mkdirSync, mkdtempSync, readdirSync, rmSync, writeFileSync } from "node:fs";
@@ -2000,6 +2000,88 @@ describe("SelfHealingManager", () => {
expect(agentStore.syncExecutionTaskLink).not.toHaveBeenCalled();
managerWithAgents.stop();
});
it("FN-8919: thrown task misses recover later stale running agents without aborting", async () => {
const agents = [
{ id: "agent-poison", state: "running", taskId: "ERR-024", updatedAt: new Date().toISOString() } as Agent,
{ id: "agent-stale", state: "running", taskId: "FN-stale", updatedAt: new Date().toISOString() } as Agent,
];
const agentStore = {
listAgents: vi.fn(async () => agents),
getActiveHeartbeatRun: vi.fn(async () => null),
updateAgentState: vi.fn(async (id: string, state: Agent["state"]) => { agents.find((agent) => agent.id === id)!.state = state; }),
syncExecutionTaskLink: vi.fn(async (id: string, taskId?: string) => { agents.find((agent) => agent.id === id)!.taskId = taskId; }),
} as unknown as AgentStore;
const getTask = vi.fn(async (taskId: string) => {
if (taskId === "ERR-024") throw new TaskNotFoundError("ERR-024");
return null;
});
const recoveryStore = createMockStore({ getTask });
const managerWithAgents = new SelfHealingManager(recoveryStore, { rootDir: "/tmp/test-project", agentStore });
await expect(managerWithAgents.recoverAgentsRunningOnInactiveTasks()).resolves.toBe(2);
expect(agentStore.syncExecutionTaskLink).toHaveBeenCalledWith("agent-poison", undefined);
expect(agentStore.syncExecutionTaskLink).toHaveBeenCalledWith("agent-stale", undefined);
expect(recoveryStore.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({
target: "agent-poison",
metadata: expect.objectContaining({
reason: "running durable agent linked to missing task without live execution proof",
}),
}));
managerWithAgents.stop();
});
it("FN-8919: transient lookup errors preserve one running link while later links recover", async () => {
const agents = [
{ id: "agent-transient", state: "running", taskId: "FN-connection", updatedAt: new Date().toISOString() } as Agent,
{ id: "agent-deleted", state: "running", taskId: "KB-1", updatedAt: new Date().toISOString() } as Agent,
{ id: "agent-stale", state: "running", taskId: "FN-stale", updatedAt: new Date().toISOString() } as Agent,
];
const agentStore = {
listAgents: vi.fn(async () => agents),
getActiveHeartbeatRun: vi.fn(async () => null),
updateAgentState: vi.fn(async (id: string, state: Agent["state"]) => { agents.find((agent) => agent.id === id)!.state = state; }),
syncExecutionTaskLink: vi.fn(async (id: string, taskId?: string) => { agents.find((agent) => agent.id === id)!.taskId = taskId; }),
} as unknown as AgentStore;
const getTask = vi.fn(async (taskId: string) => {
if (taskId === "FN-connection") throw new Error("connection terminated unexpectedly");
if (taskId === "KB-1") throw new TaskDeletedError("KB-1", "2026-08-10T00:00:00.000Z");
return null;
});
const managerWithAgents = new SelfHealingManager(createMockStore({ getTask }), { rootDir: "/tmp/test-project", agentStore });
await expect(managerWithAgents.recoverAgentsRunningOnInactiveTasks()).resolves.toBe(2);
expect(agents[0].taskId).toBe("FN-connection");
expect(agentStore.syncExecutionTaskLink).not.toHaveBeenCalledWith("agent-transient", undefined);
expect(agentStore.syncExecutionTaskLink).toHaveBeenCalledWith("agent-deleted", undefined);
expect(agentStore.syncExecutionTaskLink).toHaveBeenCalledWith("agent-stale", undefined);
managerWithAgents.stop();
});
it("FN-8919: archived, live, ephemeral, and unlinked agents keep existing running-sweep behavior", async () => {
const agents = [
{ id: "agent-archive", state: "running", taskId: "FN-archive", updatedAt: new Date().toISOString() } as Agent,
{ id: "agent-live", state: "running", taskId: "FN-live", updatedAt: new Date().toISOString() } as Agent,
{ id: "agent-ephemeral", state: "running", taskId: "FN-stale", metadata: { type: "spawned" }, updatedAt: new Date().toISOString() } as Agent,
{ id: "agent-unlinked", state: "running", updatedAt: new Date().toISOString() } as Agent,
];
const agentStore = {
listAgents: vi.fn(async () => agents),
getActiveHeartbeatRun: vi.fn(async () => null),
updateAgentState: vi.fn(),
syncExecutionTaskLink: vi.fn(),
} as unknown as AgentStore;
const getTask = vi.fn(async (taskId: string) => ({ id: taskId, column: taskId === "FN-archive" ? "archived" : "in-progress" } as Task));
const managerWithAgents = new SelfHealingManager(createMockStore({ getTask }), { rootDir: "/tmp/test-project", agentStore });
await expect(managerWithAgents.recoverAgentsRunningOnInactiveTasks()).resolves.toBe(0);
expect(agentStore.syncExecutionTaskLink).not.toHaveBeenCalled();
expect(getTask).toHaveBeenCalledTimes(2);
managerWithAgents.stop();
});
});
describe("recoverStaleHeartbeatRuns", () => {

View File

@@ -2,7 +2,12 @@
* FNXC:CodeOrganization 2026-07-16-14:00:
* Pure path/error/scope helpers peeled from self-healing.ts.
*/
import type { Task } from "@fusion/core";
import {
isTaskNotFoundError as isCoreTaskNotFoundError,
TaskDeletedError,
type Task,
type TaskStore,
} from "@fusion/core";
export function extractTaskIdFromTempMergeDir(dirname: string): string | null {
const match = /^fusion-ai-merge-(fn-\d+)-[a-z0-9]+$/i.exec(dirname);
@@ -17,6 +22,33 @@ export function isTaskNotFoundError(err: unknown): boolean {
return /\btask\s+fn-\d+\s+not found\b/i.test(getErrorMessage(err));
}
/*
FNXC:SelfHealing 2026-08-10-09:51:
Runfusion/Fusion#3397 requires durable-agent link recovery to treat PostgreSQL
`getTask` misses as ordinary stale-link input. The older FN-only regex above
must remain unchanged for temp-merge paths, because production ids also include
`ERR-024` and other prefixes. Archive snapshots are resolved tasks in a
terminal column, never misses, so this guard only normalizes thrown lookup
errors and not successful task reads.
*/
export function isMissingTaskLookupError(err: unknown): boolean {
if (isCoreTaskNotFoundError(err) || err instanceof TaskDeletedError) return true;
if (err && typeof err === "object" && (err as { name?: unknown }).name === "TaskDeletedError") return true;
return err instanceof Error && /\btask\s+\S+\s+not found\b/i.test(err.message);
}
export async function readLinkedTaskOrUndefined(
store: Pick<TaskStore, "getTask">,
taskId: string,
): Promise<Task | undefined> {
try {
return (await store.getTask(taskId)) ?? undefined;
} catch (err) {
if (isMissingTaskLookupError(err)) return undefined;
throw err;
}
}
export function buildResumeLimboStepSignature(task: Task): string {
return JSON.stringify({
currentStep: task.currentStep ?? null,

View File

@@ -171,6 +171,8 @@ export {
extractTaskIdFromTempMergeDir,
getErrorMessage,
isTaskNotFoundError,
isMissingTaskLookupError,
readLinkedTaskOrUndefined,
buildResumeLimboStepSignature,
formatRecoveryTimestamp,
matchGlob,
@@ -180,6 +182,7 @@ import {
extractTaskIdFromTempMergeDir,
getErrorMessage,
isTaskNotFoundError,
readLinkedTaskOrUndefined,
buildResumeLimboStepSignature,
formatRecoveryTimestamp,
matchesScope,
@@ -13527,7 +13530,17 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
continue;
}
const linkedTask = await this.store.getTask(agent.taskId);
const linkedTaskId = agent.taskId;
try {
/*
FNXC:SelfHealing 2026-08-10-09:52:
Runfusion/Fusion#3397 makes a thrown task-miss ordinary stale-link input:
PostgreSQL `getTask` throws instead of returning null. Archive snapshots
still resolve to terminal-column tasks and are skipped below; non-miss
failures rethrow into this per-agent boundary so they never clear a link
and one poisoned row cannot disable the remaining recovery sweep.
*/
const linkedTask = await readLinkedTaskOrUndefined(this.store, linkedTaskId);
/* FNXC:WorkflowResolvedColumns 2026-07-31-16:50 (fleet, round 2): WIP u REVIEW u TERMINAL. Keyed on
ids none matched on a renamed board, so this sweep evaluated agents whose task was still executing. */
if (linkedTask && (agentLinkLiveColumns.has(linkedTask.column) || agentLinkTerminalColumns.has(linkedTask.column))) {
@@ -13584,14 +13597,18 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
await agentStore.syncExecutionTaskLink(agent.id, undefined);
await this.emitStaleAgentAssignmentAudit({
agent: { id: agent.id, state: staleAgentState },
taskId: agent.taskId,
linkedTask,
hadFreshRun: proof.hasFreshRun,
hadActiveExecution: proof.hasActiveExecution,
reason,
});
recoveredAgentIds.add(agent.id);
log.log(`Recovered running durable agent ${agent.id} on inactive task ${agent.taskId}; file-scope lease preserved when present`);
taskId: linkedTaskId,
linkedTask,
hadFreshRun: proof.hasFreshRun,
hadActiveExecution: proof.hasActiveExecution,
reason,
});
recoveredAgentIds.add(agent.id);
log.log(`Recovered running durable agent ${agent.id} on inactive task ${linkedTaskId}; file-scope lease preserved when present`);
} catch (err) {
log.warn(`Failed to recover running durable agent ${agent.id} on ${linkedTaskId}: ${getErrorMessage(err)}`);
continue;
}
}
return recoveredAgentIds.size;
@@ -13615,8 +13632,18 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
}
const linkedTaskId = agent.taskId;
const linkedTask = await this.store.getTask(linkedTaskId);
let shouldClear = false;
try {
/*
FNXC:SelfHealing 2026-08-10-09:53:
Runfusion/Fusion#3397 applies the same fail-open-per-candidate principle
as `finalizeOrphanedPlanningSegments` (#3386): missing-task throws are
stale-link input, but non-miss failures leave this agent linked. An
archive snapshot resolves to its terminal column and deliberately keeps
this sweep's existing terminal-column clearing behavior; one bad row
must not abort later durable-agent reconciliation.
*/
const linkedTask = await readLinkedTaskOrUndefined(this.store, linkedTaskId);
let shouldClear = false;
let reason = "";
let hadFreshRun = false;
let hadActiveExecution = false;
@@ -13678,8 +13705,12 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
hadActiveExecution,
reason,
});
clearedAgentIds.add(agent.id);
log.log(`Cleared drifted durable agent task link for ${agent.id} (${linkedTaskId}): ${reason}; file-scope lease preserved when present`);
clearedAgentIds.add(agent.id);
log.log(`Cleared drifted durable agent task link for ${agent.id} (${linkedTaskId}): ${reason}; file-scope lease preserved when present`);
} catch (err) {
log.warn(`Failed to reconcile drifted durable agent task link for ${agent.id} (${linkedTaskId}): ${getErrorMessage(err)}`);
continue;
}
}
/*