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:
7
.changeset/fn-8919-agent-link-sweep-fail-open.md
Normal file
7
.changeset/fn-8919-agent-link-sweep-fail-open.md
Normal 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.
|
||||||
@@ -1,6 +1,6 @@
|
|||||||
import { describe, expect, it, vi } from "vitest";
|
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";
|
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", () => {
|
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 = {
|
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 () => {}),
|
recordRunAuditEvent: vi.fn(async () => {}),
|
||||||
} as any;
|
} 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 () => {
|
it("FN-4296: durable agent linked to archived task is cleared by sweep", async () => {
|
||||||
const agents = [makeAgent("agent-1", "FN-1")];
|
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();
|
await manager.recoverDriftedAgentTaskLinks();
|
||||||
expect(agents[0].taskId).toBeUndefined();
|
expect(agents[0].taskId).toBeUndefined();
|
||||||
|
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({
|
||||||
|
metadata: expect.objectContaining({ reason: "linked task in terminal column archived" }),
|
||||||
|
}));
|
||||||
manager.stop();
|
manager.stop();
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -185,6 +192,49 @@ describe("FN-4296: self-healing agent link drift", () => {
|
|||||||
manager.stop();
|
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 () => {
|
it("FN-4296: ephemeral agents are not touched", async () => {
|
||||||
const durable = makeAgent("agent-1", "FN-1");
|
const durable = makeAgent("agent-1", "FN-1");
|
||||||
const ephemeral = makeAgent("temp-worker", "FN-2");
|
const ephemeral = makeAgent("temp-worker", "FN-2");
|
||||||
|
|||||||
@@ -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);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -143,7 +143,7 @@ vi.mock("../merger.js", () => ({
|
|||||||
|
|
||||||
import { SelfHealingManager, isBranchAheadOfBase, MAX_AUTO_MERGE_RETRIES } from "../self-healing.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 { 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 { EventEmitter } from "node:events";
|
||||||
import { execSync } from "node:child_process";
|
import { execSync } from "node:child_process";
|
||||||
import { existsSync, mkdirSync, mkdtempSync, readdirSync, rmSync, writeFileSync } from "node:fs";
|
import { existsSync, mkdirSync, mkdtempSync, readdirSync, rmSync, writeFileSync } from "node:fs";
|
||||||
@@ -2000,6 +2000,88 @@ describe("SelfHealingManager", () => {
|
|||||||
expect(agentStore.syncExecutionTaskLink).not.toHaveBeenCalled();
|
expect(agentStore.syncExecutionTaskLink).not.toHaveBeenCalled();
|
||||||
managerWithAgents.stop();
|
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", () => {
|
describe("recoverStaleHeartbeatRuns", () => {
|
||||||
|
|||||||
@@ -2,7 +2,12 @@
|
|||||||
* FNXC:CodeOrganization 2026-07-16-14:00:
|
* FNXC:CodeOrganization 2026-07-16-14:00:
|
||||||
* Pure path/error/scope helpers peeled from self-healing.ts.
|
* 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 {
|
export function extractTaskIdFromTempMergeDir(dirname: string): string | null {
|
||||||
const match = /^fusion-ai-merge-(fn-\d+)-[a-z0-9]+$/i.exec(dirname);
|
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));
|
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 {
|
export function buildResumeLimboStepSignature(task: Task): string {
|
||||||
return JSON.stringify({
|
return JSON.stringify({
|
||||||
currentStep: task.currentStep ?? null,
|
currentStep: task.currentStep ?? null,
|
||||||
|
|||||||
@@ -171,6 +171,8 @@ export {
|
|||||||
extractTaskIdFromTempMergeDir,
|
extractTaskIdFromTempMergeDir,
|
||||||
getErrorMessage,
|
getErrorMessage,
|
||||||
isTaskNotFoundError,
|
isTaskNotFoundError,
|
||||||
|
isMissingTaskLookupError,
|
||||||
|
readLinkedTaskOrUndefined,
|
||||||
buildResumeLimboStepSignature,
|
buildResumeLimboStepSignature,
|
||||||
formatRecoveryTimestamp,
|
formatRecoveryTimestamp,
|
||||||
matchGlob,
|
matchGlob,
|
||||||
@@ -180,6 +182,7 @@ import {
|
|||||||
extractTaskIdFromTempMergeDir,
|
extractTaskIdFromTempMergeDir,
|
||||||
getErrorMessage,
|
getErrorMessage,
|
||||||
isTaskNotFoundError,
|
isTaskNotFoundError,
|
||||||
|
readLinkedTaskOrUndefined,
|
||||||
buildResumeLimboStepSignature,
|
buildResumeLimboStepSignature,
|
||||||
formatRecoveryTimestamp,
|
formatRecoveryTimestamp,
|
||||||
matchesScope,
|
matchesScope,
|
||||||
@@ -13527,7 +13530,17 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
|
|||||||
continue;
|
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
|
/* 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. */
|
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))) {
|
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 agentStore.syncExecutionTaskLink(agent.id, undefined);
|
||||||
await this.emitStaleAgentAssignmentAudit({
|
await this.emitStaleAgentAssignmentAudit({
|
||||||
agent: { id: agent.id, state: staleAgentState },
|
agent: { id: agent.id, state: staleAgentState },
|
||||||
taskId: agent.taskId,
|
taskId: linkedTaskId,
|
||||||
linkedTask,
|
linkedTask,
|
||||||
hadFreshRun: proof.hasFreshRun,
|
hadFreshRun: proof.hasFreshRun,
|
||||||
hadActiveExecution: proof.hasActiveExecution,
|
hadActiveExecution: proof.hasActiveExecution,
|
||||||
reason,
|
reason,
|
||||||
});
|
});
|
||||||
recoveredAgentIds.add(agent.id);
|
recoveredAgentIds.add(agent.id);
|
||||||
log.log(`Recovered running durable agent ${agent.id} on inactive task ${agent.taskId}; file-scope lease preserved when present`);
|
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;
|
return recoveredAgentIds.size;
|
||||||
@@ -13615,8 +13632,18 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
|
|||||||
}
|
}
|
||||||
|
|
||||||
const linkedTaskId = agent.taskId;
|
const linkedTaskId = agent.taskId;
|
||||||
const linkedTask = await this.store.getTask(linkedTaskId);
|
try {
|
||||||
let shouldClear = false;
|
/*
|
||||||
|
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 reason = "";
|
||||||
let hadFreshRun = false;
|
let hadFreshRun = false;
|
||||||
let hadActiveExecution = false;
|
let hadActiveExecution = false;
|
||||||
@@ -13678,8 +13705,12 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
|
|||||||
hadActiveExecution,
|
hadActiveExecution,
|
||||||
reason,
|
reason,
|
||||||
});
|
});
|
||||||
clearedAgentIds.add(agent.id);
|
clearedAgentIds.add(agent.id);
|
||||||
log.log(`Cleared drifted durable agent task link for ${agent.id} (${linkedTaskId}): ${reason}; file-scope lease preserved when present`);
|
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;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/*
|
/*
|
||||||
|
|||||||
Reference in New Issue
Block a user