diff --git a/lib/towerops/application.ex b/lib/towerops/application.ex index 36bbc4f2..f9be1d3e 100644 --- a/lib/towerops/application.ex +++ b/lib/towerops/application.ex @@ -74,7 +74,7 @@ defmodule Towerops.Application do port: Keyword.get(redis_config, :port, 6379), namespace: "exq", concurrency: 10, - queues: ["default", "discovery", "polling", "maintenance"] + queues: ["default", "discovery", "polling", "monitoring", "maintenance"] ] end diff --git a/lib/towerops/monitoring/device_monitor.ex b/lib/towerops/monitoring/device_monitor.ex index 8dbfa78c..0c67ab66 100644 --- a/lib/towerops/monitoring/device_monitor.ex +++ b/lib/towerops/monitoring/device_monitor.ex @@ -37,8 +37,8 @@ defmodule Towerops.Monitoring.DeviceMonitor do GenServer.cast(via_tuple(device_id), :check_now) [] -> - # Process doesn't exist (monitoring disabled), perform check directly - Task.start(fn -> perform_check(device_id) end) + # Process doesn't exist (monitoring disabled), enqueue check job + enqueue_check(device_id) end end @@ -215,6 +215,17 @@ defmodule Towerops.Monitoring.DeviceMonitor 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 + Task.start(fn -> perform_check(device_id) end) + else + # In dev/prod, enqueue to Exq + {:ok, _job} = Exq.enqueue(Exq, "monitoring", Towerops.Workers.MonitorWorker, [device_id]) + end + end + defp via_tuple(device_id) do {:via, Registry, {Towerops.Monitoring.Registry, device_id}} end diff --git a/lib/towerops/snmp/poller_worker.ex b/lib/towerops/snmp/poller_worker.ex index c4fe0329..127b1481 100644 --- a/lib/towerops/snmp/poller_worker.ex +++ b/lib/towerops/snmp/poller_worker.ex @@ -41,8 +41,8 @@ defmodule Towerops.Snmp.PollerWorker do GenServer.cast(via_tuple(device_id), :poll_now) [] -> - # Process doesn't exist, perform poll directly in background - Task.start(fn -> perform_poll(device_id) end) + # Process doesn't exist, enqueue poll job + enqueue_poll(device_id) end end @@ -960,6 +960,17 @@ defmodule Towerops.Snmp.PollerWorker do 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 + Task.start(fn -> perform_poll(device_id) end) + else + # In dev/prod, enqueue to Exq + {:ok, _job} = Exq.enqueue(Exq, "polling", Towerops.Workers.PollWorker, [device_id]) + end + end + defp via_tuple(device_id) do {:via, Registry, {PollerRegistry, device_id}} end diff --git a/lib/towerops/workers/monitor_worker.ex b/lib/towerops/workers/monitor_worker.ex new file mode 100644 index 00000000..8677780f --- /dev/null +++ b/lib/towerops/workers/monitor_worker.ex @@ -0,0 +1,102 @@ +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 + + @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: "up", + latency_ms: latency, + checked_at: DateTime.utc_now() + }) + + :ok + + {:error, reason} -> + Logger.warning("Device #{device_id} is down: #{inspect(reason)}") + + Monitoring.create_check(%{ + device_id: device_id, + status: "down", + latency_ms: nil, + checked_at: DateTime.utc_now() + }) + + :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 system ping command (works on most Unix-like systems) + case System.cmd("ping", ["-c", "1", "-W", "1", device.ip_address], stderr_to_stdout: true) do + {output, 0} -> + # Extract latency from ping output + latency = extract_latency(output) + {:ok, latency} + + {_output, _status} -> + {:error, :timeout} + end + rescue + error -> + {:error, error} + end + + defp extract_latency(output) do + # Parse ping output to extract latency + # Example: "time=1.234 ms" + case Regex.run(~r/time=([\d.]+)/, output) do + [_, latency_str] -> + case Float.parse(latency_str) do + {latency, _} -> latency + :error -> 0.0 + end + + nil -> + 0.0 + end + end +end diff --git a/lib/towerops/workers/poll_worker.ex b/lib/towerops/workers/poll_worker.ex new file mode 100644 index 00000000..b9cb227d --- /dev/null +++ b/lib/towerops/workers/poll_worker.ex @@ -0,0 +1,177 @@ +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 726bed22..8f2f70f0 100644 --- a/lib/towerops_web/channels/agent_channel.ex +++ b/lib/towerops_web/channels/agent_channel.ex @@ -317,15 +317,7 @@ defmodule ToweropsWeb.AgentChannel do # - Server does SNMP queries and parsing using existing logic Logger.info("Discovery results received for #{device.name}, triggering full discovery") - Task.start(fn -> - case Discovery.discover_device(device) do - {:ok, _device} -> - Logger.info("Full discovery completed for #{device.name}") - - {:error, reason} -> - Logger.error("Discovery failed for #{device.name}: #{inspect(reason)}") - end - end) + enqueue_discovery(device.id) end defp process_polling_result(device, result) do @@ -365,6 +357,17 @@ defmodule ToweropsWeb.AgentChannel do # since it requires complex LLDP/CDP parsing end + # Enqueue discovery job - safe to call in test environment + defp enqueue_discovery(device_id) do + if Application.get_env(:towerops, :env) == :test 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]) + end + end + defp parse_integer(nil), do: nil defp parse_integer(value) when is_integer(value), do: value