Continuum.Workflow (continuum v0.8.5)

Copy Markdown View Source

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
end

Determinism

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

activity(call, opts \\ [])

(macro)

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]

activity_all(calls)

(since 0.8.0) (macro)

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))}
end

The 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.

activity_map(items, activity, opts)

(macro)

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.

await(arg)

(macro)

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})

await_child(ref)

(since 0.3.0) (macro)

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.

compensate(ref)

(since 0.3.0) (macro)

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.

compensate_all()

(since 0.3.0) (macro)

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.

compensate_all(opts)

(since 0.4.0) (macro)

continue_as_new(input)

(since 0.3.0) (macro)

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
end

The 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.

days(n)

(macro)

hours(n)

(macro)

minutes(n)

(macro)

seconds(n)

(macro)

Returns a duration in milliseconds.

start_child(workflow, input, opts \\ [])

(since 0.3.0) (macro)

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.

timer(duration)

(macro)

Macro: durable timer.

timer(hours(24))