From d0946c3cd0cae5831434318d9076100bb974d6ac Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Sat, 24 Jan 2026 16:36:57 -0600 Subject: [PATCH] refactor: simplify job architecture from Oban coordinators to direct workers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- lib/towerops/application.ex | 13 - lib/towerops/devices.ex | 74 ++- lib/towerops/exq_supervisor.ex | 107 ---- lib/towerops/monitoring/device_monitor.ex | 253 ---------- lib/towerops/monitoring/supervisor.ex | 72 +-- .../workers/device_monitor_coordinator.ex | 92 ---- lib/towerops/workers/device_monitor_worker.ex | 213 ++++++++ .../workers/device_poller_coordinator.ex | 96 ---- .../device_poller_worker.ex} | 396 ++++----------- lib/towerops/workers/discovery_worker.ex | 27 +- lib/towerops/workers/monitor_worker.ex | 96 ---- lib/towerops/workers/poll_worker.ex | 177 ------- lib/towerops_web/channels/agent_channel.ex | 4 +- lib/towerops_web/live/device_live/form.ex | 4 +- lib/towerops_web/live/device_live/index.ex | 2 +- lib/towerops_web/live/site_live/show.ex | 4 +- mix.exs | 9 +- mix.lock | 2 - test/integration/snmp_integration_test.exs | 9 +- test/towerops/exq_supervisor_test.exs | 95 ---- .../monitoring/device_monitor_test.exs | 457 ------------------ test/towerops/monitoring/supervisor_test.exs | 155 ------ .../workers/discovery_worker_test.exs | 12 +- test/towerops/workers/monitor_worker_test.exs | 154 ------ test/towerops/workers/poll_worker_test.exs | 271 ----------- 25 files changed, 404 insertions(+), 2390 deletions(-) delete mode 100644 lib/towerops/exq_supervisor.ex delete mode 100644 lib/towerops/monitoring/device_monitor.ex delete mode 100644 lib/towerops/workers/device_monitor_coordinator.ex create mode 100644 lib/towerops/workers/device_monitor_worker.ex delete mode 100644 lib/towerops/workers/device_poller_coordinator.ex rename lib/towerops/{snmp/poller_worker.ex => workers/device_poller_worker.ex} (77%) delete mode 100644 lib/towerops/workers/monitor_worker.ex delete mode 100644 lib/towerops/workers/poll_worker.ex delete mode 100644 test/towerops/exq_supervisor_test.exs delete mode 100644 test/towerops/monitoring/device_monitor_test.exs delete mode 100644 test/towerops/monitoring/supervisor_test.exs delete mode 100644 test/towerops/workers/monitor_worker_test.exs delete mode 100644 test/towerops/workers/poll_worker_test.exs diff --git a/lib/towerops/application.ex b/lib/towerops/application.ex index 641048cb..f1db35ec 100644 --- a/lib/towerops/application.ex +++ b/lib/towerops/application.ex @@ -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 diff --git a/lib/towerops/devices.ex b/lib/towerops/devices.ex index 1b5fd239..3f0070e9 100644 --- a/lib/towerops/devices.ex +++ b/lib/towerops/devices.ex @@ -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 """ diff --git a/lib/towerops/exq_supervisor.ex b/lib/towerops/exq_supervisor.ex deleted file mode 100644 index 9b4886b7..00000000 --- a/lib/towerops/exq_supervisor.ex +++ /dev/null @@ -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 diff --git a/lib/towerops/monitoring/device_monitor.ex b/lib/towerops/monitoring/device_monitor.ex deleted file mode 100644 index bc1f3d54..00000000 --- a/lib/towerops/monitoring/device_monitor.ex +++ /dev/null @@ -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 diff --git a/lib/towerops/monitoring/supervisor.ex b/lib/towerops/monitoring/supervisor.ex index e40c192f..48f76391 100644 --- a/lib/towerops/monitoring/supervisor.ex +++ b/lib/towerops/monitoring/supervisor.ex @@ -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, []) diff --git a/lib/towerops/workers/device_monitor_coordinator.ex b/lib/towerops/workers/device_monitor_coordinator.ex deleted file mode 100644 index 613d9188..00000000 --- a/lib/towerops/workers/device_monitor_coordinator.ex +++ /dev/null @@ -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 diff --git a/lib/towerops/workers/device_monitor_worker.ex b/lib/towerops/workers/device_monitor_worker.ex new file mode 100644 index 00000000..512f3f4f --- /dev/null +++ b/lib/towerops/workers/device_monitor_worker.ex @@ -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 diff --git a/lib/towerops/workers/device_poller_coordinator.ex b/lib/towerops/workers/device_poller_coordinator.ex deleted file mode 100644 index b6972f01..00000000 --- a/lib/towerops/workers/device_poller_coordinator.ex +++ /dev/null @@ -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 diff --git a/lib/towerops/snmp/poller_worker.ex b/lib/towerops/workers/device_poller_worker.ex similarity index 77% rename from lib/towerops/snmp/poller_worker.ex rename to lib/towerops/workers/device_poller_worker.ex index 5203bc34..3896c79b 100644 --- a/lib/towerops/snmp/poller_worker.ex +++ b/lib/towerops/workers/device_poller_worker.ex @@ -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) <> = 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 diff --git a/lib/towerops/workers/discovery_worker.ex b/lib/towerops/workers/discovery_worker.ex index e5a54cfd..3d8765cb 100644 --- a/lib/towerops/workers/discovery_worker.ex +++ b/lib/towerops/workers/discovery_worker.ex @@ -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 diff --git a/lib/towerops/workers/monitor_worker.ex b/lib/towerops/workers/monitor_worker.ex deleted file mode 100644 index e95bf524..00000000 --- a/lib/towerops/workers/monitor_worker.ex +++ /dev/null @@ -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 diff --git a/lib/towerops/workers/poll_worker.ex b/lib/towerops/workers/poll_worker.ex deleted file mode 100644 index b9cb227d..00000000 --- a/lib/towerops/workers/poll_worker.ex +++ /dev/null @@ -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 diff --git a/lib/towerops_web/channels/agent_channel.ex b/lib/towerops_web/channels/agent_channel.ex index f9ddb266..2ed5638c 100644 --- a/lib/towerops_web/channels/agent_channel.ex +++ b/lib/towerops_web/channels/agent_channel.ex @@ -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 diff --git a/lib/towerops_web/live/device_live/form.ex b/lib/towerops_web/live/device_live/form.ex index 4cdffff7..d0ab086b 100644 --- a/lib/towerops_web/live/device_live/form.ex +++ b/lib/towerops_web/live/device_live/form.ex @@ -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 diff --git a/lib/towerops_web/live/device_live/index.ex b/lib/towerops_web/live/device_live/index.ex index 5d6e0d6b..d9da2cb0 100644 --- a/lib/towerops_web/live/device_live/index.ex +++ b/lib/towerops_web/live/device_live/index.ex @@ -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 diff --git a/lib/towerops_web/live/site_live/show.ex b/lib/towerops_web/live/site_live/show.ex index e627c68e..081c6cd2 100644 --- a/lib/towerops_web/live/site_live/show.ex +++ b/lib/towerops_web/live/site_live/show.ex @@ -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 diff --git a/mix.exs b/mix.exs index a7488959..20319a05 100644 --- a/mix.exs +++ b/mix.exs @@ -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" ] diff --git a/mix.lock b/mix.lock index 000e206f..d8e41840 100644 --- a/mix.lock +++ b/mix.lock @@ -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"}, diff --git a/test/integration/snmp_integration_test.exs b/test/integration/snmp_integration_test.exs index fa5e9d73..7ee4f155 100644 --- a/test/integration/snmp_integration_test.exs +++ b/test/integration/snmp_integration_test.exs @@ -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) diff --git a/test/towerops/exq_supervisor_test.exs b/test/towerops/exq_supervisor_test.exs deleted file mode 100644 index 95fd49ee..00000000 --- a/test/towerops/exq_supervisor_test.exs +++ /dev/null @@ -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 diff --git a/test/towerops/monitoring/device_monitor_test.exs b/test/towerops/monitoring/device_monitor_test.exs deleted file mode 100644 index a129fc0c..00000000 --- a/test/towerops/monitoring/device_monitor_test.exs +++ /dev/null @@ -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 diff --git a/test/towerops/monitoring/supervisor_test.exs b/test/towerops/monitoring/supervisor_test.exs deleted file mode 100644 index 3ad27009..00000000 --- a/test/towerops/monitoring/supervisor_test.exs +++ /dev/null @@ -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 diff --git a/test/towerops/workers/discovery_worker_test.exs b/test/towerops/workers/discovery_worker_test.exs index 63a7741b..4ba5bb8f 100644 --- a/test/towerops/workers/discovery_worker_test.exs +++ b/test/towerops/workers/discovery_worker_test.exs @@ -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}" diff --git a/test/towerops/workers/monitor_worker_test.exs b/test/towerops/workers/monitor_worker_test.exs deleted file mode 100644 index ca26476e..00000000 --- a/test/towerops/workers/monitor_worker_test.exs +++ /dev/null @@ -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 diff --git a/test/towerops/workers/poll_worker_test.exs b/test/towerops/workers/poll_worker_test.exs deleted file mode 100644 index 467b0e6c..00000000 --- a/test/towerops/workers/poll_worker_test.exs +++ /dev/null @@ -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