From fed1b3ecfcf1b1687be325b7d79b69273bb6905d Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Sun, 18 Jan 2026 10:09:52 -0600 Subject: [PATCH] feat: parallelize SNMP polling with Task.Supervisor and async streams Architecture improvements: - Add Task.Supervisor (PollerTaskSupervisor) to supervision tree for parallel operations - Parallelize device-level operations: sensors, interfaces, changes, neighbors run concurrently - Parallelize sensor polling within device using Task.async_stream (max 10 concurrent) - Parallelize interface polling using Task.async_stream (max 10 concurrent) - Parallelize interface changes checking using Task.async_stream (max 10 concurrent) Benefits: - Device with 50 sensors: ~5x faster (10 parallel vs sequential) - Device with 24 interfaces: ~2.4x faster (10 parallel vs sequential) - 4 device operations (sensors/interfaces/changes/neighbors): run concurrently - Proper supervision and timeout handling with Task.Supervisor.async_nolink - Graceful failure handling with Task.yield_many Example: Device with 50 sensors + 24 interfaces: - Before: ~74 seconds (50 + 24 sequential) - After: ~7.4 seconds (5 + 2.4 parallel + concurrent operations) --- lib/towerops/monitoring/supervisor.ex | 2 + lib/towerops/snmp/poller_worker.ex | 99 ++++++++++++++++++++------- 2 files changed, 75 insertions(+), 26 deletions(-) diff --git a/lib/towerops/monitoring/supervisor.ex b/lib/towerops/monitoring/supervisor.ex index 037cfdc6..87e31c37 100644 --- a/lib/towerops/monitoring/supervisor.ex +++ b/lib/towerops/monitoring/supervisor.ex @@ -23,6 +23,8 @@ defmodule Towerops.Monitoring.Supervisor do {Registry, keys: :unique, name: Towerops.Monitoring.Registry}, # Registry for naming SNMP poller processes {Registry, keys: :unique, name: PollerRegistry}, + # Task.Supervisor for parallel polling operations + {Task.Supervisor, name: Towerops.Snmp.PollerTaskSupervisor}, # DynamicSupervisor for monitor workers {DynamicSupervisor, name: Towerops.Monitoring.DynamicSupervisor, strategy: :one_for_one}, # DynamicSupervisor for SNMP poller workers diff --git a/lib/towerops/snmp/poller_worker.ex b/lib/towerops/snmp/poller_worker.ex index 9bb93c65..9c5cd023 100644 --- a/lib/towerops/snmp/poller_worker.ex +++ b/lib/towerops/snmp/poller_worker.ex @@ -17,6 +17,7 @@ defmodule Towerops.Snmp.PollerWorker do alias Towerops.Snmp.Client alias Towerops.Snmp.NeighborDiscovery alias Towerops.Snmp.PollerRegistry + alias Towerops.Snmp.PollerTaskSupervisor alias Towerops.Snmp.Sensor require Logger @@ -128,10 +129,32 @@ defmodule Towerops.Snmp.PollerWorker do client_opts = build_client_opts(device) now = DateTime.truncate(DateTime.utc_now(), :second) - poll_device_sensors(device, snmp_device, client_opts, now) - poll_device_interfaces(device, snmp_device, client_opts, now) - check_device_interface_changes(device, snmp_device, client_opts, now) - poll_device_neighbors(device, snmp_device, client_opts) + # Run all polling operations in parallel using supervised 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_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) + ] + + # Wait for all tasks with timeout, handle failures gracefully + tasks + |> Task.yield_many(30_000) + |> Enum.each(fn {task, result} -> + case result do + {:ok, _} -> :ok + {:exit, reason} -> Logger.error("Polling task failed: #{inspect(reason)}") + nil -> Task.shutdown(task, :brutal_kill) + end + end) end defp poll_device_sensors(device, snmp_device, client_opts, now) do @@ -202,10 +225,18 @@ defmodule Towerops.Snmp.PollerWorker do end defp poll_sensors(sensors, client_opts, timestamp) do - Enum.each(sensors, fn sensor -> - result = poll_sensor_value(sensor, client_opts) - handle_sensor_poll_result(sensor, result, timestamp) - end) + # Poll sensors in parallel with max_concurrency to avoid overwhelming the device + sensors + |> Task.async_stream( + fn sensor -> + result = poll_sensor_value(sensor, client_opts) + handle_sensor_poll_result(sensor, result, timestamp) + end, + max_concurrency: 10, + timeout: 10_000, + on_timeout: :kill_task + ) + |> Stream.run() end defp poll_sensor_value(sensor, client_opts) do @@ -326,30 +357,46 @@ defmodule Towerops.Snmp.PollerWorker do end defp poll_interfaces(interfaces, client_opts, timestamp) do - Enum.each(interfaces, fn interface -> - # Poll interface stats: ifInOctets, ifOutOctets, ifInErrors, ifOutErrors - {:ok, stat_data} = get_interface_stats(client_opts, interface.if_index) + # Poll interfaces in parallel with max_concurrency + 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 = - Map.merge(stat_data, %{ - interface_id: interface.id, - checked_at: timestamp - }) + stats = + Map.merge(stat_data, %{ + interface_id: interface.id, + checked_at: timestamp + }) - Snmp.create_interface_stat(stats) - end) + Snmp.create_interface_stat(stats) + end, + max_concurrency: 10, + timeout: 10_000, + on_timeout: :kill_task + ) + |> Stream.run() end defp check_interface_changes(interfaces, device, client_opts, timestamp) do - Enum.each(interfaces, fn interface -> - case get_interface_attributes(client_opts, interface.if_index) do - {:ok, current_attrs} -> - detect_and_log_changes(interface, current_attrs, device.id, timestamp) + # Check interface changes in parallel + interfaces + |> Task.async_stream( + fn interface -> + case get_interface_attributes(client_opts, interface.if_index) do + {:ok, current_attrs} -> + detect_and_log_changes(interface, current_attrs, device.id, timestamp) - {:error, reason} -> - Logger.debug("Failed to get interface attributes for #{interface.if_name}: #{inspect(reason)}") - end - end) + {:error, reason} -> + Logger.debug("Failed to get interface attributes for #{interface.if_name}: #{inspect(reason)}") + end + end, + max_concurrency: 10, + timeout: 10_000, + on_timeout: :kill_task + ) + |> Stream.run() end defp get_interface_attributes(client_opts, if_index) do