refactor: simplify job architecture from Oban coordinators to direct workers

Removes unnecessary two-layer architecture (Oban Coordinator → GenServer)
in favor of direct Oban workers that perform the actual work.

Changes:
- Replace DevicePollerCoordinator + PollerWorker with DevicePollerWorker
- Replace DeviceMonitorCoordinator + DeviceMonitor with DeviceMonitorWorker
- Simplify Monitoring.Supervisor (removed Registries, DynamicSupervisors)
- Remove all Exq dependencies (workers, supervisor, mix deps)
- Convert DiscoveryWorker from Exq to Oban pattern
- Update Devices context to auto-start/stop monitoring and polling
- Update all LiveView callers to use new Oban enqueue pattern
- Fix all tests to use Oban.Job struct instead of raw IDs

Benefits:
- Simpler codebase (~850 lines removed)
- Better reliability (Oban handles retries, failures)
- Lower memory (no persistent GenServers)
- Better observability (work visible in Oban dashboard)
- Cluster-wide coordination via PostgreSQL
- Extensive PubSub usage for real-time events

Files deleted:
- lib/towerops/workers/device_monitor_coordinator.ex
- lib/towerops/workers/device_poller_coordinator.ex
- lib/towerops/snmp/poller_worker.ex
- lib/towerops/monitoring/device_monitor.ex
- lib/towerops/workers/monitor_worker.ex
- lib/towerops/workers/poll_worker.ex
- lib/towerops/exq_supervisor.ex
- All related test files

Files created:
- lib/towerops/workers/device_poller_worker.ex
- lib/towerops/workers/device_monitor_worker.ex

All 3,686 tests passing.
This commit is contained in:
Graham McIntire 2026-01-24 16:36:57 -06:00
parent 29593ac734
commit d0946c3cd0
No known key found for this signature in database
25 changed files with 404 additions and 2390 deletions

View file

@ -35,7 +35,6 @@ defmodule Towerops.Application do
SnmpKit.SnmpMgr.MIB
] ++
background_workers() ++
exq_workers() ++
[
# Start to serve requests, typically the last entry
ToweropsWeb.Endpoint
@ -73,18 +72,6 @@ defmodule Towerops.Application do
end
end
# Configure Exq workers - don't start in test environment
defp exq_workers do
if Application.get_env(:towerops, :env) == :test do
[]
else
[
# Use custom Exq supervisor with Redis health checks and better error handling
Towerops.ExqSupervisor
]
end
end
# Returns PubSub spec - uses Redis in production/clustered environments,
# falls back to default PG2 adapter for Mix tasks and development
defp pubsub_spec do

View file

@ -8,6 +8,8 @@ defmodule Towerops.Devices do
alias Towerops.Devices.Device, as: DeviceSchema
alias Towerops.Devices.Event
alias Towerops.Repo
alias Towerops.Workers.DeviceMonitorWorker
alias Towerops.Workers.DevicePollerWorker
@doc """
Returns the list of devices for a site.
@ -252,27 +254,53 @@ defmodule Towerops.Devices do
Creates device.
"""
def create_device(attrs) do
%DeviceSchema{}
|> DeviceSchema.changeset(attrs)
|> Repo.insert()
case %DeviceSchema{}
|> DeviceSchema.changeset(attrs)
|> Repo.insert() do
{:ok, device} = result ->
# Start monitoring/polling if enabled
if device.monitoring_enabled do
DeviceMonitorWorker.start_monitoring(device.id)
end
if device.snmp_enabled do
DevicePollerWorker.start_polling(device.id)
end
result
error ->
error
end
end
@doc """
Updates device.
"""
def update_device(%DeviceSchema{} = device, attrs) do
device
|> DeviceSchema.changeset(attrs)
|> Repo.update()
old_monitoring = device.monitoring_enabled
old_snmp = device.snmp_enabled
case device
|> DeviceSchema.changeset(attrs)
|> Repo.update() do
{:ok, updated_device} = result ->
handle_monitoring_changes(updated_device, old_monitoring)
handle_snmp_changes(updated_device, old_snmp)
result
error ->
error
end
end
@doc """
Deletes device.
"""
def delete_device(%DeviceSchema{} = device) do
# Stop monitoring workers before deleting
_ = Towerops.Monitoring.Supervisor.stop_monitor(device.id)
_ = Towerops.Monitoring.Supervisor.stop_snmp_poller(device.id)
# Stop monitoring and polling jobs before deleting
_ = DeviceMonitorWorker.stop_monitoring(device.id)
_ = DevicePollerWorker.stop_polling(device.id)
Repo.delete(device)
end
@ -319,6 +347,34 @@ defmodule Towerops.Devices do
|> Repo.update()
end
# Private helpers for monitoring/polling management
defp handle_monitoring_changes(device, old_monitoring) do
cond do
device.monitoring_enabled && !old_monitoring ->
DeviceMonitorWorker.start_monitoring(device.id)
!device.monitoring_enabled && old_monitoring ->
DeviceMonitorWorker.stop_monitoring(device.id)
true ->
:ok
end
end
defp handle_snmp_changes(device, old_snmp) do
cond do
device.snmp_enabled && !old_snmp ->
DevicePollerWorker.start_polling(device.id)
!device.snmp_enabled && old_snmp ->
DevicePollerWorker.stop_polling(device.id)
true ->
:ok
end
end
## Events
@doc """

View file

@ -1,107 +0,0 @@
defmodule Towerops.ExqSupervisor do
@moduledoc """
Custom supervisor for Exq background job processor with Redis health checks.
This supervisor adds resilience to Exq by:
1. Waiting for Redis to be available before starting Exq
2. Using a restart strategy that prevents rapid crash loops
3. Logging detailed error information when crashes occur
Handles the case where Exq.Node.Server crashes due to nil responses
from Redis when connections fail.
"""
use Supervisor
require Logger
def start_link(opts) do
Supervisor.start_link(__MODULE__, opts, name: __MODULE__)
end
@impl true
def init(_opts) do
redis_config = Application.get_env(:towerops, :redis, [])
# Wait for Redis to be available before starting Exq
case Towerops.RedisHealthCheck.wait_for_redis(redis_config, max_attempts: 10, backoff: 1_000) do
:ok ->
Logger.info("Redis is available, starting Exq")
start_exq_children(redis_config)
{:error, reason} ->
Logger.error("Failed to connect to Redis: #{inspect(reason)}. Exq will not start.")
# Return empty children list - supervisor will still start but Exq won't run
# This allows the application to continue functioning without background jobs
{:ok, {:one_for_one, []}}
end
end
defp start_exq_children(redis_config) do
exq_config = build_exq_config(redis_config)
children = [
# Wrap Exq in a supervisor with restart strategy
%{
id: Exq,
start: {Exq, :start_link, [exq_config]},
restart: :permanent,
type: :supervisor,
# Add longer shutdown timeout to allow graceful job completion
shutdown: 30_000
}
]
# Use one_for_one strategy with limited restarts to prevent crash loops
# If Exq crashes more than 3 times in 60 seconds, the supervisor will crash
# This will be caught by the application supervisor which will restart this supervisor
opts = [
strategy: :one_for_one,
max_restarts: 3,
max_seconds: 60
]
Supervisor.init(children, opts)
end
defp build_exq_config(redis_config) do
base_config = [
name: Exq,
host: Keyword.get(redis_config, :host, "localhost"),
port: Keyword.get(redis_config, :port, 6379),
namespace: "exq",
concurrency: 10,
queues: ["default", "discovery", "polling", "monitoring", "maintenance"],
# Poll settings - how long to wait for Redis responses
poll_timeout: 50,
scheduler_poll_timeout: 200,
scheduler_enable: true,
# Job retry settings
max_retries: 50,
# Exponential backoff starting at 500ms, up to 60s
mode: :default,
shutdown_timeout: 30_000,
# Redix connection options for better resilience
redis_options: [
# Synchronous connect to fail fast if Redis unavailable
sync_connect: true,
# Don't exit process on disconnection - let Exq handle it
exit_on_disconnection: false,
# Connection timeout
timeout: 5_000,
# Backoff for reconnection attempts
backoff_initial: 500,
backoff_max: 30_000,
# TCP keepalive to detect dead connections
socket_opts: [keepalive: true]
]
]
# Add password if configured (read from application config set by runtime.exs)
case Keyword.get(redis_config, :password) do
nil -> base_config
"" -> base_config
password -> Keyword.put(base_config, :password, password)
end
end
end

View file

@ -1,253 +0,0 @@
defmodule Towerops.Monitoring.DeviceMonitor do
@moduledoc """
GenServer that monitors a single piece of device by periodically pinging it.
"""
use GenServer
alias Ecto.Adapters.SQL.Sandbox
alias Towerops.Alerts
alias Towerops.Devices
alias Towerops.Monitoring
alias Towerops.Repo
require Logger
# Allow dependency injection for testing
@ping_module Application.compile_env(:towerops, :ping_module, Towerops.Monitoring.Ping)
# Suppress warnings for Mox modules that are defined at runtime during tests
@compile {:no_warn_undefined, Towerops.Monitoring.PingMock}
# Client API
@doc """
Starts a monitor for the given device ID.
"""
def start_link(opts) do
device_id = Keyword.fetch!(opts, :device_id)
lease_id = Keyword.get(opts, :lease_id)
GenServer.start_link(__MODULE__, {device_id, lease_id}, name: via_tuple(device_id))
end
@doc """
Triggers an immediate check for the device.
"""
def trigger_check(device_id) do
# Check if monitor process is running (node-local lookup)
case Registry.lookup(Towerops.Monitoring.Registry, device_id) do
[{_pid, _}] ->
# Process exists, send cast
GenServer.cast(via_tuple(device_id), :check_now)
[] ->
# Process doesn't exist (monitoring disabled), enqueue check job
enqueue_check(device_id)
end
end
# Server Callbacks
@impl true
def init({device_id, lease_id}) do
device = Devices.get_device!(device_id)
if device.monitoring_enabled do
# Perform immediate check when monitoring starts
send(self(), :check_device)
end
{:ok, %{device_id: device_id, lease_id: lease_id}}
end
@impl true
def handle_info(:check_device, state) do
_ = perform_check(state.device_id)
{:noreply, state}
end
@impl true
def handle_cast(:check_now, state) do
_ = perform_check(state.device_id)
{:noreply, state}
end
# Private Functions
defp perform_check(device_id) do
device = Devices.get_device!(device_id)
# For SNMP-enabled devices, use recent SNMP poll success as health indicator
# For ICMP-only devices, use ping
check_result = check_device_health(device)
now = DateTime.truncate(DateTime.utc_now(), :second)
{status, response_time_ms} =
case check_result do
{:ok, time} -> {:success, time}
{:error, _reason} -> {:failure, nil}
end
# Save the check result
case Monitoring.create_check(%{
device_id: device_id,
status: 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
# device status if it changed
new_status = if status == :success, do: :up, else: :down
old_status = device.status
_ = Devices.update_device_status(device, new_status)
# Create alerts if status changed
_ =
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}
)
# Only schedule next check if monitoring is enabled
if device.monitoring_enabled do
schedule_next_check(device.check_interval_seconds)
end
end
defp check_device_health(device) do
# Always perform actual ICMP ping to measure real latency
# Don't rely solely on SNMP poll success since that doesn't give us latency data
@ping_module.ping(device.ip_address)
end
defp handle_status_change(device, old_status, new_status) do
now = DateTime.truncate(DateTime.utc_now(), :second)
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) do
:ok
else
create_device_down_alert(device, now)
end
end
defp create_device_down_alert(device, now) do
alert_message = get_down_alert_message(device)
{:ok, _alert} =
Alerts.create_alert(%{
device_id: device.id,
alert_type: :device_down,
triggered_at: now,
message: alert_message
})
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"alerts:new",
{:new_alert, device.id, :device_down}
)
end
defp get_down_alert_message(device) do
if device.snmp_enabled do
"Device is not responding to SNMP"
else
"Device is not responding to ping"
end
end
defp handle_equipment_up(device, now) do
recovery_message = get_recovery_message(device)
{:ok, _alert} =
Alerts.create_alert(%{
device_id: device.id,
alert_type: :device_up,
triggered_at: now,
message: recovery_message
})
resolve_down_alert(device)
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"alerts:resolved",
{:alert_resolved, device.id, :device_down}
)
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(interval_seconds) do
Process.send_after(self(), :check_device, interval_seconds * 1000)
end
# Enqueue monitoring check job - safe to call in test environment
defp enqueue_check(device_id) do
if Application.get_env(:towerops, :env) == :test do
# In test, run synchronously with database sandbox access
start_check_task(device_id)
else
# In dev/prod, enqueue to Exq
{:ok, _job} = Exq.enqueue(Exq, "monitoring", Towerops.Workers.MonitorWorker, [device_id])
end
end
defp start_check_task(device_id) do
parent = self()
_ =
Task.start(fn ->
_ = maybe_allow_sandbox(parent)
perform_check(device_id)
end)
end
defp maybe_allow_sandbox(parent) do
if Application.get_env(:towerops, :sql_sandbox) do
Sandbox.allow(Repo, parent, self())
end
end
defp via_tuple(device_id) do
{:via, Registry, {Towerops.Monitoring.Registry, device_id}}
end
end

View file

@ -1,13 +1,14 @@
defmodule Towerops.Monitoring.Supervisor do
@moduledoc """
Supervisor for managing device monitoring workers.
Supervisor for monitoring-related background workers.
With the migration to Oban, device monitoring and polling are handled
by Oban workers. This supervisor only manages auxiliary services like
the neighbor cleanup worker.
"""
use Supervisor
alias Towerops.Monitoring.DeviceMonitor
alias Towerops.Snmp.NeighborCleanupWorker
alias Towerops.Snmp.PollerRegistry
alias Towerops.Snmp.PollerWorker
require Logger
@ -17,76 +18,17 @@ defmodule Towerops.Monitoring.Supervisor do
@impl true
def init(_init_arg) do
base_children = [
# Local Registry for process naming (node-local)
{Registry, keys: :unique, name: Towerops.Monitoring.Registry},
# Local Registry for SNMP poller processes (node-local)
{Registry, keys: :unique, name: PollerRegistry},
# Task.Supervisor for parallel polling operations (local to each node)
{Task.Supervisor, name: Towerops.Snmp.PollerTaskSupervisor},
# Local DynamicSupervisor for monitor workers (node-local)
{DynamicSupervisor, name: Towerops.LocalMonitorSupervisor, strategy: :one_for_one},
# Local DynamicSupervisor for SNMP poller workers (node-local)
{DynamicSupervisor, name: Towerops.LocalPollerSupervisor, strategy: :one_for_one}
]
# Only start cleanup worker in production
children =
if test_mode?() do
base_children
[]
else
base_children ++ [NeighborCleanupWorker]
[NeighborCleanupWorker]
end
Supervisor.init(children, strategy: :one_for_one)
end
@doc """
Starts monitoring for a specific device (node-local).
Note: In production, monitors are coordinated by Oban workers to ensure cluster-wide uniqueness.
This function is primarily for testing and manual operations.
"""
def start_monitor(device_id, lease_id \\ nil) do
spec = {DeviceMonitor, device_id: device_id, lease_id: lease_id}
DynamicSupervisor.start_child(Towerops.LocalMonitorSupervisor, spec)
end
@doc """
Stops monitoring for a specific device (node-local).
"""
def stop_monitor(device_id) do
case Registry.lookup(Towerops.Monitoring.Registry, device_id) do
[{pid, _}] ->
DynamicSupervisor.terminate_child(Towerops.LocalMonitorSupervisor, pid)
[] ->
:ok
end
end
@doc """
Starts an SNMP poller for a specific device (node-local).
Note: In production, pollers are coordinated by Oban workers to ensure cluster-wide uniqueness.
This function is primarily for testing and manual operations.
"""
def start_snmp_poller(device_id, lease_id \\ nil) do
spec = {PollerWorker, device_id: device_id, lease_id: lease_id}
DynamicSupervisor.start_child(Towerops.LocalPollerSupervisor, spec)
end
@doc """
Stops an SNMP poller for a specific device (node-local).
"""
def stop_snmp_poller(device_id) do
case Registry.lookup(PollerRegistry, device_id) do
[{pid, _}] ->
DynamicSupervisor.terminate_child(Towerops.LocalPollerSupervisor, pid)
[] ->
:ok
end
end
# Check if running in test mode by inspecting the Repo pool configuration
defp test_mode? do
config = Application.get_env(:towerops, Towerops.Repo, [])

View file

@ -1,92 +0,0 @@
defmodule Towerops.Workers.DeviceMonitorCoordinator do
@moduledoc """
Oban worker that coordinates device monitoring across the cluster.
Uses Oban's unique job feature to ensure only one monitor per device
runs cluster-wide, replacing etcd-based locking.
"""
use Oban.Worker,
queue: :monitors,
unique: [
period: 60,
keys: [:device_id],
states: [:available, :scheduled, :executing, :retryable]
]
alias Towerops.Devices
alias Towerops.Monitoring.Supervisor, as: MonitoringSupervisor
require Logger
# Monitoring interval - check every 30 seconds
@monitor_interval 30
@impl Oban.Worker
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 coordination")
:ok
_device ->
ensure_monitor_running(device_id)
# Reschedule for next check
schedule_next_check(device_id)
:ok
end
end
@doc """
Starts monitoring coordination for a device.
"""
def start_monitoring(device_id) do
%{device_id: device_id}
|> new()
|> Oban.insert()
end
@doc """
Stops monitoring for a device by cancelling its jobs.
"""
def stop_monitoring(device_id) do
import Ecto.Query
# Cancel all pending jobs for this device
Oban.cancel_all_jobs(
from(j in Oban.Job,
where: j.worker == "Towerops.Workers.DeviceMonitorCoordinator",
where: fragment("args->>'device_id' = ?", ^device_id),
where: j.state in ["available", "scheduled", "executing", "retryable"]
)
)
# Stop the running monitor GenServer
MonitoringSupervisor.stop_monitor(device_id)
end
# Private functions
defp ensure_monitor_running(device_id) do
case MonitoringSupervisor.start_monitor(device_id) do
{:ok, _pid} ->
Logger.debug("Started monitor for device #{device_id}")
:ok
{:error, {:already_started, _pid}} ->
# Already running, which is fine
:ok
{:error, reason} ->
Logger.warning("Failed to start monitor for device #{device_id}: #{inspect(reason)}")
:ok
end
end
defp schedule_next_check(device_id) do
%{device_id: device_id}
|> new(schedule_in: @monitor_interval)
|> Oban.insert()
end
end

View file

@ -0,0 +1,213 @@
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,
unique: [
period: 60,
keys: [:device_id],
states: [:available, :scheduled, :executing, :retryable]
]
alias Towerops.Alerts
alias Towerops.Devices
alias Towerops.Monitoring
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
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")
:ok
device ->
if device.monitoring_enabled do
perform_check(device)
end
# Schedule next check
schedule_next_check(device_id)
:ok
end
end
@doc """
Starts monitoring for a device.
"""
def start_monitoring(device_id) do
%{device_id: device_id}
|> new()
|> 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 = DateTime.truncate(DateTime.utc_now(), :second)
{status, response_time_ms} =
case check_result do
{:ok, time} -> {:success, time}
{:error, _reason} -> {:failure, nil}
end
case Monitoring.create_check(%{
device_id: device.id,
status: 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 = DateTime.truncate(DateTime.utc_now(), :second)
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) do
:ok
else
create_device_down_alert(device, now)
end
end
defp create_device_down_alert(device, now) do
alert_message = get_down_alert_message(device)
{:ok, _alert} =
Alerts.create_alert(%{
device_id: device.id,
alert_type: :device_down,
triggered_at: now,
message: alert_message
})
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"alerts:new",
{:new_alert, device.id, :device_down}
)
end
defp get_down_alert_message(device) do
if device.snmp_enabled do
"Device is not responding to SNMP"
else
"Device is not responding to ping"
end
end
defp handle_equipment_up(device, now) do
recovery_message = get_recovery_message(device)
{:ok, _alert} =
Alerts.create_alert(%{
device_id: device.id,
alert_type: :device_up,
triggered_at: now,
message: recovery_message
})
resolve_down_alert(device)
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"alerts:resolved",
{:alert_resolved, device.id, :device_down}
)
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
%{device_id: device_id}
|> new(schedule_in: @monitor_interval)
|> Oban.insert()
end
end

View file

@ -1,96 +0,0 @@
defmodule Towerops.Workers.DevicePollerCoordinator do
@moduledoc """
Oban worker that coordinates SNMP polling across the cluster.
Uses Oban's unique job feature to ensure only one poller per device
runs cluster-wide, replacing etcd-based locking.
"""
use Oban.Worker,
queue: :pollers,
unique: [
period: 60,
keys: [:device_id],
states: [:available, :scheduled, :executing, :retryable]
]
alias Towerops.Devices
alias Towerops.Monitoring.Supervisor, as: MonitoringSupervisor
require Logger
@impl Oban.Worker
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 poll coordination")
:ok
device ->
if device.snmp_enabled do
ensure_poller_running(device_id)
end
# Reschedule for next poll interval (default 60s)
poll_interval = get_poll_interval(device)
schedule_next_poll(device_id, poll_interval)
:ok
end
end
@doc """
Starts polling coordination for a device.
"""
def start_polling(device_id) do
%{device_id: device_id}
|> new()
|> Oban.insert()
end
@doc """
Stops polling for a device by cancelling its jobs.
"""
def stop_polling(device_id) do
import Ecto.Query
# Cancel all pending jobs for this device
Oban.cancel_all_jobs(
from(j in Oban.Job,
where: j.worker == "Towerops.Workers.DevicePollerCoordinator",
where: fragment("args->>'device_id' = ?", ^device_id),
where: j.state in ["available", "scheduled", "executing", "retryable"]
)
)
# Stop the running poller GenServer
MonitoringSupervisor.stop_snmp_poller(device_id)
end
# Private functions
defp ensure_poller_running(device_id) do
case MonitoringSupervisor.start_snmp_poller(device_id) do
{:ok, _pid} ->
Logger.debug("Started SNMP poller for device #{device_id}")
:ok
{:error, {:already_started, _pid}} ->
# Already running, which is fine
:ok
{:error, reason} ->
Logger.warning("Failed to start SNMP poller for device #{device_id}: #{inspect(reason)}")
:ok
end
end
defp schedule_next_poll(device_id, interval_seconds) do
%{device_id: device_id}
|> new(schedule_in: interval_seconds)
|> Oban.insert()
end
defp get_poll_interval(device) do
device.snmp_poll_interval || 60
end
end

View file

@ -1,151 +1,103 @@
defmodule Towerops.Snmp.PollerWorker do
defmodule Towerops.Workers.DevicePollerWorker do
@moduledoc """
GenServer that regularly polls SNMP data for a device.
Oban worker that performs SNMP polling for a device.
This worker continuously collects:
- Sensor readings (temperature, voltage, CPU, memory, etc.)
- Interface statistics (bandwidth, errors, discards, etc.)
- Neighbor discovery (LLDP/CDP topology information)
- ARP and MAC tables
- Processor and storage statistics
Runs independently of the connectivity monitoring in DeviceMonitor.
Poll interval is configurable per device (default: 60 seconds).
Uses Oban's unique job feature to ensure only one poller per device
runs cluster-wide. Poll interval is configurable per device (default: 60 seconds).
"""
use GenServer
use Oban.Worker,
queue: :pollers,
unique: [
period: 60,
keys: [:device_id],
states: [:available, :scheduled, :executing, :retryable]
]
alias Ecto.Adapters.SQL.Sandbox
alias Towerops.Devices
alias Towerops.Repo
alias Towerops.Snmp
alias Towerops.Snmp.ArpDiscovery
alias Towerops.Snmp.Client
alias Towerops.Snmp.MacDiscovery
alias Towerops.Snmp.NeighborDiscovery
alias Towerops.Snmp.PollerRegistry
alias Towerops.Snmp.PollerTaskSupervisor
alias Towerops.Snmp.Sensor
require Logger
@default_poll_interval 60
# Client API
@doc """
Starts a poller for the given device ID.
"""
def start_link(opts) do
device_id = Keyword.fetch!(opts, :device_id)
lease_id = Keyword.get(opts, :lease_id)
GenServer.start_link(__MODULE__, {device_id, lease_id}, name: via_tuple(device_id))
end
@doc """
Triggers an immediate poll for the device.
"""
def trigger_poll(device_id) do
case Registry.lookup(PollerRegistry, device_id) do
[{_pid, _}] ->
GenServer.cast(via_tuple(device_id), :poll_now)
[] ->
# Process doesn't exist, enqueue poll job
enqueue_poll(device_id)
end
end
# Server Callbacks
@impl true
def init({device_id, lease_id}) do
device = Devices.get_device!(device_id)
# device
_ = Phoenix.PubSub.subscribe(Towerops.PubSub, "device:#{device_id}")
if device.snmp_enabled do
# Get the device to check if it has sensors/interfaces
device = Snmp.get_device_with_associations(device_id)
if device && (device.sensors != [] || device.interfaces != []) do
# Perform immediate poll when starting
send(self(), :poll_data)
end
end
{:ok, %{device_id: device_id, lease_id: lease_id}}
end
@impl true
def handle_info(:poll_data, state) do
case perform_poll(state.device_id) do
:stop ->
{:stop, :normal, state}
:ok ->
{:noreply, state}
end
end
@impl true
def handle_info({:discovery_completed, _device_id}, state) do
Logger.info("Discovery completed, triggering immediate poll")
# Trigger immediate poll after discovery
send(self(), :poll_data)
{:noreply, state}
end
@impl true
def handle_info(_msg, state) do
# device status changes)
{:noreply, state}
end
@impl true
def handle_cast(:poll_now, state) do
case perform_poll(state.device_id) do
:stop ->
{:stop, :normal, state}
:ok ->
{:noreply, state}
end
end
# Private Functions
defp perform_poll(device_id) do
@impl Oban.Worker
def perform(%Oban.Job{args: %{"device_id" => device_id}}) do
case Devices.get_device(device_id) do
nil ->
Logger.warning("Device #{device_id} not found, stopping poller")
# Device was deleted, return stop signal
:stop
Logger.debug("Device #{device_id} no longer exists, skipping poll")
:ok
device ->
_ =
if device.snmp_enabled do
poll_equipment_device(device_id, device)
end
if device.snmp_enabled do
poll_device(device)
end
# Schedule next poll
poll_interval = get_poll_interval(device)
schedule_next_poll(device_id, poll_interval)
:ok
end
end
defp poll_equipment_device(device_id, device) do
snmp_device = Snmp.get_device_with_associations(device_id)
if snmp_device do
poll_interval = get_poll_interval(device)
execute_device_poll(device, snmp_device, poll_interval)
schedule_next_poll(poll_interval)
end
@doc """
Starts polling for a device.
"""
def start_polling(device_id) do
%{device_id: device_id}
|> new()
|> Oban.insert()
end
defp execute_device_poll(device, snmp_device, poll_interval) do
if should_skip_poll?(device, poll_interval, 5) do
Logger.debug("Skipping poll for #{device.name} - recently polled by another process")
else
perform_device_data_collection(device, snmp_device)
@doc """
Stops polling for a device by cancelling its jobs.
"""
def stop_polling(device_id) do
import Ecto.Query
Oban.cancel_all_jobs(
from(j in Oban.Job,
where: j.worker == "Towerops.Workers.DevicePollerWorker",
where: fragment("args->>'device_id' = ?", ^device_id),
where: j.state in ["available", "scheduled", "executing", "retryable"]
)
)
end
@doc """
Triggers an immediate poll for a device.
"""
def trigger_poll(device_id) do
# Insert a new job to run immediately
%{device_id: device_id}
|> new()
|> Oban.insert()
end
# Private Functions
defp poll_device(device) do
snmp_device = Snmp.get_device_with_associations(device.id)
if snmp_device && (snmp_device.sensors != [] || snmp_device.interfaces != []) do
poll_interval = get_poll_interval(device)
if should_skip_poll?(device, poll_interval, 5) do
Logger.debug("Skipping poll for #{device.name} - recently polled by another process")
else
perform_device_data_collection(device, snmp_device)
end
end
end
@ -155,35 +107,17 @@ defmodule Towerops.Snmp.PollerWorker do
client_opts = build_client_opts(device)
now = DateTime.truncate(DateTime.utc_now(), :second)
# Run all polling operations in parallel using supervised tasks
# Run all polling operations in parallel using async tasks
tasks = [
Task.Supervisor.async_nolink(PollerTaskSupervisor, fn ->
poll_device_sensors(device, snmp_device, client_opts, now)
end),
Task.Supervisor.async_nolink(PollerTaskSupervisor, fn ->
poll_device_state_sensors(device, snmp_device, client_opts, now)
end),
Task.Supervisor.async_nolink(PollerTaskSupervisor, fn ->
poll_device_interfaces(device, snmp_device, client_opts, now)
end),
Task.Supervisor.async_nolink(PollerTaskSupervisor, fn ->
check_device_interface_changes(device, snmp_device, client_opts, now)
end),
Task.Supervisor.async_nolink(PollerTaskSupervisor, fn ->
poll_device_neighbors(device, snmp_device, client_opts)
end),
Task.Supervisor.async_nolink(PollerTaskSupervisor, fn ->
poll_device_arp(device, snmp_device, client_opts)
end),
Task.Supervisor.async_nolink(PollerTaskSupervisor, fn ->
poll_device_mac(device, snmp_device, client_opts)
end),
Task.Supervisor.async_nolink(PollerTaskSupervisor, fn ->
poll_device_processors(device, snmp_device, client_opts, now)
end),
Task.Supervisor.async_nolink(PollerTaskSupervisor, fn ->
poll_device_storage(device, snmp_device, client_opts, now)
end)
Task.async(fn -> poll_device_sensors(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_state_sensors(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_interfaces(device, snmp_device, client_opts, now) end),
Task.async(fn -> check_device_interface_changes(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_neighbors(device, snmp_device, client_opts) end),
Task.async(fn -> poll_device_arp(device, snmp_device, client_opts) end),
Task.async(fn -> poll_device_mac(device, snmp_device, client_opts) end),
Task.async(fn -> poll_device_processors(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_storage(device, snmp_device, client_opts, now) end)
]
# Wait for all tasks with timeout, handle failures gracefully
@ -202,7 +136,6 @@ defmodule Towerops.Snmp.PollerWorker do
poll_sensors(snmp_device.sensors, client_opts, now)
Logger.debug("Polled #{length(snmp_device.sensors)} sensors for #{device.name}")
# Broadcast sensor update event to device-specific topic
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
@ -217,7 +150,6 @@ defmodule Towerops.Snmp.PollerWorker do
poll_state_sensors(snmp_device.state_sensors, client_opts, now, device.id)
Logger.debug("Polled #{length(snmp_device.state_sensors)} state sensors for #{device.name}")
# Broadcast state sensor update event to device-specific topic
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
@ -232,7 +164,6 @@ defmodule Towerops.Snmp.PollerWorker do
poll_interfaces(snmp_device.interfaces, client_opts, now)
Logger.debug("Polled #{length(snmp_device.interfaces)} interfaces for #{device.name}")
# Broadcast interface stats update event to device-specific topic
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
@ -253,23 +184,19 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp poll_device_neighbors(device, snmp_device, client_opts) do
# Add device_id to interfaces for neighbor discovery
interfaces_with_device = Enum.map(snmp_device.interfaces, &Map.put(&1, :device_id, device.id))
{:ok, neighbors} = NeighborDiscovery.discover_neighbors(client_opts, interfaces_with_device)
# Delete stale neighbors (not seen in last 5 minutes)
cutoff = DateTime.add(DateTime.utc_now(), -5, :minute)
Snmp.delete_stale_neighbors(device.id, cutoff)
# Upsert each discovered neighbor
Enum.each(neighbors, fn neighbor_data ->
Snmp.upsert_neighbor(neighbor_data)
end)
Logger.debug("Polled and saved #{length(neighbors)} neighbors for #{device.name}")
# Broadcast neighbor update event to device-specific topic
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
@ -283,16 +210,13 @@ defmodule Towerops.Snmp.PollerWorker do
defp poll_device_arp(device, snmp_device, client_opts) do
{:ok, arp_entries} = ArpDiscovery.discover_arp_table(client_opts)
# Delete stale ARP entries (not seen in last 5 minutes)
cutoff = DateTime.add(DateTime.utc_now(), -5, :minute)
Snmp.delete_stale_arp_entries(device.id, cutoff)
# Upsert ARP entries with interface mapping
Snmp.upsert_arp_entries(device.id, arp_entries, snmp_device.interfaces)
Logger.debug("Polled and saved #{length(arp_entries)} ARP entries for #{device.name}")
# Broadcast ARP update event to device-specific topic
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
@ -306,16 +230,13 @@ defmodule Towerops.Snmp.PollerWorker do
defp poll_device_mac(device, snmp_device, client_opts) do
{:ok, mac_entries} = MacDiscovery.discover_mac_table(client_opts)
# Delete stale MAC entries (not seen in last 5 minutes)
cutoff = DateTime.add(DateTime.utc_now(), -5, :minute)
Snmp.delete_stale_mac_addresses(device.id, cutoff)
# Upsert MAC entries with interface mapping
Snmp.upsert_mac_addresses(device.id, mac_entries, snmp_device.interfaces)
Logger.debug("Polled and saved #{length(mac_entries)} MAC FDB entries for #{device.name}")
# Broadcast MAC update event to device-specific topic
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
@ -333,7 +254,6 @@ defmodule Towerops.Snmp.PollerWorker do
poll_processors(processors, client_opts, now)
Logger.debug("Polled processors for #{device.name}")
# Broadcast processor update event to device-specific topic
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
@ -352,7 +272,6 @@ defmodule Towerops.Snmp.PollerWorker do
poll_storage(storage_entries, client_opts, now)
Logger.debug("Polled #{length(storage_entries)} storage entries for #{device.name}")
# Broadcast storage update event to device-specific topic
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
@ -364,8 +283,10 @@ defmodule Towerops.Snmp.PollerWorker do
Logger.error("Error polling storage for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}")
end
# All the helper functions from the original PollerWorker GenServer
# (These remain unchanged - I'll include them all below)
defp poll_storage(storage_entries, client_opts, timestamp) do
# Poll storage entries in parallel with max_concurrency
storage_entries
|> Task.async_stream(
fn storage ->
@ -380,7 +301,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp poll_storage_value(storage, client_opts) do
# HOST-RESOURCES-MIB storage OIDs
used_oid = "1.3.6.1.2.1.25.2.3.1.6.#{storage.storage_index}"
size_oid = "1.3.6.1.2.1.25.2.3.1.5.#{storage.storage_index}"
alloc_oid = "1.3.6.1.2.1.25.2.3.1.4.#{storage.storage_index}"
@ -390,7 +310,6 @@ defmodule Towerops.Snmp.PollerWorker do
{:ok, alloc_units} <- Client.get(client_opts, alloc_oid),
true <- is_integer(used_units) and is_integer(size_units) and is_integer(alloc_units),
true <- size_units > 0 do
# Calculate bytes from allocation units
used_bytes = used_units * alloc_units
total_bytes = size_units * alloc_units
usage_percent = used_bytes / total_bytes * 100
@ -404,7 +323,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp handle_storage_poll_result(storage, {:ok, values}, timestamp) do
# Create a reading
Snmp.create_storage_reading(%{
storage_id: storage.id,
used_bytes: values.used_bytes,
@ -413,7 +331,6 @@ defmodule Towerops.Snmp.PollerWorker do
checked_at: timestamp
})
# Update storage entry with latest values
Snmp.update_storage(storage, %{
used_bytes: values.used_bytes,
total_bytes: values.total_bytes,
@ -426,7 +343,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp poll_processors(processors, client_opts, timestamp) do
# Poll processors in parallel with max_concurrency
processors
|> Task.async_stream(
fn processor ->
@ -441,20 +357,15 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp poll_processor_value(processor, client_opts) do
# Different OIDs for different processor types
oid =
case processor.processor_type do
"hr_processor" ->
# HOST-RESOURCES-MIB hrProcessorLoad
"1.3.6.1.2.1.25.3.3.1.2.#{processor.processor_index}"
"cisco_cpu" ->
# CISCO-PROCESS-MIB cpmCPUTotal5minRev
"1.3.6.1.4.1.9.9.109.1.1.1.1.8.#{processor.processor_index}"
"ucd_cpu" ->
# UCD-SNMP-MIB - need to calculate from user/system/idle
# For simplicity, use ssCpuUser + ssCpuSystem as approximation
poll_ucd_cpu(processor, client_opts)
_ ->
@ -482,7 +393,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp poll_ucd_cpu(_processor, client_opts) do
# UCD-SNMP-MIB: ssCpuUser + ssCpuSystem = CPU usage
user_oid = "1.3.6.1.4.1.2021.11.9.0"
system_oid = "1.3.6.1.4.1.2021.11.10.0"
@ -496,7 +406,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp handle_processor_poll_result(processor, {:ok, load_percent}, timestamp) do
# Create a reading
Snmp.create_processor_reading(%{
processor_id: processor.id,
load_percent: load_percent,
@ -504,7 +413,6 @@ defmodule Towerops.Snmp.PollerWorker do
checked_at: timestamp
})
# Update processor with latest value
Snmp.update_processor(processor, %{
load_percent: load_percent,
last_checked_at: timestamp
@ -523,8 +431,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp poll_sensors(sensors, client_opts, timestamp) do
# Poll sensors in parallel with max_concurrency to avoid overwhelming the device
# Timeout must be longer than SNMP timeout (30s) to allow slow devices
sensors
|> Task.async_stream(
fn sensor ->
@ -539,7 +445,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp poll_state_sensors(state_sensors, client_opts, timestamp, device_id) do
# Poll state sensors in parallel
state_sensors
|> Task.async_stream(
fn state_sensor ->
@ -571,7 +476,6 @@ defmodule Towerops.Snmp.PollerWorker do
new_status = entity_state_to_status(new_state_value)
new_descr = entity_state_to_descr(new_state_value)
# Update the state sensor
Snmp.update_state_sensor(state_sensor, %{
state_value: new_state_value,
state_descr: new_descr,
@ -579,7 +483,6 @@ defmodule Towerops.Snmp.PollerWorker do
last_checked_at: timestamp
})
# Detect and broadcast state changes
if old_status != new_status do
broadcast_state_sensor_change(state_sensor, old_status, new_status, device_id, timestamp)
end
@ -613,10 +516,7 @@ defmodule Towerops.Snmp.PollerWorker do
occurred_at: timestamp
}
# Broadcast to device-specific topic
_ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:#{device_id}", {:device_event, event})
# Also broadcast to global events topic
_ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:events", {:device_event, event})
Logger.info("State sensor change: #{event.message}")
@ -633,8 +533,6 @@ defmodule Towerops.Snmp.PollerWorker do
defp determine_state_change_event_type("ok"), do: "state_sensor_normal"
defp determine_state_change_event_type(_), do: "state_sensor_unknown"
# ENTITY-STATE-MIB operational status mapping
# 1=unknown, 2=disabled, 3=enabled, 4=testing
defp entity_state_to_status(1), do: "unknown"
defp entity_state_to_status(2), do: "ok"
defp entity_state_to_status(3), do: "warning"
@ -715,7 +613,6 @@ defmodule Towerops.Snmp.PollerWorker do
decoded_value = decode_snmp_value(raw_value)
if is_number(decoded_value) do
# Calculate actual value using divisor
value = decoded_value / sensor.sensor_divisor
{:ok, value}
else
@ -728,7 +625,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp poll_percentage_sensor(sensor, client_opts) do
# For percentage sensors (like storage), fetch both used and size
size_oid = sensor.metadata["size_oid"]
used_oid = sensor.sensor_oid
@ -739,9 +635,6 @@ defmodule Towerops.Snmp.PollerWorker do
true <- size > 0 do
percentage = used / size * 100
# Validate percentage and warn if it exceeds 100%
# This can happen if the device reports incorrect SNMP values or if
# hrStorageAllocationUnits differs between used/size OIDs
capped_percentage =
if percentage > 100 do
Logger.warning(
@ -765,12 +658,9 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp poll_interfaces(interfaces, client_opts, timestamp) do
# Poll interfaces in parallel with max_concurrency
# Timeout must be longer than SNMP timeout (30s) to allow slow devices
interfaces
|> Task.async_stream(
fn interface ->
# Poll interface stats: ifInOctets, ifOutOctets, ifInErrors, ifOutErrors
{:ok, stat_data} = get_interface_stats(client_opts, interface.if_index)
stats =
@ -789,8 +679,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp check_interface_changes(interfaces, device, client_opts, timestamp) do
# Check interface changes in parallel
# Timeout must be longer than SNMP timeout (30s) to allow slow devices
interfaces
|> Task.async_stream(
fn interface ->
@ -811,13 +699,9 @@ defmodule Towerops.Snmp.PollerWorker do
defp get_interface_attributes(client_opts, if_index) do
oids = [
# ifSpeed
"1.3.6.1.2.1.2.2.1.5.#{if_index}",
# ifPhysAddress
"1.3.6.1.2.1.2.2.1.6.#{if_index}",
# ifAdminStatus
"1.3.6.1.2.1.2.2.1.7.#{if_index}",
# ifOperStatus
"1.3.6.1.2.1.2.2.1.8.#{if_index}"
]
@ -837,12 +721,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
end
@spec detect_and_log_changes(
Towerops.Snmp.Interface.t(),
map(),
Ecto.UUID.t(),
DateTime.t()
) :: :ok | nil
defp detect_and_log_changes(interface, current_attrs, device_id, timestamp) do
events =
[]
@ -996,7 +874,6 @@ defmodule Towerops.Snmp.PollerWorker do
defp broadcast_interface_events(events) do
Enum.each(events, fn {:event, event_attrs} ->
# Broadcast to device-specific topic for real-time updates
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
@ -1004,7 +881,6 @@ defmodule Towerops.Snmp.PollerWorker do
{:device_event, event_attrs}
)
# Also broadcast to generic events topic for global event logger
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
@ -1055,21 +931,16 @@ defmodule Towerops.Snmp.PollerWorker do
defp format_speed(_), do: "Unknown"
defp get_interface_stats(client_opts, if_index) do
# Try High-Capacity 64-bit counters first (ifHCInOctets, ifHCOutOctets)
# These are essential for 10G+ interfaces to avoid counter wrapping
# HC counters are from IF-MIB::ifXTable (RFC 2863)
hc_oids = [
if_in_octets: "1.3.6.1.2.1.31.1.1.1.6.#{if_index}",
if_out_octets: "1.3.6.1.2.1.31.1.1.1.10.#{if_index}"
]
# Standard 32-bit counters (fallback)
std_oids = [
if_in_octets: "1.3.6.1.2.1.2.2.1.10.#{if_index}",
if_out_octets: "1.3.6.1.2.1.2.2.1.16.#{if_index}"
]
# Error/discard counters (always 32-bit, no HC versions)
error_oids = [
if_in_errors: "1.3.6.1.2.1.2.2.1.14.#{if_index}",
if_out_errors: "1.3.6.1.2.1.2.2.1.20.#{if_index}",
@ -1077,10 +948,8 @@ defmodule Towerops.Snmp.PollerWorker do
if_out_discards: "1.3.6.1.2.1.2.2.1.19.#{if_index}"
]
# Try HC counters first for octet counts
octet_results = fetch_octet_counters(client_opts, hc_oids, std_oids)
# Fetch error/discard counters
error_results =
Enum.map(error_oids, fn {key, oid} ->
case Client.get(client_opts, oid) do
@ -1092,7 +961,6 @@ defmodule Towerops.Snmp.PollerWorker do
{:ok, Map.new(octet_results ++ error_results)}
end
# Try HC (64-bit) counters first, fall back to standard 32-bit if not available
defp fetch_octet_counters(client_opts, hc_oids, std_oids) do
hc_results = Enum.map(hc_oids, &try_fetch_hc_counter(client_opts, &1))
@ -1124,37 +992,29 @@ defmodule Towerops.Snmp.PollerWorker do
end
end
# Decode SNMP values, handling Counter64 and other binary types
# Note: Client.get already unwraps SNMPKit tuples via extract_snmp_value
defp decode_snmp_value(value) when is_number(value) do
value
end
# Handle binary Counter64 values (8 bytes, big-endian unsigned integer)
defp decode_snmp_value(value) when is_binary(value) do
size = byte_size(value)
result =
case size do
8 ->
# Standard Counter64 (8 bytes, big-endian unsigned integer)
<<counter::unsigned-big-integer-size(64)>> = value
counter
16 ->
# Some SNMP implementations return 16-byte values
# Try reading last 8 bytes as Counter64
<<_prefix::binary-size(8), counter::unsigned-big-integer-size(64)>> = value
counter
s when s > 8 ->
# Large binary, try reading last 8 bytes
offset = s - 8
<<_prefix::binary-size(^offset), counter::unsigned-big-integer-size(64)>> = value
counter
_ ->
# Unknown binary format or too small
Logger.warning("Unknown SNMP binary value format, size: #{size}")
nil
end
@ -1172,7 +1032,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp build_client_opts(device) do
# Get SNMP config with hierarchical fallback (device -> site -> organization)
snmp_config = Devices.get_snmp_config(device)
[
@ -1184,13 +1043,14 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp get_poll_interval(device) do
# Use check_interval_seconds, with configurable minimum for SNMP polling
min_interval = Application.get_env(:towerops, :snmp_min_poll_interval, 300)
max(device.check_interval_seconds || @default_poll_interval, min_interval)
end
defp schedule_next_poll(interval_seconds) do
Process.send_after(self(), :poll_data, interval_seconds * 1000)
defp schedule_next_poll(device_id, interval_seconds) do
%{device_id: device_id}
|> new(schedule_in: interval_seconds)
|> Oban.insert()
end
defp should_skip_poll?(device, poll_interval_seconds, grace_period_seconds) do
@ -1203,14 +1063,11 @@ defmodule Towerops.Snmp.PollerWorker do
seconds_since_poll = DateTime.diff(now, last_poll, :second)
min_interval = poll_interval_seconds - grace_period_seconds
# Skip if polled within the minimum interval
seconds_since_poll < min_interval
end
end
@spec detect_sensor_changes(Sensor.t(), float(), DateTime.t()) :: :ok
defp detect_sensor_changes(sensor, current_value, timestamp) do
# Skip change detection if there's no previous value to compare against
if sensor.last_value == nil do
:ok
else
@ -1226,12 +1083,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
end
@spec extract_thresholds(map()) :: %{
warning_high: float() | nil,
critical_high: float() | nil,
warning_low: float() | nil,
critical_low: float() | nil
}
defp extract_thresholds(metadata) do
%{
warning_high: metadata["warning_high"],
@ -1247,13 +1098,6 @@ defmodule Towerops.Snmp.PollerWorker do
if threshold_event, do: [threshold_event | events], else: events
end
@spec check_threshold_violation(
Sensor.t(),
float(),
DateTime.t(),
Ecto.UUID.t(),
map()
) :: map() | nil
defp check_threshold_violation(sensor, current_value, timestamp, device_id, thresholds) do
threshold_checks = [
&check_critical_high/5,
@ -1267,13 +1111,6 @@ defmodule Towerops.Snmp.PollerWorker do
end) || check_returned_to_normal(sensor, current_value, timestamp, device_id, thresholds)
end
@spec check_critical_high(
Sensor.t(),
float(),
DateTime.t(),
Ecto.UUID.t(),
map()
) :: map() | nil
defp check_critical_high(sensor, current_value, timestamp, device_id, thresholds) do
if thresholds.critical_high && current_value >= thresholds.critical_high do
build_threshold_event(
@ -1289,13 +1126,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
end
@spec check_critical_low(
Sensor.t(),
float(),
DateTime.t(),
Ecto.UUID.t(),
map()
) :: map() | nil
defp check_critical_low(sensor, current_value, timestamp, device_id, thresholds) do
if thresholds.critical_low && current_value <= thresholds.critical_low do
build_threshold_event(
@ -1311,13 +1141,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
end
@spec check_warning_high(
Sensor.t(),
float(),
DateTime.t(),
Ecto.UUID.t(),
map()
) :: map() | nil
defp check_warning_high(sensor, current_value, timestamp, device_id, thresholds) do
if thresholds.warning_high && current_value >= thresholds.warning_high do
build_threshold_event(
@ -1333,13 +1156,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
end
@spec check_warning_low(
Sensor.t(),
float(),
DateTime.t(),
Ecto.UUID.t(),
map()
) :: map() | nil
defp check_warning_low(sensor, current_value, timestamp, device_id, thresholds) do
if thresholds.warning_low && current_value <= thresholds.warning_low do
build_threshold_event(
@ -1415,7 +1231,6 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp maybe_add_change_event(events, sensor, current_value, timestamp, device_id) do
# Only check for spikes/drops on percentage sensors
if sensor.sensor_unit == "%" do
change_event = check_significant_change(sensor, current_value, timestamp, device_id)
if change_event, do: [change_event | events], else: events
@ -1428,7 +1243,6 @@ defmodule Towerops.Snmp.PollerWorker do
change_percent = abs(current_value - sensor.last_value)
cond do
# Spike: increase of 30% or more
change_percent >= 30 && current_value > sensor.last_value ->
build_sensor_event(
device_id,
@ -1440,7 +1254,6 @@ defmodule Towerops.Snmp.PollerWorker do
timestamp
)
# Drop: decrease of 30% or more
change_percent >= 30 && current_value < sensor.last_value ->
build_sensor_event(
device_id,
@ -1459,10 +1272,7 @@ defmodule Towerops.Snmp.PollerWorker do
defp broadcast_sensor_events(events) do
Enum.each(events, fn event ->
# Broadcast to device-specific topic for real-time updates
_ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:#{event.device_id}", {:device_event, event})
# Also broadcast to generic events topic for global event logger
_ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:events", {:device_event, event})
Logger.debug("Sensor event: #{event.message}")
@ -1488,8 +1298,7 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp get_device_id_from_sensor(sensor) do
# Get the device to find device_id
case Repo.get(Towerops.Snmp.Device, sensor.snmp_device_id) do
case Towerops.Repo.get(Towerops.Snmp.Device, sensor.snmp_device_id) do
nil -> nil
device -> device.device_id
end
@ -1500,35 +1309,4 @@ defmodule Towerops.Snmp.PollerWorker do
end
defp format_sensor_value(_, _), do: "N/A"
# Enqueue poll job - safe to call in test environment
defp enqueue_poll(device_id) do
if Application.get_env(:towerops, :env) == :test do
# In test, run synchronously with database sandbox access
start_poll_task(device_id)
else
# In dev/prod, enqueue to Exq
{:ok, _job} = Exq.enqueue(Exq, "polling", Towerops.Workers.PollWorker, [device_id])
end
end
defp start_poll_task(device_id) do
parent = self()
_ =
Task.start(fn ->
_ = maybe_allow_sandbox(parent)
perform_poll(device_id)
end)
end
defp maybe_allow_sandbox(parent) do
if Application.get_env(:towerops, :sql_sandbox) do
Sandbox.allow(Repo, parent, self())
end
end
defp via_tuple(device_id) do
{:via, Registry, {PollerRegistry, device_id}}
end
end

View file

@ -1,31 +1,21 @@
defmodule Towerops.Workers.DiscoveryWorker do
@moduledoc """
Background worker for SNMP device discovery.
Oban worker for SNMP device discovery.
Enqueued when:
- A new device is created with SNMP enabled
- An existing device's SNMP configuration is changed
- User manually triggers discovery via "Rediscover Device" button
Queue: discovery
"""
use Oban.Worker, queue: :discovery
alias Towerops.Devices
alias Towerops.Snmp
require Logger
@doc """
Performs SNMP discovery for a device.
## Parameters
- device_id: The UUID of the device to discover
## Returns
- :ok on success
- {:error, reason} on failure
"""
def perform(device_id) do
@impl Oban.Worker
def perform(%Oban.Job{args: %{"device_id" => device_id}}) do
Logger.info("Starting SNMP discovery for device #{device_id}")
device = Devices.get_device!(device_id)
@ -44,4 +34,13 @@ defmodule Towerops.Workers.DiscoveryWorker do
Logger.error("Device #{device_id} not found")
{:error, :device_not_found}
end
@doc """
Enqueues a discovery job for a device.
"""
def enqueue(device_id) do
%{device_id: device_id}
|> new()
|> Oban.insert()
end
end

View file

@ -1,96 +0,0 @@
defmodule Towerops.Workers.MonitorWorker do
@moduledoc """
Background worker for device monitoring checks (ICMP pings).
Enqueued when:
- trigger_check is called but the DeviceMonitor GenServer doesn't exist
- Manual monitoring check trigger from UI
Queue: monitoring
"""
alias Towerops.Devices
alias Towerops.Monitoring
require Logger
defp ping_adapter do
Application.get_env(:towerops, :ping_module, Towerops.Monitoring.Ping)
end
@doc """
Performs a monitoring check for a device.
## Parameters
- device_id: The UUID of the device to check
## Returns
- :ok on success
- {:error, reason} on failure
"""
def perform(device_id) do
Logger.info("Starting monitoring check for device #{device_id}")
device = Devices.get_device!(device_id)
if device.monitoring_enabled do
case ping_device(device) do
{:ok, latency} ->
Logger.debug("Device #{device_id} is up, latency: #{latency}ms")
Monitoring.create_check(%{
device_id: device_id,
status: :success,
response_time_ms: latency,
checked_at: DateTime.utc_now()
})
# Broadcast monitoring check update to device-specific topic
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device_id}",
{:monitoring_check_updated, device_id}
)
:ok
{:error, reason} ->
Logger.warning("Device #{device_id} is down: #{inspect(reason)}")
Monitoring.create_check(%{
device_id: device_id,
status: :failure,
response_time_ms: nil,
checked_at: DateTime.utc_now()
})
# Broadcast monitoring check update to device-specific topic
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device_id}",
{:monitoring_check_updated, device_id}
)
:ok
end
else
Logger.debug("Device #{device_id} does not have monitoring enabled, skipping check")
:ok
end
rescue
Ecto.NoResultsError ->
Logger.error("Device #{device_id} not found")
{:error, :device_not_found}
error ->
Logger.error("Monitoring check failed for device #{device_id}: #{inspect(error)}")
{:error, error}
end
defp ping_device(device) do
# Use configured ping adapter (mockable in tests)
ping_adapter().ping(device.ip_address, 1000)
end
end

View file

@ -1,177 +0,0 @@
defmodule Towerops.Workers.PollWorker do
@moduledoc """
Background worker for SNMP polling.
Enqueued when:
- trigger_poll is called but the PollerWorker GenServer doesn't exist
- Manual poll trigger from UI
Queue: polling
"""
alias Towerops.Devices
alias Towerops.Snmp
alias Towerops.Snmp.Client
alias Towerops.Snmp.NeighborDiscovery
require Logger
@doc """
Performs SNMP polling for a device.
## Parameters
- device_id: The UUID of the device to poll
## Returns
- :ok on success
- {:error, reason} on failure
"""
def perform(device_id) do
Logger.info("Starting SNMP poll for device #{device_id}")
device = Devices.get_device!(device_id)
if device.snmp_enabled do
poll_device(device_id, device)
:ok
else
Logger.debug("Device #{device_id} does not have SNMP enabled, skipping poll")
:ok
end
rescue
Ecto.NoResultsError ->
Logger.error("Device #{device_id} not found")
{:error, :device_not_found}
error ->
Logger.error("SNMP poll failed for device #{device_id}: #{inspect(error)}")
{:error, error}
end
defp poll_device(device_id, device) do
snmp_device = Snmp.get_device_with_associations(device_id)
if snmp_device do
perform_device_data_collection(device, snmp_device)
else
Logger.debug("No SNMP device found for #{device_id}, skipping poll")
end
end
defp perform_device_data_collection(device, snmp_device) do
Devices.update_snmp_poll_time(device)
client_opts = build_client_opts(device)
now = DateTime.truncate(DateTime.utc_now(), :second)
poll_sensors(snmp_device.sensors, client_opts, now)
poll_interfaces(snmp_device.interfaces, client_opts, now)
poll_neighbors(device, snmp_device, client_opts)
Logger.info("SNMP poll completed for device #{device.id}")
end
defp poll_sensors(sensors, client_opts, timestamp) do
Enum.each(sensors, fn sensor ->
case poll_sensor(sensor, client_opts) do
{:ok, value} ->
Snmp.create_sensor_reading(%{
sensor_id: sensor.id,
value: value,
status: "ok",
checked_at: timestamp
})
Snmp.update_sensor(sensor, %{
last_value: value,
last_checked_at: timestamp
})
{:error, _reason} ->
Snmp.create_sensor_reading(%{
sensor_id: sensor.id,
value: nil,
status: "error",
checked_at: timestamp
})
end
end)
end
defp poll_sensor(sensor, client_opts) do
case Client.get(client_opts, sensor.sensor_oid) do
{:ok, raw_value} when is_number(raw_value) ->
value = raw_value / sensor.sensor_divisor
{:ok, value}
{:error, reason} ->
{:error, reason}
_ ->
{:error, :non_numeric}
end
end
defp poll_interfaces(interfaces, client_opts, timestamp) do
Enum.each(interfaces, fn interface ->
{:ok, stat_data} = get_interface_stats(client_opts, interface.if_index)
stats =
Map.merge(stat_data, %{
interface_id: interface.id,
checked_at: timestamp
})
Snmp.create_interface_stat(stats)
end)
end
defp get_interface_stats(client_opts, if_index) do
oids = [
if_in_octets: "1.3.6.1.2.1.2.2.1.10.#{if_index}",
if_out_octets: "1.3.6.1.2.1.2.2.1.16.#{if_index}",
if_in_errors: "1.3.6.1.2.1.2.2.1.14.#{if_index}",
if_out_errors: "1.3.6.1.2.1.2.2.1.20.#{if_index}",
if_in_discards: "1.3.6.1.2.1.2.2.1.13.#{if_index}",
if_out_discards: "1.3.6.1.2.1.2.2.1.19.#{if_index}"
]
results =
Enum.map(oids, fn {key, oid} ->
case Client.get(client_opts, oid) do
{:ok, value} when is_number(value) -> {key, value}
_ -> {key, nil}
end
end)
{:ok, Map.new(results)}
end
defp poll_neighbors(device, snmp_device, client_opts) do
interfaces_with_device = Enum.map(snmp_device.interfaces, &Map.put(&1, :device_id, snmp_device.id))
{:ok, neighbors} = NeighborDiscovery.discover_neighbors(client_opts, interfaces_with_device)
cutoff = DateTime.add(DateTime.utc_now(), -5, :minute)
Snmp.delete_stale_neighbors(device.id, cutoff)
Enum.each(neighbors, fn neighbor_data ->
Snmp.upsert_neighbor(neighbor_data)
end)
rescue
error ->
Logger.error("Error polling neighbors: #{inspect(error)}")
end
defp build_client_opts(device) do
snmp_config = Devices.get_snmp_config(device)
[
ip: device.ip_address,
community: snmp_config.community,
version: snmp_config.version,
port: device.snmp_port || 161,
timeout: 5000
]
end
end

View file

@ -427,8 +427,8 @@ defmodule ToweropsWeb.AgentChannel do
# In test, run synchronously
Task.start(fn -> Discovery.discover_device(Devices.get_device!(device_id)) end)
else
# In dev/prod, enqueue to Exq
{:ok, _job} = Exq.enqueue(Exq, "discovery", Towerops.Workers.DiscoveryWorker, [device_id])
# In dev/prod, enqueue to Oban
Towerops.Workers.DiscoveryWorker.enqueue(device_id)
end
end

View file

@ -340,8 +340,8 @@ defmodule ToweropsWeb.DeviceLive.Form do
# In test, skip discovery (tests can call Snmp.discover_device directly if needed)
:ok
else
# In dev/prod, enqueue to Exq
{:ok, _job} = Exq.enqueue(Exq, "discovery", DiscoveryWorker, [device_id])
# In dev/prod, enqueue to Oban
DiscoveryWorker.enqueue(device_id)
end
end

View file

@ -51,7 +51,7 @@ defmodule ToweropsWeb.DeviceLive.Index do
if Application.get_env(:towerops, :env) == :test do
_ = Task.start(fn -> Snmp.discover_device(Devices.get_device!(device_id)) end)
else
{:ok, _job} = Exq.enqueue(Exq, "discovery", DiscoveryWorker, [device_id])
DiscoveryWorker.enqueue(device_id)
end
end
end

View file

@ -94,8 +94,8 @@ defmodule ToweropsWeb.SiteLive.Show do
# In test, run synchronously
_ = Task.start(fn -> Snmp.discover_device(Devices.get_device!(device_id)) end)
else
# In dev/prod, enqueue to Exq
{:ok, _job} = Exq.enqueue(Exq, "discovery", DiscoveryWorker, [device_id])
# In dev/prod, enqueue to Oban
DiscoveryWorker.enqueue(device_id)
end
end
end

View file

@ -5,7 +5,7 @@ defmodule Towerops.MixProject do
[
app: :towerops,
version: "0.1.0",
elixir: "~> 1.15",
elixir: "~> 1.18",
elixirc_paths: elixirc_paths(Mix.env()),
start_permanent: Mix.env() == :prod,
aliases: aliases(),
@ -22,8 +22,7 @@ defmodule Towerops.MixProject do
def application do
[
mod: {Towerops.Application, []},
extra_applications: [:logger, :runtime_tools, :os_mon],
included_applications: [:exq]
extra_applications: [:logger, :runtime_tools, :os_mon]
]
end
@ -61,7 +60,6 @@ defmodule Towerops.MixProject do
{:cbor, "~> 1.0"},
{:protobuf, "~> 0.12"},
{:req, "~> 0.5"},
# snmpkit is now vendored in lib/snmpkit
{:yaml_elixir, "~> 2.9"},
{:telemetry_metrics, "~> 1.0"},
{:telemetry_poller, "~> 1.0"},
@ -73,7 +71,6 @@ defmodule Towerops.MixProject do
{:bandit, "~> 1.5"},
{:phoenix_pubsub_redis, "~> 3.0"},
{:ecto_psql_extras, "~> 0.6"},
{:exq, "~> 0.19"},
{:mox, "~> 1.0", only: :test},
{:honeybadger, "~> 0.24"},
{:stream_data, "~> 1.1", only: :test},
@ -90,7 +87,7 @@ defmodule Towerops.MixProject do
defp dialyzer do
[
plt_file: {:no_warn, "priv/plts/dialyzer.plt"},
plt_add_apps: [:mix, :ex_unit, :exq],
plt_add_apps: [:mix, :ex_unit],
flags: [:unmatched_returns, :error_handling, :unknown],
ignore_warnings: ".dialyzer_ignore.exs"
]

View file

@ -14,11 +14,9 @@
"ecto_psql_extras": {:hex, :ecto_psql_extras, "0.8.8", "aa02529c97f69aed5722899f5dc6360128735a92dd169f23c5d50b1f7fdede08", [:mix], [{:ecto_sql, "~> 3.7", [hex: :ecto_sql, repo: "hexpm", optional: false]}, {:postgrex, "> 0.16.0", [hex: :postgrex, repo: "hexpm", optional: false]}, {:table_rex, "~> 3.1.1 or ~> 4.0", [hex: :table_rex, repo: "hexpm", optional: false]}], "hexpm", "04c63d92b141723ad6fed2e60a4b461ca00b3594d16df47bbc48f1f4534f2c49"},
"ecto_sql": {:hex, :ecto_sql, "3.13.4", "b6e9d07557ddba62508a9ce4a484989a5bb5e9a048ae0e695f6d93f095c25d60", [:mix], [{:db_connection, "~> 2.4.1 or ~> 2.5", [hex: :db_connection, repo: "hexpm", optional: false]}, {:ecto, "~> 3.13.0", [hex: :ecto, repo: "hexpm", optional: false]}, {:myxql, "~> 0.7", [hex: :myxql, repo: "hexpm", optional: true]}, {:postgrex, "~> 0.19 or ~> 1.0", [hex: :postgrex, repo: "hexpm", optional: true]}, {:tds, "~> 2.1.1 or ~> 2.2", [hex: :tds, repo: "hexpm", optional: true]}, {:telemetry, "~> 0.4.0 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "2b38cf0749ca4d1c5a8bcbff79bbe15446861ca12a61f9fba604486cb6b62a14"},
"elixir_make": {:hex, :elixir_make, "0.9.0", "6484b3cd8c0cee58f09f05ecaf1a140a8c97670671a6a0e7ab4dc326c3109726", [:mix], [], "hexpm", "db23d4fd8b757462ad02f8aa73431a426fe6671c80b200d9710caf3d1dd0ffdb"},
"elixir_uuid": {:hex, :elixir_uuid, "1.2.1", "dce506597acb7e6b0daeaff52ff6a9043f5919a4c3315abb4143f0b00378c097", [:mix], [], "hexpm", "f7eba2ea6c3555cea09706492716b0d87397b88946e6380898c2889d68585752"},
"erlex": {:hex, :erlex, "0.2.8", "cd8116f20f3c0afe376d1e8d1f0ae2452337729f68be016ea544a72f767d9c12", [:mix], [], "hexpm", "9d66ff9fedf69e49dc3fd12831e12a8a37b76f8651dd21cd45fcf5561a8a7590"},
"esbuild": {:hex, :esbuild, "0.10.0", "b0aa3388a1c23e727c5a3e7427c932d89ee791746b0081bbe56103e9ef3d291f", [:mix], [{:jason, "~> 1.4", [hex: :jason, repo: "hexpm", optional: false]}], "hexpm", "468489cda427b974a7cc9f03ace55368a83e1a7be12fba7e30969af78e5f8c70"},
"expo": {:hex, :expo, "1.1.1", "4202e1d2ca6e2b3b63e02f69cfe0a404f77702b041d02b58597c00992b601db5", [:mix], [], "hexpm", "5fb308b9cb359ae200b7e23d37c76978673aa1b06e2b3075d814ce12c5811640"},
"exq": {:hex, :exq, "0.23.0", "143079d49cec13673849704e9407521047f4016c72117d0067405115a65247da", [:mix], [{:elixir_uuid, ">= 1.2.0", [hex: :elixir_uuid, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:poison, ">= 1.2.0 and < 7.0.0", [hex: :poison, repo: "hexpm", optional: true]}, {:redix, ">= 0.9.0", [hex: :redix, repo: "hexpm", optional: false]}], "hexpm", "38bc3216051d8278aa32b147f1eca6a74dd84283b7785834c7a3609c00f5114a"},
"file_system": {:hex, :file_system, "1.1.1", "31864f4685b0148f25bd3fbef2b1228457c0c89024ad67f7a81a3ffbc0bbad3a", [:mix], [], "hexpm", "7a15ff97dfe526aeefb090a7a9d3d03aa907e100e262a0f8f7746b78f8f87a5d"},
"finch": {:hex, :finch, "0.20.0", "5330aefb6b010f424dcbbc4615d914e9e3deae40095e73ab0c1bb0968933cadf", [:mix], [{:mime, "~> 1.0 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:mint, "~> 1.6.2 or ~> 1.7", [hex: :mint, repo: "hexpm", optional: false]}, {:nimble_options, "~> 0.4 or ~> 1.0", [hex: :nimble_options, repo: "hexpm", optional: false]}, {:nimble_pool, "~> 1.1", [hex: :nimble_pool, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "2658131a74d051aabfcba936093c903b8e89da9a1b63e430bee62045fa9b2ee2"},
"fine": {:hex, :fine, "0.1.4", "b19a89c1476c7c57afb5f9314aed5960b5bc95d5277de4cb5ee8e1d1616ce379", [:mix], [], "hexpm", "be3324cc454a42d80951cf6023b9954e9ff27c6daa255483b3e8d608670303f5"},

View file

@ -11,9 +11,9 @@ defmodule Towerops.Integration.SnmpIntegrationTest do
"""
use ExUnit.Case, async: false
alias Towerops.Monitoring.DeviceMonitor
alias Towerops.Snmp.Client
alias Towerops.Snmp.Poller
alias Towerops.Workers.DeviceMonitorWorker
@moduletag :integration
@ -126,13 +126,8 @@ defmodule Towerops.Integration.SnmpIntegrationTest do
# Set to up first so status change will trigger
{:ok, device} = Towerops.Devices.update_device_status(device, :up)
{:ok, state} = DeviceMonitor.init(device.id)
# Trigger check - will timeout trying to reach 192.0.2.1
{:noreply, _state} =
DeviceMonitor.handle_info(:check_device, state)
Process.sleep(200)
DeviceMonitorWorker.perform(%Oban.Job{args: %{"device_id" => device.id}})
alerts = Towerops.Alerts.list_devices_alerts(device.id)
refute Enum.empty?(alerts)

View file

@ -1,95 +0,0 @@
defmodule Towerops.ExqSupervisorTest do
use ExUnit.Case, async: false
alias Towerops.ExqSupervisor
describe "init/1" do
@tag :skip
test "starts Exq when Redis is available" do
# This test is skipped because it would interfere with the running application
# In a real test environment, you would:
# 1. Stop the main ExqSupervisor
# 2. Start a test instance with a different name
# 3. Verify Exq child process starts
# 4. Clean up
assert true
end
@tag :skip
test "handles Redis unavailability gracefully" do
# This test is skipped because it would require mocking Redis connection
# In a real test environment, you would:
# 1. Configure test Redis to an unreachable host
# 2. Verify supervisor starts without crashing
# 3. Verify Exq is not started
# 4. Verify appropriate error logs
assert true
end
end
describe "supervisor behavior" do
test "module defines start_link/1" do
# Verify the module has been compiled and loaded
Code.ensure_loaded!(ExqSupervisor)
assert function_exported?(ExqSupervisor, :start_link, 1)
end
test "module uses Supervisor behavior" do
Code.ensure_loaded!(ExqSupervisor)
behaviours =
:attributes
|> ExqSupervisor.module_info()
|> Keyword.get(:behaviour, [])
assert Supervisor in behaviours
end
test "init/1 callback is defined" do
Code.ensure_loaded!(ExqSupervisor)
# The init/1 callback is defined by the module
assert function_exported?(ExqSupervisor, :init, 1)
end
end
describe "configuration" do
test "reads Redis config from application environment" do
redis_config = Application.get_env(:towerops, :redis, [])
assert Keyword.has_key?(redis_config, :host) || redis_config == []
assert Keyword.has_key?(redis_config, :port) || redis_config == []
end
test "handles missing Redis config" do
# Should default to localhost:6379 if config is missing
original_config = Application.get_env(:towerops, :redis)
try do
Application.delete_env(:towerops, :redis)
# Verify the module can handle missing config
# (We can't actually test init without starting processes)
assert Application.get_env(:towerops, :redis) == nil
after
if original_config do
Application.put_env(:towerops, :redis, original_config)
end
end
end
end
describe "restart strategy" do
test "uses one_for_one strategy" do
# The supervisor uses one_for_one strategy
# This means if Exq crashes, only Exq will be restarted
# Other children (if any were added) would continue running
assert true
end
test "has limited restart policy" do
# The supervisor has max_restarts: 3, max_seconds: 60
# This prevents crash loops from overwhelming the system
assert true
end
end
end

View file

@ -1,457 +0,0 @@
defmodule Towerops.Monitoring.DeviceMonitorTest do
use Towerops.DataCase, async: false
import Mox
import Towerops.AccountsFixtures
alias Towerops.Alerts
alias Towerops.Devices
alias Towerops.Monitoring
alias Towerops.Monitoring.DeviceMonitor
alias Towerops.Monitoring.PingMock
setup :verify_on_exit!
setup :set_mox_global
setup do
# Use stub_with for GenServer tests since monitors may ping multiple times
stub_with(PingMock, Towerops.Monitoring.PingStub)
# Start the Monitoring.Supervisor for these tests
start_supervised!(Towerops.Monitoring.Supervisor)
# Create test data
user = user_fixture()
{:ok, organization} = Towerops.Organizations.create_organization(%{name: "Test Org"}, user.id)
{:ok, site} =
Towerops.Sites.create_site(%{
name: "Test Site",
organization_id: organization.id
})
{:ok, device} =
Devices.create_device(%{
name: "Test Router",
ip_address: "192.168.1.1",
site_id: site.id,
monitoring_enabled: true,
check_interval_seconds: 30
})
%{device: device, site: site, organization: organization}
end
describe "start_link/1" do
test "starts a monitor GenServer for a device", %{device: device} do
assert {:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
assert is_pid(pid)
assert Process.alive?(pid)
# Cleanup
GenServer.stop(pid)
end
end
describe "init/1" do
test "performs immediate check when monitoring is enabled", %{device: device} do
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
# Give it time to perform the initial check
Process.sleep(50)
# Verify check was created
checks = Monitoring.list_devices_checks(device.id)
assert checks != [], "Expected at least 1 check"
assert hd(checks).status == :success
assert hd(checks).response_time_ms == 10
# Cleanup
GenServer.stop(pid)
end
test "does not check when monitoring is disabled", %{device: device} do
# Disable monitoring
{:ok, device} = Devices.update_device(device, %{monitoring_enabled: false})
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
# Verify no checks were created
checks = Monitoring.list_devices_checks(device.id)
assert checks == []
# Cleanup
GenServer.stop(pid)
end
end
describe "trigger_check/1" do
test "triggers immediate check when monitor is running", %{device: device} do
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
# Count initial checks
initial_count = length(Monitoring.list_devices_checks(device.id))
DeviceMonitor.trigger_check(device.id)
Process.sleep(50)
# Should have at least one more check
final_count = length(Monitoring.list_devices_checks(device.id))
assert final_count > initial_count
# Cleanup
GenServer.stop(pid)
end
test "enqueues check job when monitor is not running", %{device: device} do
# Disable monitoring so monitor doesn't start
{:ok, device} = Devices.update_device(device, %{monitoring_enabled: false})
DeviceMonitor.trigger_check(device.id)
Process.sleep(100)
# Verify check was created
checks = Monitoring.list_devices_checks(device.id)
assert length(checks) == 1
assert hd(checks).response_time_ms == 10
end
end
describe "perform_check/1 - successful ping" do
test "creates successful monitoring check", %{device: device} do
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
check = Monitoring.get_latest_check(device.id)
assert check.status == :success
assert check.response_time_ms == 10
assert check.device_id == device.id
# Cleanup
GenServer.stop(pid)
end
test "updates device status to :up", %{device: device} do
# Start with device down
{:ok, device} = Devices.update_device_status(device, :down)
assert device.status == :down
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
# Device should be up now
device = Devices.get_device!(device.id)
assert device.status == :up
# Cleanup
GenServer.stop(pid)
end
test "broadcasts device status change via PubSub", %{device: device} do
# Subscribe to device status changes
Phoenix.PubSub.subscribe(Towerops.PubSub, "device:#{device.id}")
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
# Wait for broadcast
assert_receive {:device_status_changed, device_id, :up, nil}, 200
assert device_id == device.id
# Cleanup
GenServer.stop(pid)
end
end
describe "perform_check/1 - failed ping" do
setup do
# Override stub to return errors for this describe block
stub(PingMock, :ping, fn _ -> {:error, :timeout} end)
:ok
end
test "creates failed monitoring check", %{device: device} do
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
check = Monitoring.get_latest_check(device.id)
assert check.status == :failure
assert is_nil(check.response_time_ms)
# Cleanup
GenServer.stop(pid)
end
test "updates device status to :down", %{device: device} do
# Start with device up
{:ok, device} = Devices.update_device_status(device, :up)
assert device.status == :up
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
# Device should be down now
device = Devices.get_device!(device.id)
assert device.status == :down
# Cleanup
GenServer.stop(pid)
end
end
describe "status change - up to down" do
setup do
stub(PingMock, :ping, fn _ -> {:error, :timeout} end)
:ok
end
test "creates device_down alert", %{device: device} do
# Start with device up
{:ok, device} = Devices.update_device_status(device, :up)
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
# Check alert was created
assert Alerts.has_active_alert?(device.id, :device_down)
alert = Alerts.get_active_alert(device.id, :device_down)
assert alert.alert_type == :device_down
assert alert.message =~ "not responding"
# Cleanup
GenServer.stop(pid)
end
test "broadcasts new_alert via PubSub", %{device: device} do
{:ok, _device} = Devices.update_device_status(device, :up)
# Subscribe to alerts
Phoenix.PubSub.subscribe(Towerops.PubSub, "alerts:new")
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
# Wait for broadcast
assert_receive {:new_alert, device_id, :device_down}, 200
assert device_id == device.id
# Cleanup
GenServer.stop(pid)
end
test "does not create duplicate alerts when already down", %{device: device} do
{:ok, device} = Devices.update_device_status(device, :up)
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
# Second check - should not create duplicate
GenServer.cast(pid, :check_now)
Process.sleep(50)
# Only one active alert should exist
alerts = Alerts.list_devices_alerts(device.id)
down_alerts = Enum.filter(alerts, &(&1.alert_type == :device_down and is_nil(&1.resolved_at)))
assert length(down_alerts) == 1
# Cleanup
GenServer.stop(pid)
end
end
describe "status change - down to up" do
test "creates device_up alert", %{device: device} do
# Start with device down and create down alert
{:ok, device} = Devices.update_device_status(device, :down)
{:ok, _alert} =
Alerts.create_alert(%{
device_id: device.id,
alert_type: :device_down,
triggered_at: DateTime.utc_now(),
message: "Device is down"
})
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
# Check recovery alert was created
alerts = Alerts.list_devices_alerts(device.id)
up_alerts = Enum.filter(alerts, &(&1.alert_type == :device_up))
assert up_alerts != [], "Expected at least 1 up alert"
up_alert = hd(up_alerts)
assert up_alert.message =~ "now responding"
# Cleanup
GenServer.stop(pid)
end
test "resolves existing device_down alert", %{device: device} do
{:ok, device} = Devices.update_device_status(device, :down)
# Create down alert
{:ok, down_alert} =
Alerts.create_alert(%{
device_id: device.id,
alert_type: :device_down,
triggered_at: DateTime.utc_now(),
message: "Device is down"
})
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
# Down alert should be resolved
refute Alerts.has_active_alert?(device.id, :device_down)
resolved_alert = Alerts.get_alert!(down_alert.id)
assert resolved_alert.resolved_at
# Cleanup
GenServer.stop(pid)
end
test "broadcasts alert_resolved via PubSub", %{device: device} do
{:ok, _device} = Devices.update_device_status(device, :down)
# Create existing down alert
{:ok, _alert} =
Alerts.create_alert(%{
device_id: device.id,
alert_type: :device_down,
triggered_at: DateTime.utc_now(),
message: "Device is down"
})
# Subscribe to resolved alerts
Phoenix.PubSub.subscribe(Towerops.PubSub, "alerts:resolved")
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
# Wait for broadcast
assert_receive {:alert_resolved, device_id, :device_down}, 200
assert device_id == device.id
# Cleanup
GenServer.stop(pid)
end
end
describe "alert messages" do
setup do
stub(PingMock, :ping, fn _ -> {:error, :timeout} end)
:ok
end
test "uses ICMP message when SNMP is disabled", %{device: device} do
{:ok, device} = Devices.update_device(device, %{snmp_enabled: false})
{:ok, _device} = Devices.update_device_status(device, :up)
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
alert = Alerts.get_active_alert(device.id, :device_down)
assert alert.message == "Device is not responding to ping"
# Cleanup
GenServer.stop(pid)
end
test "uses SNMP message when SNMP is enabled", %{device: device} do
{:ok, device} = Devices.update_device(device, %{snmp_enabled: true})
{:ok, _device} = Devices.update_device_status(device, :up)
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
alert = Alerts.get_active_alert(device.id, :device_down)
assert alert.message == "Device is not responding to SNMP"
# Cleanup
GenServer.stop(pid)
end
end
describe "alert messages - recovery" do
test "recovery message mentions ICMP when SNMP disabled", %{device: device} do
{:ok, device} = Devices.update_device(device, %{snmp_enabled: false})
{:ok, device} = Devices.update_device_status(device, :down)
{:ok, _alert} =
Alerts.create_alert(%{
device_id: device.id,
alert_type: :device_down,
triggered_at: DateTime.utc_now(),
message: "Device is down"
})
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
alerts = Alerts.list_devices_alerts(device.id)
up_alert = Enum.find(alerts, &(&1.alert_type == :device_up))
assert up_alert.message == "Device is now responding to ping"
# Cleanup
GenServer.stop(pid)
end
test "recovery message mentions SNMP when SNMP enabled", %{device: device} do
{:ok, device} = Devices.update_device(device, %{snmp_enabled: true})
{:ok, device} = Devices.update_device_status(device, :down)
{:ok, _alert} =
Alerts.create_alert(%{
device_id: device.id,
alert_type: :device_down,
triggered_at: DateTime.utc_now(),
message: "Device is down"
})
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
alerts = Alerts.list_devices_alerts(device.id)
up_alert = Enum.find(alerts, &(&1.alert_type == :device_up))
assert up_alert.message == "Device is now responding to SNMP"
# Cleanup
GenServer.stop(pid)
end
end
describe "scheduling" do
test "schedules next check based on check_interval_seconds", %{device: device} do
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
# Wait for initial check
Process.sleep(50)
initial_checks = length(Monitoring.list_devices_checks(device.id, 10))
# Manually trigger the next check instead of waiting for timer
# This verifies the scheduling mechanism works without slow waits
send(pid, :check_device)
Process.sleep(50)
# Should have more checks
checks = Monitoring.list_devices_checks(device.id, 10)
assert length(checks) > initial_checks
# Cleanup
GenServer.stop(pid)
end
test "does not schedule next check when monitoring disabled", %{device: device} do
{:ok, device} = Devices.update_device(device, %{monitoring_enabled: false})
{:ok, pid} = DeviceMonitor.start_link(device_id: device.id)
Process.sleep(50)
# No checks should exist
checks = Monitoring.list_devices_checks(device.id)
assert checks == []
# Cleanup
GenServer.stop(pid)
end
end
end

View file

@ -1,155 +0,0 @@
defmodule Towerops.Monitoring.SupervisorTest do
use Towerops.DataCase, async: false
import Mox
import Towerops.AccountsFixtures
alias Towerops.Monitoring.PingMock
alias Towerops.Monitoring.Supervisor, as: MonitoringSupervisor
setup :verify_on_exit!
setup do
# Stub ping calls for any monitors that start during tests
# Allow any number of ping calls (monitors ping on init)
stub_with(PingMock, Towerops.Monitoring.PingStub)
# Start Monitoring.Supervisor for these tests
start_supervised!(MonitoringSupervisor)
user = user_fixture()
{:ok, organization} = Towerops.Organizations.create_organization(%{name: "Test Org"}, user.id)
{:ok, site} =
Towerops.Sites.create_site(%{
name: "Test Site",
organization_id: organization.id
})
{:ok, device} =
Towerops.Devices.create_device(%{
name: "Router 1",
ip_address: "192.168.1.1",
site_id: site.id,
monitoring_enabled: true
})
%{device: device, site: site}
end
describe "start_monitor/1" do
test "starts a monitor for device", %{device: device} do
assert {:ok, pid} = MonitoringSupervisor.start_monitor(device.id)
assert is_pid(pid)
assert Process.alive?(pid)
# Cleanup
MonitoringSupervisor.stop_monitor(device.id)
end
test "returns error when monitor already running", %{device: device} do
{:ok, _pid} = MonitoringSupervisor.start_monitor(device.id)
assert {:error, {:already_started, _pid}} =
MonitoringSupervisor.start_monitor(device.id)
# Cleanup
MonitoringSupervisor.stop_monitor(device.id)
end
end
describe "stop_monitor/1" do
test "stops a running monitor", %{device: device} do
{:ok, pid} = MonitoringSupervisor.start_monitor(device.id)
assert Process.alive?(pid)
assert :ok = MonitoringSupervisor.stop_monitor(device.id)
# Give it a moment to terminate
Process.sleep(10)
refute Process.alive?(pid)
end
test "returns :ok when monitor not running", %{device: device} do
assert :ok = MonitoringSupervisor.stop_monitor(device.id)
end
end
describe "start_snmp_poller/1" do
test "starts an SNMP poller for device", %{device: device} do
# Enable SNMP for the equipment
{:ok, device} =
Towerops.Devices.update_device(device, %{
snmp_enabled: true,
snmp_community: "public"
})
assert {:ok, pid} = MonitoringSupervisor.start_snmp_poller(device.id)
assert is_pid(pid)
assert Process.alive?(pid)
# Cleanup
MonitoringSupervisor.stop_snmp_poller(device.id)
end
test "returns error when poller already running", %{device: device} do
{:ok, device} =
Towerops.Devices.update_device(device, %{
snmp_enabled: true,
snmp_community: "public"
})
{:ok, _pid} = MonitoringSupervisor.start_snmp_poller(device.id)
assert {:error, {:already_started, _pid}} =
MonitoringSupervisor.start_snmp_poller(device.id)
# Cleanup
MonitoringSupervisor.stop_snmp_poller(device.id)
end
end
describe "stop_snmp_poller/1" do
test "stops a running SNMP poller", %{device: device} do
{:ok, device} =
Towerops.Devices.update_device(device, %{
snmp_enabled: true,
snmp_community: "public"
})
{:ok, pid} = MonitoringSupervisor.start_snmp_poller(device.id)
assert Process.alive?(pid)
assert :ok = MonitoringSupervisor.stop_snmp_poller(device.id)
# Give it a moment to terminate
Process.sleep(10)
refute Process.alive?(pid)
end
test "returns :ok when poller not running", %{device: device} do
assert :ok = MonitoringSupervisor.stop_snmp_poller(device.id)
end
end
describe "concurrent operations" do
test "handles multiple concurrent start requests", %{device: device} do
# Start same monitor concurrently
results =
1..5
|> Enum.map(fn _ ->
Task.async(fn -> MonitoringSupervisor.start_monitor(device.id) end)
end)
|> Enum.map(&Task.await/1)
# One should succeed, others should get already_started
successful = Enum.count(results, &match?({:ok, _}, &1))
already_started = Enum.count(results, &match?({:error, {:already_started, _}}, &1))
assert successful == 1
assert already_started == 4
# Cleanup
MonitoringSupervisor.stop_monitor(device.id)
end
end
end

View file

@ -61,7 +61,7 @@ defmodule Towerops.Workers.DiscoveryWorkerTest do
{:ok, []}
end)
assert :ok = DiscoveryWorker.perform(device.id)
assert :ok = DiscoveryWorker.perform(%Oban.Job{args: %{"device_id" => device.id}})
# Verify SNMP device was created
snmp_device = Snmp.get_device(device.id)
@ -73,7 +73,8 @@ defmodule Towerops.Workers.DiscoveryWorkerTest do
test "returns error when device not found" do
non_existent_id = Ecto.UUID.generate()
assert {:error, :device_not_found} = DiscoveryWorker.perform(non_existent_id)
assert {:error, :device_not_found} =
DiscoveryWorker.perform(%Oban.Job{args: %{"device_id" => non_existent_id}})
end
test "returns error when discovery fails", %{device: device} do
@ -82,7 +83,8 @@ defmodule Towerops.Workers.DiscoveryWorkerTest do
{:error, :timeout}
end)
assert {:error, _reason} = DiscoveryWorker.perform(device.id)
assert {:error, _reason} =
DiscoveryWorker.perform(%Oban.Job{args: %{"device_id" => device.id}})
end
test "successfully completes discovery with proper mocks", %{device: device} do
@ -106,7 +108,7 @@ defmodule Towerops.Workers.DiscoveryWorkerTest do
# Allow any number of walk calls (interface + neighbor discovery)
stub(SnmpMock, :walk, fn _target, _oid, _opts -> {:ok, []} end)
assert :ok = DiscoveryWorker.perform(device.id)
assert :ok = DiscoveryWorker.perform(%Oban.Job{args: %{"device_id" => device.id}})
# Verify SNMP device was created
snmp_device = Snmp.get_device(device.id)
@ -122,7 +124,7 @@ defmodule Towerops.Workers.DiscoveryWorkerTest do
log =
capture_log(fn ->
DiscoveryWorker.perform(device.id)
DiscoveryWorker.perform(%Oban.Job{args: %{"device_id" => device.id}})
end)
assert log =~ "SNMP discovery failed for device #{device.id}"

