From 3e5c7098899951a094ba03b6f9715527a26b5ae7 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:46:10 +0700 Subject: [PATCH 01/12] feat(workflows): add exact ready-step dispatch CAS --- server/src/workflows/store.ts | 170 +++++++++++++++++++++++++++++++++- 1 file changed, 168 insertions(+), 2 deletions(-) diff --git a/server/src/workflows/store.ts b/server/src/workflows/store.ts index f00ea2db0..5fd8c79d7 100644 --- a/server/src/workflows/store.ts +++ b/server/src/workflows/store.ts @@ -134,6 +134,13 @@ export type DueWorkflowWait = WorkflowIdentity & { attempts: number; }; +export type ReadyWorkflowStep = WorkflowIdentity & { + workflowId: string; + stepKey: string; + readyAt: Date; + attempts: number; +}; + export class WorkflowNotFoundError extends Error { constructor() { super("That workflow does not exist."); @@ -281,6 +288,13 @@ export type WorkflowStore = { id: string, key: string, ): Promise; + startReadyStep( + identity: WorkflowIdentity, + id: string, + key: string, + expectedReadyAt: Date, + expectedAttempt: number, + ): Promise; waitStep( identity: WorkflowIdentity, id: string, @@ -311,6 +325,7 @@ export type WorkflowStore = { id: string, key: string, ): Promise; + readySteps(limit: number): Promise; dueWaitingSteps(limit: number): Promise; addAsset( identity: WorkflowIdentity, @@ -609,6 +624,7 @@ export function createWorkflowStore(database: Database): WorkflowStore { id: `workflow_step_${crypto.randomUUID()}`, workflowId: id, ...step, + updatedAt: sql`date_trunc('milliseconds', now())`, })), ); }); @@ -682,6 +698,23 @@ export function createWorkflowStore(database: Database): WorkflowStore { lte(workflowSteps.waitUntil, sql`now()`), ), ); + + // Ready work has the same deterministic-key problem as an overdue wait: if a queued ready + // item was finished while the workflow was paused, resuming with the same stamp would + // collide with that finished queue row forever. Give every still-ready step a fresh, + // millisecond-stable dispatch stamp on Resume so the next sweep can offer it again. + await transaction + .update(workflowSteps) + .set({ + resumedFromWaitUntil: null, + updatedAt: sql`date_trunc('milliseconds', now())`, + }) + .where( + and( + eq(workflowSteps.workflowId, id), + eq(workflowSteps.status, "ready"), + ), + ); }); return (await planFor(identity, id)) as WorkflowPlan; }, @@ -732,6 +765,99 @@ export function createWorkflowStore(database: Database): WorkflowStore { }); }, + async startReadyStep( + identity, + id, + key, + expectedReadyAt, + expectedAttempt, + ) { + if ( + !(expectedReadyAt instanceof Date) || + Number.isNaN(expectedReadyAt.getTime()) || + !Number.isInteger(expectedAttempt) || + expectedAttempt < 0 + ) { + throw new WorkflowRefusedError( + "A queued ready step needs an exact ready timestamp and attempt.", + ); + } + + return await database.transaction(async (transaction) => { + await lockWorkflow(transaction, id); + const run = await loadOwned(transaction, identity, id); + if (run.status !== "active") { + throw new WorkflowRefusedError( + "Only an active workflow can start a ready step.", + ); + } + + const wantedKey = stepKey(key); + const [current] = await transaction + .select() + .from(workflowSteps) + .where( + and( + eq(workflowSteps.workflowId, id), + eq(workflowSteps.key, wantedKey), + ), + ) + .limit(1); + if (!current) { + throw new WorkflowRefusedError( + "That workflow step does not exist for this workflow.", + ); + } + + // A released queue item may come back after it already performed the ready->running CAS. + // Only the exact item that wrote this stamp may re-enter that same running attempt. + if ( + current.status === "running" && + current.attempts === expectedAttempt + 1 && + current.resumedFromWaitUntil?.getTime() === expectedReadyAt.getTime() + ) { + return toStep(current); + } + + if (current.attempts !== expectedAttempt) { + throw new WorkflowRefusedError( + "That workflow step moved to another attempt after this ready dispatch was scheduled.", + ); + } + + const [row] = await transaction + .update(workflowSteps) + .set({ + status: "running", + attempts: sql`${workflowSteps.attempts} + 1`, + startedAt: sql`now()`, + finishedAt: null, + failureReason: null, + waitUntil: null, + // Reuse the durable per-attempt dispatch stamp. waitStep clears it, and a later wait + // wake replaces it with its own exact wait timestamp. + resumedFromWaitUntil: expectedReadyAt, + updatedAt: sql`now()`, + }) + .where( + and( + eq(workflowSteps.workflowId, id), + eq(workflowSteps.key, wantedKey), + eq(workflowSteps.status, "ready"), + eq(workflowSteps.attempts, expectedAttempt), + eq(workflowSteps.updatedAt, expectedReadyAt), + ), + ) + .returning(); + if (!row) { + throw new WorkflowRefusedError( + "That ready step changed after this dispatch was scheduled.", + ); + } + return toStep(row); + }); + }, + async waitStep(identity, id, key, input) { if ( !(input.waitUntil instanceof Date) || @@ -903,7 +1029,10 @@ export function createWorkflowStore(database: Database): WorkflowStore { ) { await transaction .update(workflowSteps) - .set({ status: "ready", updatedAt: sql`now()` }) + .set({ + status: "ready", + updatedAt: sql`date_trunc('milliseconds', now())`, + }) .where( and( eq(workflowSteps.id, step.id), @@ -1020,7 +1149,8 @@ export function createWorkflowStore(database: Database): WorkflowStore { status: "ready", failureReason: null, finishedAt: null, - updatedAt: sql`now()`, + resumedFromWaitUntil: null, + updatedAt: sql`date_trunc('milliseconds', now())`, }) .where( and( @@ -1150,6 +1280,42 @@ export function createWorkflowStore(database: Database): WorkflowStore { }); }, + async readySteps(limit) { + if (!Number.isInteger(limit) || limit <= 0 || limit > 500) { + throw new WorkflowRefusedError( + "A ready-step scan limit must be between 1 and 500.", + ); + } + const rows = await database + .select({ + ownerUserId: workflowRuns.ownerUserId, + agentId: workflowRuns.agentId, + workflowId: workflowRuns.id, + stepKey: workflowSteps.key, + readyAt: workflowSteps.updatedAt, + attempts: workflowSteps.attempts, + }) + .from(workflowSteps) + .innerJoin(workflowRuns, eq(workflowRuns.id, workflowSteps.workflowId)) + .where( + and( + eq(workflowRuns.status, "active"), + eq(workflowSteps.status, "ready"), + ), + ) + .orderBy(asc(workflowSteps.updatedAt), asc(workflowSteps.id)) + .limit(limit); + + return rows.map((row) => ({ + ownerUserId: row.ownerUserId, + agentId: row.agentId, + workflowId: row.workflowId, + stepKey: row.stepKey, + readyAt: row.readyAt, + attempts: row.attempts, + })); + }, + async dueWaitingSteps(limit) { if (!Number.isInteger(limit) || limit <= 0 || limit > 500) { throw new WorkflowRefusedError( From 77c9eb21239b15ffecb044444f1c0ffae2422de5 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:46:34 +0700 Subject: [PATCH 02/12] feat(workflows): queue and dispatch ready steps durably --- server/src/workflows/ready.ts | 280 ++++++++++++++++++++++++++++++++++ 1 file changed, 280 insertions(+) create mode 100644 server/src/workflows/ready.ts diff --git a/server/src/workflows/ready.ts b/server/src/workflows/ready.ts new file mode 100644 index 000000000..286521959 --- /dev/null +++ b/server/src/workflows/ready.ts @@ -0,0 +1,280 @@ +import type { WorkflowStore } from "./store"; +import { DEFAULT_MAX_ATTEMPTS, type WorkQueue } from "../work/queue"; + +export const WORKFLOW_READY_DISPATCH_KIND = "workflow_ready_dispatch"; + +const DEFAULT_LIMIT = 50; +const DEFAULT_LEASE_MS = 6 * 60_000; +const DEFAULT_RENEW_EVERY_MS = 15_000; +const DEFAULT_RETRY_DELAY_MS = 5_000; + +type WorkflowReadyStore = Pick< + WorkflowStore, + "readySteps" | "startReadyStep" | "failStep" +>; + +export type WorkflowReadyOptions = { + store: WorkflowReadyStore; + queue: WorkQueue; + owner: string; + dispatch?: (input: { + ownerUserId: string; + agentId: string; + workflowId: string; + stepKey: string; + expectedAttempt: number; + }) => Promise; + limit?: number; + leaseMs?: number; + retryDelayMs?: number; + renewEveryMs?: number; + maxAttempts?: number; +}; + +export type WorkflowReadyReport = { + considered: number; + started: Array<{ workflowId: string; stepKey: string }>; + skipped: Array<{ workflowId: string; stepKey: string; reason: string }>; +}; + +function readyKey(workflowId: string, stepKey: string, readyAt: Date): string { + return `${workflowId}:${stepKey}:${readyAt.toISOString()}`; +} + +export async function offerReadyWorkflowSteps( + options: WorkflowReadyOptions, +): Promise<{ queued: number; already: number }> { + const ready = await options.store.readySteps(options.limit ?? DEFAULT_LIMIT); + let queued = 0; + let already = 0; + + for (const step of ready) { + const outcome = await options.queue.offer({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: readyKey(step.workflowId, step.stepKey, step.readyAt), + payload: { + ownerUserId: step.ownerUserId, + agentId: step.agentId, + workflowId: step.workflowId, + stepKey: step.stepKey, + readyAt: step.readyAt.toISOString(), + attempts: step.attempts, + }, + }); + if (outcome === "queued") queued += 1; + if (outcome === "already") already += 1; + } + + return { queued, already }; +} + +export async function dispatchClaimedReadyWorkflowSteps( + options: WorkflowReadyOptions, +): Promise { + const leaseMs = options.leaseMs ?? DEFAULT_LEASE_MS; + const maxAttempts = options.maxAttempts ?? DEFAULT_MAX_ATTEMPTS; + const claimed = await options.queue.claim({ + kind: WORKFLOW_READY_DISPATCH_KIND, + owner: options.owner, + leaseMs, + limit: options.limit ?? DEFAULT_LIMIT, + maxAttempts, + }); + + const report: WorkflowReadyReport = { + considered: claimed.length, + started: [], + skipped: [], + }; + + for (const item of claimed) { + const workflowId = + typeof item.payload.workflowId === "string" + ? item.payload.workflowId + : ""; + const stepKey = + typeof item.payload.stepKey === "string" ? item.payload.stepKey : ""; + const ownerUserId = + typeof item.payload.ownerUserId === "string" + ? item.payload.ownerUserId + : ""; + const agentId = + typeof item.payload.agentId === "string" ? item.payload.agentId : ""; + const readyAt = + typeof item.payload.readyAt === "string" + ? new Date(item.payload.readyAt) + : null; + const expectedAttempt = + typeof item.payload.attempts === "number" && + Number.isInteger(item.payload.attempts) && + item.payload.attempts >= 0 + ? item.payload.attempts + : -1; + + if ( + !workflowId || + !stepKey || + !ownerUserId || + !agentId || + expectedAttempt < 0 || + !readyAt || + Number.isNaN(readyAt.getTime()) + ) { + const reason = "workflow ready payload is incomplete or invalid"; + await options.queue.finish({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: item.key, + owner: options.owner, + }); + report.skipped.push({ workflowId, stepKey, reason }); + continue; + } + + const renewed = await options.queue.renew({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: item.key, + owner: options.owner, + leaseMs, + }); + if (!renewed) { + report.skipped.push({ + workflowId, + stepKey, + reason: "the lease went to another replica", + }); + continue; + } + + let startedAttempt: number | undefined; + try { + const started = await options.store.startReadyStep( + { ownerUserId, agentId }, + workflowId, + stepKey, + readyAt, + expectedAttempt, + ); + startedAttempt = started.attempts; + + if (options.dispatch) { + const renewEveryMs = + options.renewEveryMs ?? + Math.min( + DEFAULT_RENEW_EVERY_MS, + Math.max(1_000, Math.floor(leaseMs / 3)), + ); + let heartbeat: ReturnType | undefined; + try { + heartbeat = setInterval(() => { + void options.queue + .renew({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: item.key, + owner: options.owner, + leaseMs, + }) + .catch(() => {}); + }, renewEveryMs); + heartbeat.unref?.(); + + const stillOurs = await options.queue.renew({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: item.key, + owner: options.owner, + leaseMs, + }); + if (!stillOurs) { + report.skipped.push({ + workflowId, + stepKey, + reason: "the lease went to another replica before continuation", + }); + continue; + } + + await options.dispatch({ + ownerUserId, + agentId, + workflowId, + stepKey, + expectedAttempt: startedAttempt, + }); + } finally { + if (heartbeat !== undefined) clearInterval(heartbeat); + } + } + + await options.queue.finish({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: item.key, + owner: options.owner, + }); + report.started.push({ workflowId, stepKey }); + } catch (error) { + const reason = + error instanceof Error + ? error.message + : "workflow ready step could not start"; + + // 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 ( + /ready step changed|another attempt|does not exist|active workflow|not ready/i.test( + reason, + ) + ) { + await options.queue.finish({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: item.key, + owner: options.owner, + }); + report.skipped.push({ workflowId, stepKey, reason }); + continue; + } + + // Once this exact ready item has performed its CAS, exhausting the queue retry budget must + // not leave the step permanently "running". Fail that exact attempt so the durable workflow + // records a visible terminal reason and can be retried deliberately. + if (startedAttempt !== undefined && item.attempts >= maxAttempts) { + try { + await options.store.failStep( + { ownerUserId, agentId }, + workflowId, + stepKey, + `Autonomous continuation exhausted its retry budget: ${reason}`, + startedAttempt, + ); + } finally { + await options.queue.finish({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: item.key, + owner: options.owner, + }); + } + report.skipped.push({ workflowId, stepKey, reason }); + continue; + } + + await options.queue.release({ + 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 }); + } + } + + return report; +} + +export async function sweepReadyWorkflowSteps( + options: WorkflowReadyOptions, +): Promise { + const offered = await offerReadyWorkflowSteps(options); + const dispatched = await dispatchClaimedReadyWorkflowSteps(options); + return { ...offered, ...dispatched }; +} + +export const WORKFLOW_READY_MAX_ATTEMPTS = DEFAULT_MAX_ATTEMPTS; From 0b711b8117b537ab8f3f038b8ff4cb867bbb6640 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:47:08 +0700 Subject: [PATCH 03/12] fix(workflows): make ready and resume stamps monotonic --- server/src/workflows/store.ts | 20 ++++++++++++++++---- 1 file changed, 16 insertions(+), 4 deletions(-) diff --git a/server/src/workflows/store.ts b/server/src/workflows/store.ts index 5fd8c79d7..72bd92840 100644 --- a/server/src/workflows/store.ts +++ b/server/src/workflows/store.ts @@ -686,7 +686,10 @@ export function createWorkflowStore(database: Database): WorkflowStore { await transaction .update(workflowSteps) .set({ - waitUntil: sql`date_trunc('milliseconds', now())`, + waitUntil: sql`greatest( + date_trunc('milliseconds', now()), + date_trunc('milliseconds', ${workflowSteps.waitUntil}) + interval '1 millisecond' + )`, resumedFromWaitUntil: null, updatedAt: sql`now()`, }) @@ -707,7 +710,10 @@ export function createWorkflowStore(database: Database): WorkflowStore { .update(workflowSteps) .set({ resumedFromWaitUntil: null, - updatedAt: sql`date_trunc('milliseconds', now())`, + updatedAt: sql`greatest( + date_trunc('milliseconds', now()), + date_trunc('milliseconds', ${workflowSteps.updatedAt}) + interval '1 millisecond' + )`, }) .where( and( @@ -1031,7 +1037,10 @@ export function createWorkflowStore(database: Database): WorkflowStore { .update(workflowSteps) .set({ status: "ready", - updatedAt: sql`date_trunc('milliseconds', now())`, + updatedAt: sql`greatest( + date_trunc('milliseconds', now()), + date_trunc('milliseconds', ${workflowSteps.updatedAt}) + interval '1 millisecond' + )`, }) .where( and( @@ -1150,7 +1159,10 @@ export function createWorkflowStore(database: Database): WorkflowStore { failureReason: null, finishedAt: null, resumedFromWaitUntil: null, - updatedAt: sql`date_trunc('milliseconds', now())`, + updatedAt: sql`greatest( + date_trunc('milliseconds', now()), + date_trunc('milliseconds', ${workflowSteps.updatedAt}) + interval '1 millisecond' + )`, }) .where( and( From 8b2805dc29d6430559c631c48e6849b97f41e940 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:47:42 +0700 Subject: [PATCH 04/12] test(workflows): cover durable ready-step dispatch --- server/tests/workflow-ready.test.ts | 213 ++++++++++++++++++++++++++++ 1 file changed, 213 insertions(+) create mode 100644 server/tests/workflow-ready.test.ts diff --git a/server/tests/workflow-ready.test.ts b/server/tests/workflow-ready.test.ts new file mode 100644 index 000000000..621e0c177 --- /dev/null +++ b/server/tests/workflow-ready.test.ts @@ -0,0 +1,213 @@ +import { describe, expect, test } from "bun:test"; +import type { WorkItem, WorkQueue } from "../src/work/queue"; +import { + dispatchClaimedReadyWorkflowSteps, + offerReadyWorkflowSteps, + WORKFLOW_READY_DISPATCH_KIND, +} from "../src/workflows/ready"; + +function queueStub(overrides: Partial = {}): WorkQueue { + return { + offer: async () => "queued", + claim: async () => [], + renew: async () => true, + finish: async () => true, + release: async () => true, + purge: async () => 0, + ...overrides, + }; +} + +function item(attempts = 1): WorkItem { + return { + kind: WORKFLOW_READY_DISPATCH_KIND, + key: "workflow-1:render:2026-09-19T16:00:00.000Z", + attempts, + payload: { + ownerUserId: "user-1", + agentId: "bot-1", + workflowId: "workflow-1", + stepKey: "render", + readyAt: "2026-09-19T16:00:00.000Z", + attempts: 0, + }, + }; +} + +describe("workflow ready-step dispatch", () => { + test("offers one idempotent item for the exact ready stamp", async () => { + const readyAt = new Date("2026-09-19T16:00:00.000Z"); + const offered: Array> = []; + const result = await offerReadyWorkflowSteps({ + owner: "worker-1", + queue: queueStub({ + offer: async (work) => { + offered.push(work as unknown as Record); + return "queued"; + }, + }), + store: { + readySteps: async () => [ + { + ownerUserId: "user-1", + agentId: "bot-1", + workflowId: "workflow-1", + stepKey: "render", + readyAt, + attempts: 0, + }, + ], + startReadyStep: async () => { + throw new Error("not used"); + }, + failStep: async () => { + throw new Error("not used"); + }, + }, + }); + + expect(result).toEqual({ queued: 1, already: 0 }); + expect(offered[0]).toMatchObject({ + kind: WORKFLOW_READY_DISPATCH_KIND, + key: "workflow-1:render:2026-09-19T16:00:00.000Z", + payload: { + ownerUserId: "user-1", + agentId: "bot-1", + workflowId: "workflow-1", + stepKey: "render", + readyAt: "2026-09-19T16:00:00.000Z", + attempts: 0, + }, + }); + }); + + test("starts the exact ready version before dispatching its new attempt", async () => { + const order: string[] = []; + let dispatchAttempt = 0; + const report = await dispatchClaimedReadyWorkflowSteps({ + owner: "worker-1", + queue: queueStub({ + claim: async () => [item()], + finish: async () => { + order.push("finish"); + return true; + }, + }), + store: { + readySteps: async () => [], + startReadyStep: async (_identity, _id, _key, readyAt, attempt) => { + order.push("start"); + expect(readyAt.toISOString()).toBe("2026-09-19T16:00:00.000Z"); + expect(attempt).toBe(0); + return { attempts: 1 } as never; + }, + failStep: async () => { + throw new Error("not used"); + }, + }, + dispatch: async (input) => { + order.push("dispatch"); + dispatchAttempt = input.expectedAttempt; + }, + }); + + expect(order).toEqual(["start", "dispatch", "finish"]); + expect(dispatchAttempt).toBe(1); + expect(report.started).toEqual([ + { workflowId: "workflow-1", stepKey: "render" }, + ]); + }); + + test("keeps the lease alive for a long headless ready-step turn", async () => { + let renewals = 0; + await dispatchClaimedReadyWorkflowSteps({ + owner: "worker-1", + leaseMs: 60, + renewEveryMs: 5, + queue: queueStub({ + claim: async () => [item()], + renew: async () => { + renewals += 1; + return true; + }, + }), + store: { + readySteps: async () => [], + startReadyStep: async () => ({ attempts: 1 }) as never, + failStep: async () => { + throw new Error("not used"); + }, + }, + dispatch: async () => { + await new Promise((resolve) => setTimeout(resolve, 25)); + }, + }); + + expect(renewals).toBeGreaterThan(1); + }); + + test("releases a transient dispatch failure so the exact CAS can be retried", async () => { + let released = 0; + await dispatchClaimedReadyWorkflowSteps({ + owner: "worker-1", + queue: queueStub({ + claim: async () => [item()], + release: async () => { + released += 1; + return true; + }, + }), + store: { + readySteps: async () => [], + startReadyStep: async () => ({ attempts: 1 }) as never, + failStep: async () => { + throw new Error("not used"); + }, + }, + dispatch: async () => { + throw new Error("gateway unavailable"); + }, + }); + + expect(released).toBe(1); + }); + + test("fails the exact running attempt when autonomous retries are exhausted", async () => { + const failures: unknown[][] = []; + let released = 0; + let finished = 0; + + await dispatchClaimedReadyWorkflowSteps({ + owner: "worker-1", + maxAttempts: 2, + queue: queueStub({ + claim: async () => [item(2)], + release: async () => { + released += 1; + return true; + }, + finish: async () => { + finished += 1; + return true; + }, + }), + store: { + readySteps: async () => [], + startReadyStep: async () => ({ attempts: 1 }) as never, + failStep: async (...args) => { + failures.push(args); + return {} as never; + }, + }, + dispatch: async () => { + throw new Error("provider unavailable"); + }, + }); + + expect(released).toBe(0); + expect(finished).toBe(1); + expect(failures).toHaveLength(1); + expect(failures[0]?.at(-1)).toBe(1); + expect(String(failures[0]?.at(-2))).toContain("retry budget"); + }); +}); From 1992fca137af492c527f0976b0aff4572c66d79e Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:47:59 +0700 Subject: [PATCH 05/12] test(workflows): verify ready-step CAS and re-arm --- .../tests/workflow-store.integration.test.ts | 79 +++++++++++++++++++ 1 file changed, 79 insertions(+) diff --git a/server/tests/workflow-store.integration.test.ts b/server/tests/workflow-store.integration.test.ts index 3ff45bd91..4315fe0d9 100644 --- a/server/tests/workflow-store.integration.test.ts +++ b/server/tests/workflow-store.integration.test.ts @@ -281,6 +281,85 @@ describe("dependency and lifecycle transitions", () => { }); }); +describe("durable ready-step dispatch", () => { + test("starts only the exact queued ready version and retries that CAS idempotently", async () => { + const { owner, agentId, channel } = await setUp(); + const who = identity(owner, agentId); + const plan = await store.create(planInput(owner, agentId, channel.id)); + const queued = (await store.readySteps(500)).find( + (candidate) => + candidate.workflowId === plan.id && candidate.stepKey === "script", + ); + expect(queued).toBeDefined(); + expect(queued?.attempts).toBe(0); + + const started = await store.startReadyStep( + who, + plan.id, + "script", + queued?.readyAt as Date, + 0, + ); + expect(started.status).toBe("running"); + expect(started.attempts).toBe(1); + + const redelivered = await store.startReadyStep( + who, + plan.id, + "script", + queued?.readyAt as Date, + 0, + ); + expect(redelivered.status).toBe("running"); + expect(redelivered.attempts).toBe(1); + }); + + test("a manual start makes an older queued ready item stale", async () => { + const { owner, agentId, channel } = await setUp(); + const who = identity(owner, agentId); + const plan = await store.create(planInput(owner, agentId, channel.id)); + const queued = (await store.readySteps(500)).find( + (candidate) => + candidate.workflowId === plan.id && candidate.stepKey === "script", + ); + expect(queued).toBeDefined(); + + await store.startStep(who, plan.id, "script"); + await expect( + store.startReadyStep( + who, + plan.id, + "script", + queued?.readyAt as Date, + 0, + ), + ).rejects.toThrow(/another attempt/); + }); + + test("pause and resume gives still-ready work a fresh deterministic key", async () => { + const { owner, agentId, channel } = await setUp(); + const who = identity(owner, agentId); + const plan = await store.create(planInput(owner, agentId, channel.id)); + const before = (await store.readySteps(500)).find( + (candidate) => + candidate.workflowId === plan.id && candidate.stepKey === "script", + ); + expect(before).toBeDefined(); + + await store.pause(who, plan.id); + await store.resume(who, plan.id); + + const after = (await store.readySteps(500)).find( + (candidate) => + candidate.workflowId === plan.id && candidate.stepKey === "script", + ); + expect(after).toBeDefined(); + expect((after?.readyAt.getTime() ?? 0) > (before?.readyAt.getTime() ?? 0)).toBe( + true, + ); + }); +}); + describe("durable waits", () => { test("resume re-arms a due wait so a wake finished during pause cannot wedge it", async () => { const { owner, agentId, channel } = await setUp(); From c44e2bcf3ba6745bf6e8de1474e3a1a75308d498 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:48:07 +0700 Subject: [PATCH 06/12] docs(schema): clarify workflow dispatch stamp semantics --- server/src/db/schema/coworker.ts | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/server/src/db/schema/coworker.ts b/server/src/db/schema/coworker.ts index 8173d668d..16a071fb7 100644 --- a/server/src/db/schema/coworker.ts +++ b/server/src/db/schema/coworker.ts @@ -257,7 +257,11 @@ export const workflowSteps = pgTable( provider: text("provider"), /** Exact durable wake target for a waiting step. */ waitUntil: timestamp("wait_until", { withTimezone: true }), - /** Exact wait stamp whose wake most recently moved this attempt back to running. */ + /** + * Exact durable dispatch stamp for the current autonomous attempt. + * A ready-step dispatch stores its ready timestamp; waitStep clears it; a later wait wake + * stores the exact wait timestamp. The column name is retained for migration compatibility. + */ resumedFromWaitUntil: timestamp("resumed_from_wait_until", { withTimezone: true, }), From 2e86e4a31c125eb848493ef605338999825c90c5 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:48:18 +0700 Subject: [PATCH 07/12] docs(workflows): make continuation prompt fit initial steps --- server/src/workflows/runner.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/server/src/workflows/runner.ts b/server/src/workflows/runner.ts index c0eab76ee..db8720ebe 100644 --- a/server/src/workflows/runner.ts +++ b/server/src/workflows/runner.ts @@ -23,7 +23,7 @@ function continuationInstruction(input: { instruction: string; }): string { return [ - "Resume this durable workflow step now.", + "Run this durable workflow step now.", "", `Workflow: ${input.workflowId}`, `Step: ${input.stepKey}`, From 82ca75b116b5ad3f29e4779b35383701d6d609ed Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:48:33 +0700 Subject: [PATCH 08/12] feat(workflows): run ready steps autonomously --- server/src/index.ts | 35 +++++++++++++++++++++++++++++++++++ 1 file changed, 35 insertions(+) diff --git a/server/src/index.ts b/server/src/index.ts index 4e8c0f28b..9d844b9a6 100644 --- a/server/src/index.ts +++ b/server/src/index.ts @@ -110,6 +110,7 @@ import { grantedSkills, grantedTools, REFUSAL_MARKER } from "./plugins/tools"; import { createTurnRunner } from "./routines/run-turn"; import { createRoutineRunner } from "./routines/runner"; import { createRoutineStore } from "./routines/store"; +import { sweepReadyWorkflowSteps } from "./workflows/ready"; import { createWorkflowRunner } from "./workflows/runner"; import { createWorkflowStore } from "./workflows/store"; import { sweepWorkflowWaits } from "./workflows/wake"; @@ -1296,6 +1297,40 @@ if (config.handoff.maxDepth > 0 && config.handoff.maxPerRun > 0) { repeatAfterEach(kick, 2_000); } +/* + * Durable workflow ready-step execution. + * + * Creating a workflow and promoting dependencies only changes durable state. This bridge turns + * exact ready versions into headless Agent turns through the same shared work_items queue used by + * other background work. The ready timestamp and attempt form a compare-and-set token, so a manual + * start, retry, pause/resume re-arm or newer dependency transition makes older queue work harmless. + */ +const workflowReady = { + store: workflowStore, + queue: createWorkQueue(database), + owner: workOwner("workflow-ready"), + dispatch: (input: Parameters[0]) => + workflowRunner.run(input), +}; + +repeatAfterEach(async () => { + try { + const report = await sweepReadyWorkflowSteps(workflowReady); + if ( + report.queued > 0 || + report.started.length > 0 || + report.skipped.length > 0 + ) { + console.info(JSON.stringify({ type: "workflow-ready", ...report })); + } + } catch (error) { + console.warn( + "[workflows] ready steps could not be dispatched:", + error instanceof Error ? error.message : error, + ); + } +}, 5_000); + /* * Durable workflow waits. * From 3d782d38a83940fd1f3b75f426e412eb130e9827 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:48:53 +0700 Subject: [PATCH 09/12] test(release): pin autonomous ready-step dispatch --- scripts/release-preflight.ts | 25 ++++++++++++++++++++++++- 1 file changed, 24 insertions(+), 1 deletion(-) diff --git a/scripts/release-preflight.ts b/scripts/release-preflight.ts index 6f10a927d..123764fc0 100644 --- a/scripts/release-preflight.ts +++ b/scripts/release-preflight.ts @@ -1476,6 +1476,7 @@ function checkDurableWorkflowState(): void { const schema = read("server/src/db/schema/coworker.ts"); const store = read("server/src/workflows/store.ts"); 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 server = read("server/src/index.ts"); const assetMigration = read("server/drizzle/0044_workflow_assets.sql"); @@ -1505,8 +1506,14 @@ function checkDurableWorkflowState(): void { "which must be an earlier step in the same workflow", "eq(workflowSteps.waitUntil, expectedWaitUntil)", "lte(workflowSteps.waitUntil, sql`now()`)", + "readySteps(limit)", + "startReadyStep(", + "expectedReadyAt", + "eq(workflowSteps.updatedAt, expectedReadyAt)", + "expectedAttempt + 1", "dueWaitingSteps(limit)", "eq(workflowSteps.attempts, expectedAttempt)", + "greatest(", "date_trunc('milliseconds', now())", ]) { if (!store.includes(evidence)) { @@ -1528,6 +1535,22 @@ function checkDurableWorkflowState(): void { fail(`workflows: durable wake bridge is missing ${evidence}`); } } + for (const evidence of [ + "WORKFLOW_READY_DISPATCH_KIND", + "readyAt.toISOString()", + "startReadyStep(", + "queue.renew({", + "queue.release({", + "item.attempts >= maxAttempts", + "Autonomous continuation exhausted its retry budget", + ]) { + if (!ready.includes(evidence)) { + fail(`workflows: durable ready-step bridge is missing ${evidence}`); + } + } + if (!server.includes("sweepReadyWorkflowSteps(workflowReady)")) { + fail("workflows: durable ready-step bridge is not running from the server"); + } if (!server.includes("sweepWorkflowWaits(workflowWake)")) { fail("workflows: durable wake bridge is not running from the server"); } @@ -1542,7 +1565,7 @@ function checkDurableWorkflowState(): void { } } for (const evidence of [ - "Resume this durable workflow step now.", + "Run this durable workflow step now.", "expectedAttempt", "checkpoint the step durably", "The autonomous continuation ended without checkpointing this workflow step.", From c2bf997f0e157a034771159d017bc94216dc71cb Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:49:10 +0700 Subject: [PATCH 10/12] docs(workflows): document autonomous ready-step execution --- docs/workflows.md | 42 +++++++++++++++++++++++++++++++++++++----- 1 file changed, 37 insertions(+), 5 deletions(-) diff --git a/docs/workflows.md b/docs/workflows.md index ebae3fb58..af71f3183 100644 --- a/docs/workflows.md +++ b/docs/workflows.md @@ -54,8 +54,15 @@ for an identity that no longer exists. ## Dependencies and retries -Root steps begin `ready`; dependent steps begin `blocked`. Starting a step is a compare-and-set from -`ready` to `running` and increments its durable attempt count. +Root steps begin `ready`; dependent steps begin `blocked`. Ready steps are discovered by a bounded +recovery sweep and offered to the existing durable `work_items` queue. The queued identity includes +the exact ready timestamp and current attempt count; starting is a compare-and-set from that exact +`ready` version to `running` and increments the durable attempt count. + +The same exact queue item may be retried after a worker/model interruption. A dispatch stamp on the +running attempt lets only that item re-enter the same attempt; a manual start, a newer retry, or any +other state transition makes the old item stale. Long headless turns renew their queue lease while +they run so a second replica cannot claim the same autonomous step mid-turn. Completing a step promotes only blocked steps whose complete dependency set has succeeded. Workflow transitions are serialized with a PostgreSQL transaction-scoped advisory lock, so two steps finishing @@ -63,7 +70,29 @@ at the same time cannot leave a dependent step permanently blocked. A running step may fail without destroying the workflow. It becomes `failed`, keeps its attempt history and failure reason, and can be moved back to `ready` only when all of its dependencies still -succeeded. The next start increments the same attempt counter. +succeeded. The retry transition mints a fresh monotonic ready stamp so an older finished queue item +cannot suppress the new attempt. The next start increments the same attempt counter. + +When a ready-step queue item has already started its exact attempt but autonomous dispatch keeps +failing until the queue retry budget is exhausted, that exact running attempt is failed with a +visible retry-budget reason instead of being left permanently `running`. + +## Autonomous step execution + +Creating a workflow no longer requires a person or Agent turn to manually start each root or newly +unblocked step. Every active `ready` step is offered idempotently to the shared durable queue and +runs through the same governed headless Agent path used for resumed waits. The Agent receives the +workflow id, step key, durable attempt number and stored instruction, then must checkpoint the step +before its turn ends: complete it, put it into a future wait, or fail it with a concrete reason. + +This is execution orchestration, not a permission shortcut. The headless turn rebuilds the Agent for +the workflow owner and Bot, uses the workflow's existing channel/thread, and receives only the tools +and permissions that Agent currently has. Browser, Files, Computer and connector policy remain +unchanged. + +Completing one step can promote dependent steps to `ready`; those new ready versions are then picked +up by the same durable bridge. This is what lets a multi-step DAG continue autonomously instead of +stopping after each dependency boundary. ## Durable waits @@ -98,8 +127,11 @@ host filesystem path into a workflow. ## Pause, resume and cancel -Pause and resume change only the workflow gate; they do not erase steps, attempts, provider state or -wait timestamps. A paused workflow cannot start or resume work. +Pause and resume keep the workflow's steps, attempts and provider state durable. A paused workflow +cannot start or resume work. On Resume, any still-ready step receives a fresh monotonic ready stamp, +and any overdue waiting step receives a fresh exact wake timestamp. That re-arms work whose old +deterministic queue item may have been finished while the workflow was paused, without replaying a +newer or manually-started attempt. Cancel marks the run terminal and cancels every still-blocked, ready, running or waiting step in one serialized transaction. Completed or previously failed steps remain as evidence of what happened. From dff7c3b7b2eb3d36d53cd8679de9838f2cf51536 Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:50:39 +0700 Subject: [PATCH 11/12] style(workflows): format ready-step CAS --- server/src/workflows/store.ts | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/server/src/workflows/store.ts b/server/src/workflows/store.ts index 72bd92840..19a738b1e 100644 --- a/server/src/workflows/store.ts +++ b/server/src/workflows/store.ts @@ -771,13 +771,7 @@ export function createWorkflowStore(database: Database): WorkflowStore { }); }, - async startReadyStep( - identity, - id, - key, - expectedReadyAt, - expectedAttempt, - ) { + async startReadyStep(identity, id, key, expectedReadyAt, expectedAttempt) { if ( !(expectedReadyAt instanceof Date) || Number.isNaN(expectedReadyAt.getTime()) || From 6d92a41005fa5019a1203e45216f528e95e6267e Mon Sep 17 00:00:00 2001 From: duc15052006-dotcom Date: Sat, 19 Sep 2026 23:50:42 +0700 Subject: [PATCH 12/12] style(tests): format workflow ready integration coverage --- server/tests/workflow-store.integration.test.ts | 14 ++++---------- 1 file changed, 4 insertions(+), 10 deletions(-) diff --git a/server/tests/workflow-store.integration.test.ts b/server/tests/workflow-store.integration.test.ts index 4315fe0d9..34c80dcbe 100644 --- a/server/tests/workflow-store.integration.test.ts +++ b/server/tests/workflow-store.integration.test.ts @@ -326,13 +326,7 @@ describe("durable ready-step dispatch", () => { await store.startStep(who, plan.id, "script"); await expect( - store.startReadyStep( - who, - plan.id, - "script", - queued?.readyAt as Date, - 0, - ), + store.startReadyStep(who, plan.id, "script", queued?.readyAt as Date, 0), ).rejects.toThrow(/another attempt/); }); @@ -354,9 +348,9 @@ describe("durable ready-step dispatch", () => { candidate.workflowId === plan.id && candidate.stepKey === "script", ); expect(after).toBeDefined(); - expect((after?.readyAt.getTime() ?? 0) > (before?.readyAt.getTime() ?? 0)).toBe( - true, - ); + expect( + (after?.readyAt.getTime() ?? 0) > (before?.readyAt.getTime() ?? 0), + ).toBe(true); }); });