fix(FN-8998): prevent planning lock reentry stalls

This commit is contained in:
gsxdsm
2026-08-11 21:44:22 -07:00
parent 11334b1249
commit d9a0ed7837
8 changed files with 384 additions and 18 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Prevent completed planning sessions from stalling before Plan Review or execution.
category: fix
dev: Avoids nested planning lifecycle locks and preserves recoverable written plans during orphan cleanup.

View File

@@ -54,4 +54,21 @@ describe("classifyPersistedPlanHandoff", () => {
awaitingApprovalReason: "require-all", awaitingApprovalReason: "require-all",
}), options)).toBeNull(); }), options)).toBeNull();
}); });
it("recovers a planning-status task with its retained planning worktree", () => {
expect(classifyPersistedPlanHandoff(approvedNullTask({
status: "planning",
workflowStepResults: [],
worktree: "/tmp/fusion-planning-worktree",
}), options)).toBe("planning");
});
it("keeps the retained-worktree fence for null-status legacy recovery", () => {
expect(classifyPersistedPlanHandoff(approvedNullTask({
status: null,
workflowStepResults: [],
steps: [{ id: "planned-step", description: "planned", status: "pending" }],
worktree: "/tmp/fusion-planning-worktree",
}), options)).toBeNull();
});
}); });

View File

@@ -23,6 +23,7 @@ import {
import { getPromptPath } from "../../execution/spec-staleness.js"; import { getPromptPath } from "../../execution/spec-staleness.js";
import { promoteHeldTask, runHoldReleaseSweep } from "../../execution/hold-release.js"; import { promoteHeldTask, runHoldReleaseSweep } from "../../execution/hold-release.js";
import { SelfHealingManager } from "../../self-healing.js"; import { SelfHealingManager } from "../../self-healing.js";
import { TriageProcessor } from "../../triage.js";
const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({
prefix: "fusion_planning_dependency_release", prefix: "fusion_planning_dependency_release",
@@ -78,6 +79,58 @@ pgDescribe("FN-8768 planning dependency release interactions", () => {
return store.getTask(id); return store.getTask(id);
} }
it("finalizes a persisted plan without reacquiring its PostgreSQL lifecycle lock", async () => {
const store = h.store();
const taskId = "FN-8998-PG";
const prompt = [
`# ${taskId}: Planning lifecycle lock regression`,
"",
"## Mission",
"",
"Finalize a persisted plan while the task planning lifecycle lock is held.",
"",
"## Steps",
"",
"### Step 1: Implement",
"- [ ] Preserve the single-lock finalization invariant.",
"",
].join("\n");
await store.updateSettings({ planApprovalMode: "auto-approve-all" });
await store.createTaskWithReservedId(
{ description: "Planning lifecycle lock regression", column: "todo" },
{ taskId, applyDefaultWorkflowSteps: true },
);
const promptPath = getPromptPath(store.getTasksDir(), taskId);
mkdirSync(dirname(promptPath), { recursive: true });
writeFileSync(promptPath, prompt, "utf8");
await store.updateTask(taskId, { status: "planning", steps: [] });
store.taskCache.delete(taskId);
const task = await store.getTask(taskId);
await expect(
new TriageProcessor(store, h.rootDir()).recoverApprovedTask(task),
).resolves.toBe(true);
store.taskCache.delete(taskId);
const released = await store.getTask(taskId);
const currentPlan = await store.getLatestCurrentPlanEvidence(taskId);
const latestLock = await store.getLatestSpecLock(taskId);
expect(released.column).toBe("todo");
expect(released.status).toBeUndefined();
expect(released.approvedPlanFingerprint).toEqual(expect.any(String));
expect(currentPlan).toMatchObject({
plan: { status: "available", contentHash: expect.any(String) },
});
expect(latestLock).toMatchObject({
approvalFingerprint: released.approvedPlanFingerprint,
currentPlanVersion: currentPlan?.version,
currentPlanHash: currentPlan?.plan.contentHash,
});
await expect(store.getActiveSpecLock(taskId)).resolves.toEqual(latestLock);
});
it.each([ it.each([
{ {
label: "dependency mutation API", label: "dependency mutation API",

View File

@@ -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 { TaskDeletedError, TaskNotFoundError, type TaskStore, type Settings, type Task, type AgentStore, type Agent, type NotificationProvider } from "@fusion/core"; import { PlanningLifecycleLockTransportError, 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";
@@ -9491,7 +9491,189 @@ describe("SelfHealingManager", () => {
}); });
}); });
describe("recoverOrphanedPlanningTasks", () => { describe("planning handoff and orphan recovery", () => {
it("finalizes a recoverable written plan before clearing it for re-planning", async () => {
const task = {
id: "FN-PLAN-HANDOFF",
column: "todo",
status: "planning",
worktree: "/tmp/fusion-planning-worktree",
paused: false,
log: [],
updatedAt: "2026-01-01T00:00:00.000Z",
} as unknown as Task;
const recoverApprovedTriageTask = vi.fn().mockResolvedValue(true);
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([task]);
const recovery = new SelfHealingManager(store, {
rootDir: "/tmp/test-project",
getPlanningTaskIds: () => new Set<string>(),
recoverApprovedTriageTask,
});
vi.setSystemTime(new Date("2026-01-01T00:05:00.000Z"));
await expect(recovery.recoverApprovedTriageTasks()).resolves.toBe(1);
expect(recoverApprovedTriageTask).toHaveBeenCalledWith(task);
expect(store.updateTask).not.toHaveBeenCalled();
recovery.stop();
});
it("backs off a planning-lock transport failure before retrying the retained handoff", async () => {
const task = {
id: "FN-PLAN-HANDOFF-RETRY",
column: "todo",
status: "planning",
paused: false,
log: [],
updatedAt: "2026-01-01T00:00:00.000Z",
} as unknown as Task & { recoveryRetryCount?: number | null; nextRecoveryAt?: string | null };
const recoverApprovedTriageTask = vi.fn()
.mockRejectedValueOnce(new PlanningLifecycleLockTransportError("lock transport unavailable"))
.mockResolvedValueOnce(true);
const retryingStore = createMockStore({
listTasks: vi.fn(async () => [task]),
getTask: vi.fn(async () => task),
updateTaskAtomic: vi.fn(async (_id: string, updater: (live: Task) => Partial<Task> | null) => {
const patch = updater(task);
if (patch) Object.assign(task, patch);
return task;
}),
});
const recovery = new SelfHealingManager(retryingStore, {
rootDir: "/tmp/test-project",
getPlanningTaskIds: () => new Set<string>(),
recoverApprovedTriageTask,
});
vi.setSystemTime(new Date("2026-01-01T00:05:00.000Z"));
await expect(recovery.recoverApprovedTriageTasks()).resolves.toBe(0);
expect(task.status).toBe("planning");
expect(task.recoveryRetryCount).toBe(1);
expect(Date.parse(task.nextRecoveryAt!)).toBeGreaterThan(Date.now());
expect(retryingStore.logEntry).toHaveBeenCalledWith(
task.id,
expect.stringContaining("Planning lifecycle lock transport failure during approved triage recovery — retry 1/3"),
);
await expect(recovery.recoverOrphanedPlanningTasks()).resolves.toBe(0);
expect(recoverApprovedTriageTask).toHaveBeenCalledTimes(1);
vi.setSystemTime(new Date(Date.parse(task.nextRecoveryAt!) + 1));
await expect(recovery.recoverApprovedTriageTasks()).resolves.toBe(1);
expect(recoverApprovedTriageTask).toHaveBeenCalledTimes(2);
recovery.stop();
});
it("parks a retained handoff after the planning-lock transport retry budget is exhausted", async () => {
const task = {
id: "FN-PLAN-HANDOFF-EXHAUSTED",
column: "todo",
status: "planning",
paused: false,
log: [],
recoveryRetryCount: 3,
nextRecoveryAt: "2026-01-01T00:04:00.000Z",
updatedAt: "2026-01-01T00:00:00.000Z",
} as unknown as Task;
const recoverApprovedTriageTask = vi.fn().mockRejectedValue(
new PlanningLifecycleLockTransportError("lock transport unavailable"),
);
const exhaustedStore = createMockStore({
listTasks: vi.fn(async () => [task]),
getTask: vi.fn(async () => task),
updateTaskAtomic: vi.fn(async (_id: string, updater: (live: Task) => Partial<Task> | null) => {
const patch = updater(task);
if (patch) Object.assign(task, patch);
return task;
}),
});
const recovery = new SelfHealingManager(exhaustedStore, {
rootDir: "/tmp/test-project",
getPlanningTaskIds: () => new Set<string>(),
recoverApprovedTriageTask,
});
vi.setSystemTime(new Date("2026-01-01T00:05:00.000Z"));
await expect(recovery.recoverApprovedTriageTasks()).resolves.toBe(0);
expect(task.status).toBe("failed");
expect(task.error).toContain("PLANNING_LIFECYCLE_LOCK_RECOVERY_EXHAUSTED");
expect(exhaustedStore.logEntry).toHaveBeenCalledWith(task.id, expect.stringContaining("PLANNING_LIFECYCLE_LOCK_RECOVERY_EXHAUSTED"));
recovery.stop();
});
it("clears and logs when canonical written-plan recovery returns false", async () => {
const task = {
id: "FN-PLAN-HANDOFF-NOT-RECOVERABLE",
column: "todo",
status: "planning",
paused: false,
log: [],
updatedAt: "2026-01-01T00:00:00.000Z",
} as unknown as Task;
const recoverApprovedTriageTask = vi.fn().mockResolvedValue(false);
const fallbackStore = createMockStore({
listTasks: vi.fn(async () => [task]),
getTask: vi.fn(async () => task),
});
const recovery = new SelfHealingManager(fallbackStore, {
rootDir: "/tmp/test-project",
getPlanningTaskIds: () => new Set<string>(),
recoverApprovedTriageTask,
});
vi.setSystemTime(new Date("2026-01-01T00:05:00.000Z"));
await expect(recovery.recoverApprovedTriageTasks()).resolves.toBe(0);
await expect(recovery.recoverOrphanedPlanningTasks()).resolves.toBe(1);
expect(recoverApprovedTriageTask).toHaveBeenCalledWith(task);
expect(recoverApprovedTriageTask).toHaveBeenCalledTimes(1);
expect(fallbackStore.updateTask).toHaveBeenCalledWith(task.id, { status: null });
expect(fallbackStore.logEntry).toHaveBeenCalledWith(
task.id,
"Auto-recovered orphaned planning task — agent session lost, cleared for re-planning",
);
recovery.stop();
});
it("does not log a clear when the guarded fallback loses the planning-stage race", async () => {
const task = {
id: "FN-PLAN-HANDOFF-GUARDED",
column: "todo",
status: "planning",
paused: false,
log: [],
updatedAt: "2026-01-01T00:00:00.000Z",
} as unknown as Task;
const recoverApprovedTriageTask = vi.fn().mockResolvedValue(false);
const guardedStore = createMockStore({
listTasks: vi.fn(async () => [task]),
updateTaskAtomic: vi.fn(async (_id: string, updater: (live: Task) => Partial<Task> | null) => updater({
...task,
column: "in-progress",
status: null,
firstExecutionAt: "2026-01-01T00:01:00.000Z",
} as Task)),
});
const recovery = new SelfHealingManager(guardedStore, {
rootDir: "/tmp/test-project",
getPlanningTaskIds: () => new Set<string>(),
recoverApprovedTriageTask,
});
vi.setSystemTime(new Date("2026-01-01T00:05:00.000Z"));
await expect(recovery.recoverApprovedTriageTasks()).resolves.toBe(0);
await expect(recovery.recoverOrphanedPlanningTasks()).resolves.toBe(0);
expect(recoverApprovedTriageTask).toHaveBeenCalledWith(task);
expect(guardedStore.logEntry).not.toHaveBeenCalledWith(
task.id,
"Auto-recovered orphaned planning task — agent session lost, cleared for re-planning",
);
recovery.stop();
});
it("clears status for orphaned planning tasks without a recoverable prompt", async () => { it("clears status for orphaned planning tasks without a recoverable prompt", async () => {
const getPlanning = vi.fn().mockReturnValue(new Set<string>()); const getPlanning = vi.fn().mockReturnValue(new Set<string>());

View File

@@ -80,6 +80,7 @@ function createStore(opts: {
task: Task; task: Task;
settings?: Partial<Settings>; settings?: Partial<Settings>;
moveTaskIfResult?: "moved" | "refused" | "absent"; moveTaskIfResult?: "moved" | "refused" | "absent";
specLock?: boolean;
} ): TaskStore { } ): TaskStore {
const { task } = opts; const { task } = opts;
const store: Record<string, unknown> = { const store: Record<string, unknown> = {
@@ -112,6 +113,16 @@ function createStore(opts: {
on: vi.fn(), on: vi.fn(),
off: vi.fn(), off: vi.fn(),
}; };
if (opts.specLock) {
store.isBackendMode = vi.fn(() => true);
store.withPlanningLifecycleLock = vi.fn(async (_id: string, fn: () => Promise<unknown>) => fn());
store.captureCurrentPlanEvidence = vi.fn(async () => {
throw new Error("planning lifecycle lock was reacquired");
});
store.captureCurrentPlanEvidenceWhilePlanningLocked = vi.fn().mockResolvedValue(undefined);
store.lockCurrentPlanWhilePlanningLocked = vi.fn().mockResolvedValue(undefined);
store.reconcileSpecDriftWhilePlanningLocked = vi.fn().mockResolvedValue(undefined);
}
if (opts.moveTaskIfResult !== "absent") { if (opts.moveTaskIfResult !== "absent") {
store.moveTaskIf = vi.fn(async (_id: string, column: string) => { store.moveTaskIf = vi.fn(async (_id: string, column: string) => {
if (opts.moveTaskIfResult === "refused") { if (opts.moveTaskIfResult === "refused") {
@@ -150,6 +161,19 @@ describe("planning handoff outcome — recoverApprovedTask reports what finalize
expect(store.moveTaskIf).toHaveBeenCalledWith("FN-001", "todo", expect.any(Function)); expect(store.moveTaskIf).toHaveBeenCalledWith("FN-001", "todo", expect.any(Function));
}); });
it("captures plan evidence without reacquiring the planning lifecycle lock", async () => {
const task = createTask();
const store = createStore({ task, moveTaskIfResult: "moved", specLock: true });
await expect(new TriageProcessor(store, rootDir).recoverApprovedTask(task)).resolves.toBe(true);
expect(store.captureCurrentPlanEvidenceWhilePlanningLocked).toHaveBeenCalledWith(
"FN-001",
expect.stringContaining("## Steps"),
);
expect(store.captureCurrentPlanEvidence).not.toHaveBeenCalled();
});
it("reports NOT recovered when the planning-stage guard refuses the release move (FN-8361)", async () => { it("reports NOT recovered when the planning-stage guard refuses the release move (FN-8361)", async () => {
// The symptom: `true` here makes handleStuckAbortRequeue stop, so the card // The symptom: `true` here makes handleStuckAbortRequeue stop, so the card
// keeps a finished spec in the planner column with its retry budget skipped. // keeps a finished spec in the planner column with its retry budget skipped.

View File

@@ -1,9 +1,14 @@
import { isPlanReviewSatisfied, type Task } from "@fusion/core"; import { isPlanReviewSatisfied, PlanningLifecycleLockTransportError, type Task } from "@fusion/core";
export const LEGACY_NULL_PLAN_HANDOFF_STALE_MS = 30 * 60 * 1000; export const LEGACY_NULL_PLAN_HANDOFF_STALE_MS = 30 * 60 * 1000;
export type PersistedPlanHandoffKind = "planning" | "approved-null" | "legacy-null"; export type PersistedPlanHandoffKind = "planning" | "approved-null" | "legacy-null";
export function isPlanningLifecycleLockTransportError(error: unknown): error is Error {
return error instanceof PlanningLifecycleLockTransportError
|| (error instanceof Error && error.name === "PlanningLifecycleLockTransportError");
}
/** /**
* Shared persisted-state classifier for planning handoff recovery. It deliberately * Shared persisted-state classifier for planning handoff recovery. It deliberately
* excludes graph work-item/step-instance evidence, which callers must check at * excludes graph work-item/step-instance evidence, which callers must check at
@@ -36,9 +41,14 @@ export function classifyPersistedPlanHandoff(
// a retained Plan Review approval must never make an operator-held or already- // a retained Plan Review approval must never make an operator-held or already-
// executing task eligible for planning-handoff recovery. // executing task eligible for planning-handoff recovery.
if (task.awaitingApprovalReason) return null; if (task.awaitingApprovalReason) return null;
if (task.worktree || task.firstExecutionAt || task.executionStartedAt) return null; if (task.firstExecutionAt || task.executionStartedAt) return null;
// A planning worktree belongs to the planner and may legitimately survive a
// crashed session. It must not hide a written plan from canonical handoff
// recovery. Null-status compatibility recovery remains fenced below because
// at that point a retained worktree is ambiguous execution evidence.
if (task.status === "planning") return "planning"; if (task.status === "planning") return "planning";
if (task.status != null) return null; if (task.status != null) return null;
if (task.worktree) return null;
if (task.workflowStepResults?.some(isPlanReviewSatisfied)) return "approved-null"; if (task.workflowStepResults?.some(isPlanReviewSatisfied)) return "approved-null";
if (task.approvedPlanFingerprint != null) return null; if (task.approvedPlanFingerprint != null) return null;
if (task.workflowStepResults?.length) return null; if (task.workflowStepResults?.length) return null;

View File

@@ -78,7 +78,11 @@ import { finalizeProvenAutoMergeTask, validateWorkflowDoneMergeProof } from "./m
import { AutoRecoveryDispatcher } from "./healing/auto-recovery.js"; import { AutoRecoveryDispatcher } from "./healing/auto-recovery.js";
import { activeSessionRegistry, executingTaskLock } from "./agents/active-session-registry.js"; import { activeSessionRegistry, executingTaskLock } from "./agents/active-session-registry.js";
import { isTaskStillInPlanningStage } from "./execution/replan-target.js"; import { isTaskStillInPlanningStage } from "./execution/replan-target.js";
import { classifyPersistedPlanHandoff, LEGACY_NULL_PLAN_HANDOFF_STALE_MS } from "./planning-handoff-recovery.js"; import {
classifyPersistedPlanHandoff,
isPlanningLifecycleLockTransportError,
LEGACY_NULL_PLAN_HANDOFF_STALE_MS,
} from "./planning-handoff-recovery.js";
import { getPromptPath } from "./execution/spec-staleness.js"; import { getPromptPath } from "./execution/spec-staleness.js";
import { evaluateStrandedHoldContinuation, seedPreReleasePlanReviewContinuation } from "./plan-review-continuation.js"; import { evaluateStrandedHoldContinuation, seedPreReleasePlanReviewContinuation } from "./plan-review-continuation.js";
import { evaluateStrandedContinuationReclaim, RECLAIM_RETIRED_STATE } from "./workflows/stranded-continuation-reclaim.js"; import { evaluateStrandedContinuationReclaim, RECLAIM_RETIRED_STATE } from "./workflows/stranded-continuation-reclaim.js";
@@ -123,7 +127,7 @@ import {
MAX_TERMINAL_FAILURE_AUTO_RETRIES, MAX_TERMINAL_FAILURE_AUTO_RETRIES,
TERMINAL_FAILURE_CLAIM_APPLY_GRACE_MS, TERMINAL_FAILURE_CLAIM_APPLY_GRACE_MS,
} from "@fusion/core"; } from "@fusion/core";
import { BASE_DELAY_MS, computeRecoveryDecision, MAX_DELAY_MS } from "./healing/recovery-policy.js"; import { BASE_DELAY_MS, computeRecoveryDecision, formatDelay, MAX_DELAY_MS, MAX_RECOVERY_RETRIES } from "./healing/recovery-policy.js";
export { export {
COMPLETED_BLOCKED_PAUSE_REASON, COMPLETED_BLOCKED_PAUSE_REASON,
@@ -477,6 +481,13 @@ export interface SelfHealingOptions {
} }
const APPROVED_TRIAGE_RECOVERY_GRACE_MS = 60_000; const APPROVED_TRIAGE_RECOVERY_GRACE_MS = 60_000;
function isRecoveryRetryDue(task: Pick<Task, "nextRecoveryAt">, now: number): boolean {
if (!task.nextRecoveryAt) return true;
const retryAt = Date.parse(task.nextRecoveryAt);
return !Number.isFinite(retryAt) || retryAt <= now;
}
const STARVED_REFINEMENT_RECOVERY_GRACE_MS = 10 * 60_000; const STARVED_REFINEMENT_RECOVERY_GRACE_MS = 10 * 60_000;
const STARVED_PEER_PROGRESS_THRESHOLD = 3; const STARVED_PEER_PROGRESS_THRESHOLD = 3;
const STARVED_REFINEMENT_ESCALATION_COOLDOWN_MS = STARVED_REFINEMENT_RECOVERY_GRACE_MS * 4; const STARVED_REFINEMENT_ESCALATION_COOLDOWN_MS = STARVED_REFINEMENT_RECOVERY_GRACE_MS * 4;
@@ -14901,6 +14912,50 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
* This catches the mirror-image of executor recovery: planning completed, * This catches the mirror-image of executor recovery: planning completed,
* but the final transition to `todo` / `awaiting-approval` never happened. * but the final transition to `todo` / `awaiting-approval` never happened.
*/ */
private async recordPlanningHandoffTransportFailure(task: Task, error: Error): Promise<void> {
const decision = computeRecoveryDecision({
recoveryRetryCount: task.recoveryRetryCount,
nextRecoveryAt: task.nextRecoveryAt,
});
const exhaustedMessage =
`PLANNING_LIFECYCLE_LOCK_RECOVERY_EXHAUSTED: canonical planning handoff failed after `
+ `${MAX_RECOVERY_RETRIES} retries — last error: ${error.message}`;
const patch = decision.shouldRetry
? {
status: "planning" as const,
error: null,
recoveryRetryCount: decision.nextState.recoveryRetryCount,
nextRecoveryAt: decision.nextState.nextRecoveryAt,
}
: {
status: "failed" as const,
error: exhaustedMessage,
recoveryRetryCount: null,
nextRecoveryAt: null,
};
let persisted = false;
if (typeof this.store.updateTaskAtomic === "function") {
await this.store.updateTaskAtomic(task.id, (live) => {
if (!isTaskStillInPlanningStage(live)) return null;
persisted = true;
return patch;
});
} else {
const live = await this.store.getTask(task.id);
if (live && isTaskStillInPlanningStage(live)) {
await this.store.updateTask(task.id, patch);
persisted = true;
}
}
if (!persisted) return;
const action = decision.shouldRetry
? `Planning lifecycle lock transport failure during approved triage recovery — retry ${decision.nextState.recoveryRetryCount}/${MAX_RECOVERY_RETRIES} in ${formatDelay(decision.delayMs)}: ${error.message}`
: exhaustedMessage;
await this.store.logEntry(task.id, action).catch(() => undefined);
}
async recoverApprovedTriageTasks(): Promise<number> { async recoverApprovedTriageTasks(): Promise<number> {
const recoverFn = this.options.recoverApprovedTriageTask; const recoverFn = this.options.recoverApprovedTriageTask;
if (!recoverFn) return 0; if (!recoverFn) return 0;
@@ -14934,6 +14989,7 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
requirePersistedSteps: true, requirePersistedSteps: true,
}); });
return handoffKind != null return handoffKind != null
&& isRecoveryRetryDue(task, now)
&& now - new Date(task.updatedAt).getTime() >= APPROVED_TRIAGE_RECOVERY_GRACE_MS; && now - new Date(task.updatedAt).getTime() >= APPROVED_TRIAGE_RECOVERY_GRACE_MS;
}); });
@@ -14962,8 +15018,16 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
// progress, so its exact classifier owns that one compatibility shape. // progress, so its exact classifier owns that one compatibility shape.
if (handoffKind !== "legacy-null" && !isTaskStillInPlanningStage(recoveryTask)) continue; if (handoffKind !== "legacy-null" && !isTaskStillInPlanningStage(recoveryTask)) continue;
log.log(`Recovering specified triage task ${task.id}: ${task.title || task.description?.slice(0, 60) || "(untitled)"}`); log.log(`Recovering specified triage task ${task.id}: ${task.title || task.description?.slice(0, 60) || "(untitled)"}`);
const success = await recoverFn(recoveryTask); try {
if (success) recovered++; const success = await recoverFn(recoveryTask);
if (success) recovered++;
} catch (error) {
if (isPlanningLifecycleLockTransportError(error)) {
await this.recordPlanningHandoffTransportFailure(recoveryTask, error);
continue;
}
log.warn(`Recoverable planning handoff for ${task.id} could not be finalized: ${error instanceof Error ? error.message : String(error)}`);
}
} }
if (recovered > 0) { if (recovered > 0) {
@@ -15386,16 +15450,25 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
t.status === "planning" && t.status === "planning" &&
!t.paused && !t.paused &&
!planningIds.has(t.id) && !planningIds.has(t.id) &&
isRecoveryRetryDue(t, now) &&
now - new Date(t.updatedAt).getTime() >= APPROVED_TRIAGE_RECOVERY_GRACE_MS now - new Date(t.updatedAt).getTime() >= APPROVED_TRIAGE_RECOVERY_GRACE_MS
); );
if (orphaned.length === 0) return 0; if (orphaned.length === 0) return 0;
log.warn(`Found ${orphaned.length} orphaned planning triage task(s) without a recoverable prompt`); log.warn(`Found ${orphaned.length} orphaned planning triage task(s)`);
let recovered = 0; let recovered = 0;
for (const task of orphaned) { for (const task of orphaned) {
try { try {
/*
FNXC:PlanningLifecycleLockReentry 2026-08-12-03:18:
Canonical written-plan recovery runs immediately before this fallback in
both startup and steady-state pipelines. It finalizes a retained prompt
or records bounded transport backoff, which makes that task ineligible
here. A remaining candidate has no releasable prompt, so guarded clearing
for replanning is safe and does not repeat the canonical recovery call.
*/
log.log(`Recovering orphaned planning task ${task.id}: ${task.title || task.description?.slice(0, 60) || "(untitled)"}`); log.log(`Recovering orphaned planning task ${task.id}: ${task.title || task.description?.slice(0, 60) || "(untitled)"}`);
// FNXC:Triage 2026-07-29-12:00: // FNXC:Triage 2026-07-29-12:00:
// FN-8361 closes the stale listTasks → updateTask race: only clear planning // FN-8361 closes the stale listTasks → updateTask race: only clear planning
@@ -15949,4 +16022,3 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
} }
} }
} }

View File

@@ -98,11 +98,6 @@ function getPlanningLifecycleLockTransportFailure(task: Task): PlanningLifecycle
: null; : null;
} }
function isPlanningLifecycleLockTransportError(error: unknown): error is Error {
return error instanceof fusionCore.PlanningLifecycleLockTransportError
|| (error instanceof Error && error.name === "PlanningLifecycleLockTransportError");
}
/* /*
FNXC:PlanReviewReplan 2026-08-10-18:32 (TOMBSTONE — do not re-add): FNXC:PlanReviewReplan 2026-08-10-18:32 (TOMBSTONE — do not re-add):
`PLAN_REVIEW_GATE_REPLAN_CAP = 8` is DELETED. It belonged to the out-of-graph triage Plan Review gate `PLAN_REVIEW_GATE_REPLAN_CAP = 8` is DELETED. It belonged to the out-of-graph triage Plan Review gate
@@ -159,7 +154,11 @@ import type {
} from "@earendil-works/pi-coding-agent"; } from "@earendil-works/pi-coding-agent";
import { ModelFallbackExhaustedError, describeModel, formatModelMarkerDetails, promptWithFallback } from "./pi.js"; import { ModelFallbackExhaustedError, describeModel, formatModelMarkerDetails, promptWithFallback } from "./pi.js";
import { hasAdvancedPastPlanning, isTaskStillInPlanningStage, resolvePlannerLanesForTaskAsync } from "./execution/replan-target.js"; import { hasAdvancedPastPlanning, isTaskStillInPlanningStage, resolvePlannerLanesForTaskAsync } from "./execution/replan-target.js";
import { classifyPersistedPlanHandoff, LEGACY_NULL_PLAN_HANDOFF_STALE_MS } from "./planning-handoff-recovery.js"; import {
classifyPersistedPlanHandoff,
isPlanningLifecycleLockTransportError,
LEGACY_NULL_PLAN_HANDOFF_STALE_MS,
} from "./planning-handoff-recovery.js";
import { import {
createResolvedAgentSession, createResolvedAgentSession,
extractRuntimeHint, extractRuntimeHint,
@@ -4910,11 +4909,13 @@ export class TriageProcessor {
FNXC:SpecLock 2026-08-09-07:36: FNXC:SpecLock 2026-08-09-07:36:
Planning finalization writes PROMPT.md without going through updateTask({ prompt }), so it must Planning finalization writes PROMPT.md without going through updateTask({ prompt }), so it must
capture the same canonical evidence before any approval path can release the task. This remains capture the same canonical evidence before any approval path can release the task. This remains
before the release boundary: a database/parser failure leaves the planning hold intact. before the release boundary: a database/parser failure leaves the planning hold intact. The
finalizer already owns the planning lifecycle lock, so evidence capture must use the lock-assuming
seam instead of reacquiring the non-reentrant PostgreSQL advisory lock.
*/ */
const supportsSpecLock = (this.store as unknown as { isBackendMode?: () => boolean }).isBackendMode?.() === true; const supportsSpecLock = (this.store as unknown as { isBackendMode?: () => boolean }).isBackendMode?.() === true;
if (supportsSpecLock) { if (supportsSpecLock) {
await this.store.captureCurrentPlanEvidence(task.id, written); await this.store.captureCurrentPlanEvidenceWhilePlanningLocked(task.id, written);
} }
let taskIntentSignature: ReturnType<typeof extractIntentSignature> = { let taskIntentSignature: ReturnType<typeof extractIntentSignature> = {