View file

@ -1,154 +0,0 @@
defmodule Towerops.Workers.MonitorWorkerTest do
use Towerops.DataCase, async: true
import Mox
import Towerops.AccountsFixtures
alias Towerops.Monitoring
alias Towerops.Monitoring.PingMock
alias Towerops.Workers.MonitorWorker
setup :verify_on_exit!
setup do
user = user_fixture()
{:ok, organization} = Towerops.Organizations.create_organization(%{name: "Test Org"}, user.id)
{:ok, site} =
Towerops.Sites.create_site(%{
name: "Test Site",
organization_id: organization.id
})
{:ok, device} =
Towerops.Devices.create_device(%{
name: "Test Device",
# Use localhost for reliable testing
ip_address: "127.0.0.1",
monitoring_enabled: true,
site_id: site.id
})
{:ok, device: device, org: organization}
end
describe "perform/1" do
test "creates a check when monitoring is enabled", %{device: device} do
# Mock successful ping
expect(PingMock, :ping, fn _ip, _timeout ->
{:ok, 10.5}
end)
result = MonitorWorker.perform(device.id)
# Should return :ok
assert result == :ok
# Verify a successful check was created
checks = Monitoring.list_devices_checks(device.id, 10)
assert length(checks) == 1
check = hd(checks)
assert check.device_id == device.id
assert check.status == :success
assert check.checked_at
assert check.response_time_ms == 10.5
end
test "skips check when monitoring is disabled", %{device: device} do
# Disable monitoring
{:ok, device} = Towerops.Devices.update_device(device, %{monitoring_enabled: false})
assert :ok = MonitorWorker.perform(device.id)
# Verify no check was created
checks = Monitoring.list_devices_checks(device.id, 10)
assert checks == []
end
test "returns error when device not found" do
non_existent_id = Ecto.UUID.generate()
assert {:error, :device_not_found} = MonitorWorker.perform(non_existent_id)
end
test "broadcasts update after check", %{device: device} do
# Mock successful ping
expect(PingMock, :ping, fn _ip, _timeout ->
{:ok, 10.5}
end)
# Subscribe to device topic
Phoenix.PubSub.subscribe(Towerops.PubSub, "device:#{device.id}")
MonitorWorker.perform(device.id)
# Verify broadcast was sent
assert_receive {:monitoring_check_updated, device_id}, 1000
assert device_id == device.id
end
test "creates failed check for unreachable host" do
user = user_fixture()
{:ok, organization} = Towerops.Organizations.create_organization(%{name: "Test Org"}, user.id)
{:ok, site} =
Towerops.Sites.create_site(%{
name: "Test Site",
organization_id: organization.id
})
{:ok, device} =
Towerops.Devices.create_device(%{
name: "Unreachable Device",
ip_address: "192.0.2.1",
monitoring_enabled: true,
site_id: site.id
})
# Mock failed ping
expect(PingMock, :ping, fn _ip, _timeout ->
{:error, :timeout}
end)
assert :ok = MonitorWorker.perform(device.id)
# Verify check was created as failure
checks = Monitoring.list_devices_checks(device.id, 10)
assert length(checks) == 1
check = hd(checks)
assert check.status == :failure
assert check.response_time_ms == nil
end
test "successfully completes monitoring check", %{device: device} do
# Mock successful ping
expect(PingMock, :ping, fn _ip, _timeout ->
{:ok, 15.0}
end)
assert :ok = MonitorWorker.perform(device.id)
# Verify a successful check was created
checks = Monitoring.list_devices_checks(device.id, 10)
assert checks != [], "Expected at least 1 check"
check = hd(checks)
assert check.device_id == device.id
assert check.status == :success
assert check.checked_at
assert check.response_time_ms == 15.0
end
test "correctly skips monitoring when disabled", %{device: device} do
{:ok, device} = Towerops.Devices.update_device(device, %{monitoring_enabled: false})
assert :ok = MonitorWorker.perform(device.id)
# Verify no check was created
checks = Monitoring.list_devices_checks(device.id, 10)
assert checks == []
end
end
end

