Your first workflow, in one file.
Three steps, one rollback handler, no services to start. By the end you will have run it, broken it on purpose, and read back exactly what the engine recorded.
Ten minutes, one file, no services to start.
bun add bunqueueThe workflow
Section titled “The workflow”Each .step() gets a name and a handler. The handler’s return value becomes ctx.steps.<name> for every step after it, fully typed, no casts:
import { Workflow, Engine } from 'bunqueue/workflow';
const orderFlow = new Workflow<{ orderId: string; amount: number }>('order-pipeline') .step('validate', async (ctx) => { // ctx.input is typed as { orderId: string; amount: number } if (ctx.input.amount <= 0) throw new Error('Invalid amount'); return { orderId: ctx.input.orderId, validated: true }; }, { retry: 1 }) .step('charge', async (ctx) => { // ctx.steps.validate is typed from the previous step's return value const txId = await payments.charge( ctx.steps.validate.orderId, ctx.input.amount, { idempotencyKey: ctx.idempotencyKey }, ); return { transactionId: txId }; }, { compensate: async (ctx) => { // A failed charge may have committed without returning its transaction id. const charge = ctx.steps.charge ?? await payments.findByIdempotencyKey(ctx.forwardIdempotencyKey); if (charge) { await payments.refund(charge.transactionId, { idempotencyKey: ctx.idempotencyKey, }); } }, }) .step('confirm', async (ctx) => { await mailer.send( 'order-confirm', { txId: ctx.steps.charge.transactionId }, { idempotencyKey: ctx.idempotencyKey }, ); return { emailSent: true }; });The provider methods are application code, but their idempotency arguments are not decorative. They make a retry of an outcome-unknown charge or email land on the same external operation. The compensate handler also reconciles by the forward key because the charge most in need of reversal may be the one whose response never came back.
Run it
Section titled “Run it”const engine = new Engine({ embedded: true, dataPath: './data/wf.db' });engine.register(orderFlow);await engine.recover(); // after every definition is registered, before new work
const run = await engine.start('order-pipeline', { orderId: 'ORD-1', amount: 99.99 });Watch it finish
Section titled “Watch it finish”start() returns as soon as the first node is enqueued; the run continues in
the background. Poll durable state when you need a definitive answer:
const terminal = new Set(['completed', 'failed']);let exec = engine.getExecution(run.id);while (exec && !terminal.has(exec.state)) { await Bun.sleep(50); exec = engine.getExecution(run.id);}console.log(exec?.state);Event subscriptions are live notifications, not a replay log. Attach
engine.onAny() before start() if you must observe the complete event
sequence; subscribe(run.id, ...) is useful for updates after the handle is
known, but a very short workflow may already have emitted early events.
Make it fail
Section titled “Make it fail”Change confirm to throw and run it again. The engine records the failure, then walks backwards through the steps that completed and calls their compensate handlers in reverse:
const failedExecution = engine.getExecution(run.id);failedExecution?.state; // 'failed'failedExecution?.failureReason; // the error from `confirm`failedExecution?.rollbackStatus; // 'completed', unwind finishedfailedExecution?.steps.charge?.compensation?.status; // 'compensated'Two separate facts, two separate fields: why the run failed, and what the rollback then did. They are not the same question, and collapsing them makes it impossible to alert on the right one.
Shutting down cleanly
Section titled “Shutting down cleanly”import { shutdownManager } from 'bunqueue/client';
await engine.close();shutdownManager(); // stops process-wide timers and flushes pending writesWithout shutdownManager() a script that finishes its work will not exit: bunqueue’s background maintenance timers are shared across the process and keep the event loop alive.
- Steps & Control Flow, retries, branching, parallel, loops
- Rollback, what the undo actually guarantees