Semola

Workflow

Durable multi-step jobs on Redis with resumable steps

Durable workflows on Redis. Event history, deterministic replay, inline steps, multi-replica leases, and automatic orphan recovery. Workflows run multi-step processes that survive restarts. Each named step caches its result in Redis, so replay and reclaim skip work that already succeeded.

Needs a Bun.RedisClient. Workflow names must be unique in the process. Execution IDs must be unique per workflow name (Redis keys are workflow:{name}:…).

Import

import {
  defineWorkflow,
  listWorkflows,
  NotFoundError,
  type WorkflowExecution,
} from "semola/workflow";

Also exported: DuplicateWorkflowError, NonRetryableStepError, SerializationError, WorkflowStoreError, and the public workflow types (Workflow, WorkflowOptions, hooks/status/start/cancel shapes, WorkflowListItem, etc.).

Quick start

This registers embedded workers, starts one execution, then reads its current status. During replay, completed named steps are loaded from Redis instead of running again.

const onboard = defineWorkflow<{ userId: string }, { ok: true }>({
  name: "onboard-user",
  redis: redisClient,
  handler: async ({ input, step, sleep }) => {
    await step("send-email", async () => {
      await emailClient.send(input.userId);
    });

    await sleep(1000);

    await step("provision", async () => {
      await provision(input.userId);
    });

    return { ok: true };
  },
});

const { executionId } = await onboard.start({ userId: "u_1" });
const execution = await onboard.get(executionId);

start() enqueues work for embedded workers and returns pending. Poll get(executionId) for the eventual result. Call stop() on shutdown.

Durability model

Each execution has an append-only event history in Redis. A single lease owner advances the workflow by:

  1. Loading history
  2. Replaying the workflow function from the start
  3. Resolving completed step / sleep calls from history (no side effects)
  4. Running the next incomplete step inline under the same lease, or scheduling a durable timer, or completing / failing / cancelling

Side effects belong inside step. Workflow code outside step / sleep must be deterministic relative to history: no raw Date.now(), Math.random(), network, or unseeded nondeterminism.

Steps are at-least-once. step bodies must be idempotent; a crash mid-step may re-run the handler. Input, result, and step outputs are always JSON (JSON.stringify / JSON.parse).

Multi-replica and automatic recovery

N Bun processes may register the same workflow name against the same Redis. Work is distributed via one task queue and per-execution leases (lockTTL).

If a replica dies mid-run, lease expiry lets another replica (or the same process after restart) reclaim the execution. You do not need to call resume() for crash recovery; workers reclaim from Redis automatically.

History and status writes are lease-fenced (Redis compare-and-append / compare-and-set against the lease token). A writer that loses the lease cannot append; the new owner continues from history. Client paths (start / cancel / resume) append without a lease.

resume(executionId) re-queues a failed execution: persist keys, append a resume event, re-schedule failed steps, then mark active. It also finishes an interrupted resume (pending after persist, or after resume events but before active). Plain pending / running / completed / cancelled executions reject. A later failure needs a new resume event even if an older WorkflowResumed is already in history.

Examples

Start and inspect an execution

start() enqueues work and immediately returns pending. Pass a custom executionId when you want a stable key (non-empty, no :). Poll get() for the terminal result.

const { executionId, status } = await onboard.start(
  { userId: "u_1" },
  { executionId: "onboard-u_1" },
);
// status: "pending"

let execution = await onboard.get(executionId);

while (execution.status === "pending" || execution.status === "running") {
  await Bun.sleep(100);
  execution = await onboard.get(executionId);
}

console.log(execution.status, execution.result, execution.error);
console.log(execution.steps);
// steps: { name, completedAt }[] - return values stay in history, not on get()

Cancel an execution

cancel() records the request and aborts local work. The return value may still be pending or running; poll until status is cancelled.

const requested = await onboard.cancel(executionId);
console.log(requested.status);

let execution = await onboard.get(executionId);

while (execution.status !== "cancelled") {
  await Bun.sleep(100);
  execution = await onboard.get(executionId);
}

