refactor: replace Task.async_stream with Oban jobs for discovery
Changed discover_all/1 to enqueue individual Oban jobs instead of using
Task.async_stream for better error visibility and resilience.
Benefits:
- Errors are visible in Oban dashboard with full stack traces
- Automatic retries on failure via Oban
- Jobs persist across pod restarts
- No more silent failures from Task crashes
- Better concurrency control via Oban queue configuration
Changes:
- discover_all/1 now enqueues DiscoveryWorker jobs and returns immediately
- Return type changed from %{success, failed, errors} to %{enqueued, failed, errors}
- Updated tests to match new async behavior
- Added DiscoveryWorker alias to Discovery module
All tests passing (3,625 tests, 0 failures).
This commit is contained in:
parent
58883d0d58
commit
a92d003aec
2 changed files with 30 additions and 36 deletions
|
|
@ -31,6 +31,7 @@ defmodule Towerops.Snmp.Discovery do
|
|||
alias Towerops.Snmp.Sensor
|
||||
alias Towerops.Snmp.Storage
|
||||
alias Towerops.Snmp.Vlan
|
||||
alias Towerops.Workers.DiscoveryWorker
|
||||
|
||||
require Logger
|
||||
|
||||
|
|
@ -81,7 +82,7 @@ defmodule Towerops.Snmp.Discovery do
|
|||
@type profile :: module() | {:yaml, map()}
|
||||
|
||||
@type discovery_summary :: %{
|
||||
success: non_neg_integer(),
|
||||
enqueued: non_neg_integer(),
|
||||
failed: non_neg_integer(),
|
||||
errors: [term()]
|
||||
}
|
||||
|
|
@ -249,8 +250,14 @@ defmodule Towerops.Snmp.Discovery do
|
|||
end
|
||||
|
||||
@doc """
|
||||
Runs discovery for all SNMP-enabled devices in an organization.
|
||||
Returns a summary of successful and failed discoveries.
|
||||
Enqueues Oban discovery jobs for all SNMP-enabled devices in an organization.
|
||||
|
||||
Returns immediately after enqueuing jobs. Job status can be monitored via Oban dashboard.
|
||||
|
||||
## Examples
|
||||
|
||||
iex> discover_all(org_id)
|
||||
{:ok, %{enqueued: 10, failed: 0}}
|
||||
"""
|
||||
@spec discover_all(String.t()) :: {:ok, discovery_summary()}
|
||||
def discover_all(org_id) do
|
||||
|
|
@ -260,28 +267,23 @@ defmodule Towerops.Snmp.Discovery do
|
|||
|> where([e, s], s.organization_id == ^org_id and e.snmp_enabled == true)
|
||||
|> Repo.all()
|
||||
|
||||
Logger.info("Starting SNMP discovery for #{length(device_list)} devices in org #{org_id}")
|
||||
Logger.info("Enqueuing SNMP discovery for #{length(device_list)} devices in org #{org_id}")
|
||||
|
||||
# Enqueue Oban jobs for each device
|
||||
results =
|
||||
device_list
|
||||
|> Task.async_stream(
|
||||
&discover_device/1,
|
||||
max_concurrency: 5,
|
||||
timeout: 60_000,
|
||||
on_timeout: :kill_task
|
||||
)
|
||||
|> Enum.reduce(%{success: 0, failed: 0, errors: []}, fn
|
||||
{:ok, {:ok, _device}}, acc ->
|
||||
%{acc | success: acc.success + 1}
|
||||
Enum.reduce(device_list, %{enqueued: 0, failed: 0, errors: []}, fn device, acc ->
|
||||
case DiscoveryWorker.enqueue(device.id) do
|
||||
{:ok, _job} ->
|
||||
%{acc | enqueued: acc.enqueued + 1}
|
||||
|
||||
{:ok, {:error, reason}}, acc ->
|
||||
%{acc | failed: acc.failed + 1, errors: [reason | acc.errors]}
|
||||
|
||||
{:exit, :timeout}, acc ->
|
||||
%{acc | failed: acc.failed + 1, errors: [:timeout | acc.errors]}
|
||||
{:error, reason} ->
|
||||
Logger.error("Failed to enqueue discovery for device #{device.id}: #{inspect(reason)}")
|
||||
%{acc | failed: acc.failed + 1, errors: [reason | acc.errors]}
|
||||
end
|
||||
end)
|
||||
|
||||
Logger.info("Discovery completed: #{results.success} succeeded, #{results.failed} failed")
|
||||
Logger.info("Discovery jobs enqueued: #{results.enqueued} succeeded, #{results.failed} failed")
|
||||
|
||||
{:ok, results}
|
||||
end
|
||||
|
||||
|
|
|
|||
|
|
@ -795,7 +795,7 @@ defmodule Towerops.Snmp.DiscoveryTest do
|
|||
describe "discover_all/1" do
|
||||
test "discovers multiple devices concurrently", %{organization: organization, site: site} do
|
||||
# Create multiple SNMP-enabled devices
|
||||
{:ok, device1} =
|
||||
{:ok, _device1} =
|
||||
Towerops.Devices.create_device(%{
|
||||
name: "Device 1",
|
||||
ip_address: "192.168.1.10",
|
||||
|
|
@ -806,7 +806,7 @@ defmodule Towerops.Snmp.DiscoveryTest do
|
|||
snmp_port: 161
|
||||
})
|
||||
|
||||
{:ok, device2} =
|
||||
{:ok, _device2} =
|
||||
Towerops.Devices.create_device(%{
|
||||
name: "Device 2",
|
||||
ip_address: "192.168.1.11",
|
||||
|
|
@ -847,17 +847,13 @@ defmodule Towerops.Snmp.DiscoveryTest do
|
|||
|
||||
assert {:ok, summary} = Discovery.discover_all(organization.id)
|
||||
|
||||
# Should only discover SNMP-enabled devices
|
||||
assert summary.success == 2
|
||||
# Should enqueue jobs for SNMP-enabled devices
|
||||
assert summary.enqueued == 2
|
||||
assert summary.failed == 0
|
||||
|
||||
# Verify devices were created
|
||||
assert Repo.get_by(Device, device_id: device1.id)
|
||||
assert Repo.get_by(Device, device_id: device2.id)
|
||||
end
|
||||
|
||||
test "handles mixed success and failure", %{organization: organization, site: site} do
|
||||
{:ok, device1} =
|
||||
{:ok, _device1} =
|
||||
Towerops.Devices.create_device(%{
|
||||
name: "Good Device",
|
||||
ip_address: "192.168.1.20",
|
||||
|
|
@ -904,13 +900,9 @@ defmodule Towerops.Snmp.DiscoveryTest do
|
|||
|
||||
assert {:ok, summary} = Discovery.discover_all(organization.id)
|
||||
|
||||
assert summary.success == 1
|
||||
assert summary.failed == 1
|
||||
# Device that fails initial responsiveness check is marked as unresponsive
|
||||
assert :device_unresponsive in summary.errors
|
||||
|
||||
# Verify only good device was created
|
||||
assert Repo.get_by(Device, device_id: device1.id)
|
||||
# Both devices should be enqueued (Oban jobs will handle success/failure)
|
||||
assert summary.enqueued == 2
|
||||
assert summary.failed == 0
|
||||
end
|
||||
|
||||
test "updates device name from SNMP sysName when device name is empty", %{
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue