Packages
step_flow
1.1.0
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/amqp/completed_consumer.ex
defmodule StepFlow.Amqp.CompletedConsumer do
@moduledoc """
Consumer of all job with completed status.
"""
require Logger
alias StepFlow.{
Amqp.CompletedConsumer,
Jobs,
Jobs.Status,
LiveWorkers,
Metrics.JobInstrumenter,
Workflows,
Workflows.StepManager
}
use StepFlow.Amqp.CommonConsumer, %{
queue: "job_completed",
exchange: "job_response",
prefetch_count: 1,
consumer: &CompletedConsumer.consume/4
}
@doc """
Consume messages with completed topic, update Job status and continue the workflow.
"""
def consume(
channel,
tag,
_redelivered,
%{
"job_id" => job_id,
"status" => status
} = payload
) do
case Jobs.get_job(job_id) do
nil ->
Basic.reject(channel, tag, requeue: false)
job ->
if job.is_live do
case live_worker_update(job_id, payload) do
:ok ->
StepManager.check_step_status(%{job_id: job_id})
Basic.ack(channel, tag)
:error ->
Basic.reject(channel, tag, requeue: true)
end
else
workflow =
job
|> Map.get(:workflow_id)
|> Workflows.get_workflow!()
set_generated_destination_paths(payload, job)
set_output_parameters(payload, workflow)
JobInstrumenter.inc(:step_flow_jobs_completed, job.name)
{:ok, job_status} = Status.set_job_status(job_id, status)
Workflows.Status.define_workflow_status(job.workflow_id, :job_completed, job_status)
Workflows.notification_from_job(job_id)
StepManager.check_step_status(%{job_id: job_id})
Basic.ack(channel, tag)
end
end
end
def consume(channel, tag, _redelivered, payload) do
Logger.error("Job completed #{inspect(payload)}")
Basic.reject(channel, tag, requeue: false)
end
defp set_generated_destination_paths(payload, job) do
case StepFlow.Map.get_by_key_or_atom(payload, "destination_paths") do
nil ->
nil
destination_paths ->
job_parameters =
job.parameters ++
[
%{
id: "destination_paths",
type: "array_of_strings",
value: destination_paths
}
]
Jobs.update_job(job, %{parameters: job_parameters})
end
end
defp set_output_parameters(payload, workflow) do
case StepFlow.Map.get_by_key_or_atom(payload, "parameters") do
nil ->
nil
parameters ->
parameters = workflow.parameters ++ parameters
Workflows.update_workflow(workflow, %{parameters: parameters})
end
end
defp live_worker_update(job_id, payload) do
live_worker = LiveWorkers.get_by(%{"job_id" => job_id})
case live_worker do
nil ->
live_worker_creation(job_id, payload)
_ ->
case live_worker.termination_date do
nil ->
# :error
Status.set_job_status(job_id, "completed")
Workflows.notification_from_job(job_id)
:ok
_ ->
Status.set_job_status(job_id, "completed")
Workflows.notification_from_job(job_id)
:ok
end
end
end
defp live_worker_creation(job_id, payload) do
instance_id =
StepFlow.Map.get_by_key_or_atom(payload, :parameters)
|> Enum.filter(fn param ->
StepFlow.Map.get_by_key_or_atom(param, :id) == "instance_id"
end)
|> List.first()
|> StepFlow.Map.get_by_key_or_atom(:value)
host_ip =
StepFlow.Map.get_by_key_or_atom(payload, :parameters)
|> Enum.filter(fn param ->
StepFlow.Map.get_by_key_or_atom(param, :id) == "host_ip"
end)
|> List.first()
|> StepFlow.Map.get_by_key_or_atom(:value)
ports =
StepFlow.Map.get_by_key_or_atom(payload, :parameters)
|> Enum.filter(fn param ->
StepFlow.Map.get_by_key_or_atom(param, :id) == "ports"
end)
|> List.first()
|> StepFlow.Map.get_by_key_or_atom(:value)
direct_messaging_queue_name =
StepFlow.Map.get_by_key_or_atom(payload, :parameters)
|> Enum.filter(fn param ->
StepFlow.Map.get_by_key_or_atom(param, :id) == "direct_messaging_queue_name"
end)
|> List.first()
|> StepFlow.Map.get_by_key_or_atom(:value)
LiveWorkers.create_live_worker(%{
job_id: job_id,
instance_id: instance_id,
direct_messaging_queue_name: direct_messaging_queue_name,
ips: [host_ip],
ports: ports
})
:ok
end
end