fix(engine): run every lane in the task worktree; contention is a wait, not a failure
Contention prevention (why tasks shared a path at all): - Planning ran `tools: "coding"` at the repo root, so every planner had write tools in the operator's checkout and all planners shared one path. Planning now acquires the task's own worktree (TriageProcessor.acquirePlanningWorktree -> TaskExecutor.ensureTaskWorktreeForPlanning). - Graph nodes with no worktree acquired one instead of falling back to rootDir, so Plan Review / Code Review / custom gates all run isolated. Plan Review re-acquires when its recorded worktree is gone, replacing FN-7996's run-from-the-repo-root degrade. Workspace projects are unchanged. - Registration goes through acquireActiveSessionPath, which reclaims a leaked entry whose holder is provably dead and aged past the FN-5256 floor. A live holder still contends — real serialization is never clobbered. Classification (the reported symptom): - A lease held by another task is no longer a provider failure. It carries SESSION_CONTENTION_HOLD_VALUE, classifies transient, is excluded from isNonPlanDefectPlanReviewFailure, and stops burning the node's fast retries. - The executor waits it out on a 10-attempt 5s->60s ladder and then leaves the task cleanly queued. There is no terminal branch: contention always ends, so parking would only ask a human to press Retry on a condition that fixed itself. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
7
.changeset/session-contention-worktree-isolation.md
Normal file
7
.changeset/session-contention-worktree-isolation.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": minor
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Planning and every review step now run in the task's own worktree, never the shared checkout.
|
||||||
|
category: fix
|
||||||
|
dev: Planning acquires the task worktree (TriageProcessor `acquirePlanningWorktree` → `TaskExecutor.ensureTaskWorktreeForPlanning`); graph nodes with no worktree acquire one instead of falling back to `rootDir`, and Plan Review re-acquires when its recorded worktree is gone (replacing the FN-7996 repo-root degrade). Registration goes through `acquireActiveSessionPath`, which reclaims a leaked entry whose holder is provably dead and aged past the FN-5256 floor. Remaining contention gets `SESSION_CONTENTION_HOLD_VALUE`: `isSessionContentionError` classifies it transient, `isNonPlanDefectPlanReviewFailure` explicitly excludes it, and the executor waits on a 10-attempt 5s→60s ladder that ends in a benign requeue, never a park.
|
||||||
131
packages/engine/src/__tests__/node-worktree-isolation.test.ts
Normal file
131
packages/engine/src/__tests__/node-worktree-isolation.test.ts
Normal file
@@ -0,0 +1,131 @@
|
|||||||
|
/*
|
||||||
|
FNXC:NodeWorktreeIsolation 2026-07-25-22:10 (no lane runs in the shared checkout — regression):
|
||||||
|
Operator requirement: Plan Review, Code Review, and every other node run in the TASK-SPECIFIC worktree;
|
||||||
|
the shared main checkout is for merge only. Before this, read-only graph gates fell back to
|
||||||
|
`this.rootDir` because a pre-execution task has no worktree yet. That is what let two tasks share one
|
||||||
|
path (the reported FN-1398/FN-1403 Plan Review session collision) and what let reviewers read a checkout
|
||||||
|
that other tasks and the operator mutate underneath them.
|
||||||
|
|
||||||
|
Invariant under test across the node surfaces that previously degraded to the root:
|
||||||
|
- Plan Review (no worktree yet) acquires and runs in a task worktree;
|
||||||
|
- a custom read-only gate (no worktree yet) does the same — this is not Plan-Review-special;
|
||||||
|
- an existing usable worktree is REUSED, not re-acquired;
|
||||||
|
- the acquisition is skipped for workspace projects, whose sessions are browse-root-rooted by design.
|
||||||
|
*/
|
||||||
|
import { describe, expect, it, vi, beforeEach } from "vitest";
|
||||||
|
import type { TaskDetail } from "@fusion/core";
|
||||||
|
import "./executor-test-helpers.js";
|
||||||
|
import { TaskExecutor } from "../executor.js";
|
||||||
|
import {
|
||||||
|
createMockStore,
|
||||||
|
mockedExecSync,
|
||||||
|
mockedExistsSync,
|
||||||
|
resetExecutorMocks,
|
||||||
|
} from "./executor-test-helpers.js";
|
||||||
|
|
||||||
|
const ROOT = "/tmp/test";
|
||||||
|
|
||||||
|
function makeTask(overrides: Partial<TaskDetail> = {}): TaskDetail {
|
||||||
|
const now = new Date().toISOString();
|
||||||
|
return {
|
||||||
|
id: "FN-1403",
|
||||||
|
title: "Isolation",
|
||||||
|
description: "Desc",
|
||||||
|
column: "todo",
|
||||||
|
dependencies: [],
|
||||||
|
steps: [],
|
||||||
|
currentStep: 0,
|
||||||
|
log: [],
|
||||||
|
worktree: undefined,
|
||||||
|
branch: undefined,
|
||||||
|
status: null,
|
||||||
|
error: null,
|
||||||
|
paused: false,
|
||||||
|
userPaused: false,
|
||||||
|
createdAt: now,
|
||||||
|
updatedAt: now,
|
||||||
|
...overrides,
|
||||||
|
} as TaskDetail;
|
||||||
|
}
|
||||||
|
|
||||||
|
const PLAN_REVIEW_NODE = {
|
||||||
|
id: "plan-review-step",
|
||||||
|
kind: "prompt",
|
||||||
|
config: { name: "Plan Review", prompt: "Review the plan.", toolMode: "readonly" },
|
||||||
|
};
|
||||||
|
const CUSTOM_READONLY_GATE = {
|
||||||
|
id: "custom-gate",
|
||||||
|
kind: "prompt",
|
||||||
|
config: { name: "Custom Gate", prompt: "Check something.", toolMode: "readonly" },
|
||||||
|
};
|
||||||
|
|
||||||
|
describe("every workflow node runs in the task worktree, never the shared checkout", () => {
|
||||||
|
beforeEach(() => {
|
||||||
|
resetExecutorMocks();
|
||||||
|
mockedExecSync.mockReturnValue("" as any);
|
||||||
|
});
|
||||||
|
|
||||||
|
it.each([
|
||||||
|
["Plan Review", PLAN_REVIEW_NODE],
|
||||||
|
["a custom read-only gate", CUSTOM_READONLY_GATE],
|
||||||
|
])("acquires a task worktree for %s when the task has none", async (_label, node) => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const executor = new TaskExecutor(store, ROOT);
|
||||||
|
mockedExistsSync.mockReturnValue(true);
|
||||||
|
|
||||||
|
const captured: { worktreePath?: string } = {};
|
||||||
|
vi.spyOn(executor as any, "executeWorkflowStep").mockImplementation(async (...args: any[]) => {
|
||||||
|
captured.worktreePath = args[2];
|
||||||
|
return { success: true, output: "APPROVE" };
|
||||||
|
});
|
||||||
|
|
||||||
|
const live = makeTask();
|
||||||
|
store.getTask.mockResolvedValue(live as any);
|
||||||
|
await (executor as any).runGraphCustomNode(node, live, { reviewerInlineFixes: false }, undefined);
|
||||||
|
|
||||||
|
expect(captured.worktreePath).not.toBe(ROOT);
|
||||||
|
expect(captured.worktreePath).toContain(`${ROOT}/.worktrees/`);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("reuses an existing usable worktree instead of acquiring another", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const executor = new TaskExecutor(store, ROOT);
|
||||||
|
const existing = `${ROOT}/.worktrees/existing`;
|
||||||
|
mockedExistsSync.mockReturnValue(true);
|
||||||
|
|
||||||
|
const acquireSpy = vi.spyOn(executor as any, "ensureGraphCustomNodeWorktree");
|
||||||
|
const captured: { worktreePath?: string } = {};
|
||||||
|
vi.spyOn(executor as any, "executeWorkflowStep").mockImplementation(async (...args: any[]) => {
|
||||||
|
captured.worktreePath = args[2];
|
||||||
|
return { success: true, output: "APPROVE" };
|
||||||
|
});
|
||||||
|
|
||||||
|
const live = makeTask({ worktree: existing, branch: "fusion/fn-1403" });
|
||||||
|
store.getTask.mockResolvedValue(live as any);
|
||||||
|
await (executor as any).runGraphCustomNode(PLAN_REVIEW_NODE, live, {}, undefined);
|
||||||
|
|
||||||
|
expect(captured.worktreePath).toBe(existing);
|
||||||
|
expect(acquireSpy).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("leaves workspace projects on the shared browse-root (per-repo isolation is the sub-repo lease)", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const executor = new TaskExecutor(store, ROOT);
|
||||||
|
(executor as any).workspaceConfig = { repos: ["apps/web"] };
|
||||||
|
mockedExistsSync.mockReturnValue(true);
|
||||||
|
|
||||||
|
const acquireSpy = vi.spyOn(executor as any, "ensureGraphCustomNodeWorktree");
|
||||||
|
const captured: { worktreePath?: string } = {};
|
||||||
|
vi.spyOn(executor as any, "executeWorkflowStep").mockImplementation(async (...args: any[]) => {
|
||||||
|
captured.worktreePath = args[2];
|
||||||
|
return { success: true, output: "APPROVE" };
|
||||||
|
});
|
||||||
|
|
||||||
|
const live = makeTask();
|
||||||
|
store.getTask.mockResolvedValue(live as any);
|
||||||
|
await (executor as any).runGraphCustomNode(PLAN_REVIEW_NODE, live, {}, undefined);
|
||||||
|
|
||||||
|
expect(captured.worktreePath).toBe(ROOT);
|
||||||
|
expect(acquireSpy).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -296,7 +296,14 @@ describe("Plan Review missing-worktree repo-root fallback (FN-7996)", () => {
|
|||||||
mockedExecSync.mockReturnValue("" as any);
|
mockedExecSync.mockReturnValue("" as any);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("runs the Plan Review reviewer from the repo root when the recorded worktree is gone", async () => {
|
/*
|
||||||
|
FNXC:NodeWorktreeIsolation 2026-07-25-22:10:
|
||||||
|
FN-7996's invariant is unchanged — a missing recorded worktree must never terminal-park Plan Review —
|
||||||
|
but the remedy is no longer "run in the shared repo root". Every lane now runs in the TASK's own
|
||||||
|
worktree, so the reviewer RE-ACQUIRES one. The assertion below is the same symptom (stale path gone,
|
||||||
|
review still runs) with the shared checkout removed as an outcome.
|
||||||
|
*/
|
||||||
|
it("re-acquires a task worktree for Plan Review when the recorded worktree is gone (never the repo root)", async () => {
|
||||||
const store = createMockStore();
|
const store = createMockStore();
|
||||||
const executor = new TaskExecutor(store, "/tmp/test");
|
const executor = new TaskExecutor(store, "/tmp/test");
|
||||||
mockedExistsSync.mockImplementation((path: unknown) => path !== "/tmp/stale-wt");
|
mockedExistsSync.mockImplementation((path: unknown) => path !== "/tmp/stale-wt");
|
||||||
@@ -317,10 +324,13 @@ describe("Plan Review missing-worktree repo-root fallback (FN-7996)", () => {
|
|||||||
const result = await (executor as any).runGraphCustomNode(node, live, {}, undefined);
|
const result = await (executor as any).runGraphCustomNode(node, live, {}, undefined);
|
||||||
|
|
||||||
expect(result.outcome).toBe("success");
|
expect(result.outcome).toBe("success");
|
||||||
expect(captured.worktreePath).toBe("/tmp/test");
|
// Not the stale path, and — the point of the change — not the shared repo root either.
|
||||||
|
expect(captured.worktreePath).not.toBe("/tmp/stale-wt");
|
||||||
|
expect(captured.worktreePath).not.toBe("/tmp/test");
|
||||||
|
expect(captured.worktreePath).toContain("/tmp/test/.worktrees/");
|
||||||
expect(store.logEntry).toHaveBeenCalledWith(
|
expect(store.logEntry).toHaveBeenCalledWith(
|
||||||
live.id,
|
live.id,
|
||||||
expect.stringContaining("running the reviewer from the repo root"),
|
expect.stringContaining("re-acquiring a task worktree instead of running in the shared checkout"),
|
||||||
undefined,
|
undefined,
|
||||||
undefined,
|
undefined,
|
||||||
);
|
);
|
||||||
|
|||||||
@@ -0,0 +1,137 @@
|
|||||||
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30 (contention is never a provider failure — regression):
|
||||||
|
Reported symptom: FN-1403's Plan Review hit "active-session path /home/ubuntu/dev/freemap-svelte is held
|
||||||
|
by task FN-1398; task FN-1403 may not overwrite it". The engine narrated that as "Plan Review provider
|
||||||
|
failure — retrying in place (2/2)", spent the provider budget on two immediate retries that could not
|
||||||
|
possibly clear another task's hold, then left the task parked.
|
||||||
|
|
||||||
|
Invariant under test, across every surface a contention error can reach:
|
||||||
|
1. the pure classifier recognizes all THREE contention shapes (active-session registry, workspace
|
||||||
|
sub-repo acquire lease, workspace sub-repo land lease) and treats them as transient;
|
||||||
|
2. the classifier does NOT swallow unrelated failures (negative control);
|
||||||
|
3. a contention error is NOT classified as a non-plan-defect PROVIDER failure by the Plan Review
|
||||||
|
classifier — that is the exact misrouting that burned the budget;
|
||||||
|
4. registration PREVENTS contention where it is preventable: a leaked entry whose holder is dead is
|
||||||
|
reclaimed, while a LIVE holder still contends (real serialization must not be clobbered), and a
|
||||||
|
too-fresh entry is never reclaimed even if the probe says dead (FN-5256 warming floor).
|
||||||
|
*/
|
||||||
|
import { describe, expect, it } from "vitest";
|
||||||
|
import {
|
||||||
|
isSessionContentionError,
|
||||||
|
isTransientError,
|
||||||
|
} from "../transient-error-patterns.js";
|
||||||
|
import { isNonPlanDefectPlanReviewFailure } from "../transient-error-detector.js";
|
||||||
|
import {
|
||||||
|
ActiveSessionRegistry,
|
||||||
|
acquireActiveSessionPath,
|
||||||
|
} from "../active-session-registry.js";
|
||||||
|
|
||||||
|
const REPORTED_MESSAGE =
|
||||||
|
"active-session path /home/ubuntu/dev/freemap-svelte is held by task FN-1398; task FN-1403 may not overwrite it";
|
||||||
|
const ACQUIRE_BUSY_MESSAGE = "workspace sub-repo apps/web acquisition is in progress for task FN-1398";
|
||||||
|
const LAND_BUSY_MESSAGE = "workspace sub-repo apps/web land is in progress for task FN-1398";
|
||||||
|
|
||||||
|
describe("session contention classification", () => {
|
||||||
|
const contentionShapes = [
|
||||||
|
["active-session registry (the reported failure)", REPORTED_MESSAGE],
|
||||||
|
["workspace sub-repo acquire lease", ACQUIRE_BUSY_MESSAGE],
|
||||||
|
["workspace sub-repo land lease", LAND_BUSY_MESSAGE],
|
||||||
|
] as const;
|
||||||
|
|
||||||
|
for (const [label, message] of contentionShapes) {
|
||||||
|
it(`classifies ${label} as contention and as transient`, () => {
|
||||||
|
expect(isSessionContentionError(message)).toBe(true);
|
||||||
|
// Transient is what makes every generic retry path wait instead of parking the task.
|
||||||
|
expect(isTransientError(message)).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it(`does not route ${label} into the Plan Review PROVIDER-failure hold`, () => {
|
||||||
|
/*
|
||||||
|
The provider hold gives two immediate retries and then parks. Contention must be recognized
|
||||||
|
before that classifier ever sees it, so the graph publishes SESSION_CONTENTION_HOLD_VALUE and the
|
||||||
|
executor waits on a backoff ladder instead.
|
||||||
|
*/
|
||||||
|
expect(
|
||||||
|
isNonPlanDefectPlanReviewFailure({ failureValue: "session-contention-hold", errorMessage: message }),
|
||||||
|
).toBe(false);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
it("does not classify unrelated failures as contention", () => {
|
||||||
|
for (const message of [
|
||||||
|
"Plan Review returned malformed JSON",
|
||||||
|
"ENOENT: no such file or directory, open '/repo/PROMPT.md'",
|
||||||
|
"task FN-1398 is soft-deleted (deletedAt=2026-07-25) and cannot be read or mutated",
|
||||||
|
]) {
|
||||||
|
expect(isSessionContentionError(message)).toBe(false);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
describe("contention prevention at the registration seam", () => {
|
||||||
|
const OTHER = { taskId: "FN-1403", kind: "workflow-step" as const, ownerKey: "FN-1403#workflow-step" };
|
||||||
|
const PATH = "/repo/.worktrees/shared";
|
||||||
|
|
||||||
|
function registryHeldBy(taskId: string, registeredAtOffsetMs: number, now: number): ActiveSessionRegistry {
|
||||||
|
const registry = new ActiveSessionRegistry();
|
||||||
|
registry.registerPath(PATH, { taskId, kind: "workflow-step", ownerKey: `${taskId}#workflow-step` });
|
||||||
|
// Rewind registeredAt so age gating can be exercised without real time.
|
||||||
|
(registry.lookupByPath(PATH) as { registeredAt: number }).registeredAt = now - registeredAtOffsetMs;
|
||||||
|
return registry;
|
||||||
|
}
|
||||||
|
|
||||||
|
it("reclaims a leaked entry whose holder is provably dead", () => {
|
||||||
|
const now = 1_000_000;
|
||||||
|
const registry = registryHeldBy("FN-1398", 60_000, now);
|
||||||
|
|
||||||
|
const outcome = acquireActiveSessionPath(registry, PATH, OTHER, {
|
||||||
|
holderLiveProbe: () => false,
|
||||||
|
now: () => now,
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(outcome.action).toBe("reclaimed-stale-foreign");
|
||||||
|
expect(registry.lookupByPath(PATH)?.taskId).toBe("FN-1403");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("contends (never clobbers) when the holder is still live", () => {
|
||||||
|
const now = 1_000_000;
|
||||||
|
const registry = registryHeldBy("FN-1398", 60_000, now);
|
||||||
|
|
||||||
|
const outcome = acquireActiveSessionPath(registry, PATH, OTHER, {
|
||||||
|
holderLiveProbe: () => true,
|
||||||
|
now: () => now,
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(outcome).toMatchObject({ action: "contended", holderTaskId: "FN-1398" });
|
||||||
|
expect(registry.lookupByPath(PATH)?.taskId).toBe("FN-1398");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("refuses to reclaim a freshly-registered entry even when the probe reports dead (warming floor)", () => {
|
||||||
|
const now = 1_000_000;
|
||||||
|
const registry = registryHeldBy("FN-1398", 10, now);
|
||||||
|
|
||||||
|
const outcome = acquireActiveSessionPath(registry, PATH, OTHER, {
|
||||||
|
holderLiveProbe: () => false,
|
||||||
|
now: () => now,
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(outcome.action).toBe("contended");
|
||||||
|
expect(registry.lookupByPath(PATH)?.taskId).toBe("FN-1398");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("treats an absent probe as LIVE so ambiguity never clobbers a holder", () => {
|
||||||
|
const now = 1_000_000;
|
||||||
|
const registry = registryHeldBy("FN-1398", 60_000, now);
|
||||||
|
|
||||||
|
const outcome = acquireActiveSessionPath(registry, PATH, OTHER, { now: () => now });
|
||||||
|
|
||||||
|
expect(outcome.action).toBe("contended");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("registers normally on a free path and stays idempotent for the same task", () => {
|
||||||
|
const registry = new ActiveSessionRegistry();
|
||||||
|
expect(acquireActiveSessionPath(registry, PATH, OTHER).action).toBe("registered");
|
||||||
|
expect(acquireActiveSessionPath(registry, PATH, OTHER).action).toBe("registered");
|
||||||
|
expect(registry.pathsForTask("FN-1403")).toEqual([PATH]);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -181,6 +181,63 @@ export interface SelfOwnedReconcileOptions {
|
|||||||
now?: () => number;
|
now?: () => number;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30 (contention prevention — foreign-stale reclaim):
|
||||||
|
Registration contention has exactly three shapes, and only one of them is legitimate:
|
||||||
|
1. Two tasks on the SHARED repo root. Not real contention — read-only root-rooted sessions need no
|
||||||
|
path exclusivity. Eliminated by construction: `sessionRegistryPath` task-scopes the root key.
|
||||||
|
2. A LEAKED entry whose owning task is dead (crashed run, torn-down executor, engine restart that
|
||||||
|
lost the session but not the map). Waiting for that holder is waiting forever — the holder will
|
||||||
|
never release. Reclaim it, which is what this seam does.
|
||||||
|
3. A LIVE holder on the same path. This is genuine serialization (the workspace sub-repo leases are
|
||||||
|
built on it) and the caller must wait, never overwrite.
|
||||||
|
So: probe the holder for liveness, reclaim when it is provably dead AND the entry has aged past the
|
||||||
|
FN-5256 staleness floor (a just-registered entry belongs to a warming session whose maps are not
|
||||||
|
populated yet — treating it as dead would yank a live shell), and surface a typed contention error only
|
||||||
|
for case 3. `holderLiveProbe` returning true is always respected; an unknown/throwing probe must be
|
||||||
|
reported as LIVE by its caller so ambiguity refuses the reclaim.
|
||||||
|
*/
|
||||||
|
export type ForeignHolderLiveProbe = (holderTaskId: string, path: string) => boolean;
|
||||||
|
|
||||||
|
export type AcquireActiveSessionPathOutcome =
|
||||||
|
| { action: "registered" }
|
||||||
|
| { action: "reclaimed-stale-foreign"; holderTaskId: string; ageMs: number }
|
||||||
|
| { action: "contended"; holderTaskId: string; holderKind: ActiveSessionKind; ageMs: number };
|
||||||
|
|
||||||
|
export interface AcquireActiveSessionPathOptions {
|
||||||
|
/** Returns true when the foreign holder still has a live session/execution surface. */
|
||||||
|
holderLiveProbe?: ForeignHolderLiveProbe;
|
||||||
|
/** Minimum entry age before a foreign entry may be reclaimed. Defaults to `DEFAULT_SELF_OWNED_MIN_IDLE_MS`. */
|
||||||
|
minIdleMs?: number;
|
||||||
|
/** Test seam — defaults to `Date.now()`. */
|
||||||
|
now?: () => number;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function acquireActiveSessionPath(
|
||||||
|
registry: ActiveSessionRegistry,
|
||||||
|
path: string,
|
||||||
|
registration: ActiveSessionRegistration,
|
||||||
|
options: AcquireActiveSessionPathOptions = {},
|
||||||
|
): AcquireActiveSessionPathOutcome {
|
||||||
|
const existing = registry.lookupByPath(path);
|
||||||
|
if (!existing || existing.taskId === registration.taskId) {
|
||||||
|
registry.registerPath(path, registration);
|
||||||
|
return { action: "registered" };
|
||||||
|
}
|
||||||
|
|
||||||
|
const now = options.now?.() ?? Date.now();
|
||||||
|
const ageMs = now - existing.registeredAt;
|
||||||
|
const minIdleMs = options.minIdleMs ?? DEFAULT_SELF_OWNED_MIN_IDLE_MS;
|
||||||
|
const holderIsLive = options.holderLiveProbe?.(existing.taskId, path) ?? true;
|
||||||
|
if (holderIsLive || ageMs < minIdleMs) {
|
||||||
|
return { action: "contended", holderTaskId: existing.taskId, holderKind: existing.kind, ageMs };
|
||||||
|
}
|
||||||
|
|
||||||
|
registry.unregisterPath(path);
|
||||||
|
registry.registerPath(path, registration);
|
||||||
|
return { action: "reclaimed-stale-foreign", holderTaskId: existing.taskId, ageMs };
|
||||||
|
}
|
||||||
|
|
||||||
export function reconcileSelfOwnedActiveSessionForRemoval(
|
export function reconcileSelfOwnedActiveSessionForRemoval(
|
||||||
registry: ActiveSessionRegistry,
|
registry: ActiveSessionRegistry,
|
||||||
worktreePath: string,
|
worktreePath: string,
|
||||||
|
|||||||
@@ -44,6 +44,7 @@ import {
|
|||||||
import {
|
import {
|
||||||
MERGE_REGION_KINDS,
|
MERGE_REGION_KINDS,
|
||||||
PLAN_REVIEW_PROVIDER_FAILURE_HOLD_VALUE,
|
PLAN_REVIEW_PROVIDER_FAILURE_HOLD_VALUE,
|
||||||
|
SESSION_CONTENTION_HOLD_VALUE,
|
||||||
WORKFLOW_DRIFT_PARK_CONTEXT_KEY,
|
WORKFLOW_DRIFT_PARK_CONTEXT_KEY,
|
||||||
WORKFLOW_NODE_ENGINE_PAUSE_ABORT_KIND,
|
WORKFLOW_NODE_ENGINE_PAUSE_ABORT_KIND,
|
||||||
WORKFLOW_OPTIONAL_GROUP_CONTEXT_KEY,
|
WORKFLOW_OPTIONAL_GROUP_CONTEXT_KEY,
|
||||||
@@ -133,9 +134,12 @@ import { attemptBranchAutocorrect } from "./branch-autocorrect.js";
|
|||||||
import { ActiveSessionWorktreeRemovalError } from "./worktree-backend.js";
|
import { ActiveSessionWorktreeRemovalError } from "./worktree-backend.js";
|
||||||
import {canonicalizeWorktreePath, registerArchiveWorkspaceWorktreeDisposer, registerArchiveWorktreeDisposer, registerTaskMoveDisposer} from "@fusion/core";
|
import {canonicalizeWorktreePath, registerArchiveWorkspaceWorktreeDisposer, registerArchiveWorktreeDisposer, registerTaskMoveDisposer} from "@fusion/core";
|
||||||
import {
|
import {
|
||||||
|
ActiveSessionPathHeldByForeignTaskError,
|
||||||
|
acquireActiveSessionPath,
|
||||||
activeSessionRegistry,
|
activeSessionRegistry,
|
||||||
executingTaskLock,
|
executingTaskLock,
|
||||||
reconcileSelfOwnedActiveSessionForRemoval,
|
reconcileSelfOwnedActiveSessionForRemoval,
|
||||||
|
type ActiveSessionKind,
|
||||||
} from "./active-session-registry.js";
|
} from "./active-session-registry.js";
|
||||||
// CLI Agent Executor (U7): task ↔ CLI session orchestration seam.
|
// CLI Agent Executor (U7): task ↔ CLI session orchestration seam.
|
||||||
import {
|
import {
|
||||||
@@ -181,7 +185,7 @@ import { AgentLogger } from "./agent-logger.js";
|
|||||||
import { createLogger, executorLog, reviewerLog, formatError } from "./logger.js";
|
import { createLogger, executorLog, reviewerLog, formatError } from "./logger.js";
|
||||||
import { TokenCapDetector } from "./token-cap-detector.js";
|
import { TokenCapDetector } from "./token-cap-detector.js";
|
||||||
import { isUsageLimitError, checkSessionError, type UsageLimitPauser } from "./usage-limit-detector.js";
|
import { isUsageLimitError, checkSessionError, type UsageLimitPauser } from "./usage-limit-detector.js";
|
||||||
import { isNonContinuableSessionError, isNonPlanDefectPlanReviewFailure, isTransientError, isSilentTransientError } from "./transient-error-detector.js";
|
import { isNonContinuableSessionError, isNonPlanDefectPlanReviewFailure, isSessionContentionError, isTransientError, isSilentTransientError } from "./transient-error-detector.js";
|
||||||
import { withRateLimitRetry } from "./rate-limit-retry.js";
|
import { withRateLimitRetry } from "./rate-limit-retry.js";
|
||||||
import {
|
import {
|
||||||
detectExternalIntegrationEvidenceGaps,
|
detectExternalIntegrationEvidenceGaps,
|
||||||
@@ -514,6 +518,16 @@ import {
|
|||||||
|
|
||||||
const MAX_TRANSIENT_GRAPH_RESUME_RETRIES = 2;
|
const MAX_TRANSIENT_GRAPH_RESUME_RETRIES = 2;
|
||||||
const TRANSIENT_GRAPH_RESUME_RETRY_BACKOFF_MS = process.env.VITEST || process.env.NODE_ENV === "test" ? 0 : 1_000;
|
const TRANSIENT_GRAPH_RESUME_RETRY_BACKOFF_MS = process.env.VITEST || process.env.NODE_ENV === "test" ? 0 : 1_000;
|
||||||
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30:
|
||||||
|
The contention ladder is deliberately long and slow compared with the provider-failure budget (2 fast
|
||||||
|
retries): a lease is held for as long as the holder's own work takes — minutes, not milliseconds. Ten
|
||||||
|
attempts backing off 5s→60s covers ~8 minutes of waiting, after which the task is left queued for
|
||||||
|
ordinary re-dispatch rather than parked.
|
||||||
|
*/
|
||||||
|
const MAX_SESSION_CONTENTION_HOLD_RETRIES = 10;
|
||||||
|
const SESSION_CONTENTION_HOLD_BACKOFF_MS = process.env.VITEST || process.env.NODE_ENV === "test" ? 0 : 5_000;
|
||||||
|
const SESSION_CONTENTION_HOLD_MAX_BACKOFF_MS = 60_000;
|
||||||
/** How long to wait before recovering a completed task still stuck in in-progress. */
|
/** How long to wait before recovering a completed task still stuck in in-progress. */
|
||||||
const COMPLETED_TASK_WATCHDOG_MS = 60_000;
|
const COMPLETED_TASK_WATCHDOG_MS = 60_000;
|
||||||
/** How long to wait before retrying a workflow rerun handoff that never reached in-progress. */
|
/** How long to wait before retrying a workflow rerun handoff that never reached in-progress. */
|
||||||
@@ -1995,9 +2009,43 @@ export class TaskExecutor {
|
|||||||
return worktreePath;
|
return worktreePath;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30 (contention prevention at the registration seam):
|
||||||
|
Every executor session registration goes through `acquireActiveSessionPath` instead of the raw
|
||||||
|
`registerPath`, so a LEAKED entry owned by a task with no live session surface in this process is
|
||||||
|
RECLAIMED rather than throwing at the newcomer. That closes the second contention class (a dead
|
||||||
|
holder can never release, so waiting on it is waiting forever). A genuinely live holder still throws
|
||||||
|
the typed error — that case is real serialization, and callers classify it as a retryable contention
|
||||||
|
hold (SESSION_CONTENTION_HOLD_VALUE), never as a provider/model failure.
|
||||||
|
The probe reports LIVE on any uncertainty: an unknown holder with a fresh entry is treated as live by
|
||||||
|
the staleness floor, so the reclaim only ever fires on proven-dead, aged entries.
|
||||||
|
*/
|
||||||
|
private acquireSessionRegistryPath(taskId: string, registryPath: string, kind: ActiveSessionKind, ownerKey: string): void {
|
||||||
|
const outcome = acquireActiveSessionPath(activeSessionRegistry, registryPath, { taskId, kind, ownerKey }, {
|
||||||
|
holderLiveProbe: (holderTaskId) => this.hasLiveTaskSessionSurface(holderTaskId) || executingTaskLock.has(holderTaskId),
|
||||||
|
});
|
||||||
|
if (outcome.action === "contended") {
|
||||||
|
throw new ActiveSessionPathHeldByForeignTaskError(registryPath, outcome.holderTaskId, taskId);
|
||||||
|
}
|
||||||
|
if (outcome.action === "reclaimed-stale-foreign") {
|
||||||
|
executorLog.warn(
|
||||||
|
`${taskId}: reclaimed a stale active-session entry on ${registryPath} from dead task ${outcome.holderTaskId} (idle ${outcome.ageMs}ms)`,
|
||||||
|
);
|
||||||
|
void this.store.recordRunAuditEvent?.({
|
||||||
|
taskId,
|
||||||
|
agentId: "executor",
|
||||||
|
runId: generateSyntheticRunId("session-path-reclaim", taskId),
|
||||||
|
domain: "database",
|
||||||
|
mutationType: "session:reclaim-stale-foreign-path",
|
||||||
|
target: taskId,
|
||||||
|
metadata: { taskId, holderTaskId: outcome.holderTaskId, kind, ageMs: outcome.ageMs },
|
||||||
|
})?.catch?.(() => undefined);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private setActiveSession(taskId: string, sessionState: ActiveExecutorSessionState, worktreePath: string): void {
|
private setActiveSession(taskId: string, sessionState: ActiveExecutorSessionState, worktreePath: string): void {
|
||||||
this.activeSessions.set(taskId, sessionState);
|
this.activeSessions.set(taskId, sessionState);
|
||||||
activeSessionRegistry.registerPath(this.sessionRegistryPath(taskId, worktreePath), { taskId, kind: "executor", ownerKey: taskId });
|
this.acquireSessionRegistryPath(taskId, this.sessionRegistryPath(taskId, worktreePath), "executor", taskId);
|
||||||
}
|
}
|
||||||
|
|
||||||
private markGraphExecuteSelfRequeued(taskId: string): void {
|
private markGraphExecuteSelfRequeued(taskId: string): void {
|
||||||
@@ -2023,7 +2071,7 @@ export class TaskExecutor {
|
|||||||
private setActiveStepExecutor(taskId: string, stepExecutor: StepSessionExecutor, worktreePath: string, seenSteeringIds = new Set<string>()): void {
|
private setActiveStepExecutor(taskId: string, stepExecutor: StepSessionExecutor, worktreePath: string, seenSteeringIds = new Set<string>()): void {
|
||||||
this.activeStepExecutors.set(taskId, stepExecutor);
|
this.activeStepExecutors.set(taskId, stepExecutor);
|
||||||
this.activeStepExecutorSeenSteeringIds.set(taskId, seenSteeringIds);
|
this.activeStepExecutorSeenSteeringIds.set(taskId, seenSteeringIds);
|
||||||
activeSessionRegistry.registerPath(this.sessionRegistryPath(taskId, worktreePath), { taskId, kind: "step-session", ownerKey: `${taskId}#step-session` });
|
this.acquireSessionRegistryPath(taskId, this.sessionRegistryPath(taskId, worktreePath), "step-session", `${taskId}#step-session`);
|
||||||
}
|
}
|
||||||
|
|
||||||
private deleteActiveStepExecutor(taskId: string, worktreePath?: string): void {
|
private deleteActiveStepExecutor(taskId: string, worktreePath?: string): void {
|
||||||
@@ -2044,7 +2092,7 @@ export class TaskExecutor {
|
|||||||
private setActiveWorkflowStepSession(taskId: string, session: AgentSession, worktreePath: string, seenSteeringIds = new Set<string>()): void {
|
private setActiveWorkflowStepSession(taskId: string, session: AgentSession, worktreePath: string, seenSteeringIds = new Set<string>()): void {
|
||||||
this.activeWorkflowStepSessions.set(taskId, session);
|
this.activeWorkflowStepSessions.set(taskId, session);
|
||||||
this.activeWorkflowStepSessionSeenSteeringIds.set(taskId, seenSteeringIds);
|
this.activeWorkflowStepSessionSeenSteeringIds.set(taskId, seenSteeringIds);
|
||||||
activeSessionRegistry.registerPath(this.sessionRegistryPath(taskId, worktreePath), { taskId, kind: "workflow-step", ownerKey: `${taskId}#workflow-step` });
|
this.acquireSessionRegistryPath(taskId, this.sessionRegistryPath(taskId, worktreePath), "workflow-step", `${taskId}#workflow-step`);
|
||||||
}
|
}
|
||||||
|
|
||||||
private deleteActiveWorkflowStepSession(taskId: string, worktreePath?: string): void {
|
private deleteActiveWorkflowStepSession(taskId: string, worktreePath?: string): void {
|
||||||
@@ -8390,6 +8438,37 @@ export class TaskExecutor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:NodeWorktreeIsolation 2026-07-25-22:10 (planning acquires the task worktree):
|
||||||
|
Public seam for the planning/triage lane. Specification runs a CODING-tool session; pointing it at
|
||||||
|
the shared main checkout meant every planning agent had write tools in the operator's tree and every
|
||||||
|
concurrent planner shared one path. Acquire the task's own worktree up front and let the whole
|
||||||
|
lifecycle — planning, Plan Review, implementation, code review — reuse that single worktree.
|
||||||
|
Returns null (caller falls back to the root, unchanged behavior) when the project is a workspace, or
|
||||||
|
when acquisition fails: planning must never be blocked by a worktree problem.
|
||||||
|
*/
|
||||||
|
public async ensureTaskWorktreeForPlanning(taskId: string): Promise<string | null> {
|
||||||
|
try {
|
||||||
|
if (this.workspaceConfig === undefined) {
|
||||||
|
this.workspaceConfig = await loadWorkspaceConfig(this.rootDir);
|
||||||
|
}
|
||||||
|
if (this.workspaceConfig && (this.workspaceConfig.repos.length ?? 0) > 0) return null;
|
||||||
|
|
||||||
|
const live = await this.store.getTask(taskId);
|
||||||
|
if (live.worktree && existsSync(live.worktree)) return live.worktree;
|
||||||
|
|
||||||
|
const settings = await this.store.getSettings();
|
||||||
|
const acquisitionTask = live.worktree
|
||||||
|
? ({ ...live, worktree: undefined, sessionFile: undefined } as TaskDetail)
|
||||||
|
: live;
|
||||||
|
const acquired = await this.ensureGraphCustomNodeWorktree(acquisitionTask, settings, "planning");
|
||||||
|
return acquired.worktree || null;
|
||||||
|
} catch (error) {
|
||||||
|
executorLog.warn(`${taskId}: could not acquire a planning worktree — planning falls back to the repo root: ${formatError(error)}`);
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private async prepareGraphNodeExecution(
|
private async prepareGraphNodeExecution(
|
||||||
node: WorkflowIrNode,
|
node: WorkflowIrNode,
|
||||||
nodeTask: TaskDetail,
|
nodeTask: TaskDetail,
|
||||||
@@ -8605,39 +8684,59 @@ export class TaskExecutor {
|
|||||||
optionalGroupId,
|
optionalGroupId,
|
||||||
reviewerInlineFixes: (settings as Settings & { reviewerInlineFixes?: boolean }).reviewerInlineFixes,
|
reviewerInlineFixes: (settings as Settings & { reviewerInlineFixes?: boolean }).reviewerInlineFixes,
|
||||||
});
|
});
|
||||||
const executionTarget = writeCapable ? await this.store.getTask(live.id) : live;
|
let executionTarget = writeCapable ? await this.store.getTask(live.id) : live;
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:NodeWorktreeIsolation 2026-07-25-22:10 (EVERY node runs in the task's own worktree):
|
||||||
|
Operator requirement: Plan Review, Code Review — everything except merge — executes in the
|
||||||
|
task-specific worktree, never in the shared main checkout. Read-only gates used to fall back to
|
||||||
|
`this.rootDir` because a pre-execution task has no worktree yet, which is what made two tasks
|
||||||
|
share a path in the first place (the reported FN-1398/FN-1403 Plan Review collision) and what let
|
||||||
|
a reviewer read a main checkout that other tasks and the operator mutate underneath it.
|
||||||
|
ACQUIRE the worktree at planning time instead: `ensureGraphCustomNodeWorktree` is the same
|
||||||
|
acquisition the write-capable nodes already use, so the worktree/branch/baseCommitSha the
|
||||||
|
implementation session later resumes into is created once, here, and reused.
|
||||||
|
A recorded-but-missing worktree is RE-ACQUIRED (strip the stale metadata first, mirroring
|
||||||
|
prepareGraphNodeExecution) rather than degraded to the root — this replaces FN-7996's
|
||||||
|
run-Plan-Review-from-the-repo-root fallback, which is exactly the shared-path behavior being
|
||||||
|
removed. Workspace projects are unchanged: `ensureGraphCustomNodeWorktree` returns the task
|
||||||
|
untouched there, because workspace sessions are rooted at the browse-root by design and per-repo
|
||||||
|
isolation comes from the sub-repo acquire lease.
|
||||||
|
*/
|
||||||
|
const nodeDisplayName = typeof cfg.name === "string" && cfg.name.trim() ? cfg.name.trim() : node.id;
|
||||||
|
const isPlanReviewNode = node.id === "plan-review-step" || nodeDisplayName === "Plan Review" || optionalGroupId === "plan-review";
|
||||||
|
if (!this.workspaceConfig) {
|
||||||
|
const recordedWorktreeMissing = Boolean(executionTarget.worktree) && !existsSync(executionTarget.worktree!);
|
||||||
|
/*
|
||||||
|
A node with NO recorded worktree is pre-execution (planning / Plan Review): acquire one.
|
||||||
|
A node whose RECORDED worktree vanished is a different situation — for gates that review
|
||||||
|
implementation output, the work is gone with it, and handing them a fresh empty worktree would
|
||||||
|
let them review the wrong tree and pass. Those keep failing fast into the unusable-worktree
|
||||||
|
recovery (FN-7996). Plan Review is the exception: it reviews the store-injected PROMPT.md, so it
|
||||||
|
re-acquires rather than parking — this replaces its old "run from the repo root" degrade.
|
||||||
|
*/
|
||||||
|
const shouldAcquire = !executionTarget.worktree || (recordedWorktreeMissing && isPlanReviewNode);
|
||||||
|
if (shouldAcquire) {
|
||||||
|
if (recordedWorktreeMissing) {
|
||||||
|
await this.store.logEntry(
|
||||||
|
live.id,
|
||||||
|
`Plan Review worktree ${executionTarget.worktree} is missing on disk — re-acquiring a task worktree instead of running in the shared checkout`,
|
||||||
|
undefined,
|
||||||
|
this.getRunContextFor(live.id),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
const acquisitionTask = recordedWorktreeMissing
|
||||||
|
? ({ ...executionTarget, worktree: undefined, sessionFile: undefined } as TaskDetail)
|
||||||
|
: executionTarget;
|
||||||
|
executionTarget = await this.ensureGraphCustomNodeWorktree(acquisitionTask, settings, node.id);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if (writeCapable && !executionTarget.worktree && !this.workspaceConfig) {
|
if (writeCapable && !executionTarget.worktree && !this.workspaceConfig) {
|
||||||
return { outcome: "failure", value: "no-worktree-for-write-node" };
|
return { outcome: "failure", value: "no-worktree-for-write-node" };
|
||||||
}
|
}
|
||||||
|
|
||||||
/*
|
const worktreePath = executionTarget.worktree || this.rootDir;
|
||||||
FNXC:PlanReviewWorktree 2026-07-16-18:30:
|
|
||||||
FN-7996: Plan Review runs pre-execution and reviews the store-injected PROMPT.md (see
|
|
||||||
FNXC:PlanReviewSpecInjection) — it does not need worktree contents at all. But it inherited
|
|
||||||
whatever stale task.worktree metadata survived earlier park/requeue cycles, and a recycled or
|
|
||||||
pruned path made session start refuse ("Refusing to start coding agent in missing worktree"),
|
|
||||||
terminal-parking the task. When the recorded worktree is absent on disk, run Plan Review from
|
|
||||||
the repo root instead. Scoped strictly to Plan Review: other read-only gates review
|
|
||||||
implementation diffs, so silently retargeting them to the root would review the wrong tree —
|
|
||||||
they keep failing fast and route through the unusable-worktree graph-failure recovery.
|
|
||||||
*/
|
|
||||||
const nodeDisplayName = typeof cfg.name === "string" && cfg.name.trim() ? cfg.name.trim() : node.id;
|
|
||||||
const isPlanReviewNode = node.id === "plan-review-step" || nodeDisplayName === "Plan Review" || optionalGroupId === "plan-review";
|
|
||||||
let worktreePath = executionTarget.worktree || this.rootDir;
|
|
||||||
if (
|
|
||||||
isPlanReviewNode
|
|
||||||
&& !writeCapable
|
|
||||||
&& executionTarget.worktree
|
|
||||||
&& !existsSync(executionTarget.worktree)
|
|
||||||
) {
|
|
||||||
await this.store.logEntry(
|
|
||||||
live.id,
|
|
||||||
`Plan Review worktree ${executionTarget.worktree} is missing on disk — running the reviewer from the repo root (spec is store-injected)`,
|
|
||||||
undefined,
|
|
||||||
this.getRunContextFor(live.id),
|
|
||||||
);
|
|
||||||
worktreePath = this.rootDir;
|
|
||||||
}
|
|
||||||
let prompt = typeof cfg.prompt === "string" ? cfg.prompt : "";
|
let prompt = typeof cfg.prompt === "string" ? cfg.prompt : "";
|
||||||
let modelProvider = typeof cfg.modelProvider === "string" && cfg.modelProvider.trim() ? cfg.modelProvider : undefined;
|
let modelProvider = typeof cfg.modelProvider === "string" && cfg.modelProvider.trim() ? cfg.modelProvider : undefined;
|
||||||
let modelId = typeof cfg.modelId === "string" && cfg.modelId.trim() ? cfg.modelId : undefined;
|
let modelId = typeof cfg.modelId === "string" && cfg.modelId.trim() ? cfg.modelId : undefined;
|
||||||
@@ -9089,6 +9188,93 @@ export class TaskExecutor {
|
|||||||
|| latestAction === "Resuming execution after unpause";
|
|| latestAction === "Resuming execution after unpause";
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30:
|
||||||
|
Two ways in, because contention must never slip through to a park:
|
||||||
|
- the typed failure value the graph now publishes (SESSION_CONTENTION_HOLD_VALUE), and
|
||||||
|
- a message-shape fallback over the run's `:error` context patches, so contention arriving from a
|
||||||
|
path that has not been taught the typed value is still recognized.
|
||||||
|
*/
|
||||||
|
private graphFailureErrorTexts(result: WorkflowGraphTaskRunResult): string[] {
|
||||||
|
if (!result.context) return [];
|
||||||
|
const texts: string[] = [];
|
||||||
|
for (const [key, value] of Object.entries(result.context)) {
|
||||||
|
if (key.endsWith(":error") && typeof value === "string" && value.trim()) texts.push(value);
|
||||||
|
}
|
||||||
|
return texts;
|
||||||
|
}
|
||||||
|
|
||||||
|
private isSessionContentionGraphFailure(result: WorkflowGraphTaskRunResult): boolean {
|
||||||
|
if (this.graphFailureValue(result) === SESSION_CONTENTION_HOLD_VALUE) return true;
|
||||||
|
return this.graphFailureErrorTexts(result).some((text) => isSessionContentionError(text));
|
||||||
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30 (self-recovering wait — the task is never parked):
|
||||||
|
Retry the graph in place on an exponential backoff while the holder finishes. The counter is
|
||||||
|
IN-MEMORY on purpose: it needs no schema change, and an engine restart resetting it is the desired
|
||||||
|
behavior (a restart also drops the in-process registry, so the contention is gone anyway).
|
||||||
|
When the ladder is exhausted the task is left cleanly dispatchable — status/error cleared, progress
|
||||||
|
untouched — so ordinary scheduling picks it up later with a fresh budget. There is no terminal branch
|
||||||
|
here by design: lease contention always ends (the holder finishes, or self-healing sweeps it), so
|
||||||
|
parking the task would only require a human to press Retry on a condition that fixed itself.
|
||||||
|
*/
|
||||||
|
private sessionContentionHoldAttempts = new Map<string, number>();
|
||||||
|
|
||||||
|
private clearSessionContentionHold(taskId: string): void {
|
||||||
|
this.sessionContentionHoldAttempts.delete(taskId);
|
||||||
|
}
|
||||||
|
|
||||||
|
private async holdForSessionContention(
|
||||||
|
task: Task,
|
||||||
|
live: TaskDetail,
|
||||||
|
result: WorkflowGraphTaskRunResult,
|
||||||
|
): Promise<void> {
|
||||||
|
const detail = this.graphFailureErrorTexts(result).find((text) => isSessionContentionError(text));
|
||||||
|
const priorAttempts = this.sessionContentionHoldAttempts.get(task.id) ?? 0;
|
||||||
|
const attempt = priorAttempts + 1;
|
||||||
|
|
||||||
|
if (attempt > MAX_SESSION_CONTENTION_HOLD_RETRIES) {
|
||||||
|
this.clearSessionContentionHold(task.id);
|
||||||
|
const message = `Still waiting on another task to release a shared session path after ${MAX_SESSION_CONTENTION_HOLD_RETRIES} attempts — leaving the task queued for normal re-dispatch (not a failure)${detail ? `: ${detail}` : ""}`;
|
||||||
|
executorLog.warn(`${task.id}: ${message}`);
|
||||||
|
await this.store.logEntry(task.id, message, undefined, this.getRunContextFor(task.id));
|
||||||
|
if (live.status != null || live.error != null) {
|
||||||
|
await this.store.updateTask(task.id, { status: null, error: null }, this.getRunContextFor(task.id));
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
this.sessionContentionHoldAttempts.set(task.id, attempt);
|
||||||
|
const message = `Waiting on another task to release a shared session path — retrying in place (${attempt}/${MAX_SESSION_CONTENTION_HOLD_RETRIES})${detail ? `: ${detail}` : ""}`;
|
||||||
|
executorLog.warn(`${task.id}: ${message}`);
|
||||||
|
await this.store.logEntry(task.id, message, undefined, this.getRunContextFor(task.id));
|
||||||
|
// A contention hold is not a failure state: clear any stale park so the row never shows as failed
|
||||||
|
// while it is simply waiting its turn.
|
||||||
|
if (live.status != null || live.error != null) {
|
||||||
|
await this.store.updateTask(task.id, { status: null, error: null }, this.getRunContextFor(task.id));
|
||||||
|
}
|
||||||
|
|
||||||
|
const delayMs = SESSION_CONTENTION_HOLD_BACKOFF_MS === 0
|
||||||
|
? 0
|
||||||
|
: Math.min(SESSION_CONTENTION_HOLD_MAX_BACKOFF_MS, SESSION_CONTENTION_HOLD_BACKOFF_MS * 2 ** (attempt - 1));
|
||||||
|
const scheduleRetry = () => {
|
||||||
|
void (async () => {
|
||||||
|
try {
|
||||||
|
const resume = await this.store.getTask(task.id);
|
||||||
|
if (!resume || resume.deletedAt || resume.paused || resume.userPaused) {
|
||||||
|
this.clearSessionContentionHold(task.id);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
await this.execute(resume);
|
||||||
|
} catch (err) {
|
||||||
|
executorLog.error(`Failed session-contention retry for ${task.id}:`, err);
|
||||||
|
}
|
||||||
|
})();
|
||||||
|
};
|
||||||
|
setTimeout(scheduleRetry, delayMs).unref?.();
|
||||||
|
}
|
||||||
|
|
||||||
private graphFailureValue(result: WorkflowGraphTaskRunResult): string | undefined {
|
private graphFailureValue(result: WorkflowGraphTaskRunResult): string | undefined {
|
||||||
const failedNode = result.visitedNodeIds[result.visitedNodeIds.length - 1];
|
const failedNode = result.visitedNodeIds[result.visitedNodeIds.length - 1];
|
||||||
if (!failedNode || !result.context) return undefined;
|
if (!failedNode || !result.context) return undefined;
|
||||||
@@ -9943,6 +10129,18 @@ export class TaskExecutor {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
/*
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30:
|
||||||
|
Classified BEFORE every other graph-failure router. A node that could not start because another
|
||||||
|
task holds its session path or sub-repo lease is not a provider outage, not a plan defect, and not
|
||||||
|
a terminal failure — it is a wait. Route it to the self-recovering backoff hold, which never parks
|
||||||
|
the task and never consumes the provider/artifact retry budgets.
|
||||||
|
*/
|
||||||
|
if (this.isSessionContentionGraphFailure(result)) {
|
||||||
|
await this.holdForSessionContention(task, live, result);
|
||||||
|
await this.persistTokenUsage(task.id);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
/*
|
||||||
FNXC:MissingWorktreeRecovery 2026-07-16-18:25:
|
FNXC:MissingWorktreeRecovery 2026-07-16-18:25:
|
||||||
An unusable-worktree session-start refusal inside a graph node must route to the bounded
|
An unusable-worktree session-start refusal inside a graph node must route to the bounded
|
||||||
worktree-session recovery BEFORE any other classifier: FN-7977's provider-failure hold
|
worktree-session recovery BEFORE any other classifier: FN-7977's provider-failure hold
|
||||||
|
|||||||
@@ -1092,6 +1092,9 @@ export class InProcessRuntime
|
|||||||
usageLimitPauser: this.usageLimitPauser,
|
usageLimitPauser: this.usageLimitPauser,
|
||||||
agentStore: this.agentStore,
|
agentStore: this.agentStore,
|
||||||
pluginRunner: this.pluginRunner,
|
pluginRunner: this.pluginRunner,
|
||||||
|
// FNXC:NodeWorktreeIsolation 2026-07-25-22:10: planning acquires (or reuses) the task's own
|
||||||
|
// worktree through the executor's acquisition path, so no lane runs in the shared checkout.
|
||||||
|
acquirePlanningWorktree: (taskId) => this.executor.ensureTaskWorktreeForPlanning(taskId),
|
||||||
onSpecifyStart: (t) => {
|
onSpecifyStart: (t) => {
|
||||||
this.recordActivity();
|
this.recordActivity();
|
||||||
runtimeLog.log(`Specifying ${t.id}...`);
|
runtimeLog.log(`Specifying ${t.id}...`);
|
||||||
|
|||||||
@@ -24,8 +24,8 @@ now live in the import-free leaf `transient-error-patterns.ts` so the merge clas
|
|||||||
one definition of "transient" without inheriting this module's logger chain (FN-8004).
|
one definition of "transient" without inheriting this module's logger chain (FN-8004).
|
||||||
Re-exported here so every existing importer of this module keeps working unchanged.
|
Re-exported here so every existing importer of this module keeps working unchanged.
|
||||||
*/
|
*/
|
||||||
import { isTransientAuthCredentialError, isTransientError } from "./transient-error-patterns.js";
|
import { isSessionContentionError, isTransientAuthCredentialError, isTransientError } from "./transient-error-patterns.js";
|
||||||
export { TRANSIENT_ERROR_PATTERNS, isTransientAuthCredentialError, isTransientError } from "./transient-error-patterns.js";
|
export { TRANSIENT_ERROR_PATTERNS, SESSION_CONTENTION_PATTERNS, isSessionContentionError, isTransientAuthCredentialError, isTransientError } from "./transient-error-patterns.js";
|
||||||
|
|
||||||
|
|
||||||
/*
|
/*
|
||||||
@@ -49,6 +49,20 @@ export function isNonPlanDefectPlanReviewFailure(input: {
|
|||||||
}): boolean {
|
}): boolean {
|
||||||
if (input.verdict === "REVISE") return false;
|
if (input.verdict === "REVISE") return false;
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30:
|
||||||
|
Contention is NOT a provider failure and must never take the provider hold, whose budget is two
|
||||||
|
immediate retries followed by a park — retrying instantly cannot clear a lease another task holds.
|
||||||
|
Excluded at the source so a caller that forgets the dedicated contention check still cannot misroute
|
||||||
|
it; contention has its own backoff hold (SESSION_CONTENTION_HOLD_VALUE).
|
||||||
|
*/
|
||||||
|
if (
|
||||||
|
input.failureValue?.trim().toLowerCase() === "session-contention-hold"
|
||||||
|
|| (input.errorMessage ? isSessionContentionError(input.errorMessage) : false)
|
||||||
|
) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
const failureValue = input.failureValue?.trim().toLowerCase();
|
const failureValue = input.failureValue?.trim().toLowerCase();
|
||||||
if (failureValue === "exception" || failureValue === "aborted") return true;
|
if (failureValue === "exception" || failureValue === "aborted") return true;
|
||||||
|
|
||||||
|
|||||||
@@ -94,6 +94,40 @@ export const TRANSIENT_ERROR_PATTERNS: RegExp[] = [
|
|||||||
/\bACP session has no live connection\b/i,
|
/\bACP session has no live connection\b/i,
|
||||||
];
|
];
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:Reliability-ErrorClassification 2026-07-25-21:10 (session/lease contention is NOT a provider failure):
|
||||||
|
Fusion serializes work on a shared path with in-process leases: the activeSessionRegistry foreign-task
|
||||||
|
guard, the workspace sub-repo acquire lease, and the workspace sub-repo land lease. When a second task
|
||||||
|
contends, the holder throws one of these three messages. They are pure CONTENTION — another task is
|
||||||
|
mid-flight on the same path — so the only correct response is "wait and try again", never "this model /
|
||||||
|
provider / plan is broken".
|
||||||
|
Reported failure: a Plan Review collision was classified as a provider failure ("Plan Review provider
|
||||||
|
failure — retrying in place (2/2)"), retried twice with no delay against a hold retrying could not clear,
|
||||||
|
then parked the task with its budget spent. Matching these shapes as transient makes every generic
|
||||||
|
retry/requeue path (executor main session, merge classifier, durable-agent heartbeat recovery) wait it
|
||||||
|
out instead. Graph nodes get an explicit backoff hold on top — see SESSION_CONTENTION_HOLD_VALUE.
|
||||||
|
*/
|
||||||
|
export const SESSION_CONTENTION_PATTERNS: RegExp[] = [
|
||||||
|
// ActiveSessionPathHeldByForeignTaskError (active-session-registry.ts)
|
||||||
|
/active-session path\s+\S+\s+is held by task\s+\S+/i,
|
||||||
|
// WorkspaceRepoAcquireBusyError (worktree-acquisition.ts)
|
||||||
|
/workspace sub-repo\s+\S+\s+acquisition is in progress for task\s+\S+/i,
|
||||||
|
// WorkspaceRepoLandBusyError (merger-ai.ts)
|
||||||
|
/workspace sub-repo\s+\S+\s+land is in progress for task\s+\S+/i,
|
||||||
|
];
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Detect a path/lease contention failure: another task holds the session path,
|
||||||
|
* sub-repo acquire lease, or sub-repo land lease this task wants. Always
|
||||||
|
* temporary — the holder releases when its own work finishes.
|
||||||
|
*/
|
||||||
|
export function isSessionContentionError(errorMessage: string): boolean {
|
||||||
|
if (!errorMessage || typeof errorMessage !== "string") {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
return SESSION_CONTENTION_PATTERNS.some((pattern) => pattern.test(errorMessage));
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Check if an error message indicates a transient network/infrastructure error.
|
* Check if an error message indicates a transient network/infrastructure error.
|
||||||
*
|
*
|
||||||
@@ -120,6 +154,11 @@ export function isTransientError(errorMessage: string): boolean {
|
|||||||
if (isTransientAuthCredentialError(errorMessage)) {
|
if (isTransientAuthCredentialError(errorMessage)) {
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
// FNXC:Reliability-ErrorClassification 2026-07-25-21:10: lease/session contention is transient by
|
||||||
|
// construction, so every generic retry path waits for the holder instead of parking the task.
|
||||||
|
if (isSessionContentionError(errorMessage)) {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
return TRANSIENT_ERROR_PATTERNS.some((pattern) => pattern.test(errorMessage));
|
return TRANSIENT_ERROR_PATTERNS.some((pattern) => pattern.test(errorMessage));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -209,6 +209,15 @@ export interface TriageProcessorOptions {
|
|||||||
agentStore?: import("@fusion/core").AgentStore;
|
agentStore?: import("@fusion/core").AgentStore;
|
||||||
/** Plugin runner for runtime selection. When provided, enables plugin runtime lookup. */
|
/** Plugin runner for runtime selection. When provided, enables plugin runtime lookup. */
|
||||||
pluginRunner?: import("./plugin-runner.js").PluginRunner;
|
pluginRunner?: import("./plugin-runner.js").PluginRunner;
|
||||||
|
/*
|
||||||
|
FNXC:NodeWorktreeIsolation 2026-07-25-22:10:
|
||||||
|
Acquires (or reuses) the task-specific worktree so the planning session runs there instead of in the
|
||||||
|
shared main checkout. Planning uses the CODING tool surface, so running it at the repo root gave every
|
||||||
|
planner write tools in the operator's tree and made concurrent planners share one path. Optional: when
|
||||||
|
unwired (older callers, tests) or when it resolves null (workspace projects, acquisition failure),
|
||||||
|
planning falls back to the repo root exactly as before.
|
||||||
|
*/
|
||||||
|
acquirePlanningWorktree?: (taskId: string) => Promise<string | null>;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -1756,11 +1765,23 @@ export class TriageProcessor {
|
|||||||
advertising that writer in the prompt while running readonly stranded triage
|
advertising that writer in the prompt while running readonly stranded triage
|
||||||
on the original PROMPT.md stub and sent the stub into Plan Review.
|
on the original PROMPT.md stub and sent the stub into Plan Review.
|
||||||
*/
|
*/
|
||||||
|
/*
|
||||||
|
FNXC:NodeWorktreeIsolation 2026-07-25-22:10:
|
||||||
|
Planning runs in the TASK's own worktree, not the shared main checkout. This session carries the
|
||||||
|
coding tool surface (see FNXC:TriagePromptPersistence above), so rooting it at `this.rootDir`
|
||||||
|
put write tools in the operator's tree and made every concurrent planner share one path — the
|
||||||
|
same shared-path shape behind the reported Plan Review session collision. The worktree acquired
|
||||||
|
here is the one Plan Review and the implementation session then reuse.
|
||||||
|
*/
|
||||||
|
const planningCwd = (await this.options.acquirePlanningWorktree?.(task.id).catch(() => null)) || this.rootDir;
|
||||||
|
if (planningCwd !== this.rootDir) {
|
||||||
|
await this.store.logEntry(task.id, `Planning session running in task worktree ${planningCwd}`).catch(() => undefined);
|
||||||
|
}
|
||||||
const { session } = await createResolvedAgentSession({
|
const { session } = await createResolvedAgentSession({
|
||||||
sessionPurpose: "triage",
|
sessionPurpose: "triage",
|
||||||
runtimeHint: triageRuntimeHint,
|
runtimeHint: triageRuntimeHint,
|
||||||
pluginRunner: this.options.pluginRunner,
|
pluginRunner: this.options.pluginRunner,
|
||||||
cwd: this.rootDir,
|
cwd: planningCwd,
|
||||||
systemPrompt: triageSystemPromptFinal,
|
systemPrompt: triageSystemPromptFinal,
|
||||||
systemPromptLayers: triageLayers,
|
systemPromptLayers: triageLayers,
|
||||||
tools: "coding",
|
tools: "coding",
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ import type {
|
|||||||
} from "@fusion/core";
|
} from "@fusion/core";
|
||||||
import { BUILTIN_CODING_WORKFLOW_IR, PLAN_REVIEW_GROUP_ID, WorkflowIrError, getWorkflowExtensionRegistry, resolveMaxReworkCycles, isExperimentalFeatureEnabled, GRAPH_NATIVE_POST_MERGE_FLAG, isCompletionSummaryNode, classifyReviewLease, isWorkflowOptionalGroupEnabled } from "@fusion/core";
|
import { BUILTIN_CODING_WORKFLOW_IR, PLAN_REVIEW_GROUP_ID, WorkflowIrError, getWorkflowExtensionRegistry, resolveMaxReworkCycles, isExperimentalFeatureEnabled, GRAPH_NATIVE_POST_MERGE_FLAG, isCompletionSummaryNode, classifyReviewLease, isWorkflowOptionalGroupEnabled } from "@fusion/core";
|
||||||
import { isNonPlanDefectPlanReviewFailure } from "./transient-error-detector.js";
|
import { isNonPlanDefectPlanReviewFailure } from "./transient-error-detector.js";
|
||||||
|
import { isSessionContentionError } from "./transient-error-patterns.js";
|
||||||
import { isRequiredArtifactReadFailedValue, parseRequiredArtifactMissingValue } from "./required-workflow-artifacts.js";
|
import { isRequiredArtifactReadFailedValue, parseRequiredArtifactMissingValue } from "./required-workflow-artifacts.js";
|
||||||
|
|
||||||
import {
|
import {
|
||||||
@@ -66,6 +67,16 @@ interleaving impossible by construction (exactly one reviewer per gate per attem
|
|||||||
*/
|
*/
|
||||||
export const PLAN_REVIEW_LEASE_HELD_VALUE = "plan-review-lease-held";
|
export const PLAN_REVIEW_LEASE_HELD_VALUE = "plan-review-lease-held";
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30:
|
||||||
|
A node that failed because another task holds the session path / sub-repo lease it needs is CONTENDED,
|
||||||
|
not broken. It gets its own failure value so the executor can route it to a backoff hold that waits for
|
||||||
|
the holder, instead of falling into the generic `exception` bucket — where Plan Review's classifier read
|
||||||
|
it as a PROVIDER failure, retried twice with no delay, and parked the task once that budget was spent
|
||||||
|
(the reported FN-1398/FN-1403 collision). Applies to EVERY node, not just Plan Review.
|
||||||
|
*/
|
||||||
|
export const SESSION_CONTENTION_HOLD_VALUE = "session-contention-hold";
|
||||||
|
|
||||||
export type WorkflowNodeAbortKind = "engine-pause";
|
export type WorkflowNodeAbortKind = "engine-pause";
|
||||||
|
|
||||||
export const WORKFLOW_INTERRUPTED_NODE_ID_CONTEXT_KEY = "workflow:interruptedNodeId";
|
export const WORKFLOW_INTERRUPTED_NODE_ID_CONTEXT_KEY = "workflow:interruptedNodeId";
|
||||||
@@ -914,6 +925,27 @@ export class WorkflowGraphExecutor {
|
|||||||
* operator-actionable errors must remain visible in place so completed
|
* operator-actionable errors must remain visible in place so completed
|
||||||
* execution work cannot bounce back to the planner column.
|
* execution work cannot bounce back to the planner column.
|
||||||
*/
|
*/
|
||||||
|
/*
|
||||||
|
* FNXC:SessionContention 2026-07-25-21:30:
|
||||||
|
* Classified FIRST, and for EVERY optional group — not just Plan Review. A gate that could
|
||||||
|
* not start because another task holds its session path must neither traverse the failure
|
||||||
|
* edge to plan-replan (nothing about the plan is wrong) nor be labeled a provider failure
|
||||||
|
* (nothing about the provider is wrong). It becomes a contention hold the executor waits out.
|
||||||
|
*/
|
||||||
|
const sessionContentionFailure =
|
||||||
|
stepStatus === "failed"
|
||||||
|
&& (
|
||||||
|
verdictRaw === SESSION_CONTENTION_HOLD_VALUE
|
||||||
|
|| isSessionContentionError(stepOutput ?? stepNotes ?? "")
|
||||||
|
);
|
||||||
|
if (sessionContentionFailure) {
|
||||||
|
this.deps.logTaskEntry?.(
|
||||||
|
`${logPrefix} ${groupName} is waiting on another task to release its session path — holding in place`,
|
||||||
|
);
|
||||||
|
context[`node:${node.id}:outcome`] = "failure";
|
||||||
|
context[`node:${node.id}:value`] = SESSION_CONTENTION_HOLD_VALUE;
|
||||||
|
return { outcome: "failure", value: SESSION_CONTENTION_HOLD_VALUE };
|
||||||
|
}
|
||||||
const nonPlanDefectPlanReviewFailure =
|
const nonPlanDefectPlanReviewFailure =
|
||||||
node.id === PLAN_REVIEW_GROUP_ID
|
node.id === PLAN_REVIEW_GROUP_ID
|
||||||
&& stepStatus === "failed"
|
&& stepStatus === "failed"
|
||||||
@@ -1482,6 +1514,13 @@ export class WorkflowGraphExecutor {
|
|||||||
} catch (error) {
|
} catch (error) {
|
||||||
if (signal?.aborted) return this.withEnginePauseAbortContext(node, { outcome: "failure", value: "aborted" });
|
if (signal?.aborted) return this.withEnginePauseAbortContext(node, { outcome: "failure", value: "aborted" });
|
||||||
lastError = error;
|
lastError = error;
|
||||||
|
/*
|
||||||
|
FNXC:SessionContention 2026-07-25-21:30:
|
||||||
|
Immediate in-loop retries cannot clear a lease another task holds — they just burn the node's
|
||||||
|
attempt budget in milliseconds and erase the diagnostic behind a generic `exception`. Stop
|
||||||
|
retrying here and hand the executor a typed contention failure it can back off on.
|
||||||
|
*/
|
||||||
|
if (isSessionContentionError(error instanceof Error ? error.message : String(error))) break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1513,11 +1552,13 @@ export class WorkflowGraphExecutor {
|
|||||||
}
|
}
|
||||||
return degraded;
|
return degraded;
|
||||||
}
|
}
|
||||||
|
const lastErrorText = lastError instanceof Error ? lastError.message : String(lastError);
|
||||||
const failureResult: WorkflowNodeResult = {
|
const failureResult: WorkflowNodeResult = {
|
||||||
outcome: "failure",
|
outcome: "failure",
|
||||||
value: "exception",
|
// FNXC:SessionContention 2026-07-25-21:30: contention is a retryable hold, not an exception.
|
||||||
|
value: isSessionContentionError(lastErrorText) ? SESSION_CONTENTION_HOLD_VALUE : "exception",
|
||||||
contextPatch: {
|
contextPatch: {
|
||||||
[`node:${node.id}:error`]: lastError instanceof Error ? lastError.message : String(lastError),
|
[`node:${node.id}:error`]: lastErrorText,
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
if (recordProgress && this.shouldRecordNodeProgress(node)) {
|
if (recordProgress && this.shouldRecordNodeProgress(node)) {
|
||||||
|
|||||||
Reference in New Issue
Block a user