Packages
step_flow
1.3.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/step/step.ex
defmodule StepFlow.Step do
@moduledoc """
The Step context.
"""
require Logger
alias StepFlow.Artifacts
alias StepFlow.Jobs
alias StepFlow.Metrics.WorkflowInstrumenter
alias StepFlow.Repo
alias StepFlow.Step.Helpers
alias StepFlow.Step.Launch
alias StepFlow.Workflows
alias StepFlow.Workflows.Workflow
def start_next(%Workflow{id: workflow_id} = workflow) do
workflow = Repo.preload(workflow, :jobs, force: true)
is_live = workflow.is_live
jobs = Repo.preload(workflow.jobs, [:status, :progressions])
steps =
StepFlow.Map.get_by_key_or_atom(workflow, :steps)
|> Workflows.get_step_status(jobs)
{is_completed_workflow, steps_to_start} = get_steps_to_start(steps, is_live)
steps_to_start =
case {steps_to_start, jobs} do
{[], []} ->
case List.first(steps) do
nil ->
Logger.warn("#{__MODULE__}: empty workflow #{workflow_id} is completed")
{:completed_workflow, []}
step ->
{:ok, [step]}
end
{[], _} ->
{:completed_workflow, []}
{list, _} ->
{:ok, list}
end
results = start_steps(steps_to_start, workflow)
get_final_status(workflow, is_completed_workflow, Enum.uniq(results) |> Enum.sort())
end
def skip_step(workflow, step) do
step_id = StepFlow.Map.get_by_key_or_atom(step, :id)
step_name = StepFlow.Map.get_by_key_or_atom(step, :name)
Repo.preload(workflow, :jobs, force: true)
|> Jobs.create_skipped_job(step_id, step_name)
end
def abort_step_jobs(workflow, step) do
step_id = StepFlow.Map.get_by_key_or_atom(step, :id)
step_name = StepFlow.Map.get_by_key_or_atom(step, :name)
Repo.preload(workflow, :jobs, force: true)
|> Jobs.abort_jobs(step_id, step_name)
end
def skip_step_jobs(workflow, step) do
step_id = StepFlow.Map.get_by_key_or_atom(step, :id)
step_name = StepFlow.Map.get_by_key_or_atom(step, :name)
Repo.preload(workflow, :jobs, force: true)
|> Jobs.skip_jobs(step_id, step_name)
end
defp set_artifacts(workflow) do
resources = %{}
params = %{
resources: resources,
workflow_id: workflow.id
}
Artifacts.create_artifact(params)
end
defp get_steps_to_start(steps, is_live), do: iter_get_steps_to_start(steps, steps, is_live)
defp iter_get_steps_to_start(steps, all_steps, is_live, completed \\ true, result \\ [])
defp iter_get_steps_to_start([], _all_steps, _is_live, completed, result),
do: {completed, result}
defp iter_get_steps_to_start([step | steps], all_steps, true, completed, result) do
result = List.insert_at(result, -1, step)
iter_get_steps_to_start(steps, all_steps, true, completed, result)
end
defp iter_get_steps_to_start([step | steps], all_steps, false, completed, result) do
completed =
if step.status in [:completed, :skipped, :stopped] do
completed
else
false
end
result =
if step.status == :queued do
case StepFlow.Map.get_by_key_or_atom(step, :required_to_start) do
nil ->
List.insert_at(result, -1, step)
required_to_start ->
count_not_completed =
Enum.filter(all_steps, fn s ->
StepFlow.Map.get_by_key_or_atom(s, :id) in required_to_start
end)
|> Enum.map(fn s -> StepFlow.Map.get_by_key_or_atom(s, :status) end)
|> Enum.filter(fn s -> s != :completed and s != :skipped end)
|> length
if count_not_completed == 0 do
List.insert_at(result, -1, step)
else
result
end
end
else
result
end
iter_get_steps_to_start(steps, all_steps, false, completed, result)
end
defp start_steps({:completed_workflow, _}, _workflow), do: [:completed_workflow]
defp start_steps({:ok, steps}, workflow) do
dates = Helpers.get_dates()
for step <- steps do
step_name = StepFlow.Map.get_by_key_or_atom(step, :name)
step_id = StepFlow.Map.get_by_key_or_atom(step, :id)
source_paths = Launch.get_source_paths(workflow, step, dates)
Logger.warn(
"#{__MODULE__}: start to process step #{step_name} (index #{step_id}) for workflow #{
workflow.id
}"
)
{result, status} =
StepFlow.Map.get_by_key_or_atom(step, :condition)
|> case do
condition when condition in [0, nil] ->
Launch.launch_step(workflow, step)
condition ->
Helpers.template_process(
"<%= " <> condition <> "%>",
workflow,
step,
dates,
source_paths
)
|> case do
"true" ->
Launch.launch_step(workflow, step)
"false" ->
skip_step(workflow, step)
{:ok, "skipped"}
_ ->
Logger.error(
"#{__MODULE__}: cannot estimate condition for step #{step_name} (index #{
step_id
}) for workflow #{workflow.id}"
)
{:error, "bad step condition"}
end
end
Logger.info("#{step_name}: #{inspect({result, status})}")
topic = "update_workflow_" <> Integer.to_string(workflow.id)
StepFlow.Notification.send(topic, %{workflow_id: workflow.id})
status
end
end
defp get_final_status(_workflow, _is_completed_workflow, ["started"]), do: {:ok, "started"}
defp get_final_status(_workflow, _is_completed_workflow, ["created"]), do: {:ok, "started"}
defp get_final_status(_workflow, _is_completed_workflow, ["created", "started"]),
do: {:ok, "started"}
defp get_final_status(workflow, _is_completed_workflow, ["skipped"]), do: start_next(workflow)
defp get_final_status(workflow, _is_completed_workflow, ["completed"]), do: start_next(workflow)
defp get_final_status(workflow, true, [:completed_workflow]) do
Workflows.Status.define_workflow_status(workflow.id, :completed_workflow)
WorkflowInstrumenter.inc(:step_flow_workflows_completed, workflow.identifier)
set_artifacts(workflow)
Logger.warn("#{__MODULE__}: workflow #{workflow.id} is completed")
{:ok, "completed"}
end
defp get_final_status(_workflow, _is_completed_workflow, _states), do: {:ok, "still_processing"}
end