# Workflow Engine: Orchestration Without Temporal

bunqueue 2.7 ships a built-in workflow engine: saga compensation, branching and human-in-the-loop signals, all persisted to SQLite with no extra services.

Canonical: https://bunqueue.dev/blog/workflow-engine/

---

import { Aside, Card, CardGrid } from '@astrojs/starlight/components';

<div class="bq-wrap bq-hero">
  <span class="bq-eyebrow">blog · workflows</span>
  <h1 class="bq-hero-h1 bq-bench-h1">Orchestration without the <em>orchestra.</em></h1>
  <p class="bq-hero-sub">Real processes span multiple steps, and when step three fails you must undo one and two. That is orchestration, and until now it meant Temporal, Inngest, or Trigger.dev. bunqueue 2.7 builds it into the queue you already run.</p>
</div>

**The new Workflow Engine in bunqueue 2.7** lets you define multi-step processes (charge payment, reserve inventory, send confirmation, notify shipping) with a fluent TypeScript DSL, run saga compensation on failure, branch conditionally, and pause for human approval, all inside the same queue infrastructure you already use. Zero new services. Zero cloud dependencies.

## Why Build It In?

Most teams reach a point where individual jobs are not enough. An order pipeline, a CI/CD deployment, an onboarding flow, these are sequences of jobs with dependencies, rollback logic, and sometimes human decisions in the middle.

The established alternatives introduce a separate orchestration runtime or
control plane. That can be the right trade-off for multi-region coordination,
but it is another durable service, deployment and operational model to own.

bunqueue's approach: if you already have a job queue with SQLite persistence, the execution engine is a natural extension. No new databases, no new protocols, no new deployment targets.

## The API

Define a workflow with chained steps:

```typescript
import { Workflow, Engine } from 'bunqueue/workflow';

const order = new Workflow('order-pipeline')
  .step('validate', async (ctx) => {
    const { orderId } = ctx.input as { orderId: string };
    // Validate order exists and is payable
    return { orderId, validated: true };
  })
  .step(
    'charge',
    async (ctx) => {
      // Charge the customer
      return { transactionId: 'txn_abc123' };
    },
    {
      compensate: async (ctx) => {
        // Refund if a later step fails
        const { transactionId } = ctx.steps.charge as { transactionId: string };
        await refundPayment(transactionId);
      },
    }
  )
  .waitFor('manager-approval')
  .step('ship', async (ctx) => {
    return { tracking: 'TRACK-789' };
  });
```

Register the workflow and start a run:

```typescript
const engine = new Engine({ embedded: true });
engine.register(order);

const run = await engine.start('order-pipeline', { orderId: 'ORD-1' });
```

That is the entire setup. No YAML, no JSON state machines, no separate orchestrator process.

## Saga Compensation

The saga pattern is first-class. Each step can declare a `compensate` handler. When any step fails, bunqueue runs compensation handlers in reverse order, exactly like a distributed transaction rollback:

```
step1 ✅ → step2 ✅ → step3 ❌
                       ↓
              compensate(step2) → compensate(step1)
```

Compensation handlers receive the same context with all previous step results, so they have everything needed to undo their work.

## Conditional Branching

Not every workflow is linear. Use `.branch()` and `.path()` to route execution based on runtime data:

```typescript
const kyc = new Workflow('kyc-onboarding')
  .step('score', async (ctx) => {
    return { riskLevel: evaluateRisk(ctx.input) };
  })
  .branch((ctx) => {
    const score = ctx.steps.score as { riskLevel: string };
    return score.riskLevel; // 'low', 'medium', or 'high'
  })
  .path('low', (w) => w.step('auto-approve', async () => ({ approved: true })))
  .path('medium', (w) => w.step('manual-review', async () => ({ reviewRequired: true })))
  .path('high', (w) => w.step('escalate', async () => ({ escalated: true })))
  .step('finalize', async (ctx) => {
    return { completed: true };
  });
```

## Human-in-the-Loop

Some workflows need a human decision before continuing. `.waitFor()` pauses execution until an external signal arrives:

```typescript
const deploy = new Workflow('deploy')
  .step('build', async (ctx) => {
    /* ... */
  })
  .step('test', async (ctx) => {
    /* ... */
  })
  .waitFor('qa-approval')
  .step('deploy-prod', async (ctx) => {
    /* ... */
  });

engine.register(deploy);
const deployRun = await engine.start('deploy', {});

// Later, when QA approves:
await engine.signal(deployRun.id, 'qa-approval', { approvedBy: 'qa-lead' });
```

The durable cursor resumes at the same node after `recover()`. Completed records
are skipped; an interrupted record may replay because the engine cannot assume
that its external effect did or did not commit.

## Step Retry, Parallel Steps & Sub-Workflows

Steps retry automatically with exponential backoff:

```typescript
.step('call-api', async () => {
  const res = await fetch('https://api.flaky.com/data');
  if (!res.ok) throw new Error(`HTTP ${res.status}`);
  return await res.json();
}, { retry: 5, timeout: 10000 })
```

Run independent steps concurrently with `.parallel()`:

```typescript
.parallel((w) => w
  .step('fetch-orders', async () => await db.orders.list())
  .step('fetch-prefs', async () => await db.preferences.get())
  .step('fetch-activity', async () => await analytics.recent())
)
```

Compose workflows by nesting them with `.subWorkflow()`:

