diff --git a/.changeset/fair-workflow-continuation-capacity.md b/.changeset/fair-workflow-continuation-capacity.md new file mode 100644 index 0000000000..d5f7b5f9c7 --- /dev/null +++ b/.changeset/fair-workflow-continuation-capacity.md @@ -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. diff --git a/packages/engine/src/__tests__/workflow-continuation-capacity.test.ts b/packages/engine/src/__tests__/workflow-continuation-capacity.test.ts new file mode 100644 index 0000000000..4c81523f0a --- /dev/null +++ b/packages/engine/src/__tests__/workflow-continuation-capacity.test.ts @@ -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 { + 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((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((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((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((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(); + }); +}); diff --git a/packages/engine/src/runtimes/in-process-runtime.ts b/packages/engine/src/runtimes/in-process-runtime.ts index 69d9995fcb..236be8500d 100644 --- a/packages/engine/src/runtimes/in-process-runtime.ts +++ b/packages/engine/src/runtimes/in-process-runtime.ts @@ -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; /** `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; 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(); + +export async function admitPlanningContinuation(input: { + store: TaskStore; + projectId: string; + task: Task; + item: WorkflowWorkItem; + dispatch: () => Promise; +}): Promise { + 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; + 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; + onError?: (task: Task, item: WorkflowWorkItem, error: unknown) => void; +}): (task: Task, item: WorkflowWorkItem) => Promise { + 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), });