Packages
step_flow
1.4.2-rc2
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/progression_consumer.ex
defmodule StepFlow.Amqp.ProgressionConsumer do
@moduledoc """
Consumer of all progression jobs.
"""
require Logger
alias StepFlow.Amqp.ProgressionConsumer
alias StepFlow.Jobs
alias StepFlow.Jobs.Status
alias StepFlow.Progressions
alias StepFlow.Repo
alias StepFlow.Repo.Checker
alias StepFlow.Workers.WorkerStatuses
alias StepFlow.Workflows
use StepFlow.Amqp.CommonConsumer, %{
queue: "job_progression",
exchange: "job_response",
prefetch_count: 1,
consumer: &ProgressionConsumer.consume/4
}
@doc """
Consume message with job progression and save it in database.
"""
def consume(
channel,
tag,
_redelivered,
%{
"job_id" => job_id
} = payload
) do
if Checker.repo_running?() do
case Jobs.get_job(job_id) do
nil ->
Basic.reject(channel, tag, requeue: false)
job ->
job = Repo.preload(Jobs.get_job(job.id), [:status])
last_status = Status.get_last_status(job.status)
case last_status do
nil ->
add_progression(channel, tag, job, payload)
_ ->
error_status_check(channel, tag, job, payload, last_status)
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 progression #{inspect(payload)}")
Basic.reject(channel, tag, requeue: false)
end
def error_status_check(channel, tag, job, payload, last_status) do
if last_status.state == :error do
Logger.warn("Progression arrived after job #{job.id} was put in error, ignoring.")
Basic.ack(channel, tag)
else
add_progression(channel, tag, job, payload)
end
end
def add_progression(channel, tag, job, payload) do
{_, progression} = Progressions.create_progression(payload)
if progression.progression == 0 do
worker_status = WorkerStatuses.create_worker_status!(progression)
Logger.debug(
"[#{__MODULE__}] Added worker status from first progression: #{inspect(worker_status)}"
)
end
Workflows.Status.define_workflow_status(job.workflow_id, :job_progression, progression)
Workflows.notification_from_job(job.id)
Basic.ack(channel, tag)
end
end