Packages
step_flow
1.6.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/error_consumer.ex
defmodule StepFlow.Amqp.ErrorConsumer do
@moduledoc """
Consumer of all job with error status.
"""
require Logger
alias StepFlow.Amqp.ErrorConsumer
alias StepFlow.Jobs
alias StepFlow.Jobs.Status
alias StepFlow.Metrics.{JobInstrumenter, WorkflowInstrumenter}
alias StepFlow.NotificationHooks.NotificationHookManager
alias StepFlow.Repo
alias StepFlow.Repo.Checker
alias StepFlow.Step
alias StepFlow.Step.Live
alias StepFlow.Workflows
alias StepFlow.Workflows.StepManager
use StepFlow.Amqp.CommonConsumer, %{
queue: "job_error",
exchange: "job_response",
prefetch_count: 1,
consumer: &ErrorConsumer.consume/4
}
@doc """
Consume message with error topic, update Job and send a notification.
"""
def consume(channel, tag, _redelivered, %{"job_id" => job_id, "error" => description} = payload) do
if Checker.repo_running?() do
case Jobs.get_job(job_id) do
nil ->
Basic.nack(channel, tag, requeue: false)
job ->
Logger.error("Job error #{inspect(payload)}")
case Status.set_job_status(job_id, :error, %{message: description}) do
{:ok, job_status} ->
JobInstrumenter.inc(:step_flow_jobs_error, job.name)
WorkflowInstrumenter.inc(:step_flow_workflows_error, job.workflow_id)
NotificationHookManager.notification_from_job(
job_id,
description
)
if job.allow_failure do
Logger.warn("Allowing failure for job #{job_id}")
StepManager.check_step_status(%{job_id: job_id})
else
Workflows.Status.define_workflow_status(job.workflow_id, :job_error, job_status)
end
if job.is_live do
restart_workflow(job)
end
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,
%{
"job_id" => job_id,
"parameters" => [%{"id" => "message", "type" => "string", "value" => description}],
"status" => "error"
} = payload
) do
if Checker.repo_running?() do
case Jobs.get_job(job_id) do
nil ->
Basic.nack(channel, tag, requeue: false)
job ->
Logger.error("Job error #{inspect(payload)}")
JobInstrumenter.inc(:step_flow_jobs_error, job.name)
WorkflowInstrumenter.inc(:step_flow_workflows_error, job.workflow_id)
case Status.set_job_status(job_id, :error, %{message: description}) do
{:ok, job_status} ->
if job.allow_failure do
Logger.warn("Allowing failure for job #{job_id}")
NotificationHookManager.notification_from_job(
job_id,
description
)
StepManager.check_step_status(%{job_id: job_id})
else
Workflows.Status.define_workflow_status(job.workflow_id, :job_error, job_status)
NotificationHookManager.manage_notification_status(job_id, "job", "error")
NotificationHookManager.notification_from_job(
job_id,
description
)
end
if job.is_live do
restart_workflow(job)
end
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 error #{inspect(payload)}")
Basic.reject(channel, tag, requeue: false)
end
defp restart_workflow(job) do
workflow =
job
|> Map.get(:workflow_id)
|> Workflows.get_workflow!()
workflow_jobs = Repo.preload(workflow, [:jobs]).jobs
workflow_jobs
|> Live.stop_jobs()
case Live.delete_worker_from_job(job) do
{:error, "unable to publish message"} ->
Logger.error("Cannot delete worker linked to job #{job.id} in error.")
_ ->
Logger.info("Deleted worker linked to job #{job.id} in error.")
end
{:ok, restarted_workflow} =
StepFlow.Controllers.Workflows.get_attr(workflow)
|> Workflows.create_workflow()
WorkflowInstrumenter.inc(:step_flow_workflows_created, restarted_workflow.identifier)
Workflows.Status.define_workflow_status(restarted_workflow.id, :created_workflow)
Step.start_next(restarted_workflow)
StepFlow.Notification.send("new_workflow", %{workflow_id: restarted_workflow.id})
Logger.info("Workflow #{workflow.id} successfully restarted as #{restarted_workflow.id}.")
end
end