Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
d7954c0
feat(queue): defer capacity-blocked work without retries
duc15052006-dotcom Sep 20, 2026
ca55099
security(workflows): cap autonomous concurrent steps
duc15052006-dotcom Sep 20, 2026
8ac2b25
refactor(queue): keep defer compatible with queue test doubles
duc15052006-dotcom Sep 20, 2026
c8acd87
fix(workflows): defer ready work when capacity is full
duc15052006-dotcom Sep 20, 2026
0f5a581
fix(workflows): defer waits when autonomous capacity is full
duc15052006-dotcom Sep 20, 2026
086666f
test(workflows): defer capacity-blocked ready dispatches
duc15052006-dotcom Sep 20, 2026
86b0f21
test(workflows): defer capacity-blocked wait wakes
duc15052006-dotcom Sep 20, 2026
68a8964
test(queue): defer without consuming retry budget
duc15052006-dotcom Sep 20, 2026
3965833
test(workflows): enforce autonomous concurrency ceilings
duc15052006-dotcom Sep 20, 2026
56be595
test(release): pin workflow concurrency admission invariants
duc15052006-dotcom Sep 20, 2026
3bcb3d7
test(release): pin workflow concurrency admission invariants
duc15052006-dotcom Sep 20, 2026
fb229ce
docs(workflows): document autonomous concurrency ceilings
duc15052006-dotcom Sep 20, 2026
5d51e2f
style(tests): format workflow capacity coverage
duc15052006-dotcom Sep 20, 2026
d88d61a
fix(workflows): count paused in-flight turns against capacity
duc15052006-dotcom Sep 20, 2026
0170eb2
test(workflows): paused running turn still consumes capacity
duc15052006-dotcom Sep 20, 2026
ebd888b
style(release): avoid template-placeholder lint warning
duc15052006-dotcom Sep 20, 2026
cc57561
fix(release): avoid placeholder lint in queue evidence
duc15052006-dotcom Sep 20, 2026
c72920f
style(release): remove useless evidence escape
duc15052006-dotcom Sep 20, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions docs/workflows.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
18 changes: 18 additions & 0 deletions scripts/release-preflight.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down Expand Up @@ -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}`);
Expand Down Expand Up @@ -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}`);
Expand All @@ -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");
}
Expand Down
36 changes: 36 additions & 0 deletions server/src/work/queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,21 @@ export type WorkQueue = {
delayMs: number;
reason?: string;
}) => Promise<boolean>;
/**
* 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<boolean>;
/**
* Drop what is done with, older than the retention window. Returns how many went.
*
Expand Down Expand Up @@ -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,
Expand Down
19 changes: 18 additions & 1 deletion server/src/workflows/ready.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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 (
Expand Down
56 changes: 56 additions & 0 deletions server/src/workflows/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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<Parameters<Database["transaction"]>[0]>[0];
type Handle = Database | Transaction;

Expand Down Expand Up @@ -384,6 +398,43 @@ export function createWorkflowStore(database: Database): WorkflowStore {
);
}

async function admitAutonomousStep(
transaction: Transaction,
agentId: string,
): Promise<void> {
/*
* 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<number>`count(*)::int`,
agent: sql<number>`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<string> {
if (input.channelId !== undefined) {
const [row] = await database
Expand Down Expand Up @@ -890,6 +941,8 @@ export function createWorkflowStore(database: Database): WorkflowStore {
);
}

await admitAutonomousStep(transaction, identity.agentId);

const [row] = await transaction
.update(workflowSteps)
.set({
Expand Down Expand Up @@ -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({
Expand Down
19 changes: 18 additions & 1 deletion server/src/workflows/wake.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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 (
Expand Down
27 changes: 27 additions & 0 deletions server/tests/work-queue.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
36 changes: 36 additions & 0 deletions server/tests/workflow-ready.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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> = {}): WorkQueue {
return {
Expand Down Expand Up @@ -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({
Expand Down
Loading
Loading