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
80 changes: 80 additions & 0 deletions docs/workflows.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
# Durable workflows

A workflow is a durable multi-step plan owned by one person and one Bot. It is the state layer used
for work that cannot safely live in one model turn: content pipelines, long-running generation,
review chains, or plans that sleep and resume later.

This layer deliberately reuses the deployment's existing scheduler and `work_items` queue. It does
not create a second lease, timer, or coordinator.

## State model

A workflow run records its owner, Bot, destination channel, title and lifecycle:

- `active` — new steps may start and due waits may resume;
- `paused` — state remains durable but no new step starts;
- `succeeded` — every step completed;
- `failed` — the workflow was explicitly abandoned as failed;
- `cancelled` — the owner/Bot cancelled the remaining work.

Each step records a stable local key, position, instruction, dependency keys, status, attempt count,
optional provider, optional exact wait timestamp, failure reason and lifecycle timestamps.

Dependencies may point only to earlier steps supplied in the same create operation. That rule makes
the graph acyclic by construction and prevents a step from naming another workflow's state.

## Isolation

Every person-facing transition is filtered by both `owner_user_id` and `agent_id`. Two Bots owned
by the same person therefore cannot read, pause, resume, cancel or mutate each other's workflow just
because they share an owner.

The destination channel is accepted only when the owner is a member and the Bot is in that same
non-deleted channel. A missing, foreign or deleted channel is one refusal and does not reveal which
part of the check failed.

Deleting an Agent cascades its workflow rows and steps, so a pending plan cannot survive as stale work
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.

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.

## Durable waits

A running step can enter `waiting` with an exact future timestamp and an optional provider label.
The provider field is metadata only and must never contain a credential.

Due-wait discovery uses PostgreSQL's clock. Resuming a wait compares the exact persisted timestamp as
well as the step status. A stale wake therefore cannot resume a step whose wait was cancelled or
rescheduled after that wake was created.

The execution slice that turns due workflow waits into existing one-shot wake/queue items is separate
from this state foundation. Keeping the state transition separate from delivery lets the existing
`work_items` lease/idempotency mechanism remain the single durable execution queue.

## 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.

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.

## Security invariants

Workflow state never stores API keys, browser passwords or host credentials. It does not widen Browser,
Files, Computer or plugin grants. When execution wiring calls a provider later, the Agent still acts
through its existing per-person/per-Bot permission and credential paths.

The database state is a recovery ledger, not an authorization bypass: a persisted task is not evidence
that the action it describes is currently allowed.
45 changes: 45 additions & 0 deletions scripts/release-preflight.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1472,6 +1472,50 @@ function checkDurableAgentWake(): void {
}
}

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");

for (const evidence of [
'workflowRunStatus = pgEnum("workflow_run_status"',
'workflowStepStatus = pgEnum("workflow_step_status"',
"export const workflowRuns = pgTable(",
"export const workflowSteps = pgTable(",
'dependsOn: text("depends_on").array().notNull().default([])',
]) {
if (!schema.includes(evidence)) {
fail(`workflows: durable state schema is missing ${evidence}`);
}
}

for (const evidence of [
"eq(workflowRuns.ownerUserId, identity.ownerUserId)",
"eq(workflowRuns.agentId, identity.agentId)",
"pg_advisory_xact_lock",
"which must be an earlier step in the same workflow",
"eq(workflowSteps.waitUntil, expectedWaitUntil)",
"lte(workflowSteps.waitUntil, sql`now()`)",
"dueWaitingSteps(limit)",
]) {
if (!store.includes(evidence)) {
fail(`workflows: durable recovery boundary is missing ${evidence}`);
}
}

for (const evidence of [
'CREATE TYPE "public"."workflow_run_status"',
'CREATE TYPE "public"."workflow_step_status"',
'CREATE TABLE "workflow_runs"',
'CREATE TABLE "workflow_steps"',
"workflow_steps_workflow_key_idx",
]) {
if (!migration.includes(evidence)) {
fail(`workflows: durable workflow migration is missing ${evidence}`);
}
}
}

function checkVersionSources(): void {
const pkg = json("package.json");
const version = pkg.version;
Expand Down Expand Up @@ -1503,6 +1547,7 @@ checkDesktopCredentialBoundary();
checkProviderConnectionTest();
checkFirstCoworkerHandoff();
checkDurableAgentWake();
checkDurableWorkflowState();
checkVersionSources();

if (failures.length > 0) {
Expand Down
40 changes: 40 additions & 0 deletions server/drizzle/0043_workflow_state.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
CREATE TYPE "public"."workflow_run_status" AS ENUM('active', 'paused', 'succeeded', 'failed', 'cancelled');--> statement-breakpoint
CREATE TYPE "public"."workflow_step_status" AS ENUM('blocked', 'ready', 'running', 'waiting', 'succeeded', 'failed', 'cancelled');--> statement-breakpoint
CREATE TABLE "workflow_runs" (
"id" text PRIMARY KEY NOT NULL,
"owner_user_id" text NOT NULL,
"agent_id" text NOT NULL,
"channel_id" text NOT NULL,
"title" text NOT NULL,
"status" "workflow_run_status" DEFAULT 'active' NOT NULL,
"finished_at" timestamp with time zone,
"created_at" timestamp with time zone DEFAULT now() NOT NULL,
"updated_at" timestamp with time zone DEFAULT now() NOT NULL
);
--> statement-breakpoint
CREATE TABLE "workflow_steps" (
"id" text PRIMARY KEY NOT NULL,
"workflow_id" text NOT NULL,
"key" text NOT NULL,
"position" integer NOT NULL,
"instruction" text NOT NULL,
"depends_on" text[] DEFAULT '{}' NOT NULL,
"status" "workflow_step_status" DEFAULT 'blocked' NOT NULL,
"attempts" integer DEFAULT 0 NOT NULL,
"provider" text,
"wait_until" timestamp with time zone,
"failure_reason" text,
"started_at" timestamp with time zone,
"finished_at" timestamp with time zone,
"created_at" timestamp with time zone DEFAULT now() NOT NULL,
"updated_at" timestamp with time zone DEFAULT now() NOT NULL
);
--> statement-breakpoint
ALTER TABLE "workflow_runs" ADD CONSTRAINT "workflow_runs_owner_user_id_users_id_fk" FOREIGN KEY ("owner_user_id") REFERENCES "public"."users"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "workflow_runs" ADD CONSTRAINT "workflow_runs_agent_id_agents_id_fk" FOREIGN KEY ("agent_id") REFERENCES "public"."agents"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
ALTER TABLE "workflow_steps" ADD CONSTRAINT "workflow_steps_workflow_id_workflow_runs_id_fk" FOREIGN KEY ("workflow_id") REFERENCES "public"."workflow_runs"("id") ON DELETE cascade ON UPDATE no action;--> statement-breakpoint
CREATE INDEX "workflow_runs_owner_status_idx" ON "workflow_runs" USING btree ("owner_user_id","status");--> statement-breakpoint
CREATE INDEX "workflow_runs_agent_status_idx" ON "workflow_runs" USING btree ("agent_id","status");--> statement-breakpoint
CREATE UNIQUE INDEX "workflow_steps_workflow_key_idx" ON "workflow_steps" USING btree ("workflow_id","key");--> statement-breakpoint
CREATE INDEX "workflow_steps_workflow_position_idx" ON "workflow_steps" USING btree ("workflow_id","position");--> statement-breakpoint
CREATE INDEX "workflow_steps_status_wait_idx" ON "workflow_steps" USING btree ("status","wait_until");
Loading
Loading