Created three new Exq workers to handle background jobs: - PollWorker: SNMP polling operations (queue: polling) - MonitorWorker: Device monitoring checks (queue: monitoring) - DiscoveryWorker: SNMP discovery (already existed) Updated files to use Exq.enqueue instead of Task.start: - lib/towerops/snmp/poller_worker.ex - lib/towerops/monitoring/device_monitor.ex - lib/towerops_web/channels/agent_channel.ex Each module includes enqueue_* helper functions that: - Use Task.start in test environment for synchronous execution - Use Exq.enqueue in dev/prod for proper background job processing Added "monitoring" queue to Exq configuration in application.ex All tests passing (824 tests, 0 failures)
177 lines
4.6 KiB
Elixir
177 lines
4.6 KiB
Elixir
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
|