Workflow SDK
Run the Workflow SDK on Supabase with a World over Postgres and Supabase Queues, and start runs, resume hooks and protect routes as the signed-in user.
The Workflow SDK (workflow) needs a World: the storage and queue its runs,
steps, hooks and events live in. better-supabase/workflow-sdk/world is a
World on your Supabase database. Its tables sit in a workflow schema that
only the service role reads, its messages wait in a
jobs queue, and every run is copied into the
workflows block's workflow_runs, so members
list and follow runs through RLS. better-supabase/workflow-sdk has the
helpers that start runs and resume hooks as the signed-in user.
better-supabase sql add workflow-sdk-world # adds workflows and jobs as well
pnpm add workflow @workflow/world @workflow/world-postgres pgThe world entry imports pg and Node built-ins, so it runs on Node, not on
edge runtimes. The helpers in better-supabase/workflow-sdk import only
workflow.
Selecting the World
import { withWorkflow } from "workflow/next";
export default withWorkflow(nextConfig);WORKFLOW_TARGET_WORLD=better-supabase/workflow-sdk/world
SUPABASE_DB_URL=postgresql://postgres:postgres@127.0.0.1:54322/postgresThe entry's default export is createWorld(), which reads its options from
the environment, so the SDK and tools that select a World by package name
load it without code. To build one yourself, call createSupabaseWorld:
import { createSupabaseWorld } from "better-supabase/workflow-sdk/world";
export const world = createSupabaseWorld({
connectionString: process.env.SUPABASE_DB_URL,
delivery: "pg_net",
});| Variable | Option | Default |
|---|---|---|
WORKFLOW_POSTGRES_URL, SUPABASE_DB_URL, DATABASE_URL | connectionString | the first one set |
WORKFLOW_DELIVERY | delivery | poll |
WORKFLOW_FLOW_URL | flowUrl | WORKFLOW_LOCAL_BASE_URL or http://localhost:$PORT, plus /.well-known/workflow/v1/flow |
WORKFLOW_DELIVERY_SECRET | deliverySecret | the Vault secret workflow_delivery_secret |
WORKFLOW_ENCRYPTION_KEY | encryptionKey | the Vault secret workflow_encryption_key |
WORKFLOW_DELIVERY_QUEUE | queue | workflow_deliveries |
WORKFLOW_POSTGRES_WORKER_CONCURRENCY | concurrency | 50 |
WORKFLOW_POSTGRES_POLL_INTERVAL_MS | pollInterval | 250 |
The World speaks the @workflow/world protocol in
WORKFLOW_WORLD_PROTOCOL and stores its tables in the shape of the
@workflow/world-postgres version in WORLD_POSTGRES_VERSION. It passes
the @workflow/world-testing suite.
Delivery
The SDK hands each step and each workflow turn to the queue, and the queue
posts it to the flow route that withWorkflow adds.
| Mode | Who posts | Use it on |
|---|---|---|
poll (default) | the World, in the process that queues: it claims jobs and posts them | a server, a container, next dev |
pg_net | dispatch_workflow_deliveries() on a pg_cron schedule, through pg_net | serverless hosts with no long-running process |
For pg_net, set sql.modules.workflow-sdk-world.options.delivery to
pg_net (the data file then schedules the dispatcher every
options.schedule, 1 seconds by default) and store the route and a
secret in Vault:
select vault.create_secret('https://app.example.com/.well-known/workflow/v1/flow', 'workflow_flow_url');
select vault.create_secret(encode(extensions.gen_random_bytes(32), 'hex'), 'workflow_delivery_secret');In poll mode, the first message a process queues starts its poller, and
world.start() also queues the runs that were active when the last
process stopped. Call it once at startup, for example from
instrumentation.ts.
Each delivery is signed with the secret (x-bs-signature, an HMAC-SHA256
over the time, the job and the body). The World's queue handler refuses an
unsigned or stale request with a 401 Problem Details response, and
completes or fails the job itself, so a delivery whose response is lost runs
again after its lease. Leave the flow route out of your proxy's matcher: the
World checks the signature, not a session.
Poll mode needs the secret too, unless NODE_ENV is development or
test: without one, the flow route answers every delivery with 401. Store
it in Vault as above or set WORKFLOW_DELIVERY_SECRET. A secret or key that
isn't in Vault counts as not set; a Vault read that fails is not cached, so
the next delivery reads again, and until then the route answers 503 and run
data can't be read or written.
A poll delivery that the flow route doesn't answer within deliveryTimeout
seconds (300 by default, at least lease) is aborted and retried. The
poller extends the job's lease only while the request is open.
Encryption
With a 32-byte master key (base64 or hex) in encryptionKey or the Vault
secret workflow_encryption_key, the World derives an AES-256 key per run
with HKDF-SHA256, and the SDK stores the run's inputs, outputs and step
data encrypted. Without one, they are stored as JSON in the workflow
schema. When Vault can't be read, the World doesn't fall back to
unencrypted data: the read fails, and the next one tries Vault again.
Starting runs as the user
"use server";
import { startFor } from "better-supabase/workflow-sdk";
import { approveInvoice } from "@/workflows/approve-invoice";
export async function requestApproval(invoiceId: string) {
const ctx = await server.context();
const run = await startFor(ctx, approveInvoice, [invoiceId], {
idempotencyKey: `approve:${invoiceId}`,
});
return run.runId;
}startFor starts the run with the bs.tenant and bs.actor attributes
of the caller, which the World copies into workflow_runs.tenant_id and
actor_id. With idempotencyKey, it returns the run that already has the
key instead of starting another. To act as the user inside a step, pass
workflowContext(ctx) as an argument: it keeps the actor and tenant and
drops the claims, so no token ends up in the run's stored state.
| Helper | Does |
|---|---|
workflowStarter({ name: workflow }) | The start callback for schedules.tick and admission.tick, as the schedule's or request's creator |
startOnEvent(workflow, { types }) | An outbox handler that starts the workflow per matching event (invoice.paid, invoice.*, *), once |
hookMetadata(ctx, permission?) | Metadata for createHook, naming the actor and the permission a tenant member needs to resume it |
authorizeHook(token, ctx, { can }) | The hook, when the caller created it or holds its permission; HookForbiddenError (403) otherwise |
protectWebHandler(handler, key, { can }) | Answers 401 or 403 as application/problem+json before a workflow UI route or a run's stream |
Approvals with hooks
import type { RequestContext } from "better-supabase";
import { hookMetadata } from "better-supabase/workflow-sdk";
import { createHook } from "workflow";
export async function approveInvoice(invoiceId: string, ctx: RequestContext) {
"use workflow";
const hook = createHook<{ approved: boolean }>({
metadata: hookMetadata(ctx, "invoice.approve"),
});
const { approved } = await hook;
// ...
}import { authorizeHook } from "better-supabase/workflow-sdk";
import { resumeHook } from "workflow/api";
export async function POST(request: Request, { params }) {
const { token } = await params;
const ctx = await server.context(request);
await authorizeHook(token, ctx, {
can: (tenant, permission) => canInTenant(ctx, tenant, permission),
});
await resumeHook(token, await request.json());
return new Response(null, { status: 204 });
}A hook token in a link or an email is not enough to resume the run:
authorizeHook also checks who is asking. can is your permission check,
for example a call to member_can(auth.uid(), tenant, permission).
Graph workflows
better-supabase/workflow-sdk/builder runs the graphs of the
workflow builder on the Workflow SDK.
compileGraph turns a graph into dynamic workflow source: each step node
calls its own step, sleep nodes call sleep and approval nodes wait on
createHook({ token: "<runKey>:<node>" }). Pass it as compile, and
graphStarter as start, to createBuilder.
import {
executeGraph,
type GraphStepCall,
type WorkflowGraph,
} from "better-supabase/blocks/workflow-builder";
import { createHook, sleep } from "workflow";
import { steps } from "./steps";
export async function graphExecutor(
graph: WorkflowGraph,
input: unknown,
meta: { runKey: string },
) {
"use workflow";
return executeGraph(graph, input, {
runKey: meta.runKey,
step: (call: GraphStepCall) => steps[call.step](call),
sleep: (ms) => sleep(ms),
approval: (token) => createHook({ token }),
});
}The executor imports executeGraph from the block rather than from
workflow-sdk/builder, so the workflow bundle doesn't load workflow/api,
which the workflow sandbox can't run.
import {
nodeRunReporter,
type GraphStepCall,
} from "better-supabase/workflow-sdk/builder";
const report = nodeRunReporter(serviceBuilder.nodeRuns);
export async function sendEmail(call: GraphStepCall) {
"use step";
return report(call, () => email.send(call.config));
}
export const steps = { "email.send": sendEmail };Dynamic workflows are experimental in the Workflow SDK. With
WORKFLOW_EXPERIMENTAL_DYNAMIC_WORKFLOWS=1, graphStarter compiles the
published version's graph again on every start and runs that source; it
never runs the compiled form stored with the version. Source over 128 KiB
fails to publish. Node ids are 1 to 100 letters, digits, underscores or
hyphens, in validateGraph and in validate_workflow_graph. Without it, graphStarter starts executor, a static
"use workflow" function that walks the stored graph with executeGraph
and takes the same branches. Either way a run carries bs.definition,
bs.version, bs.tenant, bs.actor and bs.key, and a start with a key
that already ran returns the first run.
nodeRunReporter records each step node's status, output and error in
workflow_node_runs for the canvas. A record that fails never fails the
step.
Analytics
world.analytics lists runs (by workflow name, status, attributes and a
time window), their steps, events, hooks and waits, and the attribute keys
in use. Pages hold at most 100 runs or 1000 steps and events. startFor
uses it to find a run by its idempotency key.
Last updated on
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.
Workflow builder
Graph workflows that members edit and publish per tenant, with webhook, schedule and event triggers, credentials by reference, node-level run status and alerts on failed or slow runs.