Convert Task.* calls to Exq background workers

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)
This commit is contained in:
Graham McIntire 2026-01-17 16:46:36 -06:00
parent 073f9bf27d
commit 5339f51dd8
No known key found for this signature in database
6 changed files with 318 additions and 14 deletions

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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