Current section
Files
Jump to
Current section
Files
lib/gust_web/live/dag_live/index.ex
defmodule GustWeb.DagLive.Index do
alias Gust.Flows
use GustWeb, :live_view
@impl true
def mount(_params, _session, socket) do
Gust.PubSub.subscribe_all_files("update")
dag_defs = Gust.DAG.Loader.get_definitions()
dags =
Enum.map(dag_defs, fn {dag_id, dag_def} ->
%{id: dag_id, dag: Flows.get_dag!(dag_id), dag_def: dag_def}
end)
socket = socket |> assign(:page_title, "DAGs Listing") |> stream(:dags, dags)
{:ok, socket}
end
@impl true
def handle_info({:trigger_dag_run, dag_id, dag_def}, socket) do
{:ok, run} = Gust.DAG.Starter.start_dag_run(dag_id, dag_def)
run = Flows.get_run_with_tasks!(run.id)
{:noreply, socket |> put_flash(:info, "Run #{run.id} triggered")}
end
@impl true
def handle_info(
{:dag, :file_updated, %{action: "removed", dag_name: name, dag_def: nil}},
socket
) do
socket =
case Flows.get_dag_by_name(name) do
%Flows.Dag{} = dag ->
Flows.delete_dag!(dag)
socket |> stream_delete(:dags, dag)
end
{:noreply, socket}
end
@impl true
def handle_info(
{:dag, :file_updated, %{action: "reload", dag_def: nil}},
socket
) do
{:noreply, socket}
end
@impl true
def handle_info({:dag, :file_updated, %{action: "reload", dag_def: dag_def}}, socket) do
name = dag_def.name
socket =
with nil <- Flows.get_dag_by_name(name),
{:ok, dag} <- Flows.create_dag(%{name: name}) do
insert_dag(socket, dag, dag_def)
else
%Flows.Dag{} = dag ->
insert_dag(socket, dag, dag_def)
{:error, changeset} ->
put_flash(socket, :error, "#{name}: #{describe_error(changeset)}")
end
{:noreply, socket}
end
defp insert_dag(socket, dag, dag_def) do
stream_insert(socket, :dags, %{id: dag.id, dag: dag, dag_def: dag_def})
end
defp humanize(atom), do: atom |> to_string() |> String.replace("_", " ") |> String.capitalize()
defp describe_error(changeset) do
Ecto.Changeset.traverse_errors(changeset, fn {msg, opts} ->
Enum.reduce(opts, msg, fn {key, value}, acc ->
String.replace(acc, "%{#{key}}", to_string(value))
end)
end)
|> Enum.map_join("; ", fn {field, messages} ->
"#{humanize(field)} #{Enum.join(messages, ", ")}"
end)
end
end