Current section

Files

Jump to
gust lib gust flows.ex
Raw

lib/gust/flows.ex

defmodule Gust.Flows do
@moduledoc """
The Flows context.
"""
alias Gust.Flows.{Dag, Run, Task, Log, Secret}
import Ecto.Query, warn: false
alias Gust.Repo
def change_secret(%Secret{} = secret, attrs \\ %{}) do
Secret.changeset(secret, attrs)
end
def create_secret(attrs \\ %{}) do
%Secret{}
|> Secret.changeset(attrs)
|> Repo.insert()
end
def create_dag(attrs \\ %{}) do
%Dag{}
|> Dag.changeset(attrs)
|> Repo.insert()
end
def create_run(attrs \\ %{}) do
%Run{}
|> Run.changeset(attrs)
|> Repo.insert()
end
def create_test_run(attrs \\ %{}) do
%Run{}
|> Run.test_changeset(attrs)
|> Repo.insert()
end
def get_run!(id), do: Repo.get!(Run, id)
def get_run_with_tasks!(id) do
get_run!(id) |> Repo.preload(:tasks)
end
def get_running_runs_by_dag(dag_ids, status) do
Repo.all(
from r in Run,
where: r.dag_id in ^dag_ids and r.status == ^status
)
end
def get_log!(id), do: Repo.get!(Log, id)
def create_log(attrs \\ %{}) do
%Log{}
|> Log.changeset(attrs)
|> Repo.insert()
end
def get_task!(id), do: Repo.get!(Task, id)
def get_secret!(id), do: Repo.get!(Secret, id)
def get_secret_by_name(name), do: Repo.get_by(Secret, name: name)
def create_task(attrs \\ %{}) do
%Task{}
|> Task.changeset(attrs)
|> Repo.insert()
end
def update_secret(%Secret{} = secret, attrs) do
secret
|> Secret.changeset(attrs)
|> Repo.update()
end
def update_task_result(task, result) do
Task.changeset(task, %{result: result})
|> Repo.update()
end
def update_task_status(task, status) do
Task.changeset(task, %{status: status})
|> Repo.update()
end
def get_task_with_logs!(id) do
Task |> Repo.get!(id) |> Repo.preload(:logs)
end
def get_task_by_name_run_with_logs(name, run_id) do
get_task_by_name_run(name, run_id) |> Repo.preload(:logs)
end
def get_task_by_name_run(name, run_id) do
Task |> where(run_id: ^run_id, name: ^name) |> Repo.one()
end
def update_run_status(run, status) do
Run.changeset(run, %{status: status})
|> Repo.update()
end
def toggle_enabled(dag) do
Dag.changeset(dag, %{enabled: !dag.enabled})
|> Repo.update()
end
def get_dag!(id), do: Repo.get!(Dag, id)
def get_dag_with_runs!(id) do
Dag |> Repo.get!(id) |> Repo.preload(:runs)
end
def get_dag_with_runs_and_tasks!(name) do
runs_q =
from r in Run,
order_by: [asc: r.inserted_at],
preload: [:tasks]
Repo.one!(
from d in Dag,
where: d.name == ^name,
preload: [runs: ^runs_q]
)
end
def list_secrets do
Repo.all(Secret)
end
def list_dags do
Repo.all(Dag)
end
def get_dag_by_name(name) do
Repo.get_by(Dag, name: name)
end
def delete_dag!(dag) do
Repo.delete!(dag)
end
def delete_secret(%Secret{} = secret) do
Repo.delete(secret)
end
end