Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
42 changes: 37 additions & 5 deletions docs/workflows.md
Original file line number Diff line number Diff line change
Expand Up @@ -54,16 +54,45 @@ 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
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

Expand Down Expand Up @@ -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.
Expand Down
25 changes: 24 additions & 1 deletion scripts/release-preflight.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down Expand Up @@ -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)) {
Expand All @@ -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");
}
Expand All @@ -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.",
Expand Down
6 changes: 5 additions & 1 deletion server/src/db/schema/coworker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}),
Expand Down
35 changes: 35 additions & 0 deletions server/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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<typeof workflowRunner.run>[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.
*
Expand Down
Loading
Loading