Belay.Workflow (Belay v1.0.0-rc.6)

Copy Markdown View Source

Jobs composed with directed acyclic dependencies.

alias Belay.Workflow

{:ok, jobs} =
  Workflow.new()
  |> Workflow.add(:fetch, MyApp.Fetch.new(%{"id" => 1}))
  |> Workflow.add(:parse, MyApp.Parse.new(%{}), deps: [:fetch])
  |> Workflow.add(:store, MyApp.Store.new(%{}), deps: [:parse])
  |> Workflow.insert(MyBelay)

Dependent jobs are inserted held and released transactionally as their upstreams succeed. When an upstream fails or is cancelled, dependents cancel by default; ignore: [:failed] / ignore: [:cancelled] (per job or workflow-wide) treats that outcome as satisfied instead.

Summary

Functions

Add a named job. deps must reference previously added names.

Insert all workflow jobs atomically. Returns {:ok, %{name => job}}.

All jobs in a workflow.

Create a workflow. Options: :workflow_id, :ignore ([:failed | :cancelled]).

Status summary: %{total:, state_counts:, done?:}.

Types

t()

@type t() :: %Belay.Workflow{
  entries: term(),
  id: term(),
  ignore: term(),
  names: term()
}

Functions

add(workflow, name, arg, wf_opts \\ [])

Add a named job. deps must reference previously added names.

insert(workflow, name)

Insert all workflow jobs atomically. Returns {:ok, %{name => job}}.

jobs(name, workflow_id)

All jobs in a workflow.

new(opts \\ [])

Create a workflow. Options: :workflow_id, :ignore ([:failed | :cancelled]).

status(name, workflow_id)

Status summary: %{total:, state_counts:, done?:}.