Packages
step_flow
1.6.0-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/workers/worker_status_watcher.ex
defmodule StepFlow.Workers.WorkerStatusWatcher do
require Logger
@moduledoc false
use GenServer
alias StepFlow.Amqp.CommonEmitter
alias StepFlow.Workers.WorkerStatus
alias StepFlow.Workers.WorkerStatuses
@default_workers_status_interval 10_000
def child_spec(_) do
%{
id: StepFlow.Workers.WorkerStatusWatcher,
start: {StepFlow.Workers.WorkerStatusWatcher, :start_link, [%{}]}
}
end
def start_link(workers_status \\ %{}) do
GenServer.start_link(__MODULE__, workers_status, name: __MODULE__)
end
def update_worker_status(instance_id, worker_status) do
GenServer.cast(__MODULE__, {:push, instance_id, worker_status})
end
def get_worker_status(instance_id) do
GenServer.call(__MODULE__, :pop, instance_id)
end
@impl true
def init(workers_status) do
Logger.info("[#{__MODULE__}] Start checking workers status!")
check_workers_status()
{:ok, workers_status}
end
@impl true
def handle_info(:check, workers_status) do
Logger.debug(
"[#{__MODULE__}] Check workers status. Current workers: #{inspect(workers_status)}"
)
# Schedule once more
check_workers_status()
{:noreply, workers_status}
end
@impl true
def handle_call(:pop, instance_id, workers_status) do
Logger.debug(
"[#{__MODULE__}] Get worker #{instance_id} status from : #{inspect(workers_status)}"
)
worker_status = WorkerStatuses.get_worker_status!(instance_id)
{:reply, worker_status, workers_status}
end
@impl true
def handle_cast({:push, instance_id, worker_status}, workers_status) do
status =
case WorkerStatuses.get_worker_status(instance_id) do
nil ->
Logger.debug(
"[#{__MODULE__}] Add #{instance_id} worker to workers status: #{inspect(workers_status)}"
)
WorkerStatuses.create_worker_status!(worker_status)
status ->
Logger.debug(
"[#{__MODULE__}] Update #{instance_id} worker status: #{inspect(workers_status)}"
)
WorkerStatuses.update_worker_status!(status, worker_status)
end
update_worker_job(status)
Logger.debug("[#{__MODULE__}] Notify that workers status have been updated.")
StepFlow.Notification.send("workers_status_updated", %{})
{:noreply, workers_status}
end
defp check_workers_status do
CommonEmitter.publish(
"",
"{ \"type\": \"status\" }",
[headers: [broadcast: "true"]],
"direct_messaging"
)
interval =
Application.get_env(:step_flow, StepFlow.Workers,
workers_status_interval: @default_workers_status_interval
)
|> Keyword.get(:workers_status_interval, @default_workers_status_interval)
Process.send_after(self(), :check, interval)
end
defp update_worker_job(%WorkerStatus{} = worker_status) do
job =
case worker_status.current_job do
nil ->
nil
current_job ->
StepFlow.Jobs.get_job(current_job.job_id)
end
case job do
nil ->
{}
job ->
instance_id = StepFlow.Map.get_by_key_or_atom(worker_status, :instance_id)
if instance_id != job.last_worker_instance_id do
Logger.info("#{__MODULE__}: set instance_id #{inspect(instance_id)} to job #{job.id}.")
StepFlow.Jobs.update_job(job, %{last_worker_instance_id: instance_id})
end
end
end
end