Defines a durable workflow.
defmodule MyApp.OrderFlow do
use Continuum.Workflow, version: 1, retention: {:days, 30}
def run(%{order_id: id, items: items}) do
{:ok, validated} = activity Validation.check(items)
{:ok, _charge} = activity Payments.charge(id, validated.total),
retry: [max_attempts: 5, backoff: :exponential]
case await signal(:fraud_review, timeout: hours(24)) do
:approved -> activity Fulfillment.ship(id)
:rejected -> {:error, :rejected}
:timeout -> activity Fulfillment.ship(id)
end
end
endDeterminism
Every public and private function in the module is scanned at compile
time by Continuum.AstCheck. Calls that are known to be non-deterministic
(DateTime.utc_now/0, :rand.uniform/0, IO.puts/1, ETS access, …) are
compile errors with a remediation hint.
See Continuum.AstCheck.forbidden_calls/0 for the denylist and the
:trusted_modules config knob for extending the allowlist.
Versioning
The module's AST is hashed at compile time and exposed through
__continuum_workflow__/0. Each Postgres run stores that hash on start so
drift can be surfaced. As of v0.3, callers may pass workflow: LogicalModule
to use Continuum.Workflow to register a concrete module as a hash-specific
entrypoint for a logical workflow.
Summary
Functions
Macro: schedule an activity. The result is journaled on first execution and replayed on resume.
Macro: schedule several activities at once and wait for all of them.
Map a runtime list through a unary activity with bounded parallel scheduling.
Macro: wait for an external signal, optionally with a timeout.
Macro: suspend until a previously start_child-ed child terminates.
Macro: run the compensation of one successful compensated activity.
Macro: run all pending compensations in LIFO order (most-recent first).
Macro: tail-call continuation — complete this run and start a fresh one on the same workflow with new input.
Returns a duration in milliseconds.
Macro: start a child workflow asynchronously, returning a %Continuum.ChildRef{}.
Macro: durable timer.
Functions
Macro: schedule an activity. The result is journaled on first execution and replayed on resume.
activity Payments.charge(order_id, amount)
activity Payments.charge(order_id, amount), retry: [max_attempts: 5]
Macro: schedule several activities at once and wait for all of them.
[quote_a, quote_b] = activity_all([Pricing.quote(a), Pricing.quote(b)])Results come back in the order the activities were declared, regardless of
the order they finish in. Each element is exactly what the corresponding
activity/2 call would have returned, so a failed activity contributes an
{:error, reason} entry and the rest of the batch is unaffected. There is
deliberately no mode: option: a partial failure is an ordinary value you
pattern-match, not a control-flow mode.
case activity_all([Pricing.quote(a), Inventory.reserve(b)]) do
[{:ok, price}, {:ok, hold}] -> ship(price, hold)
results -> {:error, Enum.filter(results, &match?({:error, _}, &1))}
endThe argument must be a literal list of activity calls written at the call
site: each element's {module, function, arity} is part of the batch's
command identity, computed at macro expansion, so it cannot come from a
variable.
Batch members take no per-activity options in v0.8 — including compensate:.
Use sequential activity/2 calls when a step needs a compensation.
Map a runtime list through a unary activity with bounded parallel scheduling.
activity_map(orders, &Shipping.ship/1, concurrency: 8, key: :id)The activity must be a static &Module.function/1 capture. :concurrency
is required (1–1,000). Members are scheduled in windows of that size; the
next window starts only after all members of the current window terminate.
Results preserve input order, and failures are {:error, reason} entries.
Automatic retries keep their window slot until the member terminates.
Inputs must be a durable list. An empty list returns [] and still records
membership. By default positions identify members, so duplicate values are
distinct work. key: :field additionally requires unique durable values in
that field of each input map. Input, order, keys, and concurrency changes
during replay are drift errors, including changes in unscheduled windows.
Node and queue limits still apply and can lower actual concurrency. Run
cancellation discards outstanding tasks and prevents later windows. Manual
member retry and per-member compensation are not supported, as with
activity_all/1. For very large lists, use workflow continuations between
maps to bound accumulated history and result memory.
Macro: wait for an external signal, optionally with a timeout.
await signal(:approved)
await signal(:approved, timeout: hours(24))Or wait for a child workflow. The shorthand accepts exactly
child Mod.run(input); use start_child/3 for other setup shapes.
await child MyApp.AuditFlow.run(%{batch_id: id})
Macro: suspend until a previously start_child-ed child terminates.
Returns the child's result ({:ok, _}/{:error, _} term), the error on child
failure, or {:error, :child_cancelled} if the child was cancelled.
Macro: run the compensation of one successful compensated activity.
{:ok, charge} = activity Payments.charge(id, total), compensate: {Payments, :refund, [id]}
# ...
compensate(charge)Takes the %Continuum.ActivityRef{} (or {:ok, ref}) returned by a compensated
activity/2 call, schedules its compensation MFA through the activity worker,
and removes it from the pending compensation set so a later compensate_all/0
cannot run it twice. Returns {:ok, result} or, if the compensation fails
terminally, {:error, reason} — the run continues either way.
Macro: run all pending compensations in LIFO order (most-recent first).
rescue
e ->
compensate_all()
reraise e, __STACKTRACE__Each successful compensated activity that has not already been compensated by
compensate/1 is rolled back, newest first. Returns :ok.
Macro: tail-call continuation — complete this run and start a fresh one on the same workflow with new input.
def run(%{cycles_done: n} = state) do
activity Billing.charge(state.customer_id)
timer(days(30))
if n >= 11 do
{:ok, :year_complete}
else
continue_as_new(%{state | cycles_done: n + 1})
end
endThe current run is marked completed with result: {:continued, next_run_id};
a new run starts with the given input, sharing the chain's correlation_id,
namespace, and attributes. Use it to keep history bounded for
long-running / cron-style workflows.
Children started but not yet awaited when the run continues are re-parented
to the successor so cancelling the chain still cascades into them. The
successor cannot await them, however — their child_started events live in
the predecessor's history. Await every child you need a result from before
calling continue_as_new/1.
Returns a duration in milliseconds.
Macro: start a child workflow asynchronously, returning a %Continuum.ChildRef{}.
ref = start_child MyApp.OrderFlow, %{order_id: id}, id: "order-#{id}"
# ... do other work ...
result = await_child(ref)opts accepts id: to tie the child's deterministic run id to a key under
this parent.
Macro: durable timer.
timer(hours(24))