fix(engine): globalPause gates every graph node entry; maxWorktrees counts planning/review holders
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
7
.changeset/pause-gates-graph-and-worktree-ledger.md
Normal file
7
.changeset/pause-gates-graph-and-worktree-ledger.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
summary: Stop AI Engine now actually stops the workflow graph, and the worktree cap counts planning/review holders.
|
||||
category: fix
|
||||
dev: Two capacity-control regressions. (1) The graph interpreter never re-read settings, so globalPause did not stop node traversal — new Plan Review sessions started under pause; every node entry now polls an isPaused probe and suspends via the durable-continuation mechanism (reason "pause"), and the continuation drain refuses to dispatch while paused. (2) The scheduler's maxWorktrees ledger counted only WIP cards; under plan-in-place, planning/review lanes hold real worktrees, and the deleted global semaphore had been the accidental protection — the ledger now counts every non-terminal task holding a worktree.
|
||||
@@ -55,6 +55,75 @@ function invariantError() {
|
||||
);
|
||||
}
|
||||
|
||||
describe("workflow column boundary — global pause suspends at every node entry", () => {
|
||||
/*
|
||||
FNXC:EnginePause 2026-08-01-00:30:
|
||||
Operator regression: Stop AI Engine (globalPause) did not stop the graph — a live run started a
|
||||
fresh Plan Review model session two minutes after pause, because no node boundary ever re-read
|
||||
settings. These fail if the `isPaused` probe is removed from onNodeEntry.
|
||||
*/
|
||||
it("suspends with reason 'pause' before any move or node side effect when paused", async () => {
|
||||
const moveTask = vi.fn();
|
||||
const onSuspend = vi.fn();
|
||||
const boundary = createWorkflowColumnBoundary({
|
||||
taskId: "FN-PAUSE1",
|
||||
workflowId: "builtin:coding",
|
||||
ir: ir(),
|
||||
initialColumn: "in-review",
|
||||
moveTask,
|
||||
onSuspend,
|
||||
isPaused: async () => true,
|
||||
});
|
||||
|
||||
const result = await boundary.onNodeEntry(remediationNode());
|
||||
|
||||
expect(result).toMatchObject({
|
||||
kind: "suspended",
|
||||
reason: "pause",
|
||||
nodeId: "code-review-remediation",
|
||||
fromColumn: "in-review",
|
||||
});
|
||||
expect(onSuspend).toHaveBeenCalledTimes(1);
|
||||
// Pause must be a pure park: no move was attempted, the card stays put.
|
||||
expect(moveTask).not.toHaveBeenCalled();
|
||||
expect(boundary.currentColumn()).toBe("in-review");
|
||||
});
|
||||
|
||||
it("gates even a same-column node — each node can start a real session without moving the card", async () => {
|
||||
const onSuspend = vi.fn();
|
||||
const boundary = createWorkflowColumnBoundary({
|
||||
taskId: "FN-PAUSE2",
|
||||
workflowId: "builtin:coding",
|
||||
ir: ir(),
|
||||
initialColumn: "in-review",
|
||||
onSuspend,
|
||||
isPaused: () => true,
|
||||
});
|
||||
|
||||
const sameColumnNode = ir().nodes.find((n) => n.id === "code-review")!;
|
||||
const result = await boundary.onNodeEntry(sameColumnNode);
|
||||
|
||||
expect(result).toMatchObject({ kind: "suspended", reason: "pause", toColumn: "in-review" });
|
||||
expect(onSuspend).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("does not suspend when the probe reports unpaused", async () => {
|
||||
const moveTask = vi.fn().mockResolvedValue(undefined);
|
||||
const boundary = createWorkflowColumnBoundary({
|
||||
taskId: "FN-PAUSE3",
|
||||
workflowId: "builtin:coding",
|
||||
ir: ir(),
|
||||
initialColumn: "in-review",
|
||||
moveTask,
|
||||
isPaused: () => false,
|
||||
});
|
||||
|
||||
const result = await boundary.onNodeEntry(remediationNode());
|
||||
expect(result).toMatchObject({ kind: "entered" });
|
||||
expect(moveTask).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
|
||||
describe("workflow column boundary — capacity rejection on the remediation crossing", () => {
|
||||
it("parks (suspends) instead of failing the run when in-progress is at capacity", async () => {
|
||||
const moveTask = vi.fn().mockRejectedValue(capacityError());
|
||||
|
||||
@@ -2308,6 +2308,20 @@ export class InProcessRuntime
|
||||
*/
|
||||
private async drainWorkflowContinuations(): Promise<void> {
|
||||
if (this.workflowContinuationDrainActive || this.status !== "active") return;
|
||||
/*
|
||||
FNXC:EnginePause 2026-08-01-00:20:
|
||||
A pause-suspended run persists a runnable continuation (same mechanism as capacity). Without
|
||||
this gate the drain would re-dispatch it on the next tick and the graph would bounce
|
||||
suspend→dispatch→suspend forever while paused — and worse, dispatch genuinely new work under
|
||||
Stop AI Engine. Settings are re-read here (not event-driven) for the same reason as the
|
||||
boundary probe: the pause must bind even if `settings:updated` never reaches this instance.
|
||||
*/
|
||||
try {
|
||||
const settings = await this.taskStore.getSettings();
|
||||
if (settings.globalPause === true || settings.enginePaused === true) return;
|
||||
} catch {
|
||||
/* unreadable settings: proceed as before rather than wedging the pump */
|
||||
}
|
||||
this.workflowContinuationDrainActive = true;
|
||||
try {
|
||||
await drainDuePlanningContinuations({
|
||||
|
||||
@@ -701,6 +701,9 @@ function computeConcurrencyGateDiagnostic(params: {
|
||||
* semaphore gate uses this instead of only in-progress agentSlots.
|
||||
*/
|
||||
topLevelClaimedSlots?: number;
|
||||
/** FNXC:WorkflowScheduling 2026-07-31-23:50: every live worktree holder (wip + planning/review
|
||||
* lanes), so the maxWorktrees holders diagnostic names who actually occupies the slots. */
|
||||
worktreeHolderTaskIds?: string[];
|
||||
/** U6: additive per-column capacity gates (flag-ON only). Omitted → the legacy
|
||||
* three-gate report is byte-identical. */
|
||||
perColumnGates?: PerColumnCapacityGate[];
|
||||
@@ -758,7 +761,7 @@ function computeConcurrencyGateDiagnostic(params: {
|
||||
semaphoreGate,
|
||||
holders: {
|
||||
maxConcurrent: [...params.inProgressTaskIds],
|
||||
maxWorktrees: maxWorktreesGate ? [...params.inProgressTaskIds] : undefined,
|
||||
maxWorktrees: maxWorktreesGate ? [...(params.worktreeHolderTaskIds ?? params.inProgressTaskIds)] : undefined,
|
||||
semaphore: semaphoreGate ? [...params.inProgressTaskIds] : undefined,
|
||||
},
|
||||
// U6: additive only — present when flag-ON, omitted otherwise.
|
||||
@@ -2212,8 +2215,34 @@ export class Scheduler {
|
||||
const isReviewColumnTask = (task: Task): boolean =>
|
||||
isReviewColumnRole(columnFlagsForTask(task), task.column);
|
||||
const wipTaskIds = tasks.filter(isWipColumnTask).map((task) => task.id);
|
||||
let reservedWorktreeSlots = wipTaskIds.length;
|
||||
let reservedConcurrentSlots = reservedWorktreeSlots;
|
||||
/*
|
||||
FNXC:WorkflowScheduling 2026-07-31-23:50 (maxWorktrees counted only WIP — live board breach):
|
||||
Under plan-in-place EVERY lane's live card holds a real worktree — planning runs in the task
|
||||
worktree (triage.ts) and review/merge keeps it — but this ledger counted WIP cards only. The
|
||||
protection that used to catch the difference was the GLOBAL SEMAPHORE gate, whose FNXC below
|
||||
says exactly this ("must include every live top-level agent holder (planning triage and active
|
||||
in-review), otherwise the hold/release sweep can admit an executor on top of a full planner
|
||||
fleet"); the two-number capacity model deleted the semaphore, and this gate never learned to
|
||||
count planners. Observed live: maxWorktrees=4, four planning sessions each holding a worktree,
|
||||
and a replan dispatch admitted as the FIFTH worktree because the gate read used=0/4.
|
||||
|
||||
Count every non-terminal task that HOLDS a worktree (`task.worktree` set) in addition to WIP
|
||||
membership — wip cards without a worktree yet still reserve (they are about to acquire), and
|
||||
terminal lanes are excluded because their retained worktrees are cleanup-owned, not capacity.
|
||||
*/
|
||||
const isTerminalColumnTask = (task: Task): boolean => {
|
||||
const flags = columnFlagsForTask(task);
|
||||
if (flags) return flags.complete === true || flags.archived === true;
|
||||
return task.column === "done" || task.column === "archived";
|
||||
};
|
||||
const wipTaskIdSet = new Set(wipTaskIds);
|
||||
const nonWipWorktreeHolderIds = tasks
|
||||
.filter((task) => !wipTaskIdSet.has(task.id)
|
||||
&& !isTerminalColumnTask(task)
|
||||
&& typeof task.worktree === "string" && task.worktree.length > 0)
|
||||
.map((task) => task.id);
|
||||
let reservedWorktreeSlots = wipTaskIds.length + nonWipWorktreeHolderIds.length;
|
||||
let reservedConcurrentSlots = wipTaskIds.length;
|
||||
const inProgressTaskIds = wipTaskIds;
|
||||
const dispatchPrepByTaskId = new Map<string, {
|
||||
baseBranch: string | null;
|
||||
@@ -2828,6 +2857,7 @@ export class Scheduler {
|
||||
maxConcurrent,
|
||||
activeWorktrees: reservedWorktreeSlots,
|
||||
maxWorktrees,
|
||||
worktreeHolderTaskIds: [...inProgressTaskIds, ...nonWipWorktreeHolderIds],
|
||||
semaphore: this.options.semaphore,
|
||||
inProgressTaskIds,
|
||||
topLevelClaimedSlots,
|
||||
|
||||
@@ -61,6 +61,15 @@ export function createExecutorColumnBoundaryHooks(
|
||||
// KTD-3 drift-park loop fix (PR #2342): detectDrift clears the stale pin
|
||||
// row fields so an ordinary requeue re-resolves the CURRENT IR fresh.
|
||||
clearPin: pinPersistence.clearPin,
|
||||
/* FNXC:EnginePause 2026-08-01-00:20: settings re-read per node entry — event-independent. */
|
||||
isPaused: async () => {
|
||||
try {
|
||||
const settings = await store.getSettings();
|
||||
return settings.globalPause === true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
},
|
||||
onSuspend: async (suspension) => {
|
||||
const items = await store.listWorkflowWorkItemsForTask(task.id, { kinds: ["task"] });
|
||||
/* Only an ACTIVE row suppresses a fresh continuation. A cancelled/exhausted/manual-required
|
||||
|
||||
@@ -102,6 +102,17 @@ export interface WorkflowColumnBoundaryDeps {
|
||||
onWarn?: (message: string, detail: Record<string, unknown>) => void;
|
||||
/** Persist a durable continuation before control returns to the scheduler. */
|
||||
onSuspend?: (suspension: Extract<WorkflowColumnBoundaryEntryResult, { kind: "suspended" }>) => void | Promise<void>;
|
||||
/*
|
||||
FNXC:EnginePause 2026-08-01-00:20 (Stop AI Engine did not stop the graph):
|
||||
Operator-observed regression: with `globalPause: true` the graph runner kept crossing node
|
||||
boundaries — a live run started a NEW Plan Review step (fresh model session) two minutes after
|
||||
Stop AI Engine, and the plan-review→replan loop kept cycling "attempt N/unbounded". The legacy
|
||||
executor loop re-read settings between steps; the graph interpreter never did, and the
|
||||
event-driven abort listeners cannot be the only line of defense (they depend on
|
||||
`settings:updated` reaching this store instance). This probe is polled at EVERY node entry, so
|
||||
a pause takes effect at the next boundary even if no event ever fires.
|
||||
*/
|
||||
isPaused?: () => boolean | Promise<boolean>;
|
||||
}
|
||||
|
||||
/** The seam the graph executor consumes. */
|
||||
@@ -119,7 +130,7 @@ export type WorkflowColumnBoundaryEntryResult =
|
||||
| { kind: "entered" }
|
||||
| {
|
||||
kind: "suspended";
|
||||
reason: "capacity";
|
||||
reason: "capacity" | "pause";
|
||||
nodeId: string;
|
||||
fromColumn: string;
|
||||
toColumn: string;
|
||||
@@ -278,6 +289,36 @@ export function createWorkflowColumnBoundary(
|
||||
async onNodeEntry(node: WorkflowIrNode): Promise<WorkflowColumnBoundaryEntryResult> {
|
||||
const toColumn = node.column;
|
||||
|
||||
/*
|
||||
FNXC:EnginePause 2026-08-01-00:20:
|
||||
Pause gates EVERY node entry — columnless and same-column nodes included, because each node
|
||||
can start a real AI session regardless of whether the card moves. Suspend with the same
|
||||
durable-continuation mechanism capacity uses, so unpause resumes at exactly this node; the
|
||||
drain refuses to dispatch continuations while paused, which closes the resume loop.
|
||||
*/
|
||||
if (await deps.isPaused?.()) {
|
||||
const pauseSuspension = {
|
||||
kind: "suspended",
|
||||
reason: "pause",
|
||||
nodeId: node.id,
|
||||
fromColumn: column,
|
||||
toColumn: toColumn ?? column,
|
||||
irHash: computeWorkflowIrPin(deps.ir, node.id).irHash,
|
||||
} as const;
|
||||
await deps.onSuspend?.(pauseSuspension);
|
||||
emitWorkflowLifecycleEvent({
|
||||
type: "RunSuspended",
|
||||
taskId: deps.taskId,
|
||||
at: new Date().toISOString(),
|
||||
workflowId: deps.workflowId,
|
||||
nodeId: node.id,
|
||||
reason: "pause",
|
||||
fromColumn: column,
|
||||
toColumn: toColumn ?? column,
|
||||
});
|
||||
return pauseSuspension;
|
||||
}
|
||||
|
||||
/*
|
||||
FNXC:WorkflowEvents 2026-07-27-15:20 (U3 / R5, PR #2467 review):
|
||||
Announce the NODE ENTRY, not the column crossing — so this fires BEFORE the
|
||||
|
||||
@@ -291,7 +291,7 @@ export interface WorkflowGraphExecutorResult {
|
||||
context: Record<string, unknown>;
|
||||
visitedNodeIds: string[];
|
||||
suspended?: {
|
||||
reason: "capacity";
|
||||
reason: "capacity" | "pause";
|
||||
nodeId: string;
|
||||
fromColumn: string;
|
||||
toColumn: string;
|
||||
|
||||
@@ -177,6 +177,8 @@ export interface WorkflowColumnBoundaryHooks {
|
||||
clearPin?: () => void | Promise<void>;
|
||||
onWarn?: (message: string, detail: Record<string, unknown>) => void;
|
||||
onSuspend?: WorkflowColumnBoundaryDeps["onSuspend"];
|
||||
/** FNXC:EnginePause 2026-08-01-00:20: polled at every node entry (see boundary deps). */
|
||||
isPaused?: WorkflowColumnBoundaryDeps["isPaused"];
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -282,6 +284,7 @@ export class WorkflowGraphTaskRunner {
|
||||
clearPin: hooks.clearPin,
|
||||
onWarn: hooks.onWarn,
|
||||
onSuspend: hooks.onSuspend,
|
||||
isPaused: hooks.isPaused,
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user