```typescript
const payment = new Workflow<{ amount: number }>('payment').step('authorize', async (ctx) => ({
  authorized: ctx.input.amount,
}));

const parent = new Workflow('order')
  .step('create', async () => ({ total: 99 }))
  .subWorkflow(
    'payment',
    (ctx) => ({
      amount: (ctx.steps['create'] as { total: number }).total,
    }),
    {
      timeout: 15 * 60_000,
      pollInterval: 250,
    }
  )
  .step('confirm', async (ctx) => {
    const result = ctx.steps['sub:payment'];
    return { done: true };
  });

engine.register(payment); // register the child definition before the parent runs
engine.register(parent);
```

The parent polls the child's durable execution record. Its timeout is measured
from the child's original creation time, so restarting the parent does not grant
another full window; expiry also does not forcibly cancel a child that is still
running.

## Observability & Lifecycle Events

Subscribe to typed events for monitoring:

```typescript
engine.on('step:failed', (e) => alerting.send(e));
engine.onAny((e) => metrics.increment(`workflow.${e.type}`));
```

Fifteen event types cover the full lifecycle: workflow, step, signal and
compensation events, including `workflow:started`, `step:retry`,
`signal:timeout` and `compensation:failed`.

## Loops, forEach & Map

Iterate with `doUntil` / `doWhile`, process lists with `forEach`, and transform data with `map`:

```typescript
const pipeline = new Workflow('etl')
  .step('fetch', async () => ({ records: await db.getAll() }))
  .forEach(
    (ctx) => (ctx.steps['fetch'] as any).records,
    'enrich',
    async (ctx) => {
      const record = ctx.steps.__item;
      return await enrichment.process(record);
    }
  )
  .map('aggregate', (ctx) => {
    // Collect all forEach results
    const results = [];
    let i = 0;
    while (ctx.steps[`enrich:${i}`]) {
      results.push(ctx.steps[`enrich:${i}`]);
      i++;
    }
    return { total: results.length };
  })
  .doUntil(
    (ctx) => (ctx.steps['sync'] as any)?.synced === true,
    (w) =>
      w.step('sync', async () => {
        const ok = await externalApi.sync();
        return { synced: ok };
      }),
    { maxIterations: 10 }
  );
```

## Schema Validation & Subscribe

Validate step inputs/outputs with any `.parse()` schema (Zod, ArkType, Valibot) and monitor specific executions:

```typescript
import { z } from 'zod';

const flow = new Workflow('validated').step(
  'process',
  async (ctx) => ({ id: 'usr_1', valid: true }),
  {
    inputSchema: z.object({ email: z.string().email() }),
    outputSchema: z.object({ id: z.string(), valid: z.boolean() }),
  }
);

const run = await engine.start('validated', { email: 'user@test.com' });
const unsub = engine.subscribe(run.id, (e) => console.log(e.type));
```

`subscribe()` observes future in-process events; it is not a replay log. Read
`getExecution(run.id)` for durable state, and call `unsub()` when the live
subscription is no longer needed.

<Aside type="note">
  bunqueue does not bundle a schema library. Schema validation is duck-typed on `.parse()`, so you
  bring your own, `bun add zod` (used above), ArkType, Valibot, or anything with a compatible
  `.parse()` method. `msgpackr` is bunqueue's only runtime dependency.
</Aside>

## Cleanup & Archival

Archive old executions to keep your SQLite lean:

```typescript
engine.archive(30 * 24 * 60 * 60 * 1000); // Archive 30-day-old executions
engine.cleanup(7 * 24 * 60 * 60 * 1000); // Delete 7-day-old executions
```

## How It Works Internally

The workflow engine is built on top of bunqueue's existing Queue and Worker:

1. **Workflow** is a pure ordered graph of top-level nodes and inline step groups
2. **Engine** wraps a Queue and Worker pair, using them to schedule and process step jobs
3. **Execution state** is stored in SQLite via `WorkflowStore`, including the cursor, step records, signals, control-flow decisions and sealed definition identity
4. Each top-level node is driven by a bunqueue job; branch, parallel and loop body steps execute inside that node and persist their own records

Queue delivery is at-least-once. A completed record is skipped on re-entry, but
work left `running` has an unknown external outcome and may execute again.
Production handlers must therefore pass `ctx.idempotencyKey` to external
providers, use a durable `dataPath`, and call `recover()` after registering all
definitions on startup.

## When to Use It

The workflow engine is ideal for:

- **E-commerce pipelines**, validate, charge, fulfill, notify with rollback on failure
- **CI/CD deployments**, build, test, await approval, deploy
- **Onboarding flows**, verify identity, score risk, branch by result
- **Data pipelines**, extract, transform, load with per-step error handling
- **Approval workflows**, any process that needs human sign-off mid-execution

For simple fire-and-forget jobs, retries, or cron schedules, the standard Queue + Worker API is still the right choice. The workflow engine adds value when you need **step dependencies, rollback, or human decisions**.

## Get Started

```bash
bun add bunqueue
```

```typescript
import { Workflow, Engine } from 'bunqueue/workflow';
```

Read the full guide: [Workflow Engine Documentation](/guide/workflow/)

<Aside type="tip" title="Executable workflow examples">
  The guide scenarios run against the real workflow engine in `test/workflow-docs-examples.test.ts`;
  framework integrations have their own package-backed suite, and the durability model drives
  generated histories against a real TCP broker and SQLite.
</Aside>