Workflows
A run registry for any workflow engine that members read through RLS, cron schedules, counting semaphores, admission control for starts, and useWorkflowRuns over Realtime.
The workflows block records the runs of your durable workflows in one
table, whatever engine runs them, so a run list, a run page and a cancel
button work the same way for every engine. It also adds schedules,
semaphores and admission control. Each of these hands a start to the engine
you choose rather than running the workflow itself. The
Workflow SDK adapter is the engine that ships
with it.
better-supabase sql add workflows # adds tenant and access as well| Table | Holds |
|---|---|
workflow_runs | One row per run: engine, external_id, definition, tenant_id, actor_id, status, attributes, error and the times |
workflow_schedules | Recurring starts: name, workflow, input, cron, timezone, next_run_at, paused |
workflow_semaphores | The holders of each semaphore key, until they release it or their lease ends |
workflow_start_requests | Starts waiting for room under their key's concurrency, debounce or singleton rule |
| Permission | Lets | Default roles |
|---|---|---|
workflow.read | members read the tenant's runs and schedules | none by default |
workflow.run | members start runs, for helpers such as protectWebHandler | none by default |
workflow.admin | members cancel any run of the tenant and manage schedules | none by default |
Grant the keys to your roles in sql.modules.access.roles. A user always
reads the runs they started, without a permission.
A run's status is queued, running, waiting, completed, failed or
cancelled. Every status change sends a payload-free Realtime message on
workflow-run:<id> and, for a run with a tenant, on
workflow-runs:<tenant>. With the outbox installed,
a run that ends emits workflow_run.completed, workflow_run.failed or
workflow_run.cancelled.
Reading runs
import {
createWorkflows,
rpcTransport,
} from "better-supabase/blocks/workflows";
const workflows = createWorkflows({ transport: rpcTransport(supabase) });
const runs = await workflows.runs
.list({ tenant: organizationId, status: "running", limit: 50 })
.orThrow();
const run = await workflows.runs.get(runId).orThrow();
await workflows.runs.requestCancel(runId);runs.get and requestCancel take the row's id or the engine's run id.
list returns the newest runs first; pass the last run's
createdAt as cursor for the next page. requestCancel sets cancel_requested_at
for the run's actor, a workflow.admin of its tenant or the service role.
The engine does the cancelling. The Workflow SDK adapter calls
run.cancel() for you.
Engines write runs with runs.record (service role), and
runs.purge({ olderThan }) deletes finished runs older than 30 days by
default.
In Client Components
"use client";
import { useWorkflowRuns } from "better-supabase/blocks/workflows/react";
export function RunList({ organizationId }: { organizationId: string }) {
const { runs, error } = useWorkflowRuns({ tenant: organizationId });
if (error) return <p role="alert">{error.message}</p>;
return (
<ul>
{runs?.map((run) => (
<li key={run.id}>
{run.definition}: {run.status}
</li>
))}
</ul>
);
}useWorkflowRuns loads the list and loads it again on every message on
workflow-runs:<tenant>. Without a tenant, it loads the user's own runs
once. useWorkflowRun(id) reads one run, follows workflow-run:<id> and
returns cancel(). In Server Components, call workflows.runs.list()
instead.
Schedules
Workflow schedules sit between job schedules and per-user AI tasks; see three scheduling layers.
import {
createWorkflows,
sqlTransport,
} from "better-supabase/blocks/workflows";
import { workflowStarter } from "better-supabase/workflow-sdk";
import { weeklyDigest } from "@/workflows/digest";
const workflows = createWorkflows({ transport: sqlTransport(postgres.admin) });
await workflows.schedules.create({
name: "weekly-digest",
workflow: "weeklyDigest",
cron: "0 9 * * 1",
timezone: "Europe/Amsterdam",
input: [organizationId],
tenant: organizationId,
});
export async function GET() {
const result = await workflows.schedules
.tick({ start: workflowStarter({ weeklyDigest }) })
.orThrow();
return Response.json(result);
}cron takes five cron fields, a macro such as @daily or an interval such
as 30 seconds. tick claims the due schedules, starts each one with the
idempotency key schedule:<id>:<fire time> and moves it to its next time,
so a tick that runs twice for the same fire time starts one run. A start
that throws stays due and runs again once its lease ends. Call tick from
a cron route or a job every minute.
Members with workflow.admin create, pause and remove the schedules of
their tenant; the service role manages every schedule.
Semaphores and admission
const free = await workflows.semaphores
.acquire("stripe-sync", runId, { max: 3, ttl: 300 })
.orThrow();
await workflows.semaphores.release("stripe-sync", runId);
await workflows.admission.request({
key: `import:${organizationId}`,
workflow: "importContacts",
input: [fileId],
tenant: organizationId,
concurrency: 1,
debounce: 30,
});
await workflows.admission.tick({ start: workflowStarter({ importContacts }) });A semaphore lets at most max holders take a key at once, each until it
releases the key or its ttl runs out. Admission queues starts by key:
concurrency caps the runs of the key that are active, debounce waits
that many seconds and replaces a request that is still waiting, and
singleton drops the request while a run with the key is waiting or
active. admission.tick starts the requests whose key has room, with the
idempotency key admission:<id>. Both are service-role functions.
Functions
| Function | Granted to | Does |
|---|---|---|
workflow_runs_list(tenant, definition, status, max, before) | authenticated, service_role | The runs the caller may read, newest first |
workflow_run_get(run) | authenticated, service_role | One run by its id or its engine's id |
request_workflow_cancel(run) | authenticated, service_role | Marks a run for cancellation |
record_workflow_run(...) | service_role | Creates or updates a run by engine and external id |
purge_workflow_runs(older_than, batch) | service_role | Deletes finished runs |
create_workflow_schedule(...) | authenticated, service_role | Creates or replaces a schedule by tenant and name |
workflow_schedules_list(tenant) | authenticated, service_role | The schedules the caller may read |
pause_workflow_schedule(id, paused), remove_workflow_schedule | authenticated, service_role | Pauses, resumes or removes a schedule |
claim_due_workflow_schedules, advance_workflow_schedule | service_role | The schedule tick |
acquire_workflow_semaphore, release_workflow_semaphore | service_role | Semaphores |
request_workflow_start, claim_workflow_start_requests, mark_* | service_role | Admission control |
Last updated on