Current section

Files

Jump to
gust lib gust dag loader worker.ex
Raw

lib/gust/dag/loader/worker.ex

defmodule Gust.DAG.Loader.Worker do
@moduledoc false
alias Gust.DAG.Parser
alias Gust.DAG.Loader
use GenServer
require Logger
@impl true
def init(args) do
{:ok, args, {:continue, :load}}
end
def child_spec(arg) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [arg]},
restart: :transient,
type: :worker
}
end
def start_link(args) do
GenServer.start_link(__MODULE__, args, name: __MODULE__)
end
def get_definitions do
GenServer.call(__MODULE__, :get_definitions)
end
def reload_definition(dag_id) do
GenServer.call(__MODULE__, {:reload_definition, dag_id})
end
@impl true
def handle_call(:get_definitions, _from, state) do
{:reply, state[:dag_defs], state}
end
@impl true
def handle_call({:reload_definition, dag_id}, _from, state) do
dag_def = state[:dag_defs][dag_id]
{:ok, dag_def} = Parser.parse(dag_def.file_path)
dag_defs = Map.put(state[:dag_defs], dag_id, dag_def)
{:reply, {:ok, dag_def}, Map.put(state, :dag_defs, dag_defs)}
end
@impl true
def handle_continue(:load, %{dags_folder: folder} = state) do
{found, removed} = Loader.load(folder)
log_dag_names("+ Created DAGs", found)
log_dag_names("- Deleted DAGs", removed)
dag_defs = map_dag_def(found, folder)
Gust.DAG.Scheduler.schedule(dag_defs)
Gust.DAG.RunRestarter.restart(dag_defs)
state = Map.put(state, :dag_defs, dag_defs)
{:noreply, state}
end
defp map_dag_def(dags, folder) do
dags
|> Enum.reduce(%{}, fn dag, acc ->
file_path = "#{Path.expand(folder)}/#{dag.name}.ex"
{:ok, dag_def} = Parser.parse(file_path)
Map.put(acc, dag.id, dag_def)
end)
end
defp log_dag_names(label, dags) when is_list(dags) do
names = dags |> Enum.map_join("; ", & &1.name)
Logger.warning("#{label}: #{names}")
end
end