Skip to content
Mayura

Workflows

Durable workflows

Define and run workflows that survive restarts, wait for people and timers, and never repeat a step whose outcome is unknown.

A durable workflow is a fixed graph of steps whose progress is stored in a database. Each step is recorded before it starts and after it finishes, so a run can wait days for an approval or a date, survive a deploy or a crash, and carry on in another process. Use one when a sequence of tool calls must finish reliably and must not repeat an effect such as a payment or an email.

Durable workflows live in mayura/workflows/lifecycle. If you have not read Workflows yet, start there for the vocabulary.

A complete example

This workflow reserves stock and then sends a confirmation. It runs on SQLite (install better-sqlite3).

ts
import { defineTool, z } from 'mayura';
import { createSqliteStore } from 'mayura/storage-sqlite';
import { createWorkflowLifecycleRuntime, defineWorkflowLifecycle } from 'mayura/workflows/lifecycle';

const order = z.object({ orderId: z.string(), email: z.string() });
const reservation = z.object({ orderId: z.string(), email: z.string(), reservationId: z.string() });

const reserve = defineTool({
  id: 'orders.reserve', version: '1', description: 'Reserve stock for an order.',
  input: order, output: reservation, effects: 'write', capabilities: ['inventory:reserve'],
  execute: async request => ({ ...request, reservationId: `res_${request.orderId}` }),
});
const confirm = defineTool({
  id: 'orders.confirm', version: '1', description: 'Email the customer that the order is reserved.',
  input: reservation, output: z.object({ messageId: z.string() }), effects: 'write', capabilities: ['email:send'],
  execute: async request => ({ messageId: `msg_${request.reservationId}` }),
});

const fulfil = defineWorkflowLifecycle({
  id: 'orders.fulfil',
  version: '1',
  input: order,
  output: z.object({ messageId: z.string() }),
  nodes: [
    { kind: 'tool', id: 'reserve', tool: reserve, input: { kind: 'input', path: [] } },
    { kind: 'tool', id: 'confirm', tool: confirm, dependsOn: ['reserve'], input: { kind: 'step', stepId: 'reserve', path: [] } },
  ],
  result: { kind: 'step', stepId: 'confirm', path: [] },
});

const store = createSqliteStore({ filename: 'workflows.sqlite' });
await store.initialize();

const runtime = createWorkflowLifecycleRuntime({
  store,
  scope: { principalId: 'orders-service', projectId: 'shop' },
  permissions: { allow: ['tool:orders.reserve', 'inventory:reserve', 'tool:orders.confirm', 'email:send', 'effect:write'] },
  policyVersion: '1',
  maxCostMicros: 0,
});

try {
  const run = await runtime.submit(fulfil, { input: { orderId: 'o-1001', email: 'ada@example.com' }, idempotencyKey: 'o-1001' });
  const settled = await runtime.runUntilSettled(fulfil, run.id);
  console.log(settled.status, settled.output); // succeeded { messageId: 'msg_res_o-1001' }
} finally {
  runtime.close();
  await store.close();
}

submit validates the input and stores a new run. The run id is derived from idempotencyKey, so submitting the same key again returns the same run instead of starting a second one. runUntilSettled then advances the run in the calling process until it finishes, waits or is paused.

Runtime options

Option Required Meaning
store yes An initialized store from mayura/storage-sqlite or mayura/storage-postgres.
scope yes { principalId, projectId }. Runs are stored and looked up per scope.
permissions yes The explicit grants every step needs: tool:<id>, each capability the tool declares, and effect:<kind> for tools with effects. A step without its grants ends blocked.
policyVersion yes A label for this set of settings, recorded with each run.
maxCostMicros yes The cost budget of each run, in micros, shared by all its steps. 0 allows only free tools.
verifyHuman for approvals and respond Turns a credential into a verified person. See Approvals and human input.
approvalTtlMs no How long an approval request stays valid. Default 1 hour.
maxOutputBytes no Largest input, step output or result a run stores. Default 1 MiB.
previousPolicies no The settings of earlier releases (up to 16), so runs they started can finish. See Changing runtime settings.
now no The clock, for tests. Timers and deadlines use it.

How a run advances

runUntilSettled works in waves. Each wave starts every step whose dependencies are done, runs sibling steps in parallel, and stores each result. It returns a snapshot when there is nothing left to do right now:

