U8 PR2: the execution-policy ladder resolves its own workflow's columns (the wip literal made retry, escalation and loop protection unreachable) (#2497)
Second PR of **U8 — the graph owns execution**, independent of [#2490](https://github.com/Runfusion/Fusion/pull/2490) and of every other unit. Small, green, independently revertable. ## The defect `handleGraphFailure`'s execution-policy ladder — FN-7863/FN-7926 dispatch-loop terminalization, FN-7996 tool-failure retry, FN-7998 escalation — decided a task's own lifecycle by naming `"todo"` and `"in-progress"` **literally, at 9 sites**. U5b converted the executor's *rebounds* to `resolveReboundColumnFor`; these were left behind, each sitting somewhere an awaited resolver could not reach: inside synchronous `updateTaskAtomic` mutators, inside fire-and-forget resume closures, and in conditions evaluated before any resolution happened. **The severe one is the wip gate, and it fails silently in the worst direction:** ```ts if (live.column !== "in-progress") { // "Workflow graph run ended after task already advanced — no further action needed" return; } ``` Under a workflow that renames the implementation column, that is true of a card sitting in **its own wip column**. So the graph failure was swallowed whole — no terminal park, no status, no error, nothing on the board — and the scheduler re-dispatched the same doomed run. Every later branch sits behind that gate, which is why the retry budgets, the escalation, and the bounded terminalization were **unreachable rather than mistargeted**. This is precisely the failure the program's problem frame predicts: *a guard that stops matching disables a recovery path invisibly and the suite stays green.* I found it because my first renamed-column test for the escalation site could not reach the escalation code at all. Two further sites misbehave once the gate is passable: - **FN-7998 node escalation** wrote `column: "todo"` inside the atomic claim — parking the card where no workflow declares it, which is on the plan's **"Stop implementation if"** list and what R7 exists to clean up after. The scheduler's effective-node resolution, the entire point of a node escalation, never runs. - **FN-7863/FN-7926's `live.column === "todo"` arm** is the classic guard that stops matching. In-process the `executeNodeSelfRequeued` marker covers the same case, so this degrades only on the **durable** arm — after a restart, or for a second `TaskExecutor` instance in the process, where the column read is the only evidence the inner executor requeued. A progressing card then falls through to the terminal sink and is parked `failed`. ## The fix Resolve hold and wip **once per graph failure** through U1's `resolveTaskLifecycleColumns` and thread the pair through the ladder. Both fall back to the legacy literal when the workflow cannot be resolved, so an unresolvable workflow keeps exactly its pre-conversion behavior rather than guessing. One IR read on a terminal recovery path — not an enumeration loop. ## Red-green, measured **3 of the 8 new tests fail with this commit's executor change reverted:** ``` FAIL FN-7998 … > requeues a node escalation to the RENAMED hold column, not the literal todo FAIL FN-7998 … > still does not move the card for a MODEL-target escalation FAIL FN-7863/FN-7926 … > recognises an inner-executor requeue that landed in the RENAMED hold column Tests 3 failed | 5 passed (8) ← reverted Tests 8 passed (8) ← with the fix ``` The other **5 pass both ways by design**, and I am not claiming them as red-green — they are the regression floor: - default coding workflow still resolves hold → `todo`, wip → `in-progress` (byte-identical); - an unresolvable workflow still uses the legacy literals; - the in-process self-requeue marker still works when no workflow resolves; - and a **negative case** proving the dispatch-loop gate stays narrow — a card still in its wip column with no marker is a genuine execute failure and must NOT be swallowed as a benign recovery. Widening that gate to "any column" would have been the easy wrong fix. ## Scope Deliberately the execution-policy ladder only. **20 further column literals remain in the same method's pause-abort, merge, and in-review regions** — they belong to U5's executor slice (B4, not started) and U9's merge lane, and are untouched here. Flagging the overlap: this PR edits `executor.ts`, so whoever takes U5-B4 should rebase onto it rather than converting these 9 sites again. ## Verification - 8 new tests + the preserved-behavior suites (`executor-tool-failure-retry`, `executor-graph-requeue-gate`, `executor-task-done-blocked`, `executor-graph-boundary`, `executor-stuck-requeue-preserve-progress`, `executor-paused-abort-todo-benign`, `executor-abort-provenance`) — **9 files, 112 tests, green** - `pnpm test:gate` — green (2/10, 16/299, 1/71); `pnpm lint` clean; `tsc --noEmit` on `@fusion/engine` clean - Changeset included (`patch`, category `fix`), passes `pnpm check:changesets` 🤖 Generated with [Claude Code](https://claude.com/claude-code) <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **Bug Fixes** * Fixed execution recovery for workflows with renamed lifecycle columns so retry, escalation, and loop-protection behaviors correctly follow the workflow’s declared hold/WIP columns. * Preserved legacy behavior for default workflows and continued safe handling when lifecycle columns can’t be resolved. * **Tests** * Added a Vitest suite validating execution-policy “ladder” behavior for renamed columns, including node escalation, dispatch-loop gating, and fail-closed scenarios. <!-- end of auto-generated comment: release notes by coderabbit.ai --> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
7
.changeset/u8-execution-policy-renamed-columns.md
Normal file
7
.changeset/u8-execution-policy-renamed-columns.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Fix execution retry, escalation and loop protection silently doing nothing on renamed-column workflows.
|
||||||
|
category: fix
|
||||||
|
dev: `handleGraphFailure`'s execution-policy ladder (FN-7863/FN-7926 dispatch-loop gate, FN-7996 tool-failure retry, FN-7998 escalation) resolved hold/wip through U1's `resolveTaskLifecycleColumns` instead of the literals `"todo"`/`"in-progress"` at 9 sites. The wip literal made the whole ladder unreachable: a card in a renamed implementation column was classified "already advanced" and its graph failure was swallowed. A workflow that declares no hold/wip column resolves through KTD-10 rebound ordering or fails closed to a visible terminal park — never to an invented column; only an unreadable workflow keeps the legacy literals.
|
||||||
@@ -0,0 +1,464 @@
|
|||||||
|
/*
|
||||||
|
FNXC:WorkflowExecutionOwnership 2026-07-28-09:10 (U8 / R3, R12 — workflow-owned lifecycle):
|
||||||
|
|
||||||
|
`handleGraphFailure`'s execution-policy ladder — the FN-7863/FN-7926 dispatch-loop gate, the
|
||||||
|
FN-7996 tool-failure retry, and the FN-7998 escalation — decided a task's own lifecycle by
|
||||||
|
naming `"todo"` and `"in-progress"` literally, at 8 sites. U5b converted the executor's
|
||||||
|
*rebounds* to `resolveReboundColumnFor` and these were left behind, each sitting somewhere an
|
||||||
|
awaited resolver could not reach: inside synchronous `updateTaskAtomic` mutators, inside
|
||||||
|
fire-and-forget resume closures, and in conditions evaluated before any resolution happened.
|
||||||
|
|
||||||
|
THE SEVERE ONE IS THE WIP GATE, and it is severe in the silent direction. `if (live.column !==
|
||||||
|
"in-progress")` classifies the run as "task already advanced — no further action needed". Under
|
||||||
|
a renamed implementation column that is true of a card sitting in its OWN wip column, so the
|
||||||
|
graph failure was swallowed whole: no terminal park, no status, no error, nothing on the board —
|
||||||
|
and the scheduler re-dispatched the same doomed run. Every downstream branch is behind that
|
||||||
|
gate, which is why the tool-failure retry, the escalation, and the bounded terminalization were
|
||||||
|
all unreachable rather than merely mistargeted. It is exactly the failure the program's problem
|
||||||
|
frame predicts: a guard that stops matching disables a recovery path invisibly and the suite
|
||||||
|
stays green.
|
||||||
|
|
||||||
|
Two further sites misbehave once the gate is passable:
|
||||||
|
|
||||||
|
- FN-7998 escalation, node-target branch, wrote `column: "todo"` inside the atomic claim.
|
||||||
|
Under a workflow with no `todo` column that parks the card where no workflow declares it —
|
||||||
|
on the plan's "Stop implementation if" list, and what R7 exists to clean up after.
|
||||||
|
- FN-7863/FN-7926's `live.column === "todo"` arm is the classic guard that stops matching.
|
||||||
|
In-process the `executeNodeSelfRequeued` marker covers the same case, so a renamed workflow
|
||||||
|
degrades only on the DURABLE arm — after a restart, or for a second TaskExecutor instance in
|
||||||
|
the process, where the column read is the only evidence the inner executor requeued. A
|
||||||
|
progressing card then falls through to the terminal sink and is parked `failed`.
|
||||||
|
|
||||||
|
Every test here was written against the literal implementation and observed FAILING first
|
||||||
|
(counts in the PR). The default-workflow and unresolvable-workflow cases are the regression
|
||||||
|
floor: `builtin:coding` resolves hold -> `todo` and wip -> `in-progress`, and an unresolvable
|
||||||
|
workflow keeps the legacy literals, so both must be byte-identical after the conversion.
|
||||||
|
|
||||||
|
Roles resolve per WORKFLOW, never per role: a resolvable IR uses `resolveLifecycleColumns` for
|
||||||
|
wip and KTD-10's `resolveReboundTarget` (hold -> intake -> first) for the requeue target, so it
|
||||||
|
can only ever be handed a column it declares. The legacy literals survive for exactly one case —
|
||||||
|
a workflow that cannot be read at all — which keeps pre-conversion behavior. See the
|
||||||
|
"never given a column it does not declare" block at the bottom for why the first cut's per-role
|
||||||
|
`?? "todo"` was wrong.
|
||||||
|
*/
|
||||||
|
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||||
|
import type { TaskDetail, WorkflowIr } from "@fusion/core";
|
||||||
|
import "./executor-test-helpers.js";
|
||||||
|
import { TaskExecutor } from "../executor.js";
|
||||||
|
import { createMockStore, resetExecutorMocks } from "./executor-test-helpers.js";
|
||||||
|
|
||||||
|
const now = "2026-07-28T00:00:00.000Z";
|
||||||
|
const WF = "custom:renamed";
|
||||||
|
|
||||||
|
/** A workflow whose hold column is `drafting` and whose wip column is `building`. No `todo`. */
|
||||||
|
function renamedIr(): WorkflowIr {
|
||||||
|
return {
|
||||||
|
version: "v2",
|
||||||
|
id: WF,
|
||||||
|
nodes: [],
|
||||||
|
edges: [],
|
||||||
|
columns: [
|
||||||
|
{ id: "inbox", label: "Inbox", traits: [{ trait: "intake" }] },
|
||||||
|
{ id: "drafting", label: "Drafting", traits: [{ trait: "hold", config: { release: "capacity" } }] },
|
||||||
|
{ id: "building", label: "Building", traits: [{ trait: "wip", config: { limitSetting: "maxConcurrent" } }] },
|
||||||
|
{ id: "shipped", label: "Shipped", traits: [{ trait: "complete" }] },
|
||||||
|
],
|
||||||
|
} as unknown as WorkflowIr;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The builtin coding shape — the byte-identical regression floor. */
|
||||||
|
function defaultIr(): WorkflowIr {
|
||||||
|
return {
|
||||||
|
version: "v2",
|
||||||
|
id: WF,
|
||||||
|
nodes: [],
|
||||||
|
edges: [],
|
||||||
|
columns: [
|
||||||
|
{ id: "triage", label: "Triage", traits: [{ trait: "intake" }] },
|
||||||
|
{ id: "todo", label: "Todo", traits: [{ trait: "hold", config: { release: "capacity" } }] },
|
||||||
|
{ id: "in-progress", label: "In Progress", traits: [{ trait: "wip", config: { limitSetting: "maxConcurrent" } }] },
|
||||||
|
{ id: "done", label: "Done", traits: [{ trait: "complete" }] },
|
||||||
|
],
|
||||||
|
} as unknown as WorkflowIr;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* A VALID workflow that declares a wip column but NO hold column. `resolveReboundTarget`'s
|
||||||
|
* KTD-10 ordering (hold -> intake -> first) resolves the requeue target to the declared intake
|
||||||
|
* column; nothing may substitute the literal `todo`, which this workflow does not declare.
|
||||||
|
*/
|
||||||
|
function noHoldIr(): WorkflowIr {
|
||||||
|
return {
|
||||||
|
version: "v2",
|
||||||
|
id: WF,
|
||||||
|
nodes: [],
|
||||||
|
edges: [],
|
||||||
|
columns: [
|
||||||
|
{ id: "inbox", label: "Inbox", traits: [{ trait: "intake" }] },
|
||||||
|
{ id: "building", label: "Building", traits: [{ trait: "wip", config: { limitSetting: "maxConcurrent" } }] },
|
||||||
|
{ id: "shipped", label: "Shipped", traits: [{ trait: "complete" }] },
|
||||||
|
],
|
||||||
|
} as unknown as WorkflowIr;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A VALID workflow that declares NO wip column at all — nothing can prove a card "advanced". */
|
||||||
|
function noWipIr(): WorkflowIr {
|
||||||
|
return {
|
||||||
|
version: "v2",
|
||||||
|
id: WF,
|
||||||
|
nodes: [],
|
||||||
|
edges: [],
|
||||||
|
columns: [
|
||||||
|
{ id: "inbox", label: "Inbox", traits: [{ trait: "intake" }] },
|
||||||
|
{ id: "drafting", label: "Drafting", traits: [{ trait: "hold", config: { release: "capacity" } }] },
|
||||||
|
{ id: "shipped", label: "Shipped", traits: [{ trait: "complete" }] },
|
||||||
|
],
|
||||||
|
} as unknown as WorkflowIr;
|
||||||
|
}
|
||||||
|
|
||||||
|
function makeTask(overrides: Partial<TaskDetail> = {}): TaskDetail {
|
||||||
|
return {
|
||||||
|
id: "FN-U8-COL",
|
||||||
|
title: "Execution policy column vocabulary",
|
||||||
|
description: "renamed-column coverage for the execution-policy branches",
|
||||||
|
column: "in-progress",
|
||||||
|
dependencies: [],
|
||||||
|
steps: [{ name: "Implement", status: "in-progress" }],
|
||||||
|
currentStep: 0,
|
||||||
|
log: [],
|
||||||
|
branch: "fusion/fn-u8-col",
|
||||||
|
baseBranch: "main",
|
||||||
|
worktree: "/tmp/fusion-fn-u8-col",
|
||||||
|
status: null,
|
||||||
|
error: null,
|
||||||
|
paused: false,
|
||||||
|
userPaused: false,
|
||||||
|
toolFailureDetectorLogCursor: 0,
|
||||||
|
autoMerge: true,
|
||||||
|
mergeRetries: 0,
|
||||||
|
createdAt: now,
|
||||||
|
updatedAt: now,
|
||||||
|
...overrides,
|
||||||
|
} as TaskDetail;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* `ir: undefined` models a workflow that cannot be resolved at all — the conservative path,
|
||||||
|
* which must keep the pre-conversion literal so a rebound is never stranded.
|
||||||
|
*/
|
||||||
|
function harness(options: {
|
||||||
|
ir: WorkflowIr | undefined;
|
||||||
|
task?: Partial<TaskDetail>;
|
||||||
|
settings?: Record<string, unknown>;
|
||||||
|
entries?: Array<{ type: string }>;
|
||||||
|
}) {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = makeTask(options.task);
|
||||||
|
const selection = { workflowId: WF, stepIds: [] };
|
||||||
|
store.getTask.mockResolvedValue(task);
|
||||||
|
store.getSettings.mockResolvedValue({
|
||||||
|
maxConcurrent: 2,
|
||||||
|
maxWorktrees: 4,
|
||||||
|
pollIntervalMs: 15_000,
|
||||||
|
autoMerge: true,
|
||||||
|
executorToolFailureRetryCount: 2,
|
||||||
|
executorToolFailureRetryBackoffMs: 0,
|
||||||
|
executorToolFailureThreshold: 3,
|
||||||
|
...options.settings,
|
||||||
|
});
|
||||||
|
const entries = options.entries ?? [];
|
||||||
|
store.getAgentLogCount = vi.fn().mockResolvedValue(entries.length);
|
||||||
|
store.getAgentLogs = vi.fn().mockResolvedValue(entries);
|
||||||
|
store.claimNextToolFailureRetry = vi.fn().mockResolvedValue({ outcome: "exhausted" });
|
||||||
|
store.markToolFailureRetryExhaustedAudit = vi.fn().mockResolvedValue(true);
|
||||||
|
store.recordRunAuditEvent = vi.fn().mockResolvedValue(undefined);
|
||||||
|
store.getTaskWorkflowSelection = vi.fn(() => selection);
|
||||||
|
store.getTaskWorkflowSelectionAsync = vi.fn(async () => selection);
|
||||||
|
store.getWorkflowDefinition = vi.fn(async () => (options.ir ? { id: WF, ir: options.ir } : null));
|
||||||
|
store.updateTask.mockImplementation(async (_id: string, patch: Partial<TaskDetail>) => Object.assign(task, patch));
|
||||||
|
store.updateTaskAtomic = vi.fn(async (_id: string, updater: (current: TaskDetail) => Partial<TaskDetail> | null) => {
|
||||||
|
const updates = updater(task);
|
||||||
|
if (updates) Object.assign(task, updates);
|
||||||
|
return task;
|
||||||
|
});
|
||||||
|
const executor = new TaskExecutor(store, "/tmp/test");
|
||||||
|
(executor as any).graphToolFailureRunCursors.set(task.id, 0);
|
||||||
|
return { executor, store, task };
|
||||||
|
}
|
||||||
|
|
||||||
|
const TOOL_ERRORS = [{ type: "tool_error" }, { type: "tool_error" }, { type: "tool_error" }];
|
||||||
|
|
||||||
|
function stepExecuteFailure() {
|
||||||
|
return {
|
||||||
|
disposition: "failed" as const,
|
||||||
|
outcome: "failure" as const,
|
||||||
|
visitedNodeIds: ["steps#0:step-execute"],
|
||||||
|
context: { "node:steps#0:step-execute:value": "failure" },
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function executeNodeFailure() {
|
||||||
|
return {
|
||||||
|
disposition: "failed" as const,
|
||||||
|
outcome: "failure" as const,
|
||||||
|
visitedNodeIds: ["execute"],
|
||||||
|
context: { "node:execute:value": "recoverable" },
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("FN-7998 node escalation targets the workflow's own hold column", () => {
|
||||||
|
beforeEach(() => {
|
||||||
|
resetExecutorMocks();
|
||||||
|
vi.useFakeTimers();
|
||||||
|
});
|
||||||
|
afterEach(() => vi.useRealTimers());
|
||||||
|
|
||||||
|
it("requeues a node escalation to the RENAMED hold column, not the literal todo", async () => {
|
||||||
|
const { executor, task } = harness({
|
||||||
|
ir: renamedIr(),
|
||||||
|
task: { column: "building" },
|
||||||
|
settings: { executorModelEscalationEnabled: true, executorEscalationNodeId: "cursor-node" },
|
||||||
|
entries: TOOL_ERRORS,
|
||||||
|
});
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, stepExecuteFailure());
|
||||||
|
|
||||||
|
/* The failure this pins: the escalated card was parked in a column id the workflow does not
|
||||||
|
declare, so nothing on the board could route it and the node escalation never ran. */
|
||||||
|
expect(task).toMatchObject({
|
||||||
|
nodeId: "cursor-node",
|
||||||
|
column: "drafting",
|
||||||
|
executorEscalationAttempted: true,
|
||||||
|
status: null,
|
||||||
|
error: null,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
it("keeps the literal todo for the default coding workflow (regression floor)", async () => {
|
||||||
|
const { executor, task } = harness({
|
||||||
|
ir: defaultIr(),
|
||||||
|
settings: { executorModelEscalationEnabled: true, executorEscalationNodeId: "cursor-node" },
|
||||||
|
entries: TOOL_ERRORS,
|
||||||
|
});
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, stepExecuteFailure());
|
||||||
|
|
||||||
|
expect(task).toMatchObject({ nodeId: "cursor-node", column: "todo", executorEscalationAttempted: true });
|
||||||
|
});
|
||||||
|
|
||||||
|
it("keeps the literal todo when the workflow cannot be resolved", async () => {
|
||||||
|
const { executor, task } = harness({
|
||||||
|
ir: undefined,
|
||||||
|
settings: { executorModelEscalationEnabled: true, executorEscalationNodeId: "cursor-node" },
|
||||||
|
entries: TOOL_ERRORS,
|
||||||
|
});
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, stepExecuteFailure());
|
||||||
|
|
||||||
|
expect(task).toMatchObject({ nodeId: "cursor-node", column: "todo" });
|
||||||
|
});
|
||||||
|
|
||||||
|
it("still does not move the card for a MODEL-target escalation", async () => {
|
||||||
|
/* Only the node target requeues; a model escalation retries in place. Resolving the hold
|
||||||
|
column must not turn the model path into a move. */
|
||||||
|
const { executor, task } = harness({
|
||||||
|
ir: renamedIr(),
|
||||||
|
task: { column: "building" },
|
||||||
|
settings: {
|
||||||
|
executorModelEscalationEnabled: true,
|
||||||
|
executorEscalationProvider: "anthropic",
|
||||||
|
executorEscalationModelId: "claude-sonnet",
|
||||||
|
},
|
||||||
|
entries: TOOL_ERRORS,
|
||||||
|
});
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, stepExecuteFailure());
|
||||||
|
await vi.advanceTimersByTimeAsync(0);
|
||||||
|
|
||||||
|
expect(task).toMatchObject({ modelProvider: "anthropic", modelId: "claude-sonnet", column: "building" });
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
describe("FN-7863/FN-7926 dispatch-loop gate matches the workflow's own hold column", () => {
|
||||||
|
beforeEach(() => {
|
||||||
|
resetExecutorMocks();
|
||||||
|
vi.useFakeTimers();
|
||||||
|
});
|
||||||
|
afterEach(() => vi.useRealTimers());
|
||||||
|
|
||||||
|
it("recognises an inner-executor requeue that landed in the RENAMED hold column", async () => {
|
||||||
|
/*
|
||||||
|
No in-process `executeNodeSelfRequeued` marker: this is the durable arm of the guard, the
|
||||||
|
one that survives a restart and is the ONLY evidence available to a second TaskExecutor
|
||||||
|
instance. Before the conversion the column read did not match, the run fell through to the
|
||||||
|
terminal sink, and a progressing card was parked failed.
|
||||||
|
*/
|
||||||
|
const { executor, store, task } = harness({ ir: renamedIr(), task: { column: "drafting" } });
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, executeNodeFailure());
|
||||||
|
|
||||||
|
expect(store.logEntry).toHaveBeenCalledWith(
|
||||||
|
task.id,
|
||||||
|
expect.stringContaining("executor recovery preserved"),
|
||||||
|
undefined,
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
expect(store.updateTask).not.toHaveBeenCalledWith(
|
||||||
|
task.id,
|
||||||
|
expect.objectContaining({ status: "failed" }),
|
||||||
|
expect.anything(),
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("still recognises the default workflow's todo (regression floor)", async () => {
|
||||||
|
const { executor, store, task } = harness({ ir: defaultIr(), task: { column: "todo" } });
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, executeNodeFailure());
|
||||||
|
|
||||||
|
expect(store.logEntry).toHaveBeenCalledWith(
|
||||||
|
task.id,
|
||||||
|
expect.stringContaining("executor recovery preserved"),
|
||||||
|
undefined,
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does NOT fire for a card sitting in the workflow's wip column — it terminalizes instead", async () => {
|
||||||
|
/*
|
||||||
|
The guard must stay narrow. A card still in the implementation column with no self-requeue
|
||||||
|
marker is a genuine execute failure and belongs to the terminal sink — widening the gate to
|
||||||
|
"any column" would swallow real failures as benign recoveries.
|
||||||
|
|
||||||
|
FNXC:WorkflowExecutionOwnership 2026-07-28-19:00 (U8, PR #2497 review — coderabbit):
|
||||||
|
Asserting only the ABSENCE of the recovery log cannot
|
||||||
|
distinguish "took the terminal path" from "did nothing and returned", which is precisely the
|
||||||
|
silent-swallow failure mode this whole file exists to catch. Assert the terminal outcome.
|
||||||
|
*/
|
||||||
|
const { executor, store, task } = harness({ ir: renamedIr(), task: { column: "building" } });
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, executeNodeFailure());
|
||||||
|
|
||||||
|
expect(store.logEntry).not.toHaveBeenCalledWith(
|
||||||
|
task.id,
|
||||||
|
expect.stringContaining("executor recovery preserved"),
|
||||||
|
undefined,
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
expect(task).toMatchObject({
|
||||||
|
status: "failed",
|
||||||
|
error: "Workflow graph terminated with failure at node 'execute'",
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
it("still honours the in-process self-requeue marker when the workflow cannot be resolved", async () => {
|
||||||
|
const { executor, store, task } = harness({ ir: undefined, task: { column: "in-progress" } });
|
||||||
|
/* markGraphExecuteSelfRequeued only records while the task is graph-routed. */
|
||||||
|
(executor as any).graphRouting.add(task.id);
|
||||||
|
(executor as any).markGraphExecuteSelfRequeued(task.id);
|
||||||
|
try {
|
||||||
|
await (executor as any).handleGraphFailure(task, executeNodeFailure());
|
||||||
|
} finally {
|
||||||
|
(executor as any).graphRouting.delete(task.id);
|
||||||
|
}
|
||||||
|
|
||||||
|
expect(store.logEntry).toHaveBeenCalledWith(
|
||||||
|
task.id,
|
||||||
|
expect.stringContaining("executor recovery preserved"),
|
||||||
|
undefined,
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:WorkflowExecutionOwnership 2026-07-28-14:20 (U8 / R3, PR #2497 review — greptile P1):
|
||||||
|
The first cut resolved PER ROLE with a literal fallback (`columns?.hold ?? "todo"`), which
|
||||||
|
conflated "no workflow resolved" with "this workflow declares no hold column". The second case
|
||||||
|
is a valid workflow, and substituting `todo` there invents a column it does not declare — which
|
||||||
|
node escalation then PERSISTS. These cases pin the per-workflow fallback: a resolved workflow
|
||||||
|
may only ever be given a column it actually declares, and a role it does not declare fails
|
||||||
|
closed toward a visible failure rather than a guess.
|
||||||
|
*/
|
||||||
|
describe("a resolved workflow is never given a column it does not declare", () => {
|
||||||
|
beforeEach(() => {
|
||||||
|
resetExecutorMocks();
|
||||||
|
vi.useFakeTimers();
|
||||||
|
});
|
||||||
|
afterEach(() => vi.useRealTimers());
|
||||||
|
|
||||||
|
it("node escalation requeues to the declared intake column when the workflow has no hold column", async () => {
|
||||||
|
const { executor, task } = harness({
|
||||||
|
ir: noHoldIr(),
|
||||||
|
task: { column: "building" },
|
||||||
|
settings: { executorModelEscalationEnabled: true, executorEscalationNodeId: "cursor-node" },
|
||||||
|
entries: TOOL_ERRORS,
|
||||||
|
});
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, stepExecuteFailure());
|
||||||
|
|
||||||
|
/* KTD-10 ordering falls through hold -> intake. `todo` is not in this workflow. */
|
||||||
|
expect(task).toMatchObject({ nodeId: "cursor-node", column: "inbox", executorEscalationAttempted: true });
|
||||||
|
expect(task.column).not.toBe("todo");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("the dispatch-loop gate does not match the literal todo for a workflow that lacks it", async () => {
|
||||||
|
/*
|
||||||
|
A card parked in `todo` under a workflow that does not declare it is a stray row, not
|
||||||
|
evidence that this workflow's inner executor self-requeued. Before the per-workflow
|
||||||
|
fallback the gate matched it and reported a benign recovery for a run that had none.
|
||||||
|
*/
|
||||||
|
const { executor, store, task } = harness({ ir: noHoldIr(), task: { column: "todo" } });
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, executeNodeFailure());
|
||||||
|
|
||||||
|
expect(store.logEntry).not.toHaveBeenCalledWith(
|
||||||
|
task.id,
|
||||||
|
expect.stringContaining("executor recovery preserved"),
|
||||||
|
undefined,
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
/*
|
||||||
|
FNXC:WorkflowExecutionOwnership 2026-07-28-19:05 (U8, PR #2497 review — coderabbit):
|
||||||
|
Assert the DISPOSITION so this cannot pass on a silent return — but the disposition here is
|
||||||
|
NOT a terminal park. An earlier classifier, the execution-resume router, owns this shape: a
|
||||||
|
recoverable execute failure with incomplete steps is routed for resume, and the card is left
|
||||||
|
in place because the router's own already-there check sees it. Asserting a failed park would
|
||||||
|
encode the wrong contract; what this case pins is that the DISPATCH-LOOP gate did not claim
|
||||||
|
it, and that the run reached a real classifier rather than falling off the end.
|
||||||
|
|
||||||
|
Noted while writing this: the resume router's log says "moved back to todo" and its
|
||||||
|
already-there check is another `"todo"` literal — one of the 20 sites in this method left to
|
||||||
|
U5's executor slice. It is why the card stays put here rather than being rehomed to `inbox`.
|
||||||
|
*/
|
||||||
|
expect(store.logEntry).toHaveBeenCalledWith(
|
||||||
|
task.id,
|
||||||
|
expect.stringContaining("execution resume"),
|
||||||
|
undefined,
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
expect(store.moveTask).not.toHaveBeenCalled();
|
||||||
|
expect(task.status).not.toBe("failed");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("a workflow with no wip column terminalizes visibly instead of claiming the card advanced", async () => {
|
||||||
|
/*
|
||||||
|
The fail-closed direction that matters. With no declared wip column there is no evidence the
|
||||||
|
card moved past implementation, so the "already advanced — no further action needed" branch
|
||||||
|
must NOT fire; the operator has to see the failure.
|
||||||
|
*/
|
||||||
|
const { executor, store, task } = harness({ ir: noWipIr(), task: { column: "drafting" } });
|
||||||
|
|
||||||
|
await (executor as any).handleGraphFailure(task, stepExecuteFailure());
|
||||||
|
|
||||||
|
expect(store.logEntry).not.toHaveBeenCalledWith(
|
||||||
|
task.id,
|
||||||
|
expect.stringContaining("already advanced"),
|
||||||
|
undefined,
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
expect(task).toMatchObject({
|
||||||
|
status: "failed",
|
||||||
|
error: "Workflow graph terminated with failure at node 'steps#0:step-execute'",
|
||||||
|
});
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -14,7 +14,7 @@ import { existsSync, lstatSync, realpathSync } from "node:fs";
|
|||||||
import { readFile, rm, writeFile } from "node:fs/promises";
|
import { readFile, rm, writeFile } from "node:fs/promises";
|
||||||
import type { TaskStore, Task, TaskDetail, TaskTokenUsage, StepStatus, Settings, WorkflowStep, MissionStore, AsyncMissionStore, Slice, AgentState, AgentCapability, RunMutationContext, AgentHeartbeatConfig, Agent, AgentMemoryInclusionMode, ProjectSettings, MergeResult, WorkflowIrNode, WorkflowIrNodeKind, WorkflowStepResult as CoreWorkflowStepResult, ThinkingLevel } from "@fusion/core";
|
import type { TaskStore, Task, TaskDetail, TaskTokenUsage, StepStatus, Settings, WorkflowStep, MissionStore, AsyncMissionStore, Slice, AgentState, AgentCapability, RunMutationContext, AgentHeartbeatConfig, Agent, AgentMemoryInclusionMode, ProjectSettings, MergeResult, WorkflowIrNode, WorkflowIrNodeKind, WorkflowStepResult as CoreWorkflowStepResult, ThinkingLevel } from "@fusion/core";
|
||||||
import { getUnmetSchedulingDependencies } from "./scheduler.js";
|
import { getUnmetSchedulingDependencies } from "./scheduler.js";
|
||||||
import { RetryStormError, serializeRetryStormError, evaluateCompletedPromotionFailureProvenance, evaluateSkipBypassTaint, resolveWorkflowIrForTask, evaluateForeachMergeProof, resolveCompleteColumn, resolveMergeOrchestrationColumn, resolveReboundTarget, resolveColumnAgentBinding, resolveEffectiveAgent, instanceNodeId, getWorkflowExtensionRegistry, getBuiltinWorkflow, parseNoOpCompletionMarker, allowsAutoMergeProcessing, resolveEffectiveAutoMerge, isLiveSharedBranchGroupMemberIntegration, resolveMaxAutoMergeRetries, resolveMaxConsecutiveToolFailureRetries, resolveConsecutiveToolFailureRetryBackoffMs, resolveConsecutiveToolFailureThreshold, resolveExecutorEscalationTarget, resolveOptionalStepRevisionBudget, resolveOptionalReviewRevisionBudget, DEFAULT_MAX_POST_REVIEW_FIXES, COMPLETION_SUMMARY_NODE_ID, upsertWorkflowStepResult, AWAITING_APPROVAL_PAUSE_REASON, THINKING_LEVELS, ACTIVE_WORKFLOW_WORK_ITEM_STATES, AgentStore, resolveExecutorFallbackModel } from "@fusion/core";
|
import { RetryStormError, serializeRetryStormError, evaluateCompletedPromotionFailureProvenance, evaluateSkipBypassTaint, resolveWorkflowIrForTask, evaluateForeachMergeProof, resolveCompleteColumn, resolveMergeOrchestrationColumn, resolveReboundTarget, resolveLifecycleColumns, resolveColumnAgentBinding, resolveEffectiveAgent, instanceNodeId, getWorkflowExtensionRegistry, getBuiltinWorkflow, parseNoOpCompletionMarker, allowsAutoMergeProcessing, resolveEffectiveAutoMerge, isLiveSharedBranchGroupMemberIntegration, resolveMaxAutoMergeRetries, resolveMaxConsecutiveToolFailureRetries, resolveConsecutiveToolFailureRetryBackoffMs, resolveConsecutiveToolFailureThreshold, resolveExecutorEscalationTarget, resolveOptionalStepRevisionBudget, resolveOptionalReviewRevisionBudget, DEFAULT_MAX_POST_REVIEW_FIXES, COMPLETION_SUMMARY_NODE_ID, upsertWorkflowStepResult, AWAITING_APPROVAL_PAUSE_REASON, THINKING_LEVELS, ACTIVE_WORKFLOW_WORK_ITEM_STATES, AgentStore, resolveExecutorFallbackModel } from "@fusion/core";
|
||||||
import { finalizeProvenAutoMergeTask } from "./auto-merge-finalization.js";
|
import { finalizeProvenAutoMergeTask } from "./auto-merge-finalization.js";
|
||||||
import { mergeEffectiveSettings } from "./effective-settings.js";
|
import { mergeEffectiveSettings } from "./effective-settings.js";
|
||||||
import { generateFeatureVideo, type GenerateFeatureVideoOptions } from "./review-artifacts/feature-video.js";
|
import { generateFeatureVideo, type GenerateFeatureVideoOptions } from "./review-artifacts/feature-video.js";
|
||||||
@@ -10777,8 +10777,60 @@ export class TaskExecutor {
|
|||||||
const failedNode = result.visitedNodeIds[result.visitedNodeIds.length - 1];
|
const failedNode = result.visitedNodeIds[result.visitedNodeIds.length - 1];
|
||||||
const mergeGraphFailure = this.isMergeGraphFailure(failedNode);
|
const mergeGraphFailure = this.isMergeGraphFailure(failedNode);
|
||||||
const failureValue = this.graphFailureValue(result);
|
const failureValue = this.graphFailureValue(result);
|
||||||
|
/*
|
||||||
|
FNXC:WorkflowExecutionOwnership 2026-07-28-09:40 (U8 / R3):
|
||||||
|
The execution-policy ladder below — the FN-7863/FN-7926 dispatch-loop gate, the FN-7996
|
||||||
|
tool-failure retry, and the FN-7998 escalation — decided the task's own lifecycle by
|
||||||
|
naming `"todo"` and `"in-progress"` literally. Under any workflow that renames those
|
||||||
|
columns the whole ladder was unreachable and its failure was SILENT in the worst
|
||||||
|
direction: the `live.column !== wip` guard below classified a card sitting in its own
|
||||||
|
implementation column as "already advanced — no further action needed", so the graph
|
||||||
|
failure was swallowed, no status was written, and the scheduler re-dispatched the same
|
||||||
|
doomed run. Nothing failed; the retry budgets, the escalation, and the bounded
|
||||||
|
terminalization simply never ran.
|
||||||
|
|
||||||
|
Resolve ONCE per failure and thread the pair through the ladder. One IR read per graph
|
||||||
|
failure: this is a terminal recovery path, not an enumeration loop.
|
||||||
|
|
||||||
|
FNXC:WorkflowExecutionOwnership 2026-07-28-14:05 (U8 / R3, PR #2497 review — greptile P1):
|
||||||
|
THE FALLBACK IS PER-WORKFLOW, NEVER PER-ROLE. The first cut wrote `columns?.hold ?? "todo"`,
|
||||||
|
which conflates two different situations: "no workflow could be resolved" and "this
|
||||||
|
workflow resolved fine and simply declares no hold column". Only the first justifies the
|
||||||
|
legacy literal. For the second, substituting `todo` invents a column the workflow does not
|
||||||
|
declare — and node-target escalation then PERSISTS it, stranding the card somewhere the
|
||||||
|
board cannot route and defeating the scheduler node re-resolution the escalation exists
|
||||||
|
for. U1 returns `undefined` per missing role precisely so a caller cannot borrow an
|
||||||
|
unrelated column; `?? "todo"` threw that guarantee away one line after asking for it.
|
||||||
|
|
||||||
|
So:
|
||||||
|
- IR unresolvable -> the legacy literals, i.e. exactly pre-conversion behavior.
|
||||||
|
- IR resolved -> `resolveReboundTarget` (KTD-10: hold -> intake -> first
|
||||||
|
column), which can only ever name a DECLARED column, and
|
||||||
|
`undefined` for wip when the workflow declares none.
|
||||||
|
|
||||||
|
A `wipColumn` of `undefined` is not a wildcard — every gate below treats "I cannot prove
|
||||||
|
where the wip column is" as "do not take the shortcut", so an unprovable card terminalizes
|
||||||
|
VISIBLY rather than being swallowed by the already-advanced branch. Fail closed toward the
|
||||||
|
operator seeing the failure.
|
||||||
|
|
||||||
|
The two literals that remain are ONLY the unresolvable-workflow fallback, and they are the
|
||||||
|
same pre-conversion values `resolveReboundColumnFor` already falls back to at its ~16
|
||||||
|
executor call sites — this adds no new rule and no new reachable-by-a-valid-workflow
|
||||||
|
literal. They are legacy-compat for a task whose workflow cannot be read at all, and they
|
||||||
|
belong to the same sweep that retires `resolveReboundColumnFor`'s own `?? "todo"` when U11
|
||||||
|
removes the column; they are deliberately NOT a per-role default, which is what made the
|
||||||
|
first cut wrong.
|
||||||
|
*/
|
||||||
|
let lifecycleIr: WorkflowIr | undefined;
|
||||||
|
try {
|
||||||
|
lifecycleIr = await resolveWorkflowIrForTask(this.store, task.id);
|
||||||
|
} catch {
|
||||||
|
lifecycleIr = undefined;
|
||||||
|
}
|
||||||
|
const wipColumn = lifecycleIr ? resolveLifecycleColumns(lifecycleIr)?.wip : "in-progress";
|
||||||
|
const holdColumn = lifecycleIr ? resolveReboundTarget(lifecycleIr) : "todo";
|
||||||
const executeNodeSelfRequeued = failedNode === "execute" && this.graphExecuteSelfRequeued.has(task.id);
|
const executeNodeSelfRequeued = failedNode === "execute" && this.graphExecuteSelfRequeued.has(task.id);
|
||||||
if (failedNode === "execute" && (live.column === "todo" || executeNodeSelfRequeued)) {
|
if (failedNode === "execute" && ((holdColumn !== undefined && live.column === holdColumn) || executeNodeSelfRequeued)) {
|
||||||
/*
|
/*
|
||||||
FNXC:WorkflowLifecycle 2026-06-23-12:03:
|
FNXC:WorkflowLifecycle 2026-06-23-12:03:
|
||||||
The graph execute node delegates to the authoritative executor. If that inner executor requeues the task to todo for self-heal/retry, the outer graph failure must not override it by parking the task in review.
|
The graph execute node delegates to the authoritative executor. If that inner executor requeues the task to todo for self-heal/retry, the outer graph failure must not override it by parking the task in review.
|
||||||
@@ -10877,7 +10929,14 @@ export class TaskExecutor {
|
|||||||
if (await this.routeGraphFailureToExecutionResume(live, failedNode ?? "unknown", failureValue)) {
|
if (await this.routeGraphFailureToExecutionResume(live, failedNode ?? "unknown", failureValue)) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
if (live.column !== "in-progress") {
|
/*
|
||||||
|
FNXC:WorkflowExecutionOwnership 2026-07-28-14:10 (U8 / R3, PR #2497 review):
|
||||||
|
`wipColumn === undefined` means the workflow declares no implementation column, so there
|
||||||
|
is no evidence the card "already advanced" past one. Swallowing the failure on a guess is
|
||||||
|
the exact silent-loss this conversion exists to remove — require a KNOWN wip column before
|
||||||
|
taking the benign shortcut.
|
||||||
|
*/
|
||||||
|
if (wipColumn !== undefined && live.column !== wipColumn) {
|
||||||
const benignMessage = `Workflow graph run ended after task already advanced to '${live.column}' — no further action needed`;
|
const benignMessage = `Workflow graph run ended after task already advanced to '${live.column}' — no further action needed`;
|
||||||
executorLog.log(`${task.id}: ${benignMessage}`);
|
executorLog.log(`${task.id}: ${benignMessage}`);
|
||||||
await this.store.logEntry(task.id, benignMessage, undefined, this.getRunContextFor(task.id));
|
await this.store.logEntry(task.id, benignMessage, undefined, this.getRunContextFor(task.id));
|
||||||
@@ -10947,7 +11006,7 @@ export class TaskExecutor {
|
|||||||
const message = `Workflow graph terminated with failure at node '${failedNode ?? "unknown"}'`;
|
const message = `Workflow graph terminated with failure at node '${failedNode ?? "unknown"}'`;
|
||||||
const settings = await this.store.getSettings();
|
const settings = await this.store.getSettings();
|
||||||
const maxToolFailureRetries = resolveMaxConsecutiveToolFailureRetries(settings);
|
const maxToolFailureRetries = resolveMaxConsecutiveToolFailureRetries(settings);
|
||||||
if (maxToolFailureRetries > 0 && isExecuteFamilyNode && !live.paused && !live.userPaused && !live.deletedAt && live.column === "in-progress") {
|
if (maxToolFailureRetries > 0 && isExecuteFamilyNode && !live.paused && !live.userPaused && !live.deletedAt && live.column === wipColumn) {
|
||||||
// Prefer the execution-local boundary; recovery paths refetch durable state rather than use the stale failure snapshot.
|
// Prefer the execution-local boundary; recovery paths refetch durable state rather than use the stale failure snapshot.
|
||||||
const cursor = this.graphToolFailureRunCursors.get(task.id) ?? (await this.store.getTask(task.id))?.toolFailureDetectorLogCursor;
|
const cursor = this.graphToolFailureRunCursors.get(task.id) ?? (await this.store.getTask(task.id))?.toolFailureDetectorLogCursor;
|
||||||
const threshold = resolveConsecutiveToolFailureThreshold(settings);
|
const threshold = resolveConsecutiveToolFailureThreshold(settings);
|
||||||
@@ -10957,7 +11016,7 @@ export class TaskExecutor {
|
|||||||
await this.store.updateTask(task.id, { status: null, error: null }, this.getRunContextFor(task.id));
|
await this.store.updateTask(task.id, { status: null, error: null }, this.getRunContextFor(task.id));
|
||||||
await this.store.logEntry(task.id, `Consecutive tool-call failures — auto-retrying same model (${claim.attempt}/${maxToolFailureRetries}) instead of parking`, undefined, this.getRunContextFor(task.id));
|
await this.store.logEntry(task.id, `Consecutive tool-call failures — auto-retrying same model (${claim.attempt}/${maxToolFailureRetries}) instead of parking`, undefined, this.getRunContextFor(task.id));
|
||||||
await this.store.recordRunAuditEvent?.({ taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("tool-failure-retry", task.id), domain: "database", mutationType: "task:execution-tool-failure-retry", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", attempt: claim.attempt, maxAttempts: maxToolFailureRetries, consecutiveToolFailures: threshold, mode: "same-model" } });
|
await this.store.recordRunAuditEvent?.({ taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("tool-failure-retry", task.id), domain: "database", mutationType: "task:execution-tool-failure-retry", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", attempt: claim.attempt, maxAttempts: maxToolFailureRetries, consecutiveToolFailures: threshold, mode: "same-model" } });
|
||||||
const schedule = () => { void (async () => { const resume = await this.store.getTask(task.id); if (resume && !resume.deletedAt && !resume.paused && !resume.userPaused && resume.column === "in-progress") await this.execute(resume); })().catch((error) => executorLog.error(`${task.id}: tool-failure retry failed`, error)); };
|
const schedule = () => { void (async () => { const resume = await this.store.getTask(task.id); if (resume && !resume.deletedAt && !resume.paused && !resume.userPaused && resume.column === wipColumn) await this.execute(resume); })().catch((error) => executorLog.error(`${task.id}: tool-failure retry failed`, error)); };
|
||||||
const delay = resolveConsecutiveToolFailureRetryBackoffMs(settings);
|
const delay = resolveConsecutiveToolFailureRetryBackoffMs(settings);
|
||||||
setTimeout(schedule, delay).unref?.();
|
setTimeout(schedule, delay).unref?.();
|
||||||
return;
|
return;
|
||||||
@@ -10969,7 +11028,21 @@ export class TaskExecutor {
|
|||||||
*/
|
*/
|
||||||
const escalationTarget = resolveExecutorEscalationTarget(settings);
|
const escalationTarget = resolveExecutorEscalationTarget(settings);
|
||||||
const hasModelTarget = escalationTarget.provider !== undefined && escalationTarget.modelId !== undefined;
|
const hasModelTarget = escalationTarget.provider !== undefined && escalationTarget.modelId !== undefined;
|
||||||
const hasNodeTarget = escalationTarget.nodeId !== undefined;
|
/*
|
||||||
|
FNXC:WorkflowExecutionOwnership 2026-07-28-14:15 (U8 / R3, PR #2497 review — greptile P1):
|
||||||
|
A node escalation is a REQUEUE: it parks the card back in the hold lane so the
|
||||||
|
scheduler re-resolves the effective node. Without a declared requeue target there is
|
||||||
|
nowhere legal to put it, and persisting an invented column is worse than not
|
||||||
|
escalating — the card lands where the board cannot route it and the node is never
|
||||||
|
dispatched. Degrade to the no-node-target shape (in-place retry, which is already how
|
||||||
|
an enabled escalation with no usable target behaves) rather than writing an
|
||||||
|
undeclared column.
|
||||||
|
*/
|
||||||
|
const nodeTargetRequeueColumn = escalationTarget.nodeId !== undefined ? holdColumn : undefined;
|
||||||
|
const hasNodeTarget = escalationTarget.nodeId !== undefined && nodeTargetRequeueColumn !== undefined;
|
||||||
|
if (escalationTarget.nodeId !== undefined && nodeTargetRequeueColumn === undefined) {
|
||||||
|
await this.store.logEntry(task.id, "Node escalation downgraded to an in-place retry — this task's workflow declares no column to requeue into", undefined, this.getRunContextFor(task.id));
|
||||||
|
}
|
||||||
let claimedEscalation = false;
|
let claimedEscalation = false;
|
||||||
let priorEscalationRetryCount = 0;
|
let priorEscalationRetryCount = 0;
|
||||||
/*
|
/*
|
||||||
@@ -10980,7 +11053,7 @@ export class TaskExecutor {
|
|||||||
*/
|
*/
|
||||||
await this.store.updateTaskAtomic(task.id, (current) => {
|
await this.store.updateTaskAtomic(task.id, (current) => {
|
||||||
const ownsFailureRun = current.toolFailureDetectorLogCursor === cursor
|
const ownsFailureRun = current.toolFailureDetectorLogCursor === cursor
|
||||||
&& current.column === "in-progress"
|
&& current.column === wipColumn
|
||||||
&& !current.paused
|
&& !current.paused
|
||||||
&& !current.userPaused
|
&& !current.userPaused
|
||||||
&& !current.deletedAt;
|
&& !current.deletedAt;
|
||||||
@@ -10989,7 +11062,7 @@ export class TaskExecutor {
|
|||||||
priorEscalationRetryCount = current.consecutiveToolFailureRetryCount ?? 0;
|
priorEscalationRetryCount = current.consecutiveToolFailureRetryCount ?? 0;
|
||||||
return {
|
return {
|
||||||
...(hasModelTarget ? { modelProvider: escalationTarget.provider, modelId: escalationTarget.modelId } : {}),
|
...(hasModelTarget ? { modelProvider: escalationTarget.provider, modelId: escalationTarget.modelId } : {}),
|
||||||
...(hasNodeTarget ? { nodeId: escalationTarget.nodeId, column: "todo" as const } : {}),
|
...(hasNodeTarget ? { nodeId: escalationTarget.nodeId, column: nodeTargetRequeueColumn } : {}),
|
||||||
executorEscalationAttempted: true,
|
executorEscalationAttempted: true,
|
||||||
/* FNXC:ExecutorEscalation 2026-07-16-22:40: Invalidate the exhausted run cursor before releasing the claim so concurrent stale handlers cannot park or audit the alternate execution; the alternate captures its own cursor at startup. */
|
/* FNXC:ExecutorEscalation 2026-07-16-22:40: Invalidate the exhausted run cursor before releasing the claim so concurrent stale handlers cannot park or audit the alternate execution; the alternate captures its own cursor at startup. */
|
||||||
toolFailureDetectorLogCursor: null,
|
toolFailureDetectorLogCursor: null,
|
||||||
@@ -11001,7 +11074,7 @@ export class TaskExecutor {
|
|||||||
await this.store.logEntry(task.id, "Same-model retries exhausted — escalating to alternate model/node (one attempt) instead of parking", undefined, this.getRunContextFor(task.id));
|
await this.store.logEntry(task.id, "Same-model retries exhausted — escalating to alternate model/node (one attempt) instead of parking", undefined, this.getRunContextFor(task.id));
|
||||||
await this.store.recordRunAuditEvent?.({ taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("escalation-retry", task.id), domain: "database", mutationType: "task:execution-escalation-retry", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", hasModelTarget, hasNodeTarget, priorConsecutiveToolFailureRetryCount: priorEscalationRetryCount } });
|
await this.store.recordRunAuditEvent?.({ taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("escalation-retry", task.id), domain: "database", mutationType: "task:execution-escalation-retry", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", hasModelTarget, hasNodeTarget, priorConsecutiveToolFailureRetryCount: priorEscalationRetryCount } });
|
||||||
if (!hasNodeTarget) {
|
if (!hasNodeTarget) {
|
||||||
const scheduleEscalation = () => { void (async () => { const resumeTask = await this.store.getTask(task.id); if (resumeTask && !resumeTask.deletedAt && !resumeTask.paused && !resumeTask.userPaused && resumeTask.column === "in-progress") await this.execute(resumeTask); })().catch((error) => executorLog.error(`${task.id}: escalation retry failed`, error)); };
|
const scheduleEscalation = () => { void (async () => { const resumeTask = await this.store.getTask(task.id); if (resumeTask && !resumeTask.deletedAt && !resumeTask.paused && !resumeTask.userPaused && resumeTask.column === wipColumn) await this.execute(resumeTask); })().catch((error) => executorLog.error(`${task.id}: escalation retry failed`, error)); };
|
||||||
const handle = setTimeout(scheduleEscalation, resolveConsecutiveToolFailureRetryBackoffMs(settings));
|
const handle = setTimeout(scheduleEscalation, resolveConsecutiveToolFailureRetryBackoffMs(settings));
|
||||||
handle.unref?.();
|
handle.unref?.();
|
||||||
}
|
}
|
||||||
@@ -11023,7 +11096,7 @@ export class TaskExecutor {
|
|||||||
await this.store.updateTaskAtomic(task.id, (current) => {
|
await this.store.updateTaskAtomic(task.id, (current) => {
|
||||||
if (
|
if (
|
||||||
current.toolFailureDetectorLogCursor !== cursor
|
current.toolFailureDetectorLogCursor !== cursor
|
||||||
|| current.column !== "in-progress"
|
|| current.column !== wipColumn
|
||||||
|| current.paused
|
|| current.paused
|
||||||
|| current.userPaused
|
|| current.userPaused
|
||||||
|| current.deletedAt
|
|| current.deletedAt
|
||||||
@@ -11065,7 +11138,7 @@ export class TaskExecutor {
|
|||||||
await this.store.updateTaskAtomic(task.id, (current) => {
|
await this.store.updateTaskAtomic(task.id, (current) => {
|
||||||
if (
|
if (
|
||||||
current.toolFailureDetectorLogCursor !== failureCursor
|
current.toolFailureDetectorLogCursor !== failureCursor
|
||||||
|| current.column !== "in-progress"
|
|| current.column !== wipColumn
|
||||||
|| current.paused
|
|| current.paused
|
||||||
|| current.userPaused
|
|| current.userPaused
|
||||||
|| current.deletedAt
|
|| current.deletedAt
|
||||||
|
|||||||
Reference in New Issue
Block a user