Workflows
Durable orchestration for multi-step business processes: retries, timeouts, parallelism, Saga-style rollback and crash recovery โ in-process, with no broker and no required database.
def := workflow.New("order-processing").
Step("validate", ValidateOrder).
Step("charge", ChargeCard, workflow.WithCompensation(RefundCard)).
Step("ship", CreateShipment)
engine := workflow.NewEngine()
engine.Register(def)
res, err := engine.Run(ctx, "order-processing", order) If ship fails, RefundCard runs automatically โ the failure path is
declared next to the work, not scattered through error handling.
Why this exists
A step that charges a card and a step that ships a parcel are not the same
kind of failure: one must be undone, the other retried. Written by hand,
that logic becomes nested if err != nil blocks with rollback code only
exercised in production, at the worst possible moment. This package makes
the failure path declarative, testable and observable.
Steps and ordering
A step with no declared dependency runs after the one before it โ the common case is sequential with no ceremony. Declaring dependencies opts into a validated DAG:
workflow.New("checkout").
Step("validate", Validate).
Step("charge", Charge, workflow.WithDependsOn("validate")).
Step("reserve", Reserve, workflow.WithDependsOn("validate")).
Step("ship", Ship, workflow.WithDependsOn("charge", "reserve")) charge and reserve run concurrently; ship waits for both. The graph is
validated at registration โ cycles, duplicate names, unknown dependencies
and nil handlers are all rejected before anything runs.
Passing data between steps
func Charge(ctx *workflow.Context) error {
order, ok := workflow.Payload[Order](ctx)
if !ok {
return workflow.NonRetryable(errors.New("bad payload"))
}
id, err := billing.Charge(ctx.Ctx, order.Total)
if err != nil {
return err
}
ctx.Set("charge_id", id) // visible to every later step
return nil
} The payload is typed; metadata set with ctx.Set is shared for anything
computed along the way.
Retries
workflow.WithRetry(workflow.RetryPolicy{
MaxAttempts: 5,
Backoff: workflow.BackoffExponential,
InitialDelay: 500 * time.Millisecond,
MaxDelay: 30 * time.Second,
Jitter: 0.2,
}) Jitter is not decoration: when a dependency fails, every in-flight execution
fails at nearly the same instant, and without jitter they all retry at the
same instant too, reproducing the outage that caused the failure. Retries
stop early for errors wrapped in NonRetryable, which preserves the
original error's identity:
err := workflow.NonRetryable(ErrCardDeclined)
errors.Is(err, ErrCardDeclined) // true
errors.Is(err, workflow.ErrNonRetryable) // true Compensation (Saga)
Step("reserve", Reserve, workflow.WithCompensation(Release)).
Step("charge", Charge, workflow.WithCompensation(Refund)).
Step("ship", Ship) ship fails โ Refund, then Release โ rollback runs in reverse order
over the steps that actually succeeded. A compensation handler that itself
fails does not stop the others: the remaining side effects still need
undoing, and stopping would leave strictly more damage behind. It emits WorkflowCompensationFailed, which is a page-someone event.
Timeouts and cancellation
WithTimeout bounds one attempt; Definition.Timeout bounds the whole
execution. Both cancel through ctx.Ctx, so a step that respects its
context stops promptly. Retry backoff is cancellable too โ shutdown never
waits out a 30-second delay.
Event triggers
def := workflow.New("welcome").Step("email", SendWelcomeEmail)
workflow.OnType[UserRegistered](def)
engine.Register(def) Emitting UserRegistered now starts the workflow asynchronously, with the
event as the payload โ the emitter never blocks on it.
Durability and idempotency
engine := workflow.NewEngine(workflow.Config{Store: myStore})
n, err := engine.Resume(ctx) // continue what a crash interrupted Resume replays interrupted executions and does not re-run steps that
already completed. Store is a nine-method interface with no driver or
ORM behind it; MemoryStore is the default.
engine.Run(ctx, "order-processing", order,
workflow.WithIdempotencyKey(order.ID)) A second call with the same key returns the first execution's outcome instead of running again โ what makes at-least-once event delivery safe.
Observability
Every execution publishes framework events (WorkflowStarted, WorkflowStepFailed, WorkflowCompensationStarted, ...) and one
observability signal carrying each step as a span:
events.On(events.WorkflowFailed{}, func(_ *events.Context, e events.WorkflowFailed) error {
alert.Page("workflow %s failed at step %s: %s", e.Workflow, e.Step, e.Err)
return nil
}) The dashboard renders finished workflows through the same pipeline as HTTP requests and event dispatches, with no workflow-specific code โ but a finished-execution signal can only describe history. A long-running execution โ one that retries with backoff, or waits on a slow third party โ is exactly the one worth watching in flight, so the dashboard also assembles a live view from step events as they arrive, needing no extra configuration:
{
"execution_id": "โฆ", "workflow": "checkout", "trigger": "http",
"steps": [
{"name": "reserve", "state": "done", "duration_ms": 12.4},
{"name": "charge", "state": "retrying", "attempt": 2, "error": "timeout"},
{"name": "ship", "state": "pending"}
]
} Every step of the plan appears from the first frame, so the whole chain is
visible immediately. A step is pending, running, done, failed, retrying, or compensated. Two distinctions matter: a failure that will
be retried is retrying, not failed โ the step has no verdict yet.
And after a successful rollback, previously completed steps become compensated rather than staying done, since the work was undone.
Performance
go test -bench=Run, i5-11400F, MemoryStore, observability disabled:
| Benchmark | ns/op | B/op | allocs/op |
|---|---|---|---|
| 1 step | 5,477 | 1,752 | 24 |
| 10 steps | 41,992 | 10,846 | 89 |
| 50 steps | 164,597 | 46,661 | 337 |
| 10 parallel steps | 35,636 | 11,188 | 86 |
| concurrent (4 cores) | 4,475 | 2,491 | 32 |
Roughly 3โ4 ยตs per step, dominated by persistence calls and event publication rather than orchestration โ a step that does anything real (a query, an HTTP call) costs hundreds of microseconds to milliseconds, so orchestration overhead is not what to optimise. The DAG is topologically sorted once, at registration; definitions are immutable once registered, so execution reads them without a lock; no lock is ever held while user code runs.
Best practices
- Make steps idempotent. A step can run twice โ retries, resumes, and
at-least-once triggers all cause it. Key external calls on
ctx.ExecutionID()plus the step name. - Mark permanent failures. A declined card will be declined five times;
NonRetryableturns five doomed attempts into one. - Compensate anything with an external side effect โ money, mail, or a provisioned resource. A step that only writes to your own database inside a transaction may not need one.
- Respect the context. A step that ignores
ctx.Ctxcannot be cancelled or timed out. - Keep payloads small โ the payload is carried for the whole execution and persisted. Pass an ID, not a 2 MB struct.
Error reference
| Error | Meaning |
|---|---|
ErrWorkflowNotFound | no workflow registered under that name |
ErrDuplicateWorkflow | name already registered |
ErrInvalidWorkflow | definition rejected by validation |
ErrWorkflowCycle | dependencies form a cycle |
ErrDuplicateStep / ErrUnknownDependency | bad step graph |
ErrStepTimeout / ErrWorkflowTimeout | deadline exceeded |
ErrWorkflowCancelled | context cancelled |
ErrStepPanicked | step panicked; recovered and converted |
ErrNonRetryable | failure marked final |
ErrEngineClosed | engine is shut down |
ErrPersistenceFailure | store returned an error |
Failures are wrapped in *StepError (workflow, step, attempt), and
validation problems in *ValidationError, both reachable with errors.As.