Packages
step_flow
1.7.2
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.
"""
import Ecto.Query, warn: false
require Logger
alias StepFlow.Amqp.WorkerCreatedConsumer
alias StepFlow.Jobs.Job
alias StepFlow.Jobs.Status
alias StepFlow.LiveWorkers
alias StepFlow.NotificationHooks.NotificationHookManager
alias StepFlow.Repo
alias StepFlow.Repo.Checker
alias StepFlow.Statistics.JobsDurations
alias StepFlow.Workers.WorkerStatuses
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
if Checker.repo_running?() do
handle_created_message(channel, tag, direct_messaging_queue_name, payload)
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("Worker creation #{inspect(payload)}")
Basic.reject(channel, tag, requeue: false)
end
defp handle_created_message(channel, tag, direct_messaging_queue_name, payload) do
job =
StepFlow.Jobs.internal_list_jobs(%{
"direct_messaging_queue_name" => direct_messaging_queue_name
})
|> Map.get(:data)
|> List.first()
instance_id = Map.get(payload, "instance_id")
Logger.debug("Worker #{instance_id} creation job search result: #{inspect(job)}")
case {job, instance_id} do
{nil, nil} ->
Basic.nack(channel, tag, requeue: false)
{nil, _} ->
case declare_worker_status(payload, instance_id) do
:ok ->
Basic.ack(channel, tag)
:error ->
Basic.nack(channel, tag, requeue: false)
end
_ ->
if job.is_live do
job_id = job.id
case live_worker_update(job_id) do
:ok ->
Basic.ack(channel, tag)
:error ->
Basic.nack(channel, tag, requeue: false)
end
else
Basic.ack(channel, tag)
end
end
end
defp declare_worker_status(payload, instance_id) do
worker =
payload
|> Map.delete("parameters")
|> Map.put_new("system_info", %{"docker_container_id" => instance_id})
Logger.debug("Declare #{instance_id} worker status without job: #{inspect(worker)}")
creation_result = WorkerStatuses.create_worker_status(%{"job" => nil, "worker" => worker})
case creation_result do
{:ok, _} ->
Logger.info("Worker #{instance_id} created without job")
workers_statuses = WorkerStatuses.list_worker_statuses()
Logger.debug(
"[#{__MODULE__}] Notify that workers status have been updated: #{inspect(workers_statuses)}"
)
StepFlow.Notification.send("workers_status_updated", %{
content: StepFlow.WorkerStatusView.render("index.json", workers_statuses)
})
:ok
{:error, error} ->
Logger.error("Could not create WorkerStatus: #{inspect(error)}")
:error
end
end
defp live_worker_update(job_id) do
live_worker = LiveWorkers.get_by(%{"job_id" => job_id})
case live_worker do
nil ->
:error
_ ->
case Status.set_job_status(job_id, "ready_to_init") do
{:ok, _status} ->
LiveWorkers.update_live_worker(live_worker, %{
"creation_date" => NaiveDateTime.utc_now()
})
repo_checker()
JobsDurations.set_job_durations(job_id)
NotificationHookManager.notification_from_job(job_id)
StepManager.check_step_status(%{job_id: job_id})
:ok
{:error, message} ->
Logger.error("Cannot set job status: #{inspect(message)}")
:error
end
end
end
defp repo_checker do
query = from(job in Job, select: job.id)
stream = Repo.stream(query)
Repo.transaction(fn ->
Enum.to_list(stream)
end)
:timer.sleep(1000)
end
end