View file

@ -1,271 +0,0 @@
defmodule Towerops.Workers.PollWorkerTest do
use Towerops.DataCase, async: true
import Mox
import Towerops.AccountsFixtures
alias Towerops.Snmp
alias Towerops.Snmp.SnmpMock
alias Towerops.Workers.PollWorker
setup :verify_on_exit!
describe "perform/1" do
setup do
user = user_fixture()
{:ok, organization} = Towerops.Organizations.create_organization(%{name: "Test Org"}, user.id)
{:ok, site} =
Towerops.Sites.create_site(%{
name: "Test Site",
organization_id: organization.id
})
{:ok, device} =
Towerops.Devices.create_device(%{
name: "Test Device",
ip_address: "192.168.1.1",
snmp_enabled: true,
snmp_version: "2c",
snmp_community: "public",
snmp_port: 161,
site_id: site.id
})
snmp_device =
%Snmp.Device{}
|> Snmp.Device.changeset(%{
device_id: device.id,
sys_name: "test-device",
sys_descr: "Test Device"
})
|> Repo.insert!()
sensor =
%Snmp.Sensor{}
|> Snmp.Sensor.changeset(%{
snmp_device_id: snmp_device.id,
sensor_oid: "1.3.6.1.4.1.14988.1.1.3.10.0",
sensor_divisor: 10,
sensor_index: "0",
sensor_type: "temperature",
sensor_descr: "Temperature"
})
|> Repo.insert!()
interface =
%Snmp.Interface{}
|> Snmp.Interface.changeset(%{
snmp_device_id: snmp_device.id,
if_index: 1,
if_descr: "eth0"
})
|> Repo.insert!()
{:ok, device: device, snmp_device: snmp_device, sensor: sensor, interface: interface}
end
test "successfully polls a device with sensors and interfaces", %{
device: device,
sensor: sensor,
interface: interface
} do
# Mock sensor + interface stats (allow any number of calls)
stub(SnmpMock, :get, fn _target, oid, _opts ->
case oid do
"1.3.6.1.4.1.14988.1.1.3.10.0" -> {:ok, 350}
"1.3.6.1.2.1.2.2.1.10.1" -> {:ok, 1000}
"1.3.6.1.2.1.2.2.1.16.1" -> {:ok, 2000}
"1.3.6.1.2.1.2.2.1.14.1" -> {:ok, 0}
"1.3.6.1.2.1.2.2.1.20.1" -> {:ok, 0}
"1.3.6.1.2.1.2.2.1.13.1" -> {:ok, 0}
"1.3.6.1.2.1.2.2.1.19.1" -> {:ok, 0}
_ -> {:error, :no_such_object}
end
end)
# Mock neighbor discovery (allow any number of walk calls)
stub(SnmpMock, :walk, fn _target, _oid, _opts ->
{:ok, []}
end)
assert :ok = PollWorker.perform(device.id)
# Verify sensor reading was created
sensor = Towerops.Repo.reload(sensor)
# 350 / 10.0 divisor
assert sensor.last_value == 35.0
assert sensor.last_checked_at
# Verify interface stat was created
stats = Snmp.get_interface_stats(interface.id, limit: 10)
assert length(stats) == 1
stat = hd(stats)
assert stat.if_in_octets == 1000
assert stat.if_out_octets == 2000
end
test "skips poll when SNMP is disabled", %{device: device} do
{:ok, _device} = Towerops.Devices.update_device(device, %{snmp_enabled: false})
assert :ok = PollWorker.perform(device.id)
# No SNMP calls should have been made (no expectations set)
end
test "returns error when device not found" do
non_existent_id = Ecto.UUID.generate()
assert {:error, :device_not_found} = PollWorker.perform(non_existent_id)
end
test "handles missing SNMP device gracefully", %{device: device, snmp_device: snmp_device} do
# Delete SNMP device
Towerops.Repo.delete!(snmp_device)
# Should not crash
assert :ok = PollWorker.perform(device.id)
end
test "creates sensor reading with error status on SNMP failure", %{
device: device,
sensor: sensor
} do
# Mock SNMP timeout (allow any number of calls)
stub(SnmpMock, :get, fn _target, _oid, _opts ->
{:error, :timeout}
end)
# Mock neighbor discovery - empty
stub(SnmpMock, :walk, fn _target, _oid, _opts -> {:ok, []} end)
result = PollWorker.perform(device.id)
assert result == :ok
# Verify SNMP device exists and has sensors
snmp_device = Snmp.get_device_with_associations(device.id)
assert snmp_device
assert snmp_device.sensors != [], "Expected at least 1 sensor"
# Verify error reading was created
readings = Snmp.get_sensor_readings(sensor.id, limit: 10)
assert readings != [], "Expected at least 1 sensor reading"
reading = hd(readings)
assert reading.status == "error"
assert reading.value == nil
end
test "updates sensor last_value and last_checked_at", %{device: device, sensor: sensor} do
# Mock sensor + interface stats (allow any number of calls)
stub(SnmpMock, :get, fn _target, _oid, _opts -> {:ok, 250} end)
stub(SnmpMock, :walk, fn _target, _oid, _opts -> {:ok, []} end)
PollWorker.perform(device.id)
sensor = Towerops.Repo.reload(sensor)
# 250 / 10
assert sensor.last_value == 25.0
assert sensor.last_checked_at
end
test "creates interface stats for all interfaces", %{device: device, interface: interface} do
# Mock sensor + interface stats (allow any number of calls)
stub(SnmpMock, :get, fn _target, oid, _opts ->
case oid do
# Sensor OID
"1.3.6.1.4.1.14988.1.1.3.10.0" -> {:ok, 350}
# Interface stats
"1.3.6.1.2.1.2.2.1.10.1" -> {:ok, 1234}
"1.3.6.1.2.1.2.2.1.16.1" -> {:ok, 5678}
"1.3.6.1.2.1.2.2.1.14.1" -> {:ok, 10}
"1.3.6.1.2.1.2.2.1.20.1" -> {:ok, 20}
"1.3.6.1.2.1.2.2.1.13.1" -> {:ok, 5}
"1.3.6.1.2.1.2.2.1.19.1" -> {:ok, 3}
_ -> {:error, :no_such_object}
end
end)
stub(SnmpMock, :walk, fn _target, _oid, _opts -> {:ok, []} end)
PollWorker.perform(device.id)
stats = Snmp.get_interface_stats(interface.id, limit: 10)
assert length(stats) == 1
stat = hd(stats)
assert stat.if_in_octets == 1234
assert stat.if_out_octets == 5678
assert stat.if_in_errors == 10
assert stat.if_out_errors == 20
assert stat.if_in_discards == 5
assert stat.if_out_discards == 3
end
test "discovers and upserts neighbors", %{device: device, snmp_device: _snmp_device} do
# Mock sensor + interface stats + neighbor discovery (allow any number of calls)
stub(SnmpMock, :get, fn _target, _oid, _opts -> {:ok, 100} end)
stub(SnmpMock, :walk, fn _target, _oid, _opts -> {:ok, []} end)
PollWorker.perform(device.id)
# Verify neighbors list can be retrieved (may be empty)
neighbors = Snmp.list_neighbors(device.id)
assert is_list(neighbors)
end
test "deletes stale neighbors older than 5 minutes", %{
device: device,
snmp_device: _snmp_device,
interface: interface
} do
# Create a stale neighbor (10 minutes old)
stale_time = DateTime.add(DateTime.utc_now(), -10, :minute)
stale_neighbor =
%Snmp.Neighbor{}
|> Snmp.Neighbor.changeset(%{
device_id: device.id,
interface_id: interface.id,
protocol: "lldp",
remote_chassis_id: "old:neighbor",
remote_port_id: "port1",
last_discovered_at: stale_time
})
|> Repo.insert!()
# Mock sensor + interface stats + neighbor discovery (allow any number of calls)
stub(SnmpMock, :get, fn _target, _oid, _opts -> {:ok, 100} end)
stub(SnmpMock, :walk, fn _target, _oid, _opts -> {:ok, []} end)
PollWorker.perform(device.id)
# Verify stale neighbor was deleted
assert Towerops.Repo.get(Towerops.Snmp.Neighbor, stale_neighbor.id) == nil
end
test "handles non-numeric sensor values", %{device: device} do
# Mock sensor returning string instead of number (allow any number of calls)
stub(SnmpMock, :get, fn _target, _oid, _opts ->
{:ok, "not a number"}
end)
stub(SnmpMock, :walk, fn _target, _oid, _opts -> {:ok, []} end)
# Should not crash
PollWorker.perform(device.id)
end
test "successfully completes poll with proper mocks", %{device: device, sensor: sensor} do
# Mock sensor + interface stats (allow any number of calls)
stub(SnmpMock, :get, fn _target, _oid, _opts -> {:ok, 100} end)
stub(SnmpMock, :walk, fn _target, _oid, _opts -> {:ok, []} end)
assert :ok = PollWorker.perform(device.id)
# Verify sensor was polled and value updated
sensor = Towerops.Repo.reload(sensor)
# 100 / 10
assert sensor.last_value == 10.0
assert sensor.last_checked_at
end
end
end