@lunora/workflow gives Lunora durable workflows: chain steps into a single
unit that survives Worker restarts and redeploys and replays
deterministically on failure. Each step's result is persisted and memoized:
when a later step throws, only that step retries; the completed ones don't run
again.
It wraps Cloudflare Workflows
(a first-party, GA durable-execution engine) with type-safe authoring: you
write a defineWorkflow and codegen emits the WorkflowEntrypoint class, wires
a typed ctx.workflows handle, and reconciles the [[workflows]] binding in
wrangler.jsonc for you.
Use it for sagas, multi-stage AI pipelines, human-in-the-loop approvals, and any long-running orchestration where "re-run the whole job from scratch" isn't good enough.
Define a workflow
Workflows live in lunora/workflows.ts as named exports. Scaffold one with
vis generate lunora-workflow --name=orderPipeline, or write it by hand:
// lunora/workflows.ts
import { defineWorkflow } from "@lunora/workflow";
import { api } from "./_generated/api";
import type { Id } from "./_generated/server";
export const channelWelcome = defineWorkflow<
{ channelId: Id<"channels"> }, // params (input)
{ channelId: Id<"channels">; posted: number } // output
>({
handler: async (context) => {
const { channelId } = context.params;
// A durable, memoized, retried step.
await context.step.do("greet", () => context.run(api.messages.send, { channelId, text: "👋 Welcome!" }, { shardKey: channelId }));
// The workflow hibernates here and resumes after a minute —
// the delay survives Worker evictions and redeploys.
await context.step.sleep("settle", "1 minute");
await context.step.do("tips", () => context.run(api.messages.send, { channelId, text: "Tip: messages stream live." }, { shardKey: channelId }));
context.log.info("welcome sequence complete", { channelId });
return { channelId, posted: 2 };
},
});defineWorkflow<Params, Output> takes a single config object:
handler: the workflow body; receives the run context (below).name(optional): override the deployed workflow name. Defaults to the kebab-cased export name (orderPipeline→order-pipeline).
The run context
The handler receives one argument bundling the durable-step API and a typed function caller:
| Member | What it is |
|---|---|
context.params | The input passed to .create({ params }), typed as Params. |
context.step | The durable-step API: do / sleep / sleepUntil / waitForEvent. |
context.run | Call a Lunora function (api.*) by typed reference; { shardKey } targets a shard, { dedupId } overrides the replay-dedup id. |
context.runStep | Run a reusable, schema-validated defineStep as a durable step (see below). |
context.parallel | Fan out branches as isolated child workflow instances and await their outputs (see below). |
context.spawn | Fire-and-forget start of a declared child workflow (replay-safe; returns a live handle). |
context.waitForEvent | Hibernate until a declared event (defineWorkflowEvent) arrives; resolves with its validated payload. |
context.log | Structured logger (debug/info/warn/error), captured by wrangler tail and the Studio logs. |
context.event | The raw triggering event (instanceId, payload, timestamp, workflowName). |
context.env | The Worker environment bindings. |
Durable steps
context.step is the Cloudflare Workflows step API:
// Memoized + retried. On replay, a completed step returns its stored
// result without re-running the callback.
const order = await context.step.do("load", () => context.run(api.orders.get, { id: orderId }, { shardKey: orderId }));
// Per-step retry policy. A raw `step.do` retries its callback IN PLACE, so a
// call that writes needs an explicit `dedupId` — see "Exactly-once dispatches"
// below, or reach for `context.runStep`, which pins one for you.
await context.step.do("charge", { retries: { limit: 3, delay: "2 seconds", backoff: "exponential" }, timeout: "30 seconds" }, () =>
context.run(api.payments.charge, { orderId }, { dedupId: `charge:${orderId}` }),
);
// `delay` can also be a function of the failed attempt: return the next wait as a
// duration. It replaces `backoff`, so it can honour a provider's Retry-After.
await context.step.do(
"sync",
{ retries: { limit: 5, delay: ({ ctx, error }) => (error.message.includes("rate limit") ? `${ctx.attempt * 30} seconds` : "10 seconds") } },
() => context.run(api.customers.sync, { orderId }, { dedupId: `sync:${orderId}` }),
);
// Durable delays — minutes to weeks, not held in memory.
await context.step.sleep("cool-off", "1 minute");
await context.step.sleepUntil("run-at", new Date("2026-07-01T00:00:00Z"));
// Pause until an external event (approvals, webhooks, human-in-the-loop).
// Prefer `context.waitForEvent` with a declared event — see below.
const approval = await context.step.waitForEvent<{ approved: boolean }>("await-approval", {
type: "approval",
timeout: "7 days",
});Steps must be JSON-serialisable and idempotent, the same determinism contract Cloudflare (and Convex) impose. A step body may run more than once on retry, so don't rely on side effects outside the step's return value.
Exactly-once dispatches (dedupId)
A workflow is a replay machine: a failed step body is retried in place, and the
handler body itself re-executes from the top on every activation (after a
step.sleep, a waitForEvent, an eviction). Every context.run on those paths
is therefore issued more than once — so each one carries a deterministic
replay-dedup id, which the shard applies as its (identity, mutationId)
idempotency key. The same logical call re-issued on a retry or a replay resolves
with the first application's result instead of applying again. The id is scoped
by the workflow's export name as well as the instance id
(<workflow>/<instanceId>#body.<n>, …#step<i>.<n>), because an instance id is
unique only within one workflow: two workflows can both run as order-42.
Pinned for you on three surfaces:
- top-level
context.runin the handler body — stable across activations; context.runinside acontext.runStepbody — stable across that step's retries;context.runinside a step'srollback— stable across the rollback's own retries, and never shared with the forward call.
Three paths are not covered, and stay at-least-once:
- A bare
context.runinside a rawcontext.step.do(...)callback. The pin numbers calls in body order and cannot see thatstep.dois re-invoking its callback in place, so attempt 2 issues a new id. Usecontext.runStep, or pass your own{ dedupId }(as in thechargeexample above). - A body whose
context.runorder varies per attempt — branching onDate.now(),Math.random(), or unordered concurrent settles. The automatic ids are positional, so pass your own. - A replay more than 24 hours after the first attempt. Dedup rows are pruned past that window, so a workflow that sleeps for days re-dispatches on wake.
The full guarantee applies to mutations, which take the dedup read under the
shard's single-writer gate. An action or query is not gated, so the shard
claims its id instead while the handler runs: a second dispatch of the same id
that arrives before the first has finished is declined with
409 DISPATCH_IN_PROGRESS rather than run alongside it. The claim is released
when the first run settles, and after that the id is served from the replay
cache. It also expires after fifteen minutes, so a handler that never settles
cannot block the id forever. Past that point a new dispatch runs, even if the
old handler is somehow still going, so an action that can run longer than
fifteen minutes must be idempotent on its own.
Slow calls and the retry budget
The usual source of a decline is the workflow's own retry. context.run gives
up on a call after 30 seconds (timeoutMs overrides it), but the shard keeps
running the call. The step retries and re-issues the same id, and that id is
still claimed. A decline is not a failure of the call, so context.run does not
throw it into the engine's retry budget. It waits and re-dispatches the same id,
after pauses that start at 1 second and double up to 30 seconds, until the
first run's result is served. If that run died, the re-dispatch runs the call
itself. The wait is bounded by the claim's fifteen-minute ceiling, counted from
the first decline. A decline after that is rethrown and retried like any other
failure. The guarantee stays at-least-once: a decline is never treated as
success.
Inside context.runStep, the wait is also bounded by the step's own timeout
(Cloudflare's default is 10 minutes). It stops 2 seconds before that timeout
and rethrows the decline, so the attempt fails once, with the decline, and the
next attempt picks the wait up. Without that, the engine would time the attempt
out anyway and the abandoned wait would keep re-dispatching next to the retry.
A re-check can find the claim gone and run the call itself, so each one gets no
more timeoutMs than the attempt has left, and the wait ends while at least a
second remains. A step's rollback is bounded the same way by its own
rollbackConfig.timeout. The engine hands a rollback the forward step's
ctx.config, so the rollback's timeout is not read from there. A
context.run inside a raw context.step.do(...) cannot see its step's timeout
and is bounded by the claim alone. The Node host does not enforce step
timeouts, so there this bound charges an attempt the host would not have ended.
Two things can still spend an attempt, then: a wait that reaches the step's
timeout, and the 30-second give-up itself, which is a retryable failure like any
other. Give a slow action's step a timeout above its expected runtime, or
raise the call's timeoutMs.
Declared events (defineWorkflowEvent)
The native step.waitForEvent(name, { type }) matches instance.sendEvent({ type })
by a bare string, and the payload crosses as unknown. Both ends drift silently:
a typo'd type hibernates the instance until its timeout (24 hours by default) with
no error anywhere, and a payload whose shape changed resumes the workflow on
garbage. Declare the event once instead, so there is no literal to typo and no
second place to update on a rename, and both ends parse the payload:
// lunora/events.ts
import { defineWorkflowEvent } from "@lunora/workflow";
import { v } from "@lunora/values";
export const orderApproved = defineWorkflowEvent("order-approved", v.object({ approvedBy: v.string() }));Wait on it inside the workflow — the payload is typed and parsed:
const { approvedBy } = await context.waitForEvent(orderApproved, { name: "await approval", timeout: "7 days" });Send it from a mutation or action with the same definition:
export const approve = mutation.input({ approvedBy: v.string(), instanceId: v.string() }).mutation(async ({ ctx, args: { approvedBy, instanceId } }) => {
await ctx.workflows.get("orderPipeline").sendEvent(instanceId, orderApproved, { approvedBy });
});The payload is validated on send (a bad value fails the caller's request, before the workflow is woken) and again on receive. A payload that fails the validator on receive is a non-retryable failure — the event has already been consumed, so replaying the wait could only hibernate the instance again.
Pass a stable { name } on any wait that can outlive a deploy. The durable step is otherwise labelled event:<type>, which couples the step's identity
to the wire type: rename the event type and an instance that already recorded the wait replays into a fresh step that nothing will ever satisfy — it
hibernates to its timeout. An explicit name also tells two waits on the same event type apart in the timeline.
What this does not buy is a compile error for every mismatch: two definitions
with the same payload shape are mutually assignable, so sending orderRejected
where the workflow awaits orderApproved still type-checks. What it removes is the
hand-matched string.
Event types beginning with lunora: are reserved for the framework's own
instance-to-instance protocol (the context.parallel branch join). defineWorkflowEvent
and the typed sendEvent above reject them; the native untyped
instance.sendEvent({ type, payload }) does not go through that check.
Reusable steps (defineStep)
An inline context.step.do("name", () => …) is fine for one-offs, but
defineStep lets you define a step once, schema-validated and reusable
across workflows. Args are validated (with @lunora/values)
before the body runs, and the return value is validated after (when you
declare returns), so a bad payload fails fast instead of corrupting later
steps. Scaffold one with vis generate lunora-step --name=chargeOrder (appends
to lunora/steps.ts), or write it by hand:
// lunora/steps.ts
import { defineStep } from "@lunora/workflow";
import { v } from "@lunora/values";
import { api } from "./_generated/api";
export const charge = defineStep("charge", {
args: { orderId: v.string(), amount: v.number() },
returns: v.object({ receiptId: v.string() }),
config: { retries: { limit: 3, backoff: "exponential", delay: "10 seconds" } },
handler: async (context, { orderId, amount }) => {
if (context.attempt > 1) context.log.warn(`retrying charge for ${orderId}`);
return context.run(api.payments.charge, { orderId, amount });
},
// Compensation — runs if a *later* step fails after this one committed.
rollback: async (context) => {
await context.run(api.payments.refund, { orderId: context.args.orderId });
},
});Run it from a workflow body; context.runStep wraps it in a durable
step.do(...) for you:
export const orderPipeline = defineWorkflow<{ orderId: string }>({
handler: async (context) => {
const { receiptId } = await context.runStep(charge, { orderId: context.params.orderId, amount: 4200 });
return { receiptId };
},
});The step handler context gives you context.attempt (1-based retry counter),
context.config, context.env, context.run, context.log, and
context.step ({ name, count }). Pass
context.runStep(step, args, { config }) to override the step's durability
config for a single call.
Both context.runs above — the handler's and the rollback's — carry a pinned
replay-dedup id, so a retried charge does not charge twice and a retried
refund does not refund twice. See
Exactly-once dispatches for the paths that
still need an explicit dedupId.
Fan-out with child-DO isolation (ctx.parallel / ctx.spawn)
A workflow instance is one Durable Object: ~128 MB memory and a 5-minute CPU
budget, shared across everything it runs. Fanning out heavy work with
Promise.all(branches.map((b) => ctx.runStep(b, …))) therefore crowds every
branch onto that one budget: one OOM or timeout takes the whole batch down.
ctx.parallel runs each branch as its own child workflow instance, with its
own DO and its own memory / CPU / retry budget, and the parent hibernates (zero
cost) until the branches report back. Branches reference declared child
workflows by their lunora/workflows.ts export name;
branch<Output>(name, params) carries the result type into the returned tuple,
in declaration order:
import { branch, defineWorkflow } from "@lunora/workflow";
export const mediaPipeline = defineWorkflow<{ key: string }>({
handler: async (ctx) => {
// Each branch is a separate declared workflow running in its own DO.
const [labels, thumb, transcript] = await ctx.parallel([
branch<{ tags: string[] }>("imageTag", { key: ctx.params.key }),
branch<{ url: string }>("thumbnail", { key: ctx.params.key }),
branch<{ text: string }>("transcribe", { key: ctx.params.key }),
]);
return { labels, thumb, transcript };
},
});- Isolated resources: each branch gets a full DO budget; a heavy branch can't starve its siblings.
- Hibernating join: the parent consumes nothing while branches run; each child signals its result back when it finishes.
- Fail-fast: if any branch fails,
ctx.parallelrejects with that branch's error (non-retryable, because retrying the join can't re-run an already-failed child). Still-running siblings are left to finish (Cloudflare can't cleanly cancel a running instance). - Replay-safe: child instance ids are derived deterministically from the
parent, so a parent replay re-attaches to the existing children instead of
double-spawning. There is a cap of
MAX_BRANCHES(100) per call. - Bounded branch output: a branch reports back over Cloudflare's event
channel, which caps a payload at 1 MiB. A larger return value can never
reach the parent, so the branch is failed immediately with
BranchOutputTooLargenaming its byte count — rather than leaving the parent hibernating on its join until the branchtimeout(Cloudflare's default is 24 hours) and then compensating a branch that had succeeded. Return a reference the parent can dereference (an R2 key, a row id) instead of the payload itself.
For fire-and-forget (start a child pipeline without awaiting it), use
ctx.spawn(name, params), which returns a live instance handle:
await ctx.spawn("sendReceipt", { orderId: ctx.params.orderId });Both reach child workflows through ctx.exports, keyed by the class codegen
generates, so there is nothing extra to configure. A branch/spawn name with
no matching declared workflow throws an error.
Rollback (saga compensation)
A step's optional rollback handler is forwarded to Cloudflare's native
step rollback: it runs when a later step in the same instance fails, so you can
undo a committed side effect (refund a charge, delete an uploaded object).
The rollback context carries the original args, the error that triggered it,
the step's output (if it completed), and env / run / log. Cloudflare
owns rollback ordering and execution; defineStep just wires the handler and an
optional rollbackConfig.
That native rollback is intra-instance: it undoes steps within one
workflow. To compensate across a ctx.parallel(...) group (each branch is its
own instance), give a branch a compensateWith, the export name of a workflow to
run if a sibling branch fails after this one already completed:
const [charge, reserve] = await ctx.parallel([
branch("chargeCard", { orderId }, { compensateWith: "refundCard" }),
branch("reserveStock", { orderId }, { compensateWith: "releaseStock" }),
]);If any branch fails, every already-completed sibling with a compensateWith is
rolled back in reverse declaration order: its compensation workflow is spawned
(durable, replay-safe) with { branch, error, index, output } as its ctx.params,
so refundCard receives what chargeCard returned plus the failing sibling's
error, and can ctx.runStep(...) its own undo. A group where no branch sets
compensateWith behaves exactly as a plain fail-fast fan-out, with zero overhead
until you opt in. The failing branch itself is not group-compensated (its own per-step
rollbacks already ran inside its instance).
Failing without retries (NonRetryableError)
Throw NonRetryableError from a step or the handler to fail the instance
immediately, skipping retries. It is the portable, Node-importable mirror of
cloudflare:workflows' native error (so your workflow code stays
unit-testable). The runtime converts it to the native error at the workflow
boundary.
import { NonRetryableError } from "@lunora/workflow";
export const charge = defineStep("charge", {
args: { orderId: v.string() },
handler: async (context, { orderId }) => {
const order = await context.run(api.orders.get, { id: orderId });
if (order.status === "cancelled") {
throw new NonRetryableError("order cancelled — retrying will never succeed");
}
return context.run(api.payments.charge, { orderId });
},
});Deploy settings (schedules, limits, defaultRetention)
Three optional keys on defineWorkflow are deploy configuration: codegen reads
them statically (string and number literals only) and writes them into the
workflow's wrangler.jsonc entry.
export const nightlyReport = defineWorkflow({
handler: async (ctx) => {
/* … */
},
// Each cron tick starts a new instance with no params — no scheduled() handler needed.
schedules: ["0 3 * * *"],
// Raise the per-instance step cap (10,000 by default, up to 25,000).
limits: { steps: 25_000 },
// How long finished instances keep their state, unless create({ retention }) overrides it.
defaultRetention: { successRetention: "3 days", errorRetention: "30 days" },
});These land as schedules, limits.steps, and
default_retention.{success_retention,error_retention}. Each cron expression is
checked at codegen. When you change a setting, the existing entry is updated.
When you remove one, it stays in wrangler.jsonc and reconcile warns about it,
because a removed setting and a hand-set one look the same. Delete it by hand.
schedules needs a host that starts instances itself. It is refused with a platform_unsupported_feature diagnostic on target: "node" (nothing reads the
list there) and on celld (not verified). limits and defaultRetention are ignored by the Node host.
Start and manage instances
Codegen wires a typed ctx.workflows handle onto mutations and actions.
Start an instance from a function:
// lunora/channels.ts
import { mutation, v } from "./_generated/server";
export const create = mutation.input({ channelId: v.id("channels") }).mutation(async ({ ctx, args: { channelId } }) => {
const instance = await ctx.workflows.get("channelWelcome").create({ params: { channelId } });
return { instanceId: instance.id };
});ctx.workflows.get("<exportName>") returns a handle with:
-
create({ id?, params?, retention? }): start an instance.create()resolves once the instance is queued (not when it finishes); the workflow runs on Cloudflare's infrastructure.The queue is immediate and not transactional:
create()is an RPC that leaves the shard as it is called, and nothing defers it to the surrounding mutation's commit. A mutation that starts a workflow and then throws has already queued an instance, and that instance runs against state the mutation never wrote. Start workflows from an action that calls the mutation first (await ctx.runMutation(...), thencreate()), so the write is committed before the instance exists. -
createBatch([...]): start many in one batch. -
deleteBatch([...ids]): delete up to 100 instances and their stored state. Returns{ deleted, errors }with one entry per id; unknown ids come back as errors. -
get(id): a handle to an existing instance. -
sendEvent(instanceId, event, payload): deliver a declared event to one instance, to satisfy itscontext.waitForEvent. The payload is validated before the send.
create/get hand back the native Cloudflare instance, which exposes its own
lifecycle: status() (returns status + output/error), pause(), resume(),
restart(), terminate(), delete() (stops the instance and frees its state), subscribe(), and the untyped sendEvent({ type, payload }).
subscribe({ cursor?, filter? }) streams the instance's events without polling.
It replays the full history, then waits for new events until the instance
completes, errors, or is terminated. Each event carries eventId, instanceId,
timestamp, and a type (workflow_started, step_completed,
attempt_errored, rollback_started, …) plus that type's fields. The typed
mirror leaves those fields open so new event types still type-check:
const instance = await ctx.workflows.get("orderPipeline").get(instanceId);
const subscription = await instance.subscribe({ filter: ["step_completed", "workflow_completed"] });
for (let event = await subscription.next(); !event.done; event = await subscription.next()) {
console.log(event.value.type, event.value.stepName);
}The Node host rejects subscribe() with NOT_IMPLEMENTED, because its engine
keeps no per-instance event log.
Listing instances and step timelines
The Workflow binding can only create/get and read a single instance's
status: it has no instance list and no per-step detail. To list instances or
read a step timeline (what the Studio renders), use createWorkflowsRestClient,
which talks to Cloudflare's account-scoped Workflows REST API. The API token is
a secret, so keep this server-side:
import { createWorkflowsRestClient } from "@lunora/workflow";
const client = createWorkflowsRestClient({
accountId: env.CLOUDFLARE_ACCOUNT_ID,
apiToken: env.CLOUDFLARE_API_TOKEN, // scope: Workflows Read (Edit for setInstanceStatus)
});
const page = await client.listInstances({ workflowName: "order-pipeline", status: "running" });
const detail = await client.getInstance({ workflowName: "order-pipeline", instanceId: page.instances[0].id });getInstance returns the instance plus a normalized steps[] timeline (each
step's type, attempts, timing, output, and error).
setInstanceStatus({ action }) pauses, resumes, or terminates an instance and
requires an Edit-scoped token.
Wiring (handled by codegen)
After adding lunora/workflows.ts, re-run codegen (lunora dev does this on
save). It:
- Emits a
WorkflowEntrypointsubclass per export intolunora/_generated/workflows.ts(channelWelcome→ChannelWelcomeWorkflow). - Wires the typed
ctx.workflowshandle. - Declares each workflow in
wrangler.jsonc'sexportsmap, keyed by the generated class, with itsnameand any deploy settings.
Your worker entry must re-export the generated classes so wrangler can find them:
// src/server/index.ts
export { ChannelWelcomeWorkflow } from "../../lunora/_generated/workflows";The resulting entry:
// wrangler.jsonc (codegen-reconciled)
{
"exports": {
"ChannelWelcomeWorkflow": { "type": "workflow", "name": "channel-welcome" },
},
}There is no workflows[] binding. The shard Durable Object that runs your
mutations, the Worker that runs actions and cron triggers, and every generated
WorkflowEntrypoint (ctx.spawn, ctx.parallel) each reach the workflow as
ctx.exports.ChannelWelcomeWorkflow. Agents are declared the same way. This
needs Wrangler 4.142.0 or newer (or @cloudflare/vite-plugin 1.61.0 or newer)
to run locally.
A project with a workflows[] binding for a class Lunora declares has it moved
into exports on the next reconcile, under the binding's deployed name (a
different declared name is reported, not applied). Cloudflare shares
instances between a binding and an export with the same name, so running
instances carry over. Settings you set by hand on the binding move with it. A
binding to another Worker's workflow (script_name) is left alone. When the
installed wrangler is older than 4.142 (or @cloudflare/vite-plugin older than
1.61), reconcile leaves workflows[] untouched and warns with the upgrade
command. A workflows[] block inside an env.<name> environment is reported,
not rewritten.
Work queued before the upgrade still names the old WORKFLOW_* / AGENT_* bindings: a ctx.scheduler.runAfter(workflows.x, …) job that has not fired yet,
and a ctx.parallel branch still running. After the deploy the scheduled job fails and is dead-lettered, and the branch cannot signal its parent (the
parent waits out its timeout). Let the scheduler drain and in-flight fan-outs finish before deploying the upgrade.
On celld, which has no workflow exports, the deploy projection turns each
export back into a workflows[] binding named after the class. The runtime
looks that name up on env when ctx.exports has no match.
introspectWorkflow() in @cloudflare/vitest-plugin still needs a binding: add
a test-only workflows[] binding with the same name to introspect one.
Advisor lints
Two static advisors keep workflow wiring honest:
workflow_unknown_target(error):ctx.workflows.get("name")references a workflow that doesn't exist (typo or removed export).workflow_unused(info): a declared workflow is never started by any function (dead code that's still billed as aWorkflowEntrypoint), unless it's triggered externally.
Studio
The Studio's Workflows panel (under the Functions domain) lists every declared workflow with its export name, generated class, binding, and deployed name; lets you start an instance from a JSON-params form; and shows a table of instances with their live status (queued / running / complete / errored) and output. See the Studio guide.
Cloudflare bills for workflow instance-state storage. Keep step payloads small and set a retention policy when you don't need long-lived instance history.