Packages
step_flow
1.6.0-rc15
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,
Repo,
Repo.Checker,
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
if Checker.repo_running?() 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, 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)
case Status.set_job_status(job_id, status, payload) do
{:ok, job_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)
{:error, message} ->
Logger.error("Cannot set job status: #{inspect(message)}")
Basic.reject(channel, tag, requeue: true)
end
end
end
else
Logger.warn(
"#{__MODULE__}: The database is not available, reject and requeue consumed message..."
)
Basic.reject(channel, tag, requeue: true)
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, payload) do
job_id = job.id
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 ->
Logger.warn("No termination date found for the worker linked to job id #{job_id}")
_ ->
Logger.info("Termination date found for the worker linked to job id #{job_id}")
end
complete_live_job(job)
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
defp complete_live_job(job) do
job = Jobs.get_job!(job.id) |> Repo.preload([:status])
last_status = Status.get_last_status(job.status)
if Status.convert_to_string(last_status) == "error" do
Logger.info("Job already in error, not transiting to completed.")
:ok
else
case Status.set_job_status(job.id, "completed") do
{:ok, job_status} ->
Workflows.Status.define_workflow_status(
job.workflow_id,
:completed_workflow,
job_status
)
Workflows.notification_from_job(job.id)
:ok
{:error, message} ->
Logger.error("Cannot set job status: #{inspect(message)}")
:error
end
end
end
end