Status Meaning
running Work remains; call runUntilSettled again.
waiting Waiting for an approval, a person, a timer or a signal. nextWakeAtMs is the earliest deadline or due time.
paused An operator paused it. Nothing runs until it is resumed.
succeeded Every step is done and output holds the validated result.
failed A step failed, a human request or signal timed out, or the result did not match the output schema.
blocked A step lacked a permission or the run budget could not cover it.
cancelled Cancelled with runtime.cancel(runId).
outcome_unknown A step with effects may or may not have happened. See Unknown outcomes.

When a step fails or is blocked, the steps that depend on it are skipped and the run ends. There are no automatic retries. The snapshot also has steps (each step's status and output) and budget (spentMicros, reservedMicros and maxCostMicros). runtime.inspect(runId) returns the same snapshot without advancing the run, and runtime.events(runId) returns its stored event log.

Tool steps

Before a tool step runs, the runtime checks its grants and reserves the tool's costMicros from the run budget. If the remaining budget cannot cover it, the step is blocked. After the step, the run is charged what the tool reported with context.reportUsage, or its declared cost if it reported nothing.

How a thrown error is recorded depends on the tool's effects:

  • A none or read tool that throws simply fails the step: it changed nothing outside, so there is nothing to reconcile.
  • A write or host tool that throws is recorded as unknown, because it may have acted before failing. The run ends outcome_unknown.
  • Throw ToolRefusal (from mayura) when the tool decided not to act, for example because the order does not exist. The step fails cleanly, its reserved cost is released, and there is nothing to reconcile.

Every step's tool receives a stable context.callId (<runId>/step:<nodeId>). Pass it, or a business key such as the order id, to your provider as an idempotency key.

Agent steps

agentStep turns an agent into a tool for a step. Each time the step runs, it runs the agent in its own in-process runtime with the run's scope, and charges the run what the agent actually spent:

ts
import { agentStep } from 'mayura/workflows/lifecycle';

const summarize = agentStep(summarizer, {
  id: 'tickets.summarize',              // the workflow grants tool:tickets.summarize
  capabilities: ['tickets:summarize'],  // and each of these
  permissions: ['model:openai.responses'], // what the agent itself may use
  limits: { maxCostMicros: 50_000, maxDurationMs: 60_000 }, // the step reserves 50,000 from the run budget
});

const nodes = [{ kind: 'tool' as const, id: 'summarize', tool: summarize, input: { kind: 'input' as const, path: [] } }];

The agent's grants are separate from the workflow's: the workflow grants the step (tool:tickets.summarize, its capabilities, and effect:<kind> when the agent's tools have effects); the agent gets only permissions. The step's outcome follows the agent's:

Agent outcome Step Charged
succeeded succeeds with the agent's output what the agent spent
failed, blocked, cancelled fails; its dependents are skipped what the agent spent
outcome_unknown unknown: the run ends outcome_unknown and is left for you to reconcile its ceiling stays reserved

Cancelling the run or reaching the step's timeout cancels the agent. If the agent's tools can change things (write or host effects), the interrupted step is unknown, because the agent may have been in the middle of one. The step's timeout defaults to the agent's maxDurationMs plus 15 seconds, so the agent's own limit fires first.

Option Meaning
id, version, description The step tool's identity. version defaults to the agent's.
permissions Grants for the agent's own model and tool calls.
limits The agent's runtime limits. maxCostMicros is the step's ceiling (0 for a model that costs nothing).
capabilities Extra grants the workflow needs to run the step.
input, output, prepare, finish Give the step its own schemas: prepare turns the step input into the agent input, and finish turns the agent output into the step output, or throws to fail the step (for example after checking citations).
onRun Called with the agent's run handle, for example to trace it. A function it returns is awaited when the run settles. Neither can fail the step.
timeoutMs The step's deadline.

To build the agent for each run, for example with tools that record what this run read, pass a function instead of the agent. It then needs input, output, and effects, the strongest effect its tools may have:

ts
const investigate = agentStep((task: { question: string }) => researcherFor(task.question), {
  id: 'research.investigate', input: taskSchema, output: findingsSchema, effects: 'read',
  permissions: researcherGrants, limits: { maxCostMicros: 20_000 },
});

The model loop inside a step is not checkpointed: a crash mid-step does not resume the agent (see Restarts and unknown outcomes).

Optional steps

Any step can declare when, a binding over the input or one of its dependencies. The step runs only when the value is something other than null or false; a path that does not exist counts as null. Otherwise the step is bypassed: it costs nothing, and its dependents continue and see its output as null.

ts
import type { WorkflowLifecycleNode } from 'mayura/workflows/lifecycle';

const legalReviewStep: WorkflowLifecycleNode = {
  kind: 'tool', id: 'legal-review', tool: legalReview, dependsOn: ['classify'],
  input: { kind: 'step', stepId: 'classify', path: ['contract'] },
  when: { kind: 'step', stepId: 'classify', path: ['needsLegalReview'] },
};

The condition is decided once, when the step's dependencies are done.

Parallel slots

fanOut runs one tool per item of an array, in parallel, up to a fixed maximum. It creates slots <id>.1 to <id>.<max> and a join <id> that collects them. A slot with no item is bypassed and reserves nothing.

ts
import { defineWorkflowLifecycle, fanOut } from 'mayura/workflows/lifecycle';

const research = defineWorkflowLifecycle({
  id: 'research', version: '1', input: question, output: report,
  nodes: [
    { kind: 'tool', id: 'plan', tool: plan, input: { kind: 'input', path: [] } },
    // research.1 .. research.4 run `investigate` on plan.assignments[0..3].
    ...fanOut({ id: 'research', items: { stepId: 'plan', path: ['assignments'] }, max: 4, tool: investigate }),
    // The join outputs one entry per slot, in order, with null for a bypassed slot.
    { kind: 'tool', id: 'write', tool: write, dependsOn: ['research'], input: { kind: 'step', stepId: 'research', path: [] } },
  ],
  result: { kind: 'step', stepId: 'write', path: [] },
});

max is at most 64 and is part of the definition. Items beyond max are ignored, so cap the array in the step that produces it.

Timers

A timer step waits until an absolute time, in milliseconds since the Unix epoch, read from a binding:

ts
import type { WorkflowLifecycleNode } from 'mayura/workflows/lifecycle';

const sendAt: WorkflowLifecycleNode = { kind: 'timer', id: 'send-at', dependsOn: ['draft'], fireAtMs: { kind: 'input', path: ['sendAtMs'] } };

While it waits, the run is waiting and holds no memory, timer handle or process. Its output is { fireAtMs, firedAtMs }. Something must call runUntilSettled again at or after nextWakeAtMs; in production the worker host does this for you.

Waiting for people

Two step types stop a run until a person acts:

  • Approvals. Add approval: true to a tool step. The run waits until someone approves that exact tool call.
  • Human requests. A human step asks a typed question (information, a correction or a choice of plan) and continues with the validated answer as the step's output.
ts
import type { WorkflowLifecycleNode } from 'mayura/workflows/lifecycle';
import { z } from 'mayura';

const review: WorkflowLifecycleNode = {
  kind: 'human', id: 'review', dependsOn: ['draft'],
  request: {
    kind: 'information',
    schemaId: 'posts.review', // with a Zod response, the schema digest is derived from it
    prompt: 'Is this post ready to publish?',
    response: z.object({ publish: z.boolean(), note: z.string().optional() }),
    context: { kind: 'step', stepId: 'draft', path: [] },
    deadlineAtMs: { kind: 'input', path: ['reviewBy'] },
  },
};

Responding, verifying who may answer, and the HTTP, CLI and UI surfaces are covered in Approvals and human input.

Signals

A signal step waits for an event from another system, such as a payment arriving. The run is waiting and holds nothing while it waits. The signal's payload is validated with payload and becomes the step's output, which later steps bind to like any other output:

ts
import { defineWorkflowLifecycle } from 'mayura/workflows/lifecycle';
import { z } from 'mayura';

const checkout = defineWorkflowLifecycle({
  id: 'orders.checkout', version: '1', input: order, output: z.object({ messageId: z.string() }),
  nodes: [
    { kind: 'tool', id: 'reserve', tool: reserve, input: { kind: 'input', path: [] } },
    { kind: 'signal', id: 'paid', name: 'payment.received', dependsOn: ['reserve'],
      payload: z.object({ amountCents: z.number().int().positive() }),
      deadlineAtMs: { kind: 'input', path: ['payBy'] } },
    { kind: 'tool', id: 'ship', tool: ship, dependsOn: ['paid'], input: { kind: 'step', stepId: 'paid', path: ['amountCents'] } },
  ],
  result: { kind: 'step', stepId: 'ship', path: [] },
});

Deliver a signal from code with the runtime (use the fleet runtime in production, so the worker sees the change):

ts
await runtime.signal(checkout, {
  id: runId, name: 'payment.received', signalId: 'payment-7731', payload: { amountCents: 4_200 },
});

Operators and other services can deliver it over HTTP with client.signalWorkflow or mayura workflow-signal; see Operating workflows. Either way:

  • Idempotent. signalId identifies the signal. Delivering the same id and payload again changes nothing, so a sender can retry freely. The same id with another payload, or a second signal for a step that already has one, is refused with CONFLICT: a signal step takes exactly one signal.
  • Validated. A payload the schema rejects is refused with INVALID_INPUT and the step keeps waiting. An unknown signal name is NOT_FOUND.
  • Early signals are kept. A signal that arrives before the step starts (its dependencies are still running) is stored with the step and completes it as soon as it starts. If the step is bypassed or the run is cancelled first, the signal is dropped with it.
  • Deadlines. With deadlineAtMs, a step with no signal by then is timed_out, the run fails, and later signals are refused. A kept early signal counts only if it arrived before the deadline.
  • name defaults to the node id and is unique within a definition. The run continues on its next runUntilSettled; the worker host does this for you.

To start a new run for each event instead, see Webhooks.

Storage

Durable workflows need an initialized store: createSqliteStore from mayura/storage-sqlite for a single machine, or createPostgresStore from mayura/storage-postgres when several processes share it. Call store.initialize() once at startup. Every process that touches the same runs must use the same store and scope. The runtime never closes the store; you do. See Storage.

Workers

In production, one process submits runs and a separate worker advances them. Submit through the fleet runtime, which keeps an index of active runs, and run a host that sweeps that index:

ts
import { hostname } from 'node:os';
import { createWorkflowLeadership, createWorkflowWorker } from 'mayura/workflows';
import { createWorkflowLifecycleFleetRuntime, createWorkflowLifecycleHost } from 'mayura/workflows/lifecycle';

const options = { store, scope, permissions: { allow: grants }, policyVersion: '1', maxCostMicros: 500_000 };

// Server side: submit (and approve, respond, pause) through the fleet runtime so workers can find the run.
const runtime = createWorkflowLifecycleFleetRuntime(options);
await runtime.submit(fulfil, { input: request, idempotencyKey: request.orderId });

// Worker side: advance every indexed run on a registered definition, about once a second.
const host = createWorkflowLifecycleHost({ ...options, definitions: [fulfil], intervalMs: 1_000, maxBackoffMs: 30_000 });
const worker = createWorkflowWorker({
  units: [host],
  leadership: createWorkflowLeadership({ store, scope, role: 'orders-worker', holderId: `${hostname()}-${process.pid}` }),
});
worker.start();

// On shutdown: start nothing new, let running steps finish (up to 30 s), then hand over the lease.
await worker.drain({ timeoutMs: 30_000 });
  • The host runs waiting runs only when they are due, and skips runs whose definition is not in definitions.
  • Leadership is a durable lease: run as many worker replicas as you like, and one advances the fleet at a time. If the leader dies, another takes over when its lease expires (15 s by default).
  • In tests, call host.runOnce() instead of starting a timer.
  • With defineMayuraApplication from mayura/cli, return the worker from worker() and mayura worker runs it. See Running in production.

Restarts and unknown outcomes

Every state change is written to the store before the next one begins, so after a restart a new process continues each run from its last recorded state: waiting steps keep waiting, due timers fire, finished steps are not run again.

The one gap is a step that was running when the process died. Mayura records a step as started before it calls the tool, so it knows the tool may have acted, but not whether it did. Such a step is never run again. Once the tool's timeoutMs and a further minute have passed since it started, no live process can still finish it, so the next process that advances the run (a worker, or your own runUntilSettled) records the step as unknown, or blocked when its receipt shows the tool finished, and the run ends that way too. Until then the run stays unfinished and cannot be paused. To record it sooner, after you check the outside system, call recoverAbandoned:

ts
const snapshot = await runtime.recoverAbandoned(runId);
// The interrupted step is now `unknown` and the run ends `outcome_unknown`.

Unknown outcomes

A run ends outcome_unknown when a step with effects may or may not have happened: it threw, timed out, or was cut off by a crash. Mayura never replays it. Look at the outside system (did the refund go out?), then either finish the work by hand or submit a new run. Operators find these runs in the settled view of the run list; see Operating workflows. The best defence is an idempotent tool: pass a stable key to your provider so a repeated call is harmless.

Tracing

createWorkflowTraceExport from mayura/workflows exports each settled run as one OpenTelemetry trace, with a span per step. It keeps a durable outbox, so a restarted worker neither loses nor duplicates traces. Add its unit to the worker:

ts
import { createOtlpHttpJsonTraceExporter } from 'mayura/exporter-otlp';
import { createWorkflowTraceExport, createWorkflowWorker, lifecycleFleetTarget } from 'mayura/workflows';

const exporter = createOtlpHttpJsonTraceExporter({ endpoint: 'https://collector.example/v1/traces', serviceName: 'orders' });
const traces = createWorkflowTraceExport({
  source: host.runtime, store, scope, exportId: 'primary', definitions: [fulfil], sink: exporter.sink,
});
const worker = createWorkflowWorker({ units: [host, traces.unit({ targets: [lifecycleFleetTarget(host.runtime)] })] });

// Where you submit runs, record them so short runs are not missed:
await traces.track(runId);

Spans carry names, ids, times, statuses and budget numbers, never inputs, outputs or prompts. See Observability.

Changing runtime settings

Every run records the settings it started with: scope, permissions, policyVersion, maxCostMicros, maxOutputBytes and approvalTtlMs. A runtime continues a run only if those are its own settings or are listed in previousPolicies; any other run stops with CONFLICT. So when a release changes a setting, for example to grant a tool that a new workflow version calls, list the settings the previous release ran with:

ts
import { createWorkflowLifecycleFleetRuntime, type WorkflowLifecyclePolicy } from 'mayura/workflows/lifecycle';

// Exactly what release 1 passed. Keep it listed until release 1's runs have finished.
const release1: WorkflowLifecyclePolicy = {
  permissions: { allow: ['tool:orders.reserve', 'inventory:reserve', 'effect:write'] },
  policyVersion: '1',
  maxCostMicros: 500_000,
};

const runtime = createWorkflowLifecycleFleetRuntime({
  store,
  scope,
  permissions: { allow: ['tool:orders.reserve', 'inventory:reserve', 'tool:orders.confirm', 'email:send', 'effect:write'] },
  policyVersion: '2',
  maxCostMicros: 500_000,
  previousPolicies: [release1],
});
  • A run keeps its own settings until it finishes. Its steps are checked against the permissions it started with, and it keeps its own budget, output limit and approval lifetime. A grant added later is never given to it: a step that needs one ends blocked. An approval requested before the deploy can still be approved.
  • New runs use the current settings. A child that a saga or loop starts later uses its parent's settings.
  • To move a run onto the new settings, migrate it to a new definition version; see Operating workflows. This also works for a run whose settings are no longer listed.
  • maxOutputBytes and approvalTtlMs default as they do for the runtime, so copy exactly what the old release passed. Remove an entry once no active run uses it. Sagas, loops, fleet runtimes and hosts take the same option.

Good to know

  • Use the same settings everywhere. Give the server and the worker the same options object. When a release changes the settings, list the old ones in previousPolicies; see Changing runtime settings.
  • Never edit a definition that has runs in flight. Changing any step changes the definition's digest, and runs pinned to the old one stop. Add a new version instead; see Operating workflows.
  • Use the fleet runtime everywhere in production. Approving, responding, signalling or pausing through a plain runtime leaves the worker's index stale.
  • Errors. Every error the runtimes and stores throw is a MayuraError with a stable code. Reusing an idempotency key with other input, another definition version or other runtime settings is CONFLICT; resubmit the identical request to get the existing run. See Storage.
  • Limits. At most 128 steps, 32 human steps, 64 timers and 64 signal steps per definition; prompts up to 1 KiB.
  • runUntilSettled runs steps in the calling process. A plain runtime starts no background timers; waiting runs move only when something calls it again.