diff --git a/docs/workflows.md b/docs/workflows.md index 56c65eda7..f4bbf73cb 100644 --- a/docs/workflows.md +++ b/docs/workflows.md @@ -96,6 +96,12 @@ Completing one step can promote dependent steps to `ready`; those new ready vers up by the same durable bridge. This is what lets a multi-step DAG continue autonomously instead of stopping after each dependency boundary. +### Autonomous concurrency limits + +Before an autonomous ready step starts or a due wait resumes, the workflow store applies cluster-wide admission control under a PostgreSQL transaction advisory lock. At most **20** autonomous workflow steps may be `running` across one deployment, and at most **4** may be `running` for one Bot. + +This is an admission ceiling rather than a queue batch size, so multiple replicas cannot race past it. If capacity is full, the exact claimed work item is deferred with backoff and its claim attempt is rolled back; healthy work therefore does not exhaust its retry budget merely because the deployment stayed busy. The workflow step remains `ready` or `waiting` and can be admitted later. + ## Durable waits A running step can enter `waiting` with an exact future timestamp and an optional provider label. diff --git a/scripts/release-preflight.ts b/scripts/release-preflight.ts index ece479963..1b8810515 100644 --- a/scripts/release-preflight.ts +++ b/scripts/release-preflight.ts @@ -1498,6 +1498,7 @@ function checkDurableWorkflowState(): void { const migration = read("server/drizzle/0043_workflow_state.sql"); const ready = read("server/src/workflows/ready.ts"); const wake = read("server/src/workflows/wake.ts"); + const queue = read("server/src/work/queue.ts"); const server = read("server/src/index.ts"); const assetMigration = read("server/drizzle/0044_workflow_assets.sql"); const resumeMigration = read("server/drizzle/0045_workflow_resume_stamp.sql"); @@ -1539,6 +1540,10 @@ function checkDurableWorkflowState(): void { "greatest(", "date_trunc('milliseconds', now())", "${" + "input.waitUntil} > now()", + "MAX_CONCURRENT_WORKFLOW_STEPS = 20", + "MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT = 4", + "workflow-autonomous-capacity", + "WorkflowCapacityError", ]) { if (!store.includes(evidence)) { fail(`workflows: durable recovery boundary is missing ${evidence}`); @@ -1571,6 +1576,8 @@ function checkDurableWorkflowState(): void { "item.attempts >= maxAttempts", "Autonomous wait continuation exhausted its retry budget", "failWaitingStep(", + "WorkflowCapacityError", + "queue.defer", ]) { if (!wake.includes(evidence)) { fail(`workflows: durable wake bridge is missing ${evidence}`); @@ -1584,11 +1591,22 @@ function checkDurableWorkflowState(): void { "queue.release({", "item.attempts >= maxAttempts", "Autonomous continuation exhausted its retry budget", + "WorkflowCapacityError", + "queue.defer", ]) { if (!ready.includes(evidence)) { fail(`workflows: durable ready-step bridge is missing ${evidence}`); } } + for (const evidence of [ + "async defer({ kind, key, owner, delayMs, reason })", + "attempts: sql`greatest(", + "workItems.attempts} - 1, 0)", + ]) { + if (!queue.includes(evidence)) { + fail(`workflows: capacity-safe queue defer is missing ${evidence}`); + } + } if (!server.includes("sweepReadyWorkflowSteps(workflowReady)")) { fail("workflows: durable ready-step bridge is not running from the server"); } diff --git a/server/src/work/queue.ts b/server/src/work/queue.ts index 8676a9c12..14b7b8031 100644 --- a/server/src/work/queue.ts +++ b/server/src/work/queue.ts @@ -147,6 +147,21 @@ export type WorkQueue = { delayMs: number; reason?: string; }) => Promise; + /** + * Put claimed work back for later WITHOUT spending one of its retry attempts. + * + * This is only for an external admission condition such as a deployment concurrency ceiling: + * nothing failed, the worker simply was not allowed to begin. Claim increments attempts before + * the caller can discover that condition, so a normal release would eventually exhaust healthy + * work merely because the deployment stayed busy long enough. + */ + defer?: (input: { + kind: string; + key: string; + owner: string; + delayMs: number; + reason?: string; + }) => Promise; /** * Drop what is done with, older than the retention window. Returns how many went. * @@ -382,6 +397,27 @@ export function createWorkQueue(database: Database): WorkQueue { return Boolean(released); }, + async defer({ kind, key, owner, delayMs, reason }) { + /* + * Claim increments attempts before the worker can discover a shared-capacity refusal. Undo + * exactly that claim here, while the row is still leased to this owner, and push it out. The + * GREATEST guard is defensive against malformed/manual rows; attempts must never go negative. + */ + const [deferred] = await database + .update(workItems) + .set({ + claimedBy: null, + leaseUntil: null, + attempts: sql`greatest(${workItems.attempts} - 1, 0)`, + runAt: fromNow(delayMs), + updatedAt: sql`now()`, + ...(reason === undefined ? {} : { lastError: reason }), + }) + .where(ours(kind, key, owner)) + .returning({ key: workItems.key }); + return Boolean(deferred); + }, + async purge({ kind, olderThanMs, diff --git a/server/src/workflows/ready.ts b/server/src/workflows/ready.ts index 286521959..e1599c0a0 100644 --- a/server/src/workflows/ready.ts +++ b/server/src/workflows/ready.ts @@ -1,4 +1,4 @@ -import type { WorkflowStore } from "./store"; +import { WorkflowCapacityError, type WorkflowStore } from "./store"; import { DEFAULT_MAX_ATTEMPTS, type WorkQueue } from "../work/queue"; export const WORKFLOW_READY_DISPATCH_KIND = "workflow_ready_dispatch"; @@ -216,6 +216,23 @@ export async function dispatchClaimedReadyWorkflowSteps( ? error.message : "workflow ready step could not start"; + if (error instanceof WorkflowCapacityError) { + if (!options.queue.defer) { + throw new Error( + "work queue defer is unavailable for workflow capacity backoff", + ); + } + await options.queue.defer({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: item.key, + owner: options.owner, + delayMs: options.retryDelayMs ?? DEFAULT_RETRY_DELAY_MS, + reason, + }); + report.skipped.push({ workflowId, stepKey, reason }); + continue; + } + // A newer/manual start, pause/cancel, or changed ready stamp makes this exact queue item // permanently stale. Finishing it is safe because a later ready transition gets a new stamp. if ( diff --git a/server/src/workflows/store.ts b/server/src/workflows/store.ts index 5e951180f..2a561b479 100644 --- a/server/src/workflows/store.ts +++ b/server/src/workflows/store.ts @@ -28,6 +28,13 @@ export const MAX_WORKFLOW_FAILURE_CODE_POINTS = 500; export const MAX_WORKFLOW_PROVIDER_CODE_POINTS = 120; export const MAX_WORKFLOW_ASSET_REF_CODE_POINTS = 512; export const MAX_WORKFLOW_ASSET_LABEL_CODE_POINTS = 200; +/** + * Unattended workflow turns need hard cluster-wide ceilings just like routines. These are + * admission caps, not queue batch sizes: every replica serializes the count + state transition + * under one PostgreSQL advisory lock before a ready/waiting step becomes running. + */ +export const MAX_CONCURRENT_WORKFLOW_STEPS = 20; +export const MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT = 4; const ACTIVE_STEP_STATUSES = [ "blocked", @@ -156,6 +163,13 @@ export class WorkflowRefusedError extends Error { } } +export class WorkflowCapacityError extends WorkflowRefusedError { + constructor(message: string) { + super(message); + this.name = "WorkflowCapacityError"; + } +} + type Transaction = Parameters[0]>[0]; type Handle = Database | Transaction; @@ -384,6 +398,43 @@ export function createWorkflowStore(database: Database): WorkflowStore { ); } + async function admitAutonomousStep( + transaction: Transaction, + agentId: string, + ): Promise { + /* + * Count + admission is one serialized decision across every API replica. Row locks cannot + * protect "how many rows are running", so the advisory lock is the cluster-wide mutex for this + * small critical section. It is taken only after the exact workflow's own lock, and no path + * takes these two locks in the reverse order. + */ + await transaction.execute( + sql`select pg_advisory_xact_lock(hashtext('workflow-autonomous-capacity'))`, + ); + + const [counts] = await transaction + .select({ + deployment: sql`count(*)::int`, + agent: sql`count(*) filter (where ${workflowRuns.agentId} = ${agentId})::int`, + }) + .from(workflowSteps) + .innerJoin(workflowRuns, eq(workflowRuns.id, workflowSteps.workflowId)) + .where(eq(workflowSteps.status, "running")); + + const deploymentRunning = counts?.deployment ?? 0; + const agentRunning = counts?.agent ?? 0; + if (deploymentRunning >= MAX_CONCURRENT_WORKFLOW_STEPS) { + throw new WorkflowCapacityError( + `Workflow capacity is full: this deployment already has ${MAX_CONCURRENT_WORKFLOW_STEPS} autonomous steps running.`, + ); + } + if (agentRunning >= MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT) { + throw new WorkflowCapacityError( + `Workflow capacity is full: this Bot already has ${MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT} autonomous steps running.`, + ); + } + } + async function resolveChannel(input: WorkflowInput): Promise { if (input.channelId !== undefined) { const [row] = await database @@ -890,6 +941,8 @@ export function createWorkflowStore(database: Database): WorkflowStore { ); } + await admitAutonomousStep(transaction, identity.agentId); + const [row] = await transaction .update(workflowSteps) .set({ @@ -1025,6 +1078,9 @@ export function createWorkflowStore(database: Database): WorkflowStore { ) { return toStep(current); } + + await admitAutonomousStep(transaction, identity.agentId); + const [row] = await transaction .update(workflowSteps) .set({ diff --git a/server/src/workflows/wake.ts b/server/src/workflows/wake.ts index 37c8bd871..fa8c9e562 100644 --- a/server/src/workflows/wake.ts +++ b/server/src/workflows/wake.ts @@ -1,4 +1,4 @@ -import type { WorkflowStore } from "./store"; +import { WorkflowCapacityError, type WorkflowStore } from "./store"; import { DEFAULT_MAX_ATTEMPTS, type WorkQueue } from "../work/queue"; export const WORKFLOW_WAIT_RESUME_KIND = "workflow_wait_resume"; @@ -229,6 +229,23 @@ export async function dispatchClaimedWorkflowWaits( ? error.message : "workflow wait could not resume"; + if (error instanceof WorkflowCapacityError) { + if (!options.queue.defer) { + throw new Error( + "work queue defer is unavailable for workflow capacity backoff", + ); + } + await options.queue.defer({ + kind: WORKFLOW_WAIT_RESUME_KIND, + key: item.key, + owner: options.owner, + delayMs: options.retryDelayMs ?? DEFAULT_RETRY_DELAY_MS, + reason, + }); + report.skipped.push({ workflowId, stepKey, reason }); + continue; + } + // A changed/cancelled/stale wait is final for this exact timestamp. The // store's compare-and-set refusal is what makes an old queue item harmless. if ( diff --git a/server/tests/work-queue.integration.test.ts b/server/tests/work-queue.integration.test.ts index cb0cef984..c076b0abf 100644 --- a/server/tests/work-queue.integration.test.ts +++ b/server/tests/work-queue.integration.test.ts @@ -209,6 +209,33 @@ describe("claiming durable work", () => { ).toHaveLength(0); }); + test("deferring gives a capacity-blocked claim back without spending an attempt", async () => { + await queue.offer({ kind, key: "capacity" }); + const [first] = await queue.claim({ + kind, + owner: "replica-1", + leaseMs: 30_000, + }); + expect(first?.attempts).toBe(1); + + expect( + await queue.defer?.({ + kind, + key: "capacity", + owner: "replica-1", + delayMs: 0, + reason: "capacity full", + }), + ).toBe(true); + + const [second] = await queue.claim({ + kind, + owner: "replica-2", + leaseMs: 30_000, + }); + expect(second?.attempts).toBe(1); + }); + /* * Both of these are the same bug from two ends: the lease says who may act on an item, and only * `renew` used to ask. A replica whose lease had quietly gone could delete or reschedule work diff --git a/server/tests/workflow-ready.test.ts b/server/tests/workflow-ready.test.ts index 621e0c177..bb8c30a1c 100644 --- a/server/tests/workflow-ready.test.ts +++ b/server/tests/workflow-ready.test.ts @@ -5,6 +5,7 @@ import { offerReadyWorkflowSteps, WORKFLOW_READY_DISPATCH_KIND, } from "../src/workflows/ready"; +import { WorkflowCapacityError } from "../src/workflows/store"; function queueStub(overrides: Partial = {}): WorkQueue { return { @@ -146,6 +147,41 @@ describe("workflow ready-step dispatch", () => { expect(renewals).toBeGreaterThan(1); }); + test("defers capacity-blocked ready work without spending a retry", async () => { + let deferred = 0; + let released = 0; + const report = await dispatchClaimedReadyWorkflowSteps({ + owner: "worker-1", + retryDelayMs: 9_000, + queue: queueStub({ + claim: async () => [item(4)], + defer: async (input) => { + deferred += 1; + expect(input.delayMs).toBe(9_000); + return true; + }, + release: async () => { + released += 1; + return true; + }, + }), + store: { + readySteps: async () => [], + startReadyStep: async () => { + throw new WorkflowCapacityError("workflow capacity is full"); + }, + failStep: async () => { + throw new Error("not used"); + }, + }, + }); + + expect(deferred).toBe(1); + expect(released).toBe(0); + expect(report.started).toEqual([]); + expect(report.skipped[0]?.reason).toContain("capacity"); + }); + test("releases a transient dispatch failure so the exact CAS can be retried", async () => { let released = 0; await dispatchClaimedReadyWorkflowSteps({ diff --git a/server/tests/workflow-store.integration.test.ts b/server/tests/workflow-store.integration.test.ts index 9045b1027..1ab5aad7d 100644 --- a/server/tests/workflow-store.integration.test.ts +++ b/server/tests/workflow-store.integration.test.ts @@ -19,6 +19,9 @@ import { } from "../src/db/schema"; import { createWorkflowStore, + MAX_CONCURRENT_WORKFLOW_STEPS, + MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT, + WorkflowCapacityError, WorkflowNotFoundError, WorkflowRefusedError, } from "../src/workflows/store"; @@ -366,6 +369,120 @@ describe("durable ready-step dispatch", () => { ).rejects.toThrow(/another attempt/); }); + test("caps autonomous running steps per Bot without mutating the refused ready step", async () => { + const { owner, agentId, channel } = await setUp(); + const who = identity(owner, agentId); + const plans = []; + + for ( + let index = 0; + index <= MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT; + index += 1 + ) { + plans.push( + await store.create({ + ...planInput(owner, agentId, channel.id), + title: `Capacity ${index}`, + }), + ); + } + + for (const plan of plans.slice( + 0, + MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT, + )) { + const queued = (await store.readySteps(500)).find( + (candidate) => + candidate.workflowId === plan.id && candidate.stepKey === "script", + ); + await store.startReadyStep( + who, + plan.id, + "script", + queued?.readyAt as Date, + 0, + ); + } + + // Pausing the workflow does not stop the already-running headless turn, so it must + // continue occupying capacity until that step actually checkpoints. + await store.pause(who, plans[0]!.id); + + const refused = plans[MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT]!; + const queued = (await store.readySteps(500)).find( + (candidate) => + candidate.workflowId === refused.id && candidate.stepKey === "script", + ); + await expect( + store.startReadyStep( + who, + refused.id, + "script", + queued?.readyAt as Date, + 0, + ), + ).rejects.toBeInstanceOf(WorkflowCapacityError); + + const unchanged = await store.get(who, refused.id); + expect(unchanged?.steps[0]?.status).toBe("ready"); + expect(unchanged?.steps[0]?.attempts).toBe(0); + }); + + test("caps autonomous workflow steps across Bots in the deployment", async () => { + const owner = await createUser(); + const agentIds = []; + for (let index = 0; index < 6; index += 1) { + agentIds.push(await createAgent(owner, `Capacity Bot ${index}`)); + } + const channel = await createChannel(owner, agentIds); + let started = 0; + + for (const agentId of agentIds.slice(0, 5)) { + for ( + let index = 0; + index < MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT; + index += 1 + ) { + const plan = await store.create({ + ...planInput(owner, agentId, channel.id), + title: `Deployment capacity ${started}`, + }); + const queued = (await store.readySteps(500)).find( + (candidate) => + candidate.workflowId === plan.id && candidate.stepKey === "script", + ); + await store.startReadyStep( + identity(owner, agentId), + plan.id, + "script", + queued?.readyAt as Date, + 0, + ); + started += 1; + } + } + expect(started).toBe(MAX_CONCURRENT_WORKFLOW_STEPS); + + const extraAgent = agentIds[5]!; + const extra = await store.create({ + ...planInput(owner, extraAgent, channel.id), + title: "Deployment capacity overflow", + }); + const queued = (await store.readySteps(500)).find( + (candidate) => + candidate.workflowId === extra.id && candidate.stepKey === "script", + ); + await expect( + store.startReadyStep( + identity(owner, extraAgent), + extra.id, + "script", + queued?.readyAt as Date, + 0, + ), + ).rejects.toBeInstanceOf(WorkflowCapacityError); + }); + test("pause and resume gives still-ready work a fresh deterministic key", async () => { const { owner, agentId, channel } = await setUp(); const who = identity(owner, agentId); diff --git a/server/tests/workflow-wake.test.ts b/server/tests/workflow-wake.test.ts index ed0a14fd7..36c01ce75 100644 --- a/server/tests/workflow-wake.test.ts +++ b/server/tests/workflow-wake.test.ts @@ -5,6 +5,7 @@ import { offerDueWorkflowWaits, WORKFLOW_WAIT_RESUME_KIND, } from "../src/workflows/wake"; +import { WorkflowCapacityError } from "../src/workflows/store"; function queueStub(overrides: Partial = {}): WorkQueue { return { @@ -216,6 +217,58 @@ describe("workflow wake bridge", () => { expect(renewals).toBeGreaterThan(1); }); + test("defers capacity-blocked wait wakes without spending a retry", async () => { + let deferred = 0; + let released = 0; + const report = await dispatchClaimedWorkflowWaits({ + owner: "worker-1", + retryDelayMs: 11_000, + queue: queueStub({ + claim: async () => [ + { + kind: WORKFLOW_WAIT_RESUME_KIND, + key: "wake-capacity", + attempts: 4, + payload: { + ownerUserId: "user-1", + agentId: "bot-1", + workflowId: "workflow-1", + stepKey: "render", + waitUntil: "2026-09-20T07:30:00.000Z", + attempts: 2, + }, + }, + ], + defer: async (input) => { + deferred += 1; + expect(input.delayMs).toBe(11_000); + return true; + }, + release: async () => { + released += 1; + return true; + }, + }), + store: { + dueWaitingSteps: async () => [], + resumeWaitingStep: async () => { + throw new WorkflowCapacityError("workflow capacity is full"); + }, + failWaitingStep: async () => { + throw new Error("not used"); + }, + failStep: async () => { + throw new Error("not used"); + }, + }, + }); + + expect(deferred).toBe(1); + expect(released).toBe(0); + expect(report.resumed).toEqual([]); + expect(report.skipped[0]?.reason).toContain("capacity"); + }); + test("releases the same wake when headless dispatch fails", async () => { let released = 0; await dispatchClaimedWorkflowWaits({