console.log(execution.cancelledAt);

Resume a failed execution

resume() re-queues a failed execution and returns pending. Completed steps stay cached and are not repeated. Non-failed statuses reject unless an interrupted resume is already in progress.

const execution = await onboard.get(executionId);

if (execution.status === "failed") {
  const resumed = await onboard.resume(executionId);
  console.log(resumed.status); // "pending"
}

List executions

listWorkflows() scans executions without requiring a registered workflow instance. Filter by workflow name, status, or both (scalars or arrays).

const active = await listWorkflows(redisClient, {
  status: ["pending", "running"],
});

const failed = await listWorkflows(redisClient, {
  name: "onboard-user",
  status: "failed",
});

for (const item of failed) {
  await onboard.resume(item.id);
}

Results are lightweight WorkflowListItem snapshots. Use instance get() for input, result, error, and step details. Filter name with a string or string array; resume through the matching workflow instance.

Fail without retrying

Call fail() inside a step for a non-retryable failure.

const payment = defineWorkflow<{ orderId: string }, { charged: true }>({
  name: "charge-order",
  redis: redisClient,
  handler: async ({ input, step }) => {
    await step("charge", async ({ fail }) => {
      const charged = await charge(input.orderId);

      if (!charged) {
        fail("card declined");
      }
    });

    return { charged: true };
  },
});

Observe retries and completion

Hooks observe real lifecycle transitions. Hook errors do not fail the workflow.

const sync = defineWorkflow<{ accountId: string }, { synced: true }>({
  name: "sync-account",
  redis: redisClient,
  retries: 2,
  hooks: {
    onRetry: ({ stepName, attempt, nextRetryDelayMs }) => {
      console.log(stepName, attempt, nextRetryDelayMs);
    },
    onComplete: ({ executionId, result }) => {
      console.log(executionId, result.synced);
    },
  },
  handler: async ({ input, step }) => {
    await step("sync", () => syncAccount(input.accountId));
    return { synced: true };
  },
});

Limit concurrency by partition

Without partitionBy, concurrency is global. With partitionBy, it is per key. partitionKey on start overrides partitionBy; resume keeps the stored key.

const deploy = defineWorkflow<{ envId: string }, void>({
  name: "deploy",
  redis: redisClient,
  concurrency: 3,
  partitionBy: (input) => input.envId,
  handler: async ({ input, step }) => {
    await step("apply", () => applyDeployment(input.envId));
  },
});

await deploy.start({ envId: "production" });
await deploy.start({ envId: "staging" }, { partitionKey: "shared" });

Different keys do not share the cap. Same-key executions are limited to concurrency.

Keep terminal history

retentionTTL controls how long completed / failed / cancelled executions stay in Redis. retentionMax caps how many terminal executions remain per workflow name.

const audit = defineWorkflow<{ userId: string }, { ok: true }>({
  name: "audit-user",
  redis: redisClient,
  retentionTTL: 60_000,
  retentionMax: 100,
  handler: async ({ input, step }) => {
    await step("record", () => recordAudit(input.userId));
    return { ok: true };
  },
});

Pass Infinity to keep forever, or 0 to unlink immediately after terminal. Failed executions can be resumed only while their keys still exist.

Graceful shutdown

stop() ends polling, waits for in-flight work, and releases this process registration. Other replicas keep reclaiming and running the same workflow name.

await onboard.stop();

Reference

Handler context

  • input, executionId, signal
  • step(name, handler) - durable side effect; handler gets { input, signal, fail }
  • sleep(ms) - durable timer (survives replay / reclaim)

fail(message) inside a step marks a non-retryable failure (NonRetryableStepError).

