fix(engine): enforce capacity for workflow continuations
Route durable planning and review continuations through shared project admission, retain reservations through execution, and prevent duplicate continuations from releasing another run's slot.
This commit is contained in:
7
.changeset/fair-workflow-continuation-capacity.md
Normal file
7
.changeset/fair-workflow-continuation-capacity.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
summary: Keep resumed planning and review workflows within the configured active worktree limit.
|
||||
category: fix
|
||||
dev: Routes durable workflow continuations through shared project admission without double-counting active task handoffs.
|
||||
@@ -0,0 +1,268 @@
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import type { Task, TaskStore, WorkflowWorkItem } from "@fusion/core";
|
||||
|
||||
import { projectAdmissionCoordinator } from "../concurrency.js";
|
||||
import {
|
||||
admitPlanningContinuation,
|
||||
createPlanningContinuationDispatcher,
|
||||
drainDuePlanningContinuations,
|
||||
} from "../runtimes/in-process-runtime.js";
|
||||
|
||||
const PROJECT_ID = "/test/workflow-continuation-capacity";
|
||||
const CONTINUATION_ID = "FN-CONTINUATION";
|
||||
|
||||
function task(id: string, patch: Partial<Task> = {}): Task {
|
||||
return {
|
||||
id,
|
||||
title: id,
|
||||
description: id,
|
||||
column: "todo",
|
||||
priority: "medium",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: "2026-08-01T00:00:00.000Z",
|
||||
updatedAt: "2026-08-01T00:00:00.000Z",
|
||||
...patch,
|
||||
} as Task;
|
||||
}
|
||||
|
||||
function store(
|
||||
tasks: Task[],
|
||||
settings: { maxConcurrent: number; maxWorktrees: number; worktreeLimitEnabled: boolean } = {
|
||||
maxConcurrent: 12,
|
||||
maxWorktrees: 9,
|
||||
worktreeLimitEnabled: true,
|
||||
},
|
||||
): TaskStore {
|
||||
return {
|
||||
getSettings: vi.fn(async () => settings),
|
||||
listTasks: vi.fn(async () => tasks),
|
||||
getTaskWorkflowSelection: vi.fn(() => undefined),
|
||||
getTaskWorkflowSelectionAsync: vi.fn(async () => undefined),
|
||||
getWorkflowDefinition: vi.fn(async () => undefined),
|
||||
} as unknown as TaskStore;
|
||||
}
|
||||
|
||||
const item = {
|
||||
id: "continuation-1",
|
||||
taskId: CONTINUATION_ID,
|
||||
nodeId: "plan-review",
|
||||
kind: "task",
|
||||
state: "runnable",
|
||||
waitReason: "planning",
|
||||
createdAt: "2026-08-01T00:00:00.000Z",
|
||||
} as WorkflowWorkItem;
|
||||
|
||||
afterEach(() => {
|
||||
projectAdmissionCoordinator.releaseReservation(CONTINUATION_ID);
|
||||
projectAdmissionCoordinator.releaseReservation("FN-CONTINUATION-2");
|
||||
projectAdmissionCoordinator.releaseReservation("FN-CONTINUATION-3");
|
||||
projectAdmissionCoordinator.releaseReservation("FN-BLOCKER");
|
||||
});
|
||||
|
||||
describe("workflow continuation active-slot admission", () => {
|
||||
it("does not start a tenth task when nine active tasks already hold the worktree budget", async () => {
|
||||
const active = Array.from({ length: 8 }, (_, index) =>
|
||||
task(`FN-PLAN-${index}`, { status: "planning" }),
|
||||
);
|
||||
active.push(task("FN-REVIEW", {
|
||||
workflowStepResults: [{
|
||||
workflowStepId: "code-review",
|
||||
workflowStepName: "Code Review",
|
||||
phase: "pre-merge",
|
||||
source: "optional-group",
|
||||
status: "pending",
|
||||
startedAt: "2026-08-01T00:00:00.000Z",
|
||||
}],
|
||||
}));
|
||||
const dispatch = vi.fn(async () => {});
|
||||
|
||||
const admitted = await admitPlanningContinuation({
|
||||
store: store([...active, task(CONTINUATION_ID)]),
|
||||
projectId: PROJECT_ID,
|
||||
task: task(CONTINUATION_ID),
|
||||
item,
|
||||
dispatch,
|
||||
});
|
||||
|
||||
expect(admitted).toBe(false);
|
||||
expect(dispatch).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("starts the ninth task when only eight active tasks hold slots", async () => {
|
||||
const active = Array.from({ length: 8 }, (_, index) =>
|
||||
task(`FN-PLAN-${index}`, { status: "planning" }),
|
||||
);
|
||||
const dispatch = vi.fn(async () => {});
|
||||
|
||||
const admitted = await admitPlanningContinuation({
|
||||
store: store([...active, task(CONTINUATION_ID)]),
|
||||
projectId: PROJECT_ID,
|
||||
task: task(CONTINUATION_ID),
|
||||
item,
|
||||
dispatch,
|
||||
});
|
||||
|
||||
expect(admitted).toBe(true);
|
||||
expect(dispatch).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("loads project capacity only after an earlier coordinator handoff settles", async () => {
|
||||
const eightActive = Array.from({ length: 8 }, (_, index) =>
|
||||
task(`FN-PLAN-${index}`, { status: "planning" }),
|
||||
);
|
||||
const ninthActive = task("FN-LANDED-HANDOFF", { status: "planning" });
|
||||
const continuation = task(CONTINUATION_ID);
|
||||
let liveTasks = [...eightActive, continuation];
|
||||
const taskStore = store(liveTasks);
|
||||
vi.mocked(taskStore.listTasks).mockImplementation(async () => liveTasks);
|
||||
const dispatch = vi.fn(async () => {});
|
||||
let releaseBlocker!: () => void;
|
||||
const blockerStarted = new Promise<void>((resolveStarted) => {
|
||||
void projectAdmissionCoordinator.admitOldest({
|
||||
projectId: PROJECT_ID,
|
||||
maxConcurrent: 9,
|
||||
claimed: () => 8,
|
||||
claimedTaskIds: () => eightActive.map((candidate) => candidate.id),
|
||||
refresh: async () => [{
|
||||
taskId: "FN-BLOCKER",
|
||||
projectId: PROJECT_ID,
|
||||
createdAt: "2026-07-31T23:59:59.000Z",
|
||||
start: async () => {
|
||||
resolveStarted();
|
||||
await new Promise<void>((resolve) => { releaseBlocker = resolve; });
|
||||
liveTasks = [...eightActive, ninthActive, continuation];
|
||||
projectAdmissionCoordinator.releaseReservation("FN-BLOCKER");
|
||||
},
|
||||
}],
|
||||
});
|
||||
});
|
||||
await blockerStarted;
|
||||
|
||||
const admission = admitPlanningContinuation({
|
||||
store: taskStore,
|
||||
projectId: PROJECT_ID,
|
||||
task: continuation,
|
||||
item,
|
||||
dispatch,
|
||||
});
|
||||
await Promise.resolve();
|
||||
expect(taskStore.listTasks).not.toHaveBeenCalled();
|
||||
releaseBlocker();
|
||||
const admitted = await admission;
|
||||
|
||||
expect(taskStore.listTasks).toHaveBeenCalledOnce();
|
||||
expect(admitted).toBe(false);
|
||||
expect(dispatch).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("ignores maxWorktrees when worktree limiting is disabled", async () => {
|
||||
const active = Array.from({ length: 2 }, (_, index) =>
|
||||
task(`FN-PLAN-${index}`, { status: "planning" }),
|
||||
);
|
||||
const continuation = task(CONTINUATION_ID);
|
||||
const dispatch = vi.fn(async () => {});
|
||||
|
||||
const admitted = await admitPlanningContinuation({
|
||||
store: store([...active, continuation], {
|
||||
maxConcurrent: 3,
|
||||
maxWorktrees: 1,
|
||||
worktreeLimitEnabled: false,
|
||||
}),
|
||||
projectId: PROJECT_ID,
|
||||
task: continuation,
|
||||
item,
|
||||
dispatch,
|
||||
});
|
||||
|
||||
expect(admitted).toBe(true);
|
||||
expect(dispatch).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("allows an already-active task to resume without claiming a second slot", async () => {
|
||||
const active = Array.from({ length: 8 }, (_, index) =>
|
||||
task(`FN-PLAN-${index}`, { status: "planning" }),
|
||||
);
|
||||
const continuing = task(CONTINUATION_ID, { status: "planning" });
|
||||
const dispatch = vi.fn(async () => {});
|
||||
|
||||
const admitted = await admitPlanningContinuation({
|
||||
store: store([...active, continuing]),
|
||||
projectId: PROJECT_ID,
|
||||
task: continuing,
|
||||
item,
|
||||
dispatch,
|
||||
});
|
||||
|
||||
expect(admitted).toBe(true);
|
||||
expect(dispatch).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("holds the final slot across a same-drain handoff until the resumed run settles", async () => {
|
||||
const active = Array.from({ length: 8 }, (_, index) =>
|
||||
task(`FN-PLAN-${index}`, { status: "planning" }),
|
||||
);
|
||||
const first = task(CONTINUATION_ID);
|
||||
const second = task("FN-CONTINUATION-2", { createdAt: "2026-08-01T00:00:01.000Z" });
|
||||
const tasks = [...active, first, second];
|
||||
const taskStore = store(tasks);
|
||||
const pendingResolvers: Array<() => void> = [];
|
||||
const execute = vi.fn(() => new Promise<void>((resolve) => pendingResolvers.push(resolve)));
|
||||
const items = [
|
||||
item,
|
||||
{ ...item, id: "continuation-2", taskId: second.id, createdAt: second.createdAt },
|
||||
] as WorkflowWorkItem[];
|
||||
|
||||
await drainDuePlanningContinuations({
|
||||
listDue: async () => items,
|
||||
getTask: async (taskId) => tasks.find((candidate) => candidate.id === taskId),
|
||||
cancelOrphan: async () => {},
|
||||
defer: async () => {},
|
||||
dispatch: createPlanningContinuationDispatcher({
|
||||
store: taskStore,
|
||||
projectId: PROJECT_ID,
|
||||
execute,
|
||||
}),
|
||||
nowMs: () => Date.now(),
|
||||
warn: () => {},
|
||||
});
|
||||
|
||||
expect(execute).toHaveBeenCalledOnce();
|
||||
expect(execute).toHaveBeenCalledWith(first);
|
||||
pendingResolvers[0]?.();
|
||||
await Promise.resolve();
|
||||
});
|
||||
|
||||
it("does not dispatch duplicate due continuations for the same inactive task", async () => {
|
||||
const active = Array.from({ length: 7 }, (_, index) =>
|
||||
task(`FN-PLAN-${index}`, { status: "planning" }),
|
||||
);
|
||||
const continuation = task(CONTINUATION_ID);
|
||||
const other = task("FN-CONTINUATION-2");
|
||||
const fourth = task("FN-CONTINUATION-3");
|
||||
const taskStore = store([...active, continuation, other, fourth]);
|
||||
const settles: Array<() => void> = [];
|
||||
const execute = vi.fn(() => new Promise<void>((resolve) => { settles.push(resolve); }));
|
||||
const dispatch = createPlanningContinuationDispatcher({
|
||||
store: taskStore,
|
||||
projectId: PROJECT_ID,
|
||||
execute,
|
||||
});
|
||||
|
||||
const admissions = await Promise.all([
|
||||
dispatch(continuation, item),
|
||||
dispatch(continuation, { ...item, id: "continuation-duplicate" }),
|
||||
]);
|
||||
expect(admissions).toEqual([true, true]);
|
||||
expect(execute).toHaveBeenCalledOnce();
|
||||
expect(await dispatch(other, { ...item, id: "continuation-other", taskId: other.id })).toBe(true);
|
||||
expect(execute).toHaveBeenCalledTimes(2);
|
||||
expect(await dispatch(fourth, { ...item, id: "continuation-fourth", taskId: fourth.id })).toBe(false);
|
||||
expect(execute).toHaveBeenCalledTimes(2);
|
||||
|
||||
settles.forEach((settle) => settle());
|
||||
await Promise.resolve();
|
||||
});
|
||||
});
|
||||
@@ -66,6 +66,11 @@ import { attachAgentLinkSync } from "../task-agent-sync.js";
|
||||
import { createRunAuditor, generateSyntheticRunId } from "../run-audit.js";
|
||||
import { setImmediate as setImmediateCb } from "node:timers";
|
||||
import { seedPreReleasePlanReviewContinuation } from "../plan-review-continuation.js";
|
||||
import {
|
||||
persistedTopLevelAgentTaskIdsFromStore,
|
||||
projectAdmissionCoordinator,
|
||||
resolveActiveTaskCapacityLimit,
|
||||
} from "../concurrency.js";
|
||||
|
||||
/*
|
||||
FNXC:WorkflowResolvedColumns 2026-07-31-14:40 (fleet — long-tail fallback arms):
|
||||
@@ -365,7 +370,9 @@ export interface DuePlanningContinuationDrainDeps {
|
||||
) => Promise<void>;
|
||||
/** `item` is passed only so the caller's failure log can keep naming the work
|
||||
* item verbatim; the extraction is otherwise a byte-for-byte body move. */
|
||||
dispatch: (task: Task, item: WorkflowWorkItem) => void;
|
||||
/** Return false when shared capacity rejected this item; later FIFO items
|
||||
* cannot fit either, so the bounded pass stops without repeating snapshots. */
|
||||
dispatch: (task: Task, item: WorkflowWorkItem) => boolean | void | Promise<boolean | void>;
|
||||
nowMs: () => number;
|
||||
warn: (message: string) => void;
|
||||
}
|
||||
@@ -420,10 +427,121 @@ export async function drainDuePlanningContinuations(
|
||||
const deferral = resolveParkedContinuationDeferral(resolved, deps.nowMs());
|
||||
if (deferral) await deps.defer(deferral);
|
||||
if (resolved.kind !== "actionable") continue;
|
||||
deps.dispatch(resolved.task, resolved.item);
|
||||
if (await deps.dispatch(resolved.task, resolved.item) === false) break;
|
||||
}
|
||||
}
|
||||
|
||||
const planningContinuationRuns = new Set<string>();
|
||||
|
||||
export async function admitPlanningContinuation(input: {
|
||||
store: TaskStore;
|
||||
projectId: string;
|
||||
task: Task;
|
||||
item: WorkflowWorkItem;
|
||||
dispatch: () => Promise<void>;
|
||||
}): Promise<boolean> {
|
||||
const runKey = `${input.projectId}:${input.task.id}`;
|
||||
// A task owns one top-level slot regardless of how many durable continuation
|
||||
// rows point at it. Treat a duplicate due row as already handled; admitting it
|
||||
// would attach two releasers to one task-keyed coordinator reservation.
|
||||
if (planningContinuationRuns.has(runKey)) return true;
|
||||
const settings = await input.store.getSettings();
|
||||
let selected = false;
|
||||
let duplicateHandled = false;
|
||||
const loadClaimSnapshot = async (): Promise<{ count: number; ids: string[] }> => {
|
||||
/*
|
||||
FNXC:WorkflowContinuationCapacity 2026-08-01-06:20:
|
||||
A dependency-cleared task continuation can resume directly in a same-column Plan Review node.
|
||||
That path does not cross the scheduler-owned hold→WIP boundary, so dispatching it directly let
|
||||
the new reviewer become a tenth live task while maxWorktrees was nine. Count the exact canonical
|
||||
live population (including pending workflow-step leases) and enter through the shared project
|
||||
coordinator before the continuation starts. Full rows are intentional here: slim task snapshots
|
||||
are not a contract for workflowStepResults, while a pending optional-step lease is a live agent.
|
||||
*/
|
||||
const tasks = await input.store.listTasks({ slim: false, includeArchived: false });
|
||||
const ids = await persistedTopLevelAgentTaskIdsFromStore(input.store, tasks);
|
||||
return { count: ids.length, ids };
|
||||
};
|
||||
// Resuming another node of an already-live task is a same-slot handoff, not a
|
||||
// new admission. Check only this fully hydrated task here; the project-wide
|
||||
// snapshot belongs inside the serialized coordinator drain below.
|
||||
const taskAlreadyActive = (await persistedTopLevelAgentTaskIdsFromStore(input.store, [input.task]))
|
||||
.includes(input.task.id);
|
||||
if (taskAlreadyActive) {
|
||||
void input.dispatch().catch(() => {});
|
||||
return true;
|
||||
}
|
||||
// This snapshot is intentionally created lazily inside the coordinator drain.
|
||||
// A prior lane may have been finishing its own handoff before this task's
|
||||
// turn; a pre-drain project snapshot can admit into its newly occupied slot.
|
||||
let admissionSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
|
||||
const getAdmissionSnapshot = () => admissionSnapshot ??= loadClaimSnapshot();
|
||||
await projectAdmissionCoordinator.admitOldest({
|
||||
projectId: input.projectId,
|
||||
maxConcurrent: resolveActiveTaskCapacityLimit({
|
||||
maxConcurrent: settings.maxConcurrent ?? 2,
|
||||
maxWorktrees: settings.maxWorktrees ?? 4,
|
||||
worktreeLimitEnabled: settings.worktreeLimitEnabled,
|
||||
}),
|
||||
claimed: async () => (await getAdmissionSnapshot()).count,
|
||||
claimedTaskIds: async () => (await getAdmissionSnapshot()).ids,
|
||||
refresh: async () => [{
|
||||
taskId: input.task.id,
|
||||
projectId: input.projectId,
|
||||
createdAt: input.item.createdAt ?? input.task.createdAt,
|
||||
start: async () => {
|
||||
// The preflight above is only a fast path. This serialized check is the
|
||||
// ownership authority when concurrent drains race the same durable row.
|
||||
if (planningContinuationRuns.has(runKey)) {
|
||||
duplicateHandled = true;
|
||||
// The coordinator's task-keyed Set already contains the ORIGINAL
|
||||
// run's reservation. Accept this no-op candidate so its decline path
|
||||
// cannot release capacity owned by that still-running workflow.
|
||||
return true;
|
||||
}
|
||||
selected = true;
|
||||
planningContinuationRuns.add(runKey);
|
||||
// Keep the coordinator reservation for the whole resumed run. The task
|
||||
// can remain canonically inactive until its first workflow node writes a
|
||||
// pending lease; releasing at executor entry recreates the over-cap gap.
|
||||
let run: Promise<void>;
|
||||
try {
|
||||
run = input.dispatch();
|
||||
} catch (error) {
|
||||
planningContinuationRuns.delete(runKey);
|
||||
throw error;
|
||||
}
|
||||
void run
|
||||
.finally(() => {
|
||||
planningContinuationRuns.delete(runKey);
|
||||
projectAdmissionCoordinator.releaseReservation(input.task.id);
|
||||
})
|
||||
.catch(() => {});
|
||||
},
|
||||
}],
|
||||
});
|
||||
return selected || duplicateHandled;
|
||||
}
|
||||
|
||||
export function createPlanningContinuationDispatcher(input: {
|
||||
store: TaskStore;
|
||||
projectId: string;
|
||||
execute: (task: Task) => Promise<void>;
|
||||
onError?: (task: Task, item: WorkflowWorkItem, error: unknown) => void;
|
||||
}): (task: Task, item: WorkflowWorkItem) => Promise<boolean> {
|
||||
return (task, item) => admitPlanningContinuation({
|
||||
store: input.store,
|
||||
projectId: input.projectId,
|
||||
task,
|
||||
item,
|
||||
dispatch: async () => {
|
||||
await input.execute(task).catch((error) => {
|
||||
input.onError?.(task, item, error);
|
||||
});
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* FNXC:WorkflowScheduling 2026-07-21-12:30:
|
||||
* Select due planning continuations whose task remains dispatchable.
|
||||
@@ -2354,11 +2472,14 @@ export class InProcessRuntime
|
||||
},
|
||||
cancelOrphan: (item, reason) => this.cancelOrphanedWorkflowWorkItem(item, reason),
|
||||
defer: (deferral) => this.deferParkedWorkflowWorkItem(deferral),
|
||||
dispatch: (task, item) => {
|
||||
void this.executor.execute(task).catch((error) => {
|
||||
dispatch: createPlanningContinuationDispatcher({
|
||||
store: this.taskStore,
|
||||
projectId: this.taskStore.getRootDir(),
|
||||
execute: (task) => this.executor.execute(task),
|
||||
onError: (_task, item, error) => {
|
||||
runtimeLog.error(`Workflow continuation ${item.id} failed:`, error);
|
||||
});
|
||||
},
|
||||
},
|
||||
}),
|
||||
nowMs: () => Date.now(),
|
||||
warn: (message) => runtimeLog.warn(message),
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user