defmodule Oban.Web.Workflows.DetailComponent do
use Oban.Web, :live_component
alias Oban.Web.Components.Core
alias Oban.Web.Timing
alias Oban.Web.WorkflowQuery
@states ~w(suspended available scheduled executing retryable completed cancelled discarded)a
@impl Phoenix.LiveComponent
def update(assigns, socket) do
sub_workflows = assigns[:sub_workflows] || []
socket =
socket
|> assign(assigns)
|> assign(:sub_workflows, sub_workflows)
|> assign_new(:graph_open?, fn -> true end)
|> assign_new(:subs_open?, fn -> match?([_ | _], sub_workflows) end)
|> assign_new(:graph_data, fn -> %{jobs: [], sub_workflows: []} end)
|> push_graph_data()
{:ok, socket}
end
defp push_graph_data(socket) do
if socket.assigns[:graph_open?] do
push_event(socket, "graph-data", socket.assigns.graph_data)
else
socket
end
end
@impl Phoenix.LiveComponent
def render(assigns) do
~H"""
<%= if @workflow do %>
<.header
access={@access}
myself={@myself}
pro_available?={@pro_available?}
workflow={@workflow}
parent_workflow={@parent_workflow}
/>
<.progress_bar workflow={@workflow} subs={@sub_workflows} />
<.stats_grid workflow={@workflow} sub_workflows={@sub_workflows} />
<.graph_section myself={@myself} graph_open?={@graph_open?} graph_data={@graph_data} />
<.sub_workflows_section
myself={@myself}
subs_open?={@subs_open?}
sub_workflows={@sub_workflows}
/>
<% else %>
Workflow not found
This workflow may have been deleted or doesn't exist.
<% end %>
"""
end
# Header
attr :access, :any, required: true
attr :myself, :any, required: true
attr :pro_available?, :boolean, required: true
attr :workflow, :map, required: true
attr :parent_workflow, :map, default: nil
defp header(assigns) do
~H"""
<.parent_breadcrumb :if={@parent_workflow} parent={@parent_workflow} />
"""
end
attr :parent, :map, required: true
defp parent_breadcrumb(assigns) do
~H"""
sub-workflow of
<.link
navigate={oban_path([:workflows, @parent.id])}
class="ml-1 font-medium text-violet-600 hover:text-violet-500 dark:text-violet-400"
>
{@parent.name || @parent.id}
"""
end
# Progress Bar
attr :workflow, :map, required: true
attr :subs, :list, default: []
defp progress_bar(assigns) do
wf = assigns.workflow
sub = count_sub_states(assigns.subs)
total = Enum.reduce(@states, 0, &(Map.fetch!(wf, &1) + &2)) + length(assigns.subs)
completed = wf.completed + sub.completed
states = [
{:suspended, wf.suspended + sub.suspended, "bg-gray-400", "Suspended"},
{:scheduled, wf.scheduled + sub.scheduled, "bg-indigo-400", "Scheduled"},
{:available, wf.available + sub.available, "bg-blue-400", "Available"},
{:retryable, wf.retryable + sub.retryable, "bg-yellow-400", "Retryable"},
{:executing, wf.executing + sub.executing, "bg-emerald-400", "Executing"},
{:completed, completed, "bg-cyan-400", "Completed"},
{:cancelled, wf.cancelled + sub.cancelled, "bg-violet-400", "Cancelled"},
{:discarded, wf.discarded + sub.discarded, "bg-rose-400", "Discarded"}
]
percent = if total > 0, do: round(completed / total * 100), else: 0
assigns =
assign(assigns, states: states, total: total, completed: completed, percent: percent)
~H"""
{@percent}% Complete
{@completed}/{@total} jobs
<%= for {_state, count, color, _label} <- @states, count > 0 do %>
<% end %>
<%= for {_state, count, color, label} <- @states do %>
{label}
{count}
<% end %>
"""
end
# Graph Section
attr :myself, :any, required: true
attr :graph_open?, :boolean, required: true
attr :graph_data, :map, required: true
defp graph_section(assigns) do
~H"""
"""
end
# Stats Grid
attr :workflow, :any, required: true
attr :sub_workflows, :list, required: true
defp stats_grid(assigns) do
queues = Map.get(assigns.workflow.meta || %{}, "queues", [])
assigns = assign(assigns, queues: queues)
~H"""
Workflow ID
{@workflow.id}
Started
<.format_started_at workflow={@workflow} />
Duration
<.format_duration workflow={@workflow} />
Status
{@workflow.state}
Subs
{length(@sub_workflows)}
"""
end
attr :workflow, :any, required: true
defp format_started_at(assigns) do
wf = assigns.workflow
executed? = wf.executing + wf.completed > 0
started = if executed?, do: wf.started_at
assigns = assign(assigns, started: started)
~H"""
-
—
"""
end
attr :workflow, :any, required: true
defp format_duration(assigns) do
wf = assigns.workflow
executing? = wf.state == "executing"
started? = not is_nil(wf.started_at)
duration =
if wf.started_at && wf.completed_at do
DateTime.diff(wf.completed_at, wf.started_at, :millisecond)
end
formatted =
if is_nil(duration) or duration <= 0 do
"—"
else
duration |> div(1000) |> Timing.to_duration()
end
assigns =
assign(assigns,
executing?: executing?,
started?: started?,
formatted: formatted,
started_at: wf.started_at
)
~H"""
-
{@formatted}
"""
end
# Sub-workflows Section
attr :myself, :any, required: true
attr :subs_open?, :boolean, required: true
attr :sub_workflows, :list, required: true
defp sub_workflows_section(assigns) do
subs_count = length(assigns.sub_workflows)
assigns = assign(assigns, subs_count: subs_count)
~H"""
0} class="mt-3">
|
Name
|
ID
|
Progress
|
Started
|
Duration
|
Status
|
<.sub_workflow_row :for={sub <- @sub_workflows} workflow={sub} />
"""
end
attr :workflow, :any, required: true
defp sub_workflow_row(assigns) do
wf = assigns.workflow
total = Enum.reduce(@states, 0, &(Map.fetch!(wf, &1) + &2))
percent = if total > 0, do: round(wf.completed / total * 100), else: 0
assigns =
assign(assigns,
completed: wf.completed,
total: total,
percent: percent,
state: wf.state
)
~H"""
<.link navigate={oban_path([:workflows, @workflow.id])} class="contents">
|
{@workflow.name || @workflow.id}
|
{@workflow.id}
|
|
<.format_started_at workflow={@workflow} />
|
<.format_duration workflow={@workflow} />
|
<.status_icon state={@state} />
|
"""
end
attr :state, :string, default: nil
defp status_icon(assigns) do
~H"""
<%= case @state do %>
<% "executing" -> %>
<% "completed" -> %>
<% "retryable" -> %>
<% "cancelled" -> %>
<% "discarded" -> %>
<% _ -> %>
<% end %>
"""
end
# Event Handlers
@impl Phoenix.LiveComponent
def handle_event("cancel-workflow", _params, socket) do
send(self(), {:cancel_workflow, socket.assigns.workflow.id})
{:noreply, socket}
end
def handle_event("retry-workflow", _params, socket) do
send(self(), {:retry_workflow, socket.assigns.workflow.id})
{:noreply, socket}
end
def handle_event("toggle-graph", _params, socket) do
graph_open? = not socket.assigns[:graph_open?]
socket = assign(socket, :graph_open?, graph_open?)
socket = if graph_open?, do: push_graph_data(socket), else: socket
{:noreply, socket}
end
def handle_event("toggle-subs", _params, socket) do
{:noreply, assign(socket, :subs_open?, not socket.assigns[:subs_open?])}
end
def handle_event("navigate-to-job", %{"job_id" => job_id}, socket) do
{:noreply, push_navigate(socket, to: oban_path([:jobs, job_id]))}
end
def handle_event("navigate-to-workflow", %{"workflow_id" => workflow_id}, socket) do
{:noreply, push_navigate(socket, to: oban_path([:workflows, workflow_id]))}
end
def handle_event("expand-sub-workflow", %{"workflow_id" => sub_workflow_id}, socket) do
%{jobs: jobs, truncated: truncated} =
WorkflowQuery.get_sub_workflow_jobs(socket.assigns.conf, sub_workflow_id)
payload = %{workflow_id: sub_workflow_id, jobs: jobs, truncated: truncated}
socket = push_event(socket, "sub-workflow-jobs", payload)
{:noreply, socket}
end
# Helpers
defp cancel_tooltip(true), do: "Cancel all jobs in this workflow"
defp cancel_tooltip(false), do: "Cancel requires Oban Pro"
defp retry_tooltip(true), do: "Retry failed jobs in this workflow"
defp retry_tooltip(false), do: "Retry requires Oban Pro"
defp has_retryable?(workflow) do
workflow.retryable + workflow.discarded + workflow.cancelled > 0
end
defp count_sub_states(subs) do
init = Map.new(@states, &{&1, 0})
Enum.reduce(subs, init, fn sub, acc ->
Map.update!(acc, String.to_existing_atom(sub.state), &(&1 + 1))
end)
end
end