Packages
step_flow
1.9.0-rc1
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/controllers/step/launch_workflows.ex
defmodule StepFlow.Step.LaunchWorkflows do
@moduledoc """
The Step Workflows launching parameters.
"""
require Logger
alias StepFlow.Amqp.CommonEmitter
alias StepFlow.Jobs
alias StepFlow.Jobs.Status
alias StepFlow.Step
alias StepFlow.Step.Helpers
alias StepFlow.Step.Launch
alias StepFlow.Workflows
def generate_child_workflow_job(workflow, step) do
job_params = %{
name: StepFlow.Map.get_by_key_or_atom(step, :name),
step_id: StepFlow.Map.get_by_key_or_atom(step, :id),
workflow_id: workflow.id,
allow_failure: StepFlow.Map.get_by_key_or_atom(step, :allow_failure),
parameters: StepFlow.Map.get_by_key_or_atom(step, :parameters)
}
Jobs.create_job(job_params)
end
# CAREFUL: trick here is that source_paths can be a list or a string
def generate_child_workflow_params(workflow, step, job, dates, source_paths) do
wf_identifier = Helpers.get_value_in_parameters(step, "identifier") |> List.first()
wf_major = Helpers.get_value_in_parameters(step, "version_major") |> List.first()
wf_minor = Helpers.get_value_in_parameters(step, "version_minor") |> List.first()
wf_micro = Helpers.get_value_in_parameters(step, "version_micro") |> List.first()
wf_def =
StepFlow.WorkflowDefinitions.get_workflow_definition(
wf_identifier,
wf_major,
wf_minor,
wf_micro
)
wf_params =
Helpers.get_value_in_parameters(step, "parameters")
|> List.first()
|> Launch.filter_and_pre_compile_parameters(workflow, step, dates, source_paths)
%{
identifier: wf_identifier,
parameters: wf_params,
parent_id: job.id,
reference: "[AUTOGENERATED] Child from workflow #{workflow.id}",
schema_version: StepFlow.Map.get_by_key_or_atom(wf_def, :schema_version),
steps: StepFlow.Map.get_by_key_or_atom(wf_def, :steps),
version_major: wf_major,
version_minor: wf_minor,
version_micro: wf_micro,
user_uuid: workflow.user_uuid,
tags: StepFlow.Map.get_by_key_or_atom(wf_def, :tags)
}
end
defp start_workflow(workflow_params, workflow, job) do
case Workflows.create_workflow(workflow_params) do
{:ok, child_workflow} ->
Logger.info(
"Created child workflow #{child_workflow.id} from parent workflow #{workflow.id}"
)
Workflows.Status.define_workflow_status(child_workflow.id, :created_workflow)
{:ok, datetime} = DateTime.now("Etc/UTC")
case CommonEmitter.publish_json(
"job_progression",
0,
%{
job_id: job.id,
datetime: datetime,
docker_container_id: "workflow",
progression: 0
},
"job_response"
) do
:ok ->
Step.start_next(child_workflow)
StepFlow.Notification.send("new_workflow", %{workflow_id: child_workflow.id})
_ ->
Logger.error("[#{__MODULE__}] Could not send progression.")
end
{:error, _} ->
Logger.error("Cannot create child workflow from parent workflow #{workflow.id}")
with :error <-
CommonEmitter.publish_json(
"job_error",
0,
%{
job_id: job.id,
status: "error",
parameters: [
%{
"id" => "message",
"type" => "string",
"value" =>
"Cannot create child workflow from parent workflow #{workflow.id}"
}
]
},
"job_response"
) do
Logger.error("[#{__MODULE__}] Could not send error message.")
end
end
end
def start_child_workflow_one_for_many(workflow, step, dates, source_paths) do
datetime = NaiveDateTime.to_string(DateTime.utc_now())
{:ok, job} = generate_child_workflow_job(workflow, step)
Status.set_job_status(job.id, "queued", %{}, datetime)
generate_child_workflow_params(workflow, step, job, dates, source_paths)
|> start_workflow(workflow, job)
{:ok, "started"}
end
defp start_child_workflow_one_for_one(workflow, step, dates, source_path) do
datetime = NaiveDateTime.to_string(DateTime.utc_now())
{:ok, job} = generate_child_workflow_job(workflow, step)
Status.set_job_status(job.id, "queued", %{}, datetime)
generate_child_workflow_params(workflow, step, job, dates, source_path)
|> start_workflow(workflow, job)
end
def start_child_workflows_one_for_one(_, _, _, []),
do: {:ok, "started"}
def start_child_workflows_one_for_one(
workflow,
step,
dates,
[source_path | source_paths]
) do
start_child_workflow_one_for_one(
workflow,
step,
dates,
source_path
)
start_child_workflows_one_for_one(workflow, step, dates, source_paths)
end
end