Options

  • name (required) - unique per process
  • redis (required) - Bun.RedisClient
  • handler (required)
  • retries - step retries before workflow fails (default: 3; 0 = fail on first error)
  • retryBackoff - { baseDelay, multiplier, maxDelay } (defaults: 1000 / 2x / 30000)
  • hooks - onStart, onRetry, onError, onComplete, onCancel (see Hooks)
  • lockTTL - execution lease TTL in ms (default: 300000); also used as capacity slot TTL. While a process holds an execution, it refreshes that execution's capacity slot for the full lifetime (including sleep and retry backoff). After process death, the Redis slot remains owned until TTL; the next reclaim re-attaches via the same executionId. Differing replica concurrency values mean the effective cap is the max.
  • retentionTTL - how long terminal executions (completed / failed / cancelled) stay in Redis, in ms (default: 86400000, 24h). Infinity keeps them forever. Any other value must be a non-negative number. 0 unlinks immediately after terminal. Pending and running keys are never expired. Failed executions can be resumed only while they still exist. A background sweep also expires leftover terminal keys with no TTL.
  • retentionMax - optional cap on terminal executions per workflow name. Must be a positive integer. Oldest are UNLINKed when the cap is exceeded. Works with or without a finite retentionTTL.
  • concurrency - max parallel instances across replicas (default: 1). Without partitionBy, all executions share one Redis slot pool of size concurrency (key *) and this process runs that many pollers. With partitionBy, each key has its own pool of size concurrency (no global cap). If replicas disagree on concurrency, the effective cap is the max.
  • partitionBy - (input) => string for per-key concurrency across replicas. Empty keys throw. Cap applies for the whole execution, including durable waits. Replaces the global concurrency cap. The key * is reserved for the unpartitioned pool.
  • pollInterval - idle poll backoff ms (default: 100)

start(input, { executionId?, partitionKey? }) - partitionKey overrides partitionBy. Custom executionId must be non-empty and must not contain :. Empty partitionKey throws.

partitionKey on start overrides partitionBy when both are present. Without partitionBy, partitionKey is stored on meta but capacity stays on the global * pool. Empty keys throw. The resolved key is stored on execution meta so resume keeps the original partition.

Failed steps retry with exponential backoff before the workflow is marked failed. Default retries: 3 means 4 total attempts. retries: 0 fails on the first error.

cancel is honored during retry backoff and sleep, not only between steps. After terminal failure, resume(executionId) re-queues the execution.

Hooks

Optional lifecycle callbacks. Errors in hooks never fail the workflow. Hooks fire on real transitions, not every history replay.

  • onStart - once when the execution first moves pending → running
  • onRetry - before each step retry backoff (attempt, nextRetryDelayMs, retriesRemaining, …)
  • onError - retries exhausted, or immediately after fail()
  • onComplete - terminal success
  • onCancel - terminal cancel

Redis keys

Prefix: workflow:

KeyPurpose
workflow:{name}:history:{executionId}append-only event list
workflow:{name}:meta:{executionId}status cache for get()
workflow:{name}:lease:{executionId}owner token + TTL
workflow:{name}:queueworkflow task queue
workflow:{name}:timerssorted set of due timers / retry delays
workflow:{name}:timer-deadunparseable timer payloads (dead letter)
workflow:{name}:activenon-terminal execution ids (reclaimer)
workflow:{name}:terminalzset of terminal execution ids when retentionMax is set
workflow:{name}:partition:{key}:{slot}concurrency slots (SET NX PX, re-ownable by same execution). Key * is the unpartitioned pool; partitionBy keys are per-key pools

Terminal meta and history keys receive PEXPIRE from retentionTTL (or are UNLINKed when the TTL has already elapsed). Leases and partition slots keep using lockTTL. listWorkflows skips empty SCAN hits (expired tombstones).

Notes

  • Keep step / sleep call order and names stable across deploys; replay matches history by call sequence (a0, a1, …), and a renamed step at the same position is nondeterminism.
  • Duplicate defineWorkflow({ name }) in the same process throws DuplicateWorkflowError.
  • Workers run embedded in your Bun process against Redis, not as a separate matching service.
  • get().steps is { name, completedAt }[] from completed history events (not a separate meta cache). Step return values are not exposed on get().
  • Successful step results are written to history even if cancel arrives mid-handler; the next advance then honors cancel.

Statuses

pending | running | completed | failed | cancelled

On this page