Workflows
Operating workflows
Run durable workflows in production: list and inspect runs, cancel, pause, hold the whole fleet, and ship new definition versions safely.
Once durable workflows run in production, operators need to see what is running, stop a run that is going wrong, freeze everything during an incident, and ship a new version of a workflow while old runs are still in flight. This page covers those tasks. It assumes you run lifecycle workflows with a worker host, as in Durable workflows.
Every operation is available three ways: from code on the runtime, over HTTP through the
agent server (and so from mayura/client, the CLI and the operator console), and for fleet
operations through a dedicated fleet control.
Serve the operator API
createWorkflowOperatorTransports builds all the server-side workflow operations over your runtimes: listing,
run views, cancel, approve, signals, pause, resume, the fleet hold and migrations. Spread its result into the server
options.
import { createAgentServer } from 'mayura/server';
import {
createWorkflowCommandJournal, createWorkflowFleetControl, createWorkflowMigrationCatalog,
createWorkflowOperatorTransports, lifecycleOperatorTarget,
} from 'mayura/workflows';
import { createWorkflowLifecycleFleetRuntime } from 'mayura/workflows/lifecycle';
const runtime = createWorkflowLifecycleFleetRuntime({ store, scope, permissions, policyVersion: '1', maxCostMicros: 500_000, verifyHuman });
const operator = createWorkflowOperatorTransports({
store, scope,
journal: createWorkflowCommandJournal({ store, scope }), // each command id applies at most once
fleet: createWorkflowFleetControl({ store, scope }), // optional: the fleet hold
migrations: createWorkflowMigrationCatalog([refundsV1toV2]), // optional: reviewed migrations
targets: [lifecycleOperatorTarget({ runtime, store, scope, definitions: [refundsV1, refundsV2], approvalCredential })],
});
const server = createAgentServer({ agents, authenticate, publicOrigin: 'https://ops.example.com', ...operator });- Register every definition version that still has runs. A run whose version is not registered is reported as not found.
approvalCredentialmaps the authenticated operator's id to a credential yourverifyHumanaccepts. Without it, approvals over HTTP are unavailable. See Approvals and human input.- The transports serve one scope. Callers authenticated for another scope see nothing and change nothing.
Callers need these capabilities in their server identity:
| Capability | Allows |
|---|---|
workflows:read |
List runs, read a run, list and preview migrations, read the fleet hold. |
workflows:control |
Cancel, approve, pause, resume, send signals. |
workflows:fleet |
Hold and release the fleet, run fleet sweeps. workflows:control does not imply it. |
workflows:migrate |
Apply a migration. |
List and inspect runs
const active = await client.workflows({ limit: 20 }); // running, waiting and paused
const finished = await client.workflows({ view: 'settled', limit: 20 }); // recently finished, and runs to reconcile
const view = await client.workflow(active.items[0]!.runId);
console.log(view.definitionVersion, view.status, view.revision, view.steps);Pages come with a next cursor; pass it back as after. The settled view keeps up to 16,384 lifecycle runs per scope,
dropping the oldest finished ones first. Runs that ended outcome_unknown are kept longest, because someone still has
to reconcile them. From code, runtime.inspect(runId) and runtime.settled() give the same information. The CLI has
mayura workflow-list (add --settled for the finished view) and mayura workflow-get.
Revisions and command ids
Every run view has a revision that changes whenever the run changes. Commands that change a run carry the revision
the operator looked at and a commandId:
await client.pauseWorkflow(runId, view.revision, { commandId: crypto.randomUUID() });- If the run changed since the operator looked, the command is refused (HTTP 409). Read the run again and decide again.
- A command id applies at most once. If a request times out, send it again with the same command id: you get the recorded outcome instead of a second action. Never make up a new id just to retry.
Cancel
runtime.cancel(runId), client.cancelWorkflow(runId, revision, { commandId }) or mayura workflow-cancel stops a run:
steps that have not started are skipped, and the run ends cancelled. Cancelling never undoes an effect that already
happened, and a step that is already running is not interrupted in its external system.
Pause and resume
A pause stops a run from starting any new step without cancelling it. Use it to investigate a run, or before migrating it.
runtime.pause(runId),client.pauseWorkflow,mayura workflow-pause.runtime.resume(runId),client.resumeWorkflow,mayura workflow-resume.
A run with a step in flight cannot be paused (CONFLICT); wait for the step to finish or recover it first. Resuming
restores waiting if the run still waits for an approval, a person or a timer; it never answers the wait for you.
Approvals and answers given while a run is paused are kept and take effect after it resumes.
Hold the whole fleet
During an incident you may want every worker to stop advancing runs at once. The fleet hold is a durable flag per scope.
import { createWorkflowFleetControl, lifecycleFleetTarget, type WorkflowFleetSweepCursor } from 'mayura/workflows';
import { createWorkflowLifecycleHost } from 'mayura/workflows/lifecycle';
const fleet = createWorkflowFleetControl({ store, scope });
// Every host checks the hold before each cycle and drives nothing while it is set.
const host = createWorkflowLifecycleHost({ ...options, definitions, hold: fleet });
await fleet.hold();
// Optionally also mark each active run paused, one page at a time.
let cursor: WorkflowFleetSweepCursor | null = null;
do { cursor = (await fleet.sweepPause([lifecycleFleetTarget(host.runtime)], { cursor })).nextCursor; } while (cursor);
// Later
await fleet.release();
do { cursor = (await fleet.sweepResume([lifecycleFleetTarget(host.runtime)], { cursor })).nextCursor; } while (cursor);hold()andrelease()are idempotent. Releasing resumes nothing by itself.- A pause sweep records which runs it paused; a resume sweep resumes only those, never runs someone paused on purpose.
- Over HTTP:
client.holdWorkflowFleet(),client.sweepWorkflowFleet('pause', { cursor }),client.releaseWorkflowFleet(). On the CLI:mayura fleet-hold,mayura fleet-sweep,mayura fleet-release,mayura fleet-get. See Operations commands. - A host that cannot read the hold fails its cycle rather than drive runs.
Signals
A lifecycle run waiting at a signal step (see Durable workflows) takes its signal
over the operator API. createWorkflowOperatorTransports includes workflowSignals, which delivers it to the step
whose name is the signal name. It needs the workflows:control capability, the run's current revision and a command
id, like every other command:
const view = await client.workflow(runId);
await client.signalWorkflow(runId, {
revision: view.revision, signalName: 'payment.received', signalId: 'payment-7731', value: { amountCents: 4_200 },
}, { commandId: 'payment-7731' });On the CLI: mayura workflow-signal --id <run id> --revision <n> --command-id <id> --signal-id <id> --signal-name <name> --value-file <file>. See Operations commands.
- The value (at most 4 KiB of JSON) is validated with the step's payload schema. A value it rejects, a second signal
for a step that already has one, or a stale revision answers
409; an unknown run or signal name answers404. - Retrying with the same command id returns the recorded result. The step also remembers the signal id, so the same signal delivered again under a new command id changes nothing.
- The sender is the authenticated caller; its id is recorded in the run's event log with the signal.
- Graph and tree runs have no signal steps: a signal to one of them answers
404.
Ship a new version of a workflow
Every run is pinned to the digest of the definition it started with, and never continues under a different one. To change a workflow:
- Add a new version. Keep the same
id, give it a newversion, and submit new runs with it. Never edit a version that has runs in flight. - Keep the old version registered in
definitionsfor every host and operator target until its runs finish. The host skips runs whose version is missing (they are reported asunregistered_definition), so they wait instead of running the wrong code. - Or migrate in-flight runs to the new version after reviewing a plan (below).
- Retire the old version once nothing uses it.
Runs are also pinned to the runtime settings they started with (permissions, policyVersion, maxCostMicros,
maxOutputBytes, approvalTtlMs). If the release changes any of them, for example to grant a tool the new version
calls, pass the previous release's settings in previousPolicies on every runtime and host. Otherwise its in-flight
runs stop with CONFLICT. Listed runs finish under their own settings and never gain the new grants; migrating a run
moves it onto the new settings. See Durable workflows.
Check before you deploy
inventoryWorkflowVersions counts active runs per version. Run it against production storage before a deploy:
import { assertWorkflowVersionsRetained, inventoryWorkflowVersions, lifecycleFleetTarget } from 'mayura/workflows';
const inventory = await inventoryWorkflowVersions({
store, scope,
targets: [lifecycleFleetTarget(runtime, 'lifecycle', { includePaused: true })],
registered: nextReleaseDefinitions, // what the new release will register
});
assertWorkflowVersionsRetained(inventory); // throws if a run would be stranded, or the scan was incomplete
console.log(inventory.retirable); // registered versions with no active runsMigrate in-flight runs
A migration is a reviewed declaration that moves a paused run from one version to another. Mayura works out from the run's real state what happens to each step and refuses anything unsafe.
import { defineWorkflowMigration } from 'mayura/workflows/lifecycle';
const refundsV1toV2 = defineWorkflowMigration({
id: 'refunds-1-to-2',
from: refundsV1,
to: refundsV2,
renames: { announce: 'notify' }, // new step id -> old step id
description: 'Notify the customer once the refund is issued.',
});
await runtime.pause(runId);
const preview = await runtime.migrate(refundsV1toV2, { id: runId, actorId: 'alice', commandId: 'change-4211', dryRun: true });
if (preview.plan.allowed) {
await runtime.migrate(refundsV1toV2, { id: runId, actorId: 'alice', commandId: 'change-4211' });
await runtime.resume(runId);
}The plan lists one action per step:
| Action | When | Effect |
|---|---|---|
keep |
The step is unchanged. | Its state is carried over. |
update |
The step changed but never started. | It runs under the new definition. |
reset |
The step changed while waiting for an approval, a person or a timer. | The request is issued again; old digests stop working. |
accept |
The step changed after it finished, and you listed it in acceptCompleted. |
Its result is kept. |
add |
The step is new. | It starts pending. |
remove |
The step is gone and never started, or finished and is listed in acceptRemoved. |
It is dropped. |
A migration is always refused for a run that is not paused, a step with work in flight or an unknown outcome, a carried step whose new dependencies have not all succeeded, a rename of a finished step, and a plan made against an older revision. Migrate through the fleet runtime so the worker's index follows the run. The run's event log records the migration id, both digests, the actor and the command id.
Operators can do the same over HTTP: client.workflowMigrations(runId), client.planWorkflowMigration(runId, id) and
client.migrateWorkflow(runId, id, revision, { commandId }), backed by the migrations catalog above. The
operator console has a review dialog for it.
Runs with unknown outcomes
A run that ends outcome_unknown stays in the settled view until it drops out. A run whose step was in flight when a
process died stays active and cannot be paused until the tool's timeoutMs and a further minute have passed; the next
worker pass then records the step as unknown, without running it again. To record it sooner, after checking the
outside system, call runtime.recoverAbandoned(runId) from code. There is no HTTP or CLI command for this yet. See
Durable workflows.
Other workflow kinds and hosts
The pieces above also cover the other durable runtimes:
- Sagas and loops run on
createWorkflowCompositeHostfrommayura/workflows/composites, which also acceptshold. See Sagas and loops. - Graphs and trees (
mayura/workflows/graphs,mayura/workflows/children) are advanced by coordinators (createWorkflowGraphCoordinator,createWorkflowTreeCoordinator) that work through pages of active runs found by a discovery (createWorkflowGraphDiscovery,createWorkflowTreeDiscovery). Wrap a coordinator withcoordinatorUnitto run it in a worker, and expose its runs withgraphOperatorTargetortreeOperatorTarget. - Workers.
createWorkflowWorkersupervises up to 32 units (hosts, coordinator units, the trace export) under one leadership lease, andworker.drain()shuts them down gracefully. Give each duty its own leadershiprole.
Good to know
- Operator commands are journaled permanently, one small record per command id. Use a fresh UUID per operator action.
- Approvals are attributed to the id your
approvalCredentialandverifyHumanproduce. If every operator token maps to one service identity, so does every approval. - Hosts and coordinators only drive runs whose definitions they were given. After a migration, register the new version everywhere before resuming.