Packages
step_flow
1.8.1-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/amqp/stopped_consumer.ex
defmodule StepFlow.Amqp.StoppedConsumer do
@moduledoc """
Consumer of all job with stopped status.
"""
require Logger
alias StepFlow.{
Amqp.StoppedConsumer,
Jobs,
Jobs.Status,
Metrics.JobInstrumenter,
NotificationHooks.NotificationHookManager,
Repo.Checker,
Statistics.JobsDurations,
Workflows,
Workflows.StepManager
}
use StepFlow.Amqp.CommonConsumer, %{
queue: "job_stopped",
exchange: "job_response",
prefetch_count: 1,
consumer: &StoppedConsumer.consume/4
}
@doc """
Consume messages with stopped topic, update Job status and continue the workflow.
"""
def consume(
channel,
tag,
_redelivered,
%{
"job_id" => job_id,
"datetime" => datetime,
"status" => status
} = _payload
) do
if Checker.repo_running?() do
case Jobs.get_job(job_id) do
nil ->
Basic.nack(channel, tag, requeue: false)
job ->
JobInstrumenter.inc(:step_flow_jobs_status_total, job.name, "stopped")
case Status.set_job_status(job_id, status, %{}, datetime) do
{:ok, job_status} ->
Workflows.Status.define_workflow_status(job.workflow_id, :job_stopped, job_status)
JobsDurations.set_job_durations(job_id)
NotificationHookManager.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.nack(channel, tag, requeue: false)
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 stopped #{inspect(payload)}")
Basic.reject(channel, tag, requeue: false)
end
end