From d7954c04ff165e68162f098fd4506b5964654f69 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:16:49 +0700 Subject: [PATCH 01/18] feat(queue): defer capacity-blocked work without retries --- server/src/work/queue.ts | 36 ++++++++++++++++++++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/server/src/work/queue.ts b/server/src/work/queue.ts index 8676a9c12..aa665bec8 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, From ca5509985c164eef233479ce17430abe9ac8ef85 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:17:35 +0700 Subject: [PATCH 02/18] security(workflows): cap autonomous concurrent steps --- server/src/workflows/store.ts | 61 +++++++++++++++++++++++++++++++++++ 1 file changed, 61 insertions(+) diff --git a/server/src/workflows/store.ts b/server/src/workflows/store.ts index 5e951180f..d39bed672 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,48 @@ 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( + and( + eq(workflowRuns.status, "active"), + 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 +946,8 @@ export function createWorkflowStore(database: Database): WorkflowStore { ); } + await admitAutonomousStep(transaction, identity.agentId); + const [row] = await transaction .update(workflowSteps) .set({ @@ -1025,6 +1083,9 @@ export function createWorkflowStore(database: Database): WorkflowStore { ) { return toStep(current); } + + await admitAutonomousStep(transaction, identity.agentId); + const [row] = await transaction .update(workflowSteps) .set({ From 8ac2b25ace7892dcbd92f82fffc41a98d304062d Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:18:00 +0700 Subject: [PATCH 03/18] refactor(queue): keep defer compatible with queue test doubles --- server/src/work/queue.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/server/src/work/queue.ts b/server/src/work/queue.ts index aa665bec8..14b7b8031 100644 --- a/server/src/work/queue.ts +++ b/server/src/work/queue.ts @@ -155,7 +155,7 @@ export type WorkQueue = { * 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: { + defer?: (input: { kind: string; key: string; owner: string; From c8acd87d63a6dd9a434ac0833c149c8054e3717c Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:18:12 +0700 Subject: [PATCH 04/18] fix(workflows): defer ready work when capacity is full --- server/src/workflows/ready.ts | 19 ++++++++++++++++++- 1 file changed, 18 insertions(+), 1 deletion(-) 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 ( From 0f5a581ed3b9ddd1f6b5e977e83d651505c24e4b Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:18:15 +0700 Subject: [PATCH 05/18] fix(workflows): defer waits when autonomous capacity is full --- server/src/workflows/wake.ts | 19 ++++++++++++++++++- 1 file changed, 18 insertions(+), 1 deletion(-) 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 ( From 086666fc77be748857ab02dd8e4ec0733e33d1a2 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:18:54 +0700 Subject: [PATCH 06/18] test(workflows): defer capacity-blocked ready dispatches --- server/tests/workflow-ready.test.ts | 36 +++++++++++++++++++++++++++++ 1 file changed, 36 insertions(+) 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({ From 86b0f2141f71ce57f2f2398cb986948dd7f10ca8 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:18:59 +0700 Subject: [PATCH 07/18] test(workflows): defer capacity-blocked wait wakes --- server/tests/workflow-wake.test.ts | 53 ++++++++++++++++++++++++++++++ 1 file changed, 53 insertions(+) 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({ From 68a8964eb49bce8c3727cfd1677a4af0ff33fdcd Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:19:41 +0700 Subject: [PATCH 08/18] test(queue): defer without consuming retry budget --- server/tests/work-queue.integration.test.ts | 27 +++++++++++++++++++++ 1 file changed, 27 insertions(+) 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 From 3965833d3b3bb1330d5d309a3d5fe188052aa693 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:19:44 +0700 Subject: [PATCH 09/18] test(workflows): enforce autonomous concurrency ceilings --- .../tests/workflow-store.integration.test.ts | 102 ++++++++++++++++++ 1 file changed, 102 insertions(+) diff --git a/server/tests/workflow-store.integration.test.ts b/server/tests/workflow-store.integration.test.ts index 9045b1027..b5ad46d00 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,105 @@ 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, + ); + } + + 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); From 56be595ded1d48a66b1cee575b72bcf396fb1d7b Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:20:22 +0700 Subject: [PATCH 10/18] test(release): pin workflow concurrency admission invariants --- scripts/release-preflight.ts | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/scripts/release-preflight.ts b/scripts/release-preflight.ts index ece479963..bd7754461 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"); @@ -1584,11 +1585,21 @@ 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 })", + "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"); } From 3bcb3d761d3bbf40c05537d330661ae4798f6861 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:20:46 +0700 Subject: [PATCH 11/18] test(release): pin workflow concurrency admission invariants --- scripts/release-preflight.ts | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/scripts/release-preflight.ts b/scripts/release-preflight.ts index bd7754461..4c59b85c5 100644 --- a/scripts/release-preflight.ts +++ b/scripts/release-preflight.ts @@ -1540,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}`); @@ -1572,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}`); From fb229ce1896395225e1905d9fc6076e0eda2e598 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:20:51 +0700 Subject: [PATCH 12/18] docs(workflows): document autonomous concurrency ceilings --- docs/workflows.md | 6 ++++++ 1 file changed, 6 insertions(+) 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. From 5d51e2f49412f140293ddf391521d8b1c5444e95 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:21:44 +0700 Subject: [PATCH 13/18] style(tests): format workflow capacity coverage --- server/tests/workflow-store.integration.test.ts | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) diff --git a/server/tests/workflow-store.integration.test.ts b/server/tests/workflow-store.integration.test.ts index b5ad46d00..8f782855d 100644 --- a/server/tests/workflow-store.integration.test.ts +++ b/server/tests/workflow-store.integration.test.ts @@ -374,7 +374,11 @@ describe("durable ready-step dispatch", () => { const who = identity(owner, agentId); const plans = []; - for (let index = 0; index <= MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT; index += 1) { + for ( + let index = 0; + index <= MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT; + index += 1 + ) { plans.push( await store.create({ ...planInput(owner, agentId, channel.id), @@ -383,7 +387,10 @@ describe("durable ready-step dispatch", () => { ); } - for (const plan of plans.slice(0, MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT)) { + 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", @@ -427,7 +434,11 @@ describe("durable ready-step dispatch", () => { let started = 0; for (const agentId of agentIds.slice(0, 5)) { - for (let index = 0; index < MAX_CONCURRENT_WORKFLOW_STEPS_PER_AGENT; index += 1) { + 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}`, From d88d61afd547da1602374f0bf33739b40790ce76 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:22:14 +0700 Subject: [PATCH 14/18] fix(workflows): count paused in-flight turns against capacity --- server/src/workflows/store.ts | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/server/src/workflows/store.ts b/server/src/workflows/store.ts index d39bed672..2a561b479 100644 --- a/server/src/workflows/store.ts +++ b/server/src/workflows/store.ts @@ -419,12 +419,7 @@ export function createWorkflowStore(database: Database): WorkflowStore { }) .from(workflowSteps) .innerJoin(workflowRuns, eq(workflowRuns.id, workflowSteps.workflowId)) - .where( - and( - eq(workflowRuns.status, "active"), - eq(workflowSteps.status, "running"), - ), - ); + .where(eq(workflowSteps.status, "running")); const deploymentRunning = counts?.deployment ?? 0; const agentRunning = counts?.agent ?? 0; From 0170eb215e03cdd1a07668aa59be5d6f3ce16ce8 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:22:17 +0700 Subject: [PATCH 15/18] test(workflows): paused running turn still consumes capacity --- server/tests/workflow-store.integration.test.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/server/tests/workflow-store.integration.test.ts b/server/tests/workflow-store.integration.test.ts index 8f782855d..1ab5aad7d 100644 --- a/server/tests/workflow-store.integration.test.ts +++ b/server/tests/workflow-store.integration.test.ts @@ -404,6 +404,10 @@ describe("durable ready-step dispatch", () => { ); } + // 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) => From ebd888b038c15bc6e5d8fbe01b805634ee7f9f21 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:23:20 +0700 Subject: [PATCH 16/18] style(release): avoid template-placeholder lint warning --- scripts/release-preflight.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/scripts/release-preflight.ts b/scripts/release-preflight.ts index 4c59b85c5..5c090d830 100644 --- a/scripts/release-preflight.ts +++ b/scripts/release-preflight.ts @@ -1600,7 +1600,7 @@ function checkDurableWorkflowState(): void { } for (const evidence of [ "async defer({ kind, key, owner, delayMs, reason })", - "greatest(${workItems.attempts} - 1, 0)", + "greatest(" + "${workItems.attempts}" + " - 1, 0)", ]) { if (!queue.includes(evidence)) { fail(`workflows: capacity-safe queue defer is missing ${evidence}`); From cc57561c1b7b5e10c34a905a116c40d470c5e572 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:24:23 +0700 Subject: [PATCH 17/18] fix(release): avoid placeholder lint in queue evidence --- scripts/release-preflight.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/scripts/release-preflight.ts b/scripts/release-preflight.ts index 5c090d830..f9c307f97 100644 --- a/scripts/release-preflight.ts +++ b/scripts/release-preflight.ts @@ -1600,7 +1600,8 @@ function checkDurableWorkflowState(): void { } for (const evidence of [ "async defer({ kind, key, owner, delayMs, reason })", - "greatest(" + "${workItems.attempts}" + " - 1, 0)", + "attempts: sql\`greatest(", + "workItems.attempts} - 1, 0)", ]) { if (!queue.includes(evidence)) { fail(`workflows: capacity-safe queue defer is missing ${evidence}`); From c72920f3f96fa6f259278526e87977d2dc5639a6 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sun, 20 Sep 2026 10:25:37 +0700 Subject: [PATCH 18/18] style(release): remove useless evidence escape --- scripts/release-preflight.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/scripts/release-preflight.ts b/scripts/release-preflight.ts index f9c307f97..1b8810515 100644 --- a/scripts/release-preflight.ts +++ b/scripts/release-preflight.ts @@ -1600,7 +1600,7 @@ function checkDurableWorkflowState(): void { } for (const evidence of [ "async defer({ kind, key, owner, delayMs, reason })", - "attempts: sql\`greatest(", + "attempts: sql`greatest(", "workItems.attempts} - 1, 0)", ]) { if (!queue.includes(evidence)) {