Packages
step_flow
1.8.1
1.9.0-rc2
1.9.0-rc1
1.9.0-rc0
1.8.2
1.8.1
1.8.1-rc8
1.8.1-rc7
1.8.1-rc6
1.8.1-rc5
1.8.1-rc4
1.8.1-rc3
1.8.1-rc2
1.8.1-rc1
1.8.1-rc0
1.8.0
1.8.0-rc3
1.8.0-rc2
1.8.0-rc1
1.8.0-rc0
1.7.3
1.7.3-rc4
1.7.3-rc3
1.7.3-rc2
1.7.3-rc1
1.7.3-rc0
1.7.2
1.7.2-rc4
1.7.2-rc3
1.7.2-rc2
1.7.2-rc1
1.7.2-rc0
1.7.1
1.7.0
1.7.0-rc1
1.7.0-rc0
1.6.1
1.6.1-rc1
1.6.1-rc0
1.6.0
1.6.0-rc9
1.6.0-rc8
1.6.0-rc7
1.6.0-rc6
1.6.0-rc5
1.6.0-rc4
1.6.0-rc3
1.6.0-rc20
1.6.0-rc2
1.6.0-rc19
1.6.0-rc18
1.6.0-rc17
1.6.0-rc16
1.6.0-rc15
1.6.0-rc14
1.6.0-rc13
1.6.0-rc12
1.6.0-rc11
1.6.0-rc10
1.6.0-rc1
1.5.0
1.5.0-rc1
1.4.2-rc2
1.4.2-rc1
1.4.1
1.4.1-rc1
1.4.0
1.4.0-rc4
1.4.0-rc3
1.4.0-rc2
1.4.0-rc1
1.3.1
1.3.0
1.3.0-rc
1.2.0
1.1.0
1.0.0
1.0.0-rc9
1.0.0-rc8
1.0.0-rc7
1.0.0-rc6
1.0.0-rc5
1.0.0-rc1
0.2.13
0.2.12
0.2.11
0.2.10
0.2.9
0.2.8
0.2.7
0.2.6
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.8
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
0.0.24
0.0.23
0.0.22
0.0.21
0.0.20
0.0.19
0.0.18
0.0.17
0.0.16
0.0.15
0.0.14
0.0.13
0.0.12
0.0.11
0.0.10
0.0.9
0.0.8
0.0.7
0.0.6
0.0.4
0.0.3
0.0.2
0.0.1
Step flow manager for Elixir applications
Current section
Files
Jump to
Current section
Files
lib/step_flow/models/workflows/workflows.ex
defmodule StepFlow.Workflows do
@moduledoc """
The Workflows context.
"""
import Ecto.Query, warn: false
alias StepFlow.Artifacts.Artifact
alias StepFlow.Jobs
alias StepFlow.Jobs.Status
alias StepFlow.QueryFilter
alias StepFlow.Repo
alias StepFlow.Roles
alias StepFlow.Workflows
alias StepFlow.Workflows.Workflow
require Logger
@doc """
Returns the list of workflows.
## Examples
iex> list_workflows()
[%Workflow{}, ...]
"""
def list_workflows(params \\ %{}) do
page =
Map.get(params, "page", 0)
|> StepFlow.Integer.force()
size =
Map.get(params, "size", 10)
|> StepFlow.Integer.force()
offset = page * size
query =
from(workflow in Workflow)
|> apply_default_query_filters(params)
total_query = from(item in subquery(query), select: count(item.id))
total =
Repo.all(total_query)
|> List.first()
query =
from(
workflow in subquery(query),
order_by: [desc: :inserted_at],
offset: ^offset,
limit: ^size
)
workflows =
Repo.all(query)
|> Repo.preload([:artifacts, :status, jobs: :child_workflow])
|> preload_workflows
%{
data: workflows,
total: total,
page: page,
size: size
}
end
def apply_default_query_filters(query, params \\ %{}) do
query =
from(workflow in subquery(query))
|> search_for_reference(params, :search)
|> QueryFilter.filter_query(params, :reference)
|> QueryFilter.filter_query(params, :identifier)
|> filter_version(params)
|> QueryFilter.filter_query(params, :is_live)
|> filter_deleted(params)
|> QueryFilter.filter_query(params, :user_uuid)
|> filter_mode(params)
|> QueryFilter.apply_end_date_filter(params, :end_date)
|> QueryFilter.apply_end_date_filter(params, :before_date)
|> QueryFilter.apply_start_date_filter(params, :after_date)
|> QueryFilter.apply_start_date_filter(params, :start_date)
|> filter_status(params, :states, extract_start_date(params))
allowed_workflows = check_rights(params)
query =
case {StepFlow.Map.get_by_key_or_atom(params, :workflow_ids),
Enum.member?(allowed_workflows, "*")} do
{nil, true} ->
query
{nil, false} ->
from(
workflow in query,
where: workflow.identifier in ^allowed_workflows
)
{workflow_ids, true} ->
from(
workflow in query,
where: workflow.identifier in ^workflow_ids
)
{workflow_ids, false} ->
intersect =
workflow_ids
|> Enum.filter(fn element -> element in allowed_workflows end)
from(
workflow in query,
where: workflow.identifier in ^intersect
)
end
query =
case StepFlow.Map.get_by_key_or_atom(params, :ids) do
nil ->
query
identifiers ->
from(workflow in query, where: workflow.id in ^identifiers)
end
query
end
defp extract_start_date(params) do
cond do
start_date = StepFlow.Map.get_by_key_or_atom(params, :start_date) -> start_date
after_date = StepFlow.Map.get_by_key_or_atom(params, :after_date) -> after_date
true -> nil
end
end
defp search_for_reference(query, params, field) do
case StepFlow.Map.get_by_key_or_atom(params, field) do
nil ->
query
search ->
like = "%#{search}%"
from(
workflow in subquery(query),
where: ilike(workflow.reference, ^like)
)
end
end
def check_rights(params) do
roles =
StepFlow.Map.get_by_key_or_atom(params, :roles, [])
|> Roles.get_roles()
view_rights =
case roles do
nil ->
[]
_ ->
StepFlow.Controllers.Roles.get_rights_for_entity_type_and_action(
roles,
"workflow",
"view"
)
end
for %{entity: entity} <- view_rights,
do:
String.split(entity, "::")
|> List.last()
end
defp filter_deleted(query, params) do
case Map.get(params, "deleted") do
value when value in [nil, "none"] ->
from(
workflow in query,
where: workflow.deleted == false
)
"only" ->
from(
workflow in query,
where: workflow.deleted == true
)
"all" ->
query
end
end
defp filter_mode(query, params) do
case Map.get(params, "mode") do
nil ->
from(workflow in query)
["live", "file"] ->
from(workflow in query)
["live"] ->
from(
workflow in query,
where: workflow.is_live == true
)
["file"] ->
from(
workflow in query,
where: workflow.is_live == false
)
end
end
defp filter_version(query, params) do
case Map.get(params, "version") do
nil ->
from(workflow in query)
versions ->
from(
workflow in query,
where:
fragment("CONCAT(version_major,'.',version_minor,'.',version_micro)") in ^versions
)
end
end
def filter_status(query, params, key, start_date \\ nil) do
case StepFlow.Map.get_by_key_or_atom(params, key) do
nil ->
query
states ->
if start_date != nil do
datetime =
case NaiveDateTime.from_iso8601(start_date) do
{:ok, date} ->
date
_ ->
NaiveDateTime.new!(
Date.from_iso8601!(start_date),
Time.new!(0, 0, 0)
)
end
from(
workflow in query,
join:
workflow_status in subquery(
from(
workflow_status in Workflows.Status,
order_by: [desc: workflow_status.id, desc: workflow_status.workflow_id],
distinct: [desc: workflow_status.workflow_id],
where: fragment("?::timestamp", workflow_status.inserted_at) >= ^datetime
)
),
on: workflow.id == workflow_status.workflow_id,
where: workflow_status.state in ^states
)
else
from(
workflow in query,
join:
workflow_status in subquery(
from(
workflow_status in Workflows.Status,
order_by: [desc: workflow_status.id, desc: workflow_status.workflow_id],
distinct: [desc: workflow_status.workflow_id]
)
),
on: workflow.id == workflow_status.workflow_id,
where: workflow_status.state in ^states
)
end
end
end
@doc """
Gets a single workflows.
Raises `Ecto.NoResultsError` if the Workflow does not exist.
## Examples
iex> get_workflows!(123)
%Workflow{}
iex> get_workflows!(456)
** (Ecto.NoResultsError)
"""
def get_workflow!(id) do
Repo.get!(Workflow, id)
|> Repo.preload([:artifacts, jobs: :child_workflow])
|> preload_workflow
end
@doc """
Gets a single workflows containing the specified job ID.
Raises `Ecto.NoResultsError` if the Workflow does not exist.
## Examples
iex> get_workflow_for_job!(19)
%Workflow{}
iex> get_workflows!(456)
** (Ecto.NoResultsError)
"""
def get_workflow_for_job!(job_id) do
job = Jobs.get_job!(job_id)
get_workflow!(job.workflow_id)
end
defp preload_workflow(workflow) do
jobs = Repo.preload(workflow.jobs, [:status, :progressions])
steps =
workflow
|> Map.get(:steps)
|> StepFlow.Controllers.Workflows.get_steps_with_status(jobs)
workflow
|> Map.put(:steps, steps)
|> Map.put(:jobs, jobs)
end
def preload_workflows(workflows, result \\ [])
def preload_workflows([], result), do: result
def preload_workflows([workflow | workflows], result) do
result = List.insert_at(result, -1, workflow |> preload_workflow)
preload_workflows(workflows, result)
end
def get_step_definition(job) do
job = Repo.preload(job, workflow: [jobs: :child_workflow])
step =
Enum.filter(job.workflow.steps, fn step ->
Map.get(step, "id") == job.step_id
end)
|> List.first()
%{step: step, workflow: job.workflow}
end
@doc """
Creates a workflow.
## Examples
iex> create_workflow(%{field: value})
{:ok, %Workflow{}}
iex> create_workflow(%{field: bad_value})
{:error, %Ecto.Changeset{}}
"""
def create_workflow(attrs \\ %{}) do
%Workflow{}
|> Workflow.changeset(attrs)
|> Repo.insert()
end
@doc """
Updates a workflow.
## Examples
iex> update_workflow(workflow, %{field: new_value})
{:ok, %Workflow{}}
iex> update_workflow(workflow, %{field: bad_value})
{:error, %Ecto.Changeset{}}
"""
def update_workflow(%Workflow{} = workflow, attrs) do
workflow
|> Workflow.changeset(attrs)
|> Repo.update()
end
@doc """
Deletes a Workflow.
## Examples
iex> delete_workflow(workflow)
{:ok, %Workflow{}}
iex> delete_workflow(workflow)
{:error, %Ecto.Changeset{}}
"""
def delete_workflow(%Workflow{} = workflow) do
workflow
|> Workflow.changeset(%{deleted: true})
|> Repo.update()
end
@doc """
Duplicates a Workflow.
## Examples
iex> duplicate_workflow(workflow_id, user_uuid)
{:ok, Integer}
iex> duplicate_workflow(workflow_id, user_uuid)
{:error, %Ecto.Changeset{}}
"""
def duplicate_workflow(workflow_id, user_uuid) do
workflow = Workflows.get_workflow!(workflow_id)
attrs = StepFlow.Controllers.Workflows.get_attr(workflow)
start_params = workflow.start_parameters
attrs =
attrs
|> Map.put(
:parameters,
if start_params == [] do
workflow.parameters
else
start_params
end
)
|> Map.put(:user_uuid, user_uuid)
case Workflows.create_workflow(attrs) do
{:ok, duplicated_workflow} ->
Logger.info(
"Workflow #{workflow_id} successfully duplicated as #{duplicated_workflow.id}."
)
Workflows.Status.define_workflow_status(duplicated_workflow.id, :created_workflow)
StepFlow.Step.start_next(duplicated_workflow)
StepFlow.Notification.send("new_workflow", %{workflow_id: duplicated_workflow.id})
{:ok, duplicated_workflow.id}
{:error, message} ->
Logger.error(
"Cannot duplicate workflow #{workflow.id} #{workflow.identifier}: #{inspect(message)}."
)
{:error, message}
end
end
@doc """
Returns an `%Ecto.Changeset{}` for tracking workflow changes.
## Examples
iex> change_workflow(workflow)
%Ecto.Changeset{source: %Workflow{}}
"""
def change_workflow(%Workflow{} = workflow) do
Workflow.changeset(workflow, %{})
end
def get_completed_statistics(scale, delta) do
query =
from(
workflow in Workflow,
inner_join:
artifacts in subquery(
from(
artifacts in Artifact,
where:
artifacts.inserted_at > datetime_add(^NaiveDateTime.utc_now(), ^delta, ^scale),
group_by: artifacts.workflow_id,
select: %{
workflow_id: artifacts.workflow_id,
inserted_at: max(artifacts.inserted_at)
}
)
),
on: workflow.id == artifacts.workflow_id,
group_by: workflow.identifier,
select: %{
count: count(),
duration:
fragment(
"CAST(EXTRACT(EPOCH FROM (SELECT avg(? - ?))) AS FLOAT)",
artifacts.inserted_at,
workflow.inserted_at
),
identifier: workflow.identifier
}
)
Repo.all(query)
end
@doc """
Convert workflow version fields to a version string
"""
def get_workflow_version_as_string(workflow) do
to_string(workflow.version_major) <>
"." <> to_string(workflow.version_minor) <> "." <> to_string(workflow.version_micro)
end
def abort(workflow) do
workflow.steps
|> abort_running_step_jobs(workflow)
workflow.steps
|> skip_remaining_steps(workflow)
result = Workflows.Status.set_workflow_status(workflow.id, :stopped)
topic = "update_workflow_" <> Integer.to_string(workflow.id)
StepFlow.Notification.send(topic, %{workflow_id: workflow.id})
result
end
@doc """
Returns processing live workflow with their reference and identifier
"""
def get_processing_workflow_live do
query =
from(
workflow in Workflow,
where: workflow.is_live == true and workflow.deleted == false,
select: %{
identifier: workflow.identifier,
reference: workflow.reference,
id: workflow.id
},
order_by: [desc: :inserted_at],
limit: 1
)
query =
from(
workflow in subquery(query),
group_by: [workflow.reference, workflow.identifier],
select: %{
identifier: workflow.identifier,
reference: workflow.reference,
count: count(workflow)
}
)
|> filter_status(%{"states" => [:processing]}, :states)
Repo.all(query)
end
def pause(workflow, action, trigger_date_time) do
case action do
action when action in ["resume", "abort"] ->
description = %{
action: action,
trigger_at:
DateTime.from_unix!(trigger_date_time, :millisecond)
|> DateTime.to_naive()
}
workflow = Repo.preload(workflow, jobs: :child_workflow)
has_processing_steps =
workflow.steps
|> Enum.filter(fn step -> step.jobs.processing > 0 end)
|> Enum.empty?()
|> Kernel.not()
workflow_status =
if has_processing_steps do
:pausing
else
:paused
end
result =
Workflows.Status.set_workflow_status(workflow.id, workflow_status, nil, description)
_paused_steps = pause_remaining_steps(workflow.steps, workflow.jobs, %{})
topic = "update_workflow_" <> Integer.to_string(workflow.id)
StepFlow.Notification.send(topic, %{workflow_id: workflow.id})
result
_ ->
{:error, "Unknown action: #{action}"}
end
end
def resume(workflow) do
{:ok, _status} = Workflows.Status.set_workflow_status(workflow.id, :pending)
case StepFlow.Step.start_next(workflow) do
{:ok, _} ->
workflow = Repo.preload(workflow, jobs: :child_workflow)
_resumed_steps = resume_remaining_steps(workflow.steps, workflow.jobs)
result = Workflows.Status.set_workflow_status(workflow.id, :processing)
topic = "update_workflow_" <> Integer.to_string(workflow.id)
StepFlow.Notification.send(topic, %{workflow_id: workflow.id})
result
error ->
error
end
end
defp abort_running_step_jobs([], _workflow), do: nil
defp abort_running_step_jobs([step | steps], workflow) do
case step.status do
:processing ->
case StepFlow.Map.get_by_key_or_atom(step, :mode) do
# When nested workflows
mode when mode in ["workflow_one_for_one", "workflow_one_for_many"] ->
abort_nested_workflows(step, workflow)
_ ->
[
StepFlow.Step.abort_step_jobs(workflow, step)
]
end
_ ->
nil
end
abort_running_step_jobs(steps, workflow)
end
defp abort_nested_workflows(step, workflow) do
step_id = StepFlow.Map.get_by_key_or_atom(step, :id)
step_jobs =
workflow.jobs
|> Enum.filter(fn job -> job.step_id == step_id end)
Enum.each(step_jobs, fn job ->
get_workflow!(job.child_workflow.id)
|> abort()
end)
end
defp skip_remaining_steps([], _workflow), do: nil
defp skip_remaining_steps([step | steps], workflow) do
case step.status do
:queued -> StepFlow.Step.skip_step(workflow, step)
:paused -> StepFlow.Step.skip_step_jobs(workflow, step)
:processing -> StepFlow.Step.skip_step_jobs(workflow, step)
_ -> nil
end
skip_remaining_steps(steps, workflow)
end
defp pause_remaining_steps(_steps, _workflow_jobs, _post_action, _paused_steps \\ [])
defp pause_remaining_steps([], _workflow_jobs, _post_action, paused_steps), do: paused_steps
defp pause_remaining_steps([step | steps], workflow_jobs, post_action, paused_steps) do
# Get step jobs
step_id = StepFlow.Map.get_by_key_or_atom(step, :id)
step_jobs =
workflow_jobs
|> Enum.filter(fn job -> job.step_id == step_id end)
paused_steps =
case step.status do
status when status in [:queued, :processing] ->
Logger.debug("Pause #{inspect(status)} step: #{inspect(step)}")
case StepFlow.Map.get_by_key_or_atom(step, :mode) do
mode when mode in ["workflow_one_for_one", "workflow_one_for_many"] ->
trigger_date_time = DateTime.utc_now() |> DateTime.to_unix()
pause_nested_workflows(step_jobs, post_action, trigger_date_time)
workflow = get_workflow!(List.first(step_jobs).workflow_id)
pause(workflow, "resume", trigger_date_time)
_ ->
[
StepFlow.Step.pause_step(step, step_jobs, post_action) | paused_steps
]
end
_ ->
paused_steps
end
pause_remaining_steps(steps, workflow_jobs, post_action, paused_steps)
end
defp pause_nested_workflows(step_jobs, post_action, trigger_date_time) do
Enum.each(step_jobs, fn job ->
{:ok, _status} = Status.set_job_status(job.id, :paused, post_action)
get_workflow!(job.child_workflow.id)
|> pause("resume", trigger_date_time)
end)
end
defp resume_remaining_steps([], _workflow_jobs), do: nil
defp resume_remaining_steps([step | steps], workflow_jobs) do
case step.status do
:paused ->
Logger.debug("Resume paused step: #{inspect(step)}")
step_id = StepFlow.Map.get_by_key_or_atom(step, :id)
step_jobs =
workflow_jobs
|> Enum.filter(fn job -> job.step_id == step_id end)
case StepFlow.Map.get_by_key_or_atom(step, :mode) do
# When nested workflows
mode when mode in ["workflow_one_for_one", "workflow_one_for_many"] ->
Enum.each(step_jobs, fn job ->
get_workflow!(job.child_workflow.id)
|> resume()
end)
_ ->
StepFlow.Step.resume_step(step_jobs)
end
_ ->
nil
end
resume_remaining_steps(steps, workflow_jobs)
end
end