Packages
step_flow
1.0.0-rc8
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/worker_created_consumer.ex
defmodule StepFlow.Amqp.WorkerCreatedConsumer do
@moduledoc """
Consumer of all worker creations.
"""
require Logger
alias StepFlow.Amqp.WorkerCreatedConsumer
alias StepFlow.Jobs.Status
alias StepFlow.LiveWorkers
alias StepFlow.Workflows
alias StepFlow.Workflows.StepManager
use StepFlow.Amqp.CommonConsumer, %{
queue: "worker_created",
exchange: "worker_response",
prefetch_count: 1,
consumer: &WorkerCreatedConsumer.consume/4
}
@doc """
Consume worker created message.
"""
def consume(
channel,
tag,
_redelivered,
%{
"direct_messaging_queue_name" => direct_messaging_queue_name
} = _payload
) do
job =
StepFlow.Jobs.list_jobs(%{
"direct_messaging_queue_name" => direct_messaging_queue_name
})
|> Map.get(:data)
|> List.first()
Logger.debug("Worker Creation job search result: #{inspect(job)}")
case job do
nil ->
Basic.reject(channel, tag, requeue: false)
_ ->
job_id = job.id
case live_worker_update(job_id) do
:ok ->
Basic.ack(channel, tag)
:error ->
Basic.reject(channel, tag, requeue: true)
end
end
end
def consume(channel, tag, _redelivered, payload) do
Logger.error("Worker creation #{inspect(payload)}")
Basic.reject(channel, tag, requeue: false)
end
defp live_worker_update(job_id) do
live_worker = LiveWorkers.get_by(%{"job_id" => job_id})
case live_worker do
nil ->
:error
_ ->
LiveWorkers.update_live_worker(live_worker, %{
"creation_date" => NaiveDateTime.utc_now()
})
Status.set_job_status(job_id, "ready_to_init")
Workflows.notification_from_job(job_id)
StepManager.check_step_status(%{job_id: job_id})
:ok
end
end
end