Lightning.WorkOrders (Lightning v2.19.0)

Copy Markdown View Source

Context for creating WorkOrders.

Work Orders

Work Orders represent the entrypoint for a unit of work in Lightning. They allow you to track the status of a webhook or cron trigger.

For example if a user makes a request to a webhook endpoint, a Work Order is created with it's associated Workflow and Dataclip.

Every Work Order has at least one Run, which represents a single invocation of the Workflow. If the workflow fails, and the run is retried, a new Run is created on the Work Order.

This allows you group all the runs for a single webhook, and track the success or failure of a given dataclip.

Creating Work Orders

Work Orders can be created in three ways:

  1. Via a webhook trigger
  2. Via a cron trigger
  3. Manually by a user (via the UI or API)

Retries do not create new Work Orders, but rather new Runs on the existing Work Order.

Summary

Functions

Cancels available runs for the given work orders.

Enqueues a single background job to cancel all available runs matching the given work order IDs. Used for "cancel all matching" where the count may be large.

Enqueue multiple runs for retry in the same transaction.

Get a Work Order by id.

Recent work orders whose runs executed against the snapshot published as version_number, newest first.

Recent work orders whose runs executed against a snapshot that was never released (draft, test, or intermediate saves), newest first.

Retry a run from a given step.

Subscribes to work order events across all projects.

Subscribes to the work order events of a single project.

Updates the state of a WorkOrder based on the state of a run.

Returns a query for work orders belonging to a specific project

Types

dataclip_input()

work_order_option()

@type work_order_option() ::
  {:workflow, Lightning.Workflows.Workflow.t()}
  | {:dataclip, dataclip_input()}
  | {:created_by, Lightning.Accounts.User.t()}
  | {:project_id, Ecto.UUID.t()}
  | {:without_run, boolean()}

Functions

build_for(trigger, attrs)

cancel_many(work_orders, opts)

@spec cancel_many([Lightning.WorkOrder.t()], keyword()) ::
  {:ok, non_neg_integer()} | {:error, any()}

Cancels available runs for the given work orders.

For small batches, cancels synchronously via atomic UPDATE. Returns {:ok, count} where count is the number of runs cancelled.

cancel_many_async(work_orders, opts)

@spec cancel_many_async([Lightning.WorkOrder.t()], keyword()) ::
  {:ok, Oban.Job.t()} | {:error, Ecto.Changeset.t()}

Enqueues a single background job to cancel all available runs matching the given work order IDs. Used for "cancel all matching" where the count may be large.

create_for(manual)

create_for(trigger, multi \\ Multi.new(), opts)

@spec create_for(Lightning.Workflows.Trigger.t(), Ecto.Multi.t(), [
  work_order_option()
]) ::
  {:ok, Lightning.WorkOrder.t()}
  | {:error, Ecto.Changeset.t(Lightning.WorkOrder.t()) | :workflow_deleted}

Create a new Work Order.

For a webhook

create_for(trigger, workflow: workflow, dataclip: dataclip)

enqueue_many_for_retry(workorders_ids, creating_user_id, queue \\ "default")

Enqueue multiple runs for retry in the same transaction.

get(id, opts \\ [])

@spec get(Ecto.UUID.t(), [{:include, list()}]) :: Lightning.WorkOrder.t() | nil

Get a Work Order by id.

Optionally preload associations by passing a list to :include. Supports nested preloads.

Lightning.WorkOrders.get(id, include: [:runs])
Lightning.WorkOrders.get(id, include: [workflow: :project, runs: []])

get_last_runs_steps_with_dataclips(workorders, jobs)

get_workorders_for_version(workflow_id, version_number)

@spec get_workorders_for_version(Ecto.UUID.t(), integer()) :: [
  Lightning.WorkOrder.t()
]

Recent work orders whose runs executed against the snapshot published as version_number, newest first.

Bounded like get_workorders_with_runs/2, not the full history page. Each work order carries only the runs matching this release, so one retried across several versions shows only the runs belonging to this one. [] when no such release exists.

get_workorders_unversioned(workflow_id)

@spec get_workorders_unversioned(Ecto.UUID.t()) :: [Lightning.WorkOrder.t()]

Recent work orders whose runs executed against a snapshot that was never released (draft, test, or intermediate saves), newest first.

The draft/unversioned counterpart to get_workorders_for_version/2: same bounded cap of 20 work orders, and each work order carries only its unversioned runs.

get_workorders_with_runs(workflow_id, run_id)

limit_run_creation(project_id, runs_count \\ 1)

@spec limit_run_creation(Ecto.UUID.t(), non_neg_integer()) ::
  :ok | Lightning.Extensions.UsageLimiting.error()

retry(run_id, step_id, opts)

@spec retry(
  Lightning.Run.t() | Ecto.UUID.t(),
  Lightning.Invocation.Step.t() | Ecto.UUID.t(),
  [
    work_order_option(),
    ...
  ]
) :: {:ok, Lightning.Run.t()} | {:error, Ecto.Changeset.t() | :workflow_deleted}

Retry a run from a given step.

This will create a new Run on the Work Order, and enqueue it for processing.

When creating a new Run, a graph of the workflow is created steps that are independent from the selected step and its downstream flow are associated with this new run, but not executed again.

retry_many(workorders, opts)

@spec retry_many([Lightning.WorkOrder.t(), ...] | [Lightning.RunStep.t(), ...], [
  work_order_option(),
  ...
]) ::
  {:ok, enqueued_count :: non_neg_integer(),
   discarded_count :: non_neg_integer()}
  | Lightning.Extensions.UsageLimiting.error()
  | {:error, :enqueue_error}

retry_many(workorders, job_id, opts)

@spec retry_many([Lightning.WorkOrder.t(), ...], job_id :: Ecto.UUID.t(), [
  work_order_option(),
  ...
]) ::
  {:ok, enqueued_count :: non_neg_integer(),
   discarded_count :: non_neg_integer()}
  | Lightning.Extensions.UsageLimiting.error()
  | {:error, :enqueue_error}

subscribe()

Subscribes to work order events across all projects.

subscribe(project_id)

Subscribes to the work order events of a single project.

update_state(run)

@spec update_state(Lightning.Run.t()) :: {:ok, Lightning.WorkOrder.t()}

Updates the state of a WorkOrder based on the state of a run.

This considers the state of all runs on the WorkOrder, with the Run passed in as the latest run.

See Lightning.WorkOrders.Query.state_for/1 for more details.

work_orders_for_project_query(project)

@spec work_orders_for_project_query(Lightning.Projects.Project.t()) ::
  Ecto.Queryable.t()

Returns a query for work orders belonging to a specific project