From 08a3f2851b1aef50a31646c58234bb06e1103209 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Mon, 10 Aug 2026 03:31:50 -0700 Subject: [PATCH] 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) --- .../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 +++++++++++ .../engine/src/__tests__/self-healing.test.ts | 84 ++++++++++++++++++- .../src/healing/self-healing-path-utils.ts | 34 +++++++- packages/engine/src/self-healing.ts | 57 ++++++++++--- 6 files changed, 272 insertions(+), 19 deletions(-) create mode 100644 .changeset/fn-8919-agent-link-sweep-fail-open.md create mode 100644 packages/engine/src/__tests__/self-healing-path-utils-task-miss.test.ts diff --git a/.changeset/fn-8919-agent-link-sweep-fail-open.md b/.changeset/fn-8919-agent-link-sweep-fail-open.md new file mode 100644 index 0000000000..44088b119e --- /dev/null +++ b/.changeset/fn-8919-agent-link-sweep-fail-open.md @@ -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. diff --git a/packages/engine/src/__tests__/self-healing-agent-link-drift.test.ts b/packages/engine/src/__tests__/self-healing-agent-link-drift.test.ts index c2b6aee711..29a5c96f4e 100644 --- a/packages/engine/src/__tests__/self-healing-agent-link-drift.test.ts +++ b/packages/engine/src/__tests__/self-healing-agent-link-drift.test.ts @@ -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, hasActiveAgentExecution?: (agentId: string) => boolean) { + function buildManager(agents: Agent[], tasks: Record, 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"); diff --git a/packages/engine/src/__tests__/self-healing-path-utils-task-miss.test.ts b/packages/engine/src/__tests__/self-healing-path-utils-task-miss.test.ts new file mode 100644 index 0000000000..8ba1acd505 --- /dev/null +++ b/packages/engine/src/__tests__/self-healing-path-utils-task-miss.test.ts @@ -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); + }); +}); diff --git a/packages/engine/src/__tests__/self-healing.test.ts b/packages/engine/src/__tests__/self-healing.test.ts index 1fd0718a32..1e709b80ec 100644 --- a/packages/engine/src/__tests__/self-healing.test.ts +++ b/packages/engine/src/__tests__/self-healing.test.ts @@ -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", () => { diff --git a/packages/engine/src/healing/self-healing-path-utils.ts b/packages/engine/src/healing/self-healing-path-utils.ts index df6aaa3f32..fd7e3b2042 100644 --- a/packages/engine/src/healing/self-healing-path-utils.ts +++ b/packages/engine/src/healing/self-healing-path-utils.ts @@ -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, + taskId: string, +): Promise { + 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, diff --git a/packages/engine/src/self-healing.ts b/packages/engine/src/self-healing.ts index ab5cfb94d3..8697e77ace 100644 --- a/packages/engine/src/self-healing.ts +++ b/packages/engine/src/self-healing.ts @@ -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; + } } /*