defmodule Oban.Web.WorkflowsPage do
@behaviour Oban.Web.Page
use Oban.Web, :live_component
alias Oban.Pro.Workflow
alias Oban.Web.{Page, SearchComponent, SortComponent, Telemetry, Utils, WorkflowQuery}
alias Oban.Web.Workflows.{DetailComponent, TableComponent}
@compile {:no_warn_undefined, Oban.Pro.Workflow}
@known_params ~w(ids limit names queues sort_by sort_dir states workers)
@keep_on_mount ~w(
default_params
detail
sub_workflows
graph_data
params
parent_workflow
workflow
workflows
)a
@inc_limit 20
@max_limit 100
@min_limit 20
@impl Phoenix.LiveComponent
def render(assigns) do
~H"""
<%= cond do %>
<% not @has_workflows? -> %>
<.pro_promo />
<% @detail -> %>
<.live_component
id="detail"
access={@access}
conf={@conf}
module={DetailComponent}
pro_available?={@pro_available?}
workflow={@workflow}
parent_workflow={@parent_workflow}
sub_workflows={@sub_workflows}
graph_data={@graph_data}
/>
<% true -> %>
<.live_component id="workflows-table" module={TableComponent} workflows={@workflows} />
<.load_button label="Show Less" click="load-less" active={@show_less?} myself={@myself} />
<.load_button label="Show More" click="load-more" active={@show_more?} myself={@myself} />
<% end %>
"""
end
attr :active, :boolean, required: true
attr :click, :string, required: true
attr :label, :string, required: true
attr :myself, :any, required: true
defp load_button(assigns) do
~H"""
Workflows
Orchestrate jobs with dependencies for sequential execution, fan-out parallelization, and
fan-in convergence. Build fault-tolerant processing pipelines that scale horizontally
across all nodes.
-
Fully Distributed
— high availability and scalability across your infrastructure
-
Cascading Context
— pass cumulative context between jobs for seamless data flow
-
Nested Sub-Workflows
— compose hierarchically for better organization and reusability
<.link
href="https://oban.pro"
target="_blank"
class="inline-flex items-center px-5 py-2.5 rounded-lg bg-violet-600 hover:bg-violet-700 text-white font-medium transition-colors"
>
Learn about Oban Pro
"""
end
@impl Page
def handle_mount(socket) do
default = %{limit: @min_limit, sort_by: "inserted", sort_dir: "desc"}
assigns = Map.drop(socket.assigns, @keep_on_mount)
%{socket | assigns: assigns}
|> assign(:default_params, default)
|> assign(:has_workflows?, Utils.has_workflows?(socket.assigns.conf))
|> assign(:pro_available?, Utils.has_pro?())
|> assign_new(:detail, fn -> nil end)
|> assign_new(:sub_workflows, fn -> [] end)
|> assign_new(:params, fn -> default end)
|> assign_new(:parent_workflow, fn -> nil end)
|> assign_new(:show_less?, fn -> false end)
|> assign_new(:show_more?, fn -> false end)
|> assign_new(:workflow, fn -> nil end)
|> assign_new(:workflows, fn -> [] end)
end
@impl Page
def handle_refresh(socket) do
%{params: params, conf: conf, detail: detail, has_workflows?: has_workflows?} = socket.assigns
cond do
not has_workflows? ->
socket
detail ->
workflow = WorkflowQuery.get_workflow(conf, detail)
sub_workflows = WorkflowQuery.get_sub_workflows(conf, detail)
parent_workflow = WorkflowQuery.get_sup_workflow(conf, detail)
graph_data = WorkflowQuery.get_workflow_graph(conf, detail)
assign(socket,
workflow: workflow,
sub_workflows: sub_workflows,
parent_workflow: parent_workflow,
graph_data: graph_data
)
true ->
workflows = WorkflowQuery.all_workflows(params, conf)
limit = params.limit
assign(socket,
workflows: workflows,
show_less?: limit > @min_limit,
show_more?: limit < @max_limit and length(workflows) == limit
)
end
end
@impl Page
def handle_params(%{"id" => workflow_id}, _uri, socket) do
conf = socket.assigns.conf
workflow = WorkflowQuery.get_workflow(conf, workflow_id)
sub_workflows = WorkflowQuery.get_sub_workflows(conf, workflow_id)
parent_workflow = WorkflowQuery.get_sup_workflow(conf, workflow_id)
graph_data = WorkflowQuery.get_workflow_graph(conf, workflow_id)
title = if workflow, do: workflow.name || workflow_id, else: "Workflow"
socket =
assign(socket,
detail: workflow_id,
workflow: workflow,
sub_workflows: sub_workflows,
parent_workflow: parent_workflow,
graph_data: graph_data,
page_title: page_title(title)
)
{:noreply, socket}
end
def handle_params(params, _uri, socket) do
params =
params
|> Map.take(@known_params)
|> decode_params()
socket =
socket
|> assign(page_title: page_title("Workflows"))
|> assign(detail: nil, params: Map.merge(socket.assigns.default_params, params))
|> handle_refresh()
{:noreply, socket}
end
@impl Phoenix.LiveComponent
def handle_event("load-less", _params, socket) do
if socket.assigns.show_less? do
send(self(), {:params, :limit, -@inc_limit})
end
{:noreply, socket}
end
def handle_event("load-more", _params, socket) do
if socket.assigns.show_more? do
send(self(), {:params, :limit, @inc_limit})
end
{:noreply, socket}
end
@impl Page
def handle_info({:params, :limit, inc}, socket) when is_integer(inc) do
params =
socket.assigns.params
|> Map.update!(:limit, &(&1 + inc))
|> without_defaults(socket.assigns.default_params)
{:noreply, push_patch(socket, to: oban_path(:workflows, params), replace: true)}
end
def handle_info(:refresh, socket) do
{:noreply, handle_refresh(socket)}
end
def handle_info({:cancel_workflow, workflow_id}, socket) do
enforce_access!(:cancel_workflows, socket.assigns.access)
socket =
if Utils.has_pro?() do
Telemetry.action(:cancel_workflow, socket, [workflow_id: workflow_id], fn ->
Workflow.cancel_jobs(socket.assigns.conf.name, workflow_id)
end)
socket
|> handle_refresh()
|> put_flash_with_clear(:info, "Workflow jobs cancelled")
else
put_flash_with_clear(socket, :error, "Cancel requires Oban Pro")
end
{:noreply, socket}
end
def handle_info({:retry_workflow, workflow_id}, socket) do
enforce_access!(:retry_workflows, socket.assigns.access)
socket =
if Utils.has_pro?() do
Telemetry.action(:retry_workflow, socket, [workflow_id: workflow_id], fn ->
Workflow.retry_jobs(socket.assigns.conf.name, workflow_id)
end)
socket
|> handle_refresh()
|> put_flash_with_clear(:info, "Workflow jobs retried")
else
put_flash_with_clear(socket, :error, "Retry requires Oban Pro")
end
{:noreply, socket}
end
def handle_info(_event, socket) do
{:noreply, socket}
end
end