262 lines
7.2 KiB
Elixir
262 lines
7.2 KiB
Elixir
defmodule Towerops.Workers.DeviceMonitorWorker do
|
|
@moduledoc """
|
|
Oban worker that monitors device health by pinging it.
|
|
|
|
Uses Oban's unique job feature to ensure only one monitor per device
|
|
runs cluster-wide. Monitors check device connectivity and create alerts
|
|
on status changes.
|
|
"""
|
|
use Oban.Worker,
|
|
queue: :monitors,
|
|
max_attempts: 1,
|
|
unique: [
|
|
period: :infinity,
|
|
keys: [:device_id],
|
|
states: [:available, :scheduled, :retryable]
|
|
],
|
|
replace: [
|
|
scheduled: [:scheduled_at],
|
|
available: [:scheduled_at]
|
|
]
|
|
|
|
alias Towerops.Agents
|
|
alias Towerops.Alerts
|
|
alias Towerops.Alerts.StormDetector
|
|
alias Towerops.Devices
|
|
alias Towerops.Monitoring
|
|
alias Towerops.Snmp.Client
|
|
alias Towerops.Workers.AlertNotificationWorker
|
|
alias Towerops.Workers.PollingOffset
|
|
|
|
require Logger
|
|
|
|
@ping_module Application.compile_env(:towerops, :ping_module, Towerops.Monitoring.Ping)
|
|
@compile {:no_warn_undefined, Towerops.Monitoring.PingMock}
|
|
|
|
@monitor_interval 30
|
|
|
|
@impl Oban.Worker
|
|
@spec perform(Oban.Job.t()) :: :ok
|
|
def perform(%Oban.Job{args: %{"device_id" => device_id}}) do
|
|
case Devices.get_device(device_id) do
|
|
nil ->
|
|
Logger.debug("Device #{device_id} no longer exists, skipping monitor and not rescheduling")
|
|
:ok
|
|
|
|
device ->
|
|
# Only reschedule if device still has monitoring enabled
|
|
if should_continue_monitoring?(device) do
|
|
maybe_perform_check(device)
|
|
schedule_next_check_with_error_handling(device_id)
|
|
else
|
|
Logger.debug("Device #{device_id} monitoring disabled, not rescheduling",
|
|
device_id: device_id,
|
|
monitoring_enabled: device.monitoring_enabled
|
|
)
|
|
end
|
|
|
|
:ok
|
|
end
|
|
end
|
|
|
|
# Check if device should continue being monitored
|
|
# Returns false if Phoenix SNMP is disabled, monitoring is disabled, or device has any agent
|
|
defp should_continue_monitoring?(device) do
|
|
!Client.phoenix_snmp_disabled() && device.monitoring_enabled &&
|
|
!Agents.device_has_effective_agent?(device.id)
|
|
end
|
|
|
|
# Agent check already done by should_continue_monitoring?, no need to re-check
|
|
defp maybe_perform_check(device) do
|
|
perform_check(device)
|
|
end
|
|
|
|
defp schedule_next_check_with_error_handling(device_id) do
|
|
# Verify device still exists before scheduling next check
|
|
case Devices.get_device(device_id) do
|
|
nil ->
|
|
Logger.debug("Device deleted, not rescheduling monitoring", device_id: device_id)
|
|
:ok
|
|
|
|
_device ->
|
|
case schedule_next_check(device_id) do
|
|
{:ok, _job} ->
|
|
:ok
|
|
|
|
{:error, changeset} ->
|
|
Logger.error("Failed to schedule next monitoring check for device #{device_id}: #{inspect(changeset.errors)}")
|
|
:ok
|
|
end
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Starts monitoring for a device.
|
|
"""
|
|
def start_monitoring(device_id) do
|
|
offset = PollingOffset.calculate_offset(device_id, @monitor_interval)
|
|
|
|
%{device_id: device_id}
|
|
|> new(schedule_in: offset)
|
|
|> Oban.insert()
|
|
end
|
|
|
|
@doc """
|
|
Stops monitoring for a device by cancelling its jobs.
|
|
"""
|
|
def stop_monitoring(device_id) do
|
|
import Ecto.Query
|
|
|
|
Oban.cancel_all_jobs(
|
|
from(j in Oban.Job,
|
|
where: j.worker == "Towerops.Workers.DeviceMonitorWorker",
|
|
where: fragment("args->>'device_id' = ?", ^device_id),
|
|
where: j.state in ["available", "scheduled", "executing", "retryable"]
|
|
)
|
|
)
|
|
end
|
|
|
|
@doc """
|
|
Triggers an immediate check for a device.
|
|
"""
|
|
def trigger_check(device_id) do
|
|
%{device_id: device_id}
|
|
|> new()
|
|
|> Oban.insert()
|
|
end
|
|
|
|
# Private Functions
|
|
|
|
defp perform_check(device) do
|
|
check_result = check_device_health(device)
|
|
|
|
now = Towerops.Time.now()
|
|
|
|
{status, response_time_ms} =
|
|
case check_result do
|
|
{:ok, time} -> {:success, time}
|
|
{:error, _reason} -> {:failure, nil}
|
|
end
|
|
|
|
case Monitoring.create_monitoring_check(%{
|
|
device_id: device.id,
|
|
status: to_string(status),
|
|
response_time_ms: response_time_ms,
|
|
checked_at: now
|
|
}) do
|
|
{:ok, _check} ->
|
|
:ok
|
|
|
|
{:error, changeset} ->
|
|
Logger.error("Failed to create monitoring check for device #{device.id}: #{inspect(changeset.errors)}")
|
|
end
|
|
|
|
new_status = if status == :success, do: :up, else: :down
|
|
old_status = device.status
|
|
|
|
_ = Devices.update_device_status(device, new_status)
|
|
|
|
_ =
|
|
if old_status != new_status do
|
|
handle_status_change(device, old_status, new_status)
|
|
end
|
|
|
|
# Broadcast status change via PubSub
|
|
_ =
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"device:#{device.id}",
|
|
{:device_status_changed, device.id, new_status, nil}
|
|
)
|
|
end
|
|
|
|
defp check_device_health(device) do
|
|
@ping_module.ping(device.ip_address)
|
|
end
|
|
|
|
defp handle_status_change(device, old_status, new_status) do
|
|
now = Towerops.Time.now()
|
|
|
|
case {old_status, new_status} do
|
|
{_, :down} ->
|
|
handle_equipment_down(device, now)
|
|
|
|
{_, :up} ->
|
|
handle_equipment_up(device, now)
|
|
end
|
|
end
|
|
|
|
defp handle_equipment_down(device, now) do
|
|
if Alerts.has_active_alert?(device.id, "device_down") or
|
|
Alerts.has_active_alert?(device.id, "site_outage") do
|
|
:ok
|
|
else
|
|
# Route through the StormDetector for site-level correlation.
|
|
# The StormDetector buffers events for a short window, then either:
|
|
# - Creates a single "site_outage" alert if multiple devices at the same site are down
|
|
# - Creates individual "device_down" alerts if below the correlation threshold
|
|
# Maintenance checks are handled inside the StormDetector.
|
|
StormDetector.register_device_down(device, now)
|
|
end
|
|
end
|
|
|
|
defp handle_equipment_up(device, now) do
|
|
recovery_message = get_recovery_message(device)
|
|
|
|
# Use case instead of pattern match to handle errors gracefully
|
|
case Alerts.create_alert(%{
|
|
device_id: device.id,
|
|
alert_type: "device_up",
|
|
triggered_at: now,
|
|
resolved_at: now,
|
|
message: recovery_message
|
|
}) do
|
|
{:ok, alert} ->
|
|
# Enqueue notification job - Oban handles retries and persistence
|
|
AlertNotificationWorker.enqueue_trigger(alert.id)
|
|
|
|
resolve_down_alert(device)
|
|
|
|
_ =
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"alerts:org:#{device.organization_id}:resolved",
|
|
{:alert_resolved, device.id, :device_down}
|
|
)
|
|
|
|
:ok
|
|
|
|
{:error, reason} ->
|
|
require Logger
|
|
|
|
Logger.error("Failed to create device up alert for #{device.name}: #{inspect(reason)}",
|
|
device_id: device.id
|
|
)
|
|
|
|
{:error, reason}
|
|
end
|
|
end
|
|
|
|
defp get_recovery_message(device) do
|
|
if device.snmp_enabled do
|
|
"Device is now responding to SNMP"
|
|
else
|
|
"Device is now responding to ping"
|
|
end
|
|
end
|
|
|
|
defp resolve_down_alert(device) do
|
|
case Alerts.get_active_alert(device.id, "device_down") do
|
|
nil -> :ok
|
|
alert -> Alerts.resolve_alert(alert)
|
|
end
|
|
end
|
|
|
|
defp schedule_next_check(device_id) do
|
|
offset = PollingOffset.calculate_offset(device_id, @monitor_interval)
|
|
|
|
%{device_id: device_id}
|
|
|> new(schedule_in: @monitor_interval + offset)
|
|
|> Oban.insert()
|
|
end
|
|
end
|