fix(FN-8998): prevent planning lock reentry stalls
This commit is contained in:
7
.changeset/fix-planning-lock-reentry.md
Normal file
7
.changeset/fix-planning-lock-reentry.md
Normal 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.
|
||||
@@ -54,4 +54,21 @@ describe("classifyPersistedPlanHandoff", () => {
|
||||
awaitingApprovalReason: "require-all",
|
||||
}), 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();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -23,6 +23,7 @@ import {
|
||||
import { getPromptPath } from "../../execution/spec-staleness.js";
|
||||
import { promoteHeldTask, runHoldReleaseSweep } from "../../execution/hold-release.js";
|
||||
import { SelfHealingManager } from "../../self-healing.js";
|
||||
import { TriageProcessor } from "../../triage.js";
|
||||
|
||||
const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({
|
||||
prefix: "fusion_planning_dependency_release",
|
||||
@@ -78,6 +79,58 @@ pgDescribe("FN-8768 planning dependency release interactions", () => {
|
||||
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([
|
||||
{
|
||||
label: "dependency mutation API",
|
||||
|
||||
@@ -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 { 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 { execSync } from "node:child_process";
|
||||
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 () => {
|
||||
const getPlanning = vi.fn().mockReturnValue(new Set<string>());
|
||||
|
||||
|
||||
@@ -80,6 +80,7 @@ function createStore(opts: {
|
||||
task: Task;
|
||||
settings?: Partial<Settings>;
|
||||
moveTaskIfResult?: "moved" | "refused" | "absent";
|
||||
specLock?: boolean;
|
||||
} ): TaskStore {
|
||||
const { task } = opts;
|
||||
const store: Record<string, unknown> = {
|
||||
@@ -112,6 +113,16 @@ function createStore(opts: {
|
||||
on: 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") {
|
||||
store.moveTaskIf = vi.fn(async (_id: string, column: string) => {
|
||||
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));
|
||||
});
|
||||
|
||||
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 () => {
|
||||
// The symptom: `true` here makes handleStuckAbortRequeue stop, so the card
|
||||
// keeps a finished spec in the planner column with its retry budget skipped.
|
||||
|
||||
@@ -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 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
|
||||
* 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-
|
||||
// executing task eligible for planning-handoff recovery.
|
||||
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 != null) return null;
|
||||
if (task.worktree) return null;
|
||||
if (task.workflowStepResults?.some(isPlanReviewSatisfied)) return "approved-null";
|
||||
if (task.approvedPlanFingerprint != null) return null;
|
||||
if (task.workflowStepResults?.length) return null;
|
||||
|
||||
@@ -78,7 +78,11 @@ import { finalizeProvenAutoMergeTask, validateWorkflowDoneMergeProof } from "./m
|
||||
import { AutoRecoveryDispatcher } from "./healing/auto-recovery.js";
|
||||
import { activeSessionRegistry, executingTaskLock } from "./agents/active-session-registry.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 { evaluateStrandedHoldContinuation, seedPreReleasePlanReviewContinuation } from "./plan-review-continuation.js";
|
||||
import { evaluateStrandedContinuationReclaim, RECLAIM_RETIRED_STATE } from "./workflows/stranded-continuation-reclaim.js";
|
||||
@@ -123,7 +127,7 @@ import {
|
||||
MAX_TERMINAL_FAILURE_AUTO_RETRIES,
|
||||
TERMINAL_FAILURE_CLAIM_APPLY_GRACE_MS,
|
||||
} 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 {
|
||||
COMPLETED_BLOCKED_PAUSE_REASON,
|
||||
@@ -477,6 +481,13 @@ export interface SelfHealingOptions {
|
||||
}
|
||||
|
||||
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_PEER_PROGRESS_THRESHOLD = 3;
|
||||
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,
|
||||
* 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> {
|
||||
const recoverFn = this.options.recoverApprovedTriageTask;
|
||||
if (!recoverFn) return 0;
|
||||
@@ -14934,6 +14989,7 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
|
||||
requirePersistedSteps: true,
|
||||
});
|
||||
return handoffKind != null
|
||||
&& isRecoveryRetryDue(task, now)
|
||||
&& 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.
|
||||
if (handoffKind !== "legacy-null" && !isTaskStillInPlanningStage(recoveryTask)) continue;
|
||||
log.log(`Recovering specified triage task ${task.id}: ${task.title || task.description?.slice(0, 60) || "(untitled)"}`);
|
||||
const success = await recoverFn(recoveryTask);
|
||||
if (success) recovered++;
|
||||
try {
|
||||
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) {
|
||||
@@ -15386,16 +15450,25 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
|
||||
t.status === "planning" &&
|
||||
!t.paused &&
|
||||
!planningIds.has(t.id) &&
|
||||
isRecoveryRetryDue(t, now) &&
|
||||
now - new Date(t.updatedAt).getTime() >= APPROVED_TRIAGE_RECOVERY_GRACE_MS
|
||||
);
|
||||
|
||||
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;
|
||||
for (const task of orphaned) {
|
||||
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)"}`);
|
||||
// FNXC:Triage 2026-07-29-12:00:
|
||||
// FN-8361 closes the stale listTasks → updateTask race: only clear planning
|
||||
@@ -15949,4 +16022,3 @@ const movedTask = await this.store.moveTask(task.id, completeLane);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -98,11 +98,6 @@ function getPlanningLifecycleLockTransportFailure(task: Task): PlanningLifecycle
|
||||
: 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):
|
||||
`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";
|
||||
import { ModelFallbackExhaustedError, describeModel, formatModelMarkerDetails, promptWithFallback } from "./pi.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 {
|
||||
createResolvedAgentSession,
|
||||
extractRuntimeHint,
|
||||
@@ -4910,11 +4909,13 @@ export class TriageProcessor {
|
||||
FNXC:SpecLock 2026-08-09-07:36:
|
||||
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
|
||||
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;
|
||||
if (supportsSpecLock) {
|
||||
await this.store.captureCurrentPlanEvidence(task.id, written);
|
||||
await this.store.captureCurrentPlanEvidenceWhilePlanningLocked(task.id, written);
|
||||
}
|
||||
|
||||
let taskIntentSignature: ReturnType<typeof extractIntentSignature> = {
|
||||
|
||||
Reference in New Issue
Block a user