refactor: convert periodic workers to Oban Cron for better resilience
Replaced GenServer-based periodic workers with Oban Cron jobs to improve
pod rollover resilience and simplify architecture.
Worker Changes:
- NeighborCleanupWorker: GenServer → Oban Cron (hourly)
- Cleans stale neighbors, ARP entries, and MAC addresses
- Runs every hour via Oban.Plugins.Cron
- StaleAgentWorker: GenServer → Oban Cron (every minute)
- Detects agents that haven't checked in for 10+ minutes
- Refactored to reduce nesting (extracted helper functions)
- Removed stateful tracking (now stateless, re-evaluates each run)
- AgentLatencyEvaluator: GenServer → Oban Cron (every 5 minutes)
- Latency-based agent reassignment with 20% threshold
- Removed trigger_evaluation/0 (no longer needed)
- JobHealthCheckWorker: NEW Oban Cron worker (every 10 minutes)
- Safety net to recover missing monitor/poller jobs
- Auto-creates jobs for devices with monitoring/SNMP enabled
Infrastructure Changes:
- Removed Monitoring.Supervisor (no longer needed)
- Updated application.ex to remove GenServer workers from supervision tree
- Added Oban.Plugins.Cron to dev.exs and runtime.exs
- All workers now run cluster-wide via PostgreSQL-backed coordination
Test Updates:
- Updated all worker tests to call perform(%Oban.Job{args: %{}})
- Removed GenServer lifecycle tests (start_link, send messages, etc.)
- Removed async sleep calls (no longer needed)
Benefits:
- Better pod rollover resilience (Cron jobs run cluster-wide)
- Simpler architecture (no GenServers for periodic tasks)
- Better observability (all jobs visible in Oban dashboard)
- Safety net for missing jobs (JobHealthCheckWorker)
- Stateless workers (easier to reason about and test)
Documentation:
- Updated CLAUDE.md with Background Job Architecture section
- Documented job types, queues, and resilience features
All tests passing (3,686 tests, 0 failures).
This commit is contained in:
parent
5109f7b6a1
commit
ee8a3220c4
11 changed files with 241 additions and 324 deletions
39
CLAUDE.md
39
CLAUDE.md
|
|
@ -164,6 +164,45 @@ All LiveViews, LiveComponents, and HTML modules automatically get these imports/
|
|||
- Phoenix LiveReload watches for file changes
|
||||
- Code reloader enabled via `listeners: [Phoenix.CodeReloader]`
|
||||
|
||||
### Background Job Architecture
|
||||
|
||||
The application uses Oban for all background job processing with PostgreSQL-backed queuing for cluster-wide coordination and resilience.
|
||||
|
||||
**Job Types:**
|
||||
|
||||
1. **Self-Scheduling Workers** (per-device recurring jobs):
|
||||
- `DeviceMonitorWorker` - Health check polling (60s default interval)
|
||||
- `DevicePollerWorker` - SNMP data collection (60s default interval)
|
||||
- Schedule next run before completing via `schedule_next_poll/2`
|
||||
- One job per enabled device, automatically created/cancelled when device settings change
|
||||
|
||||
2. **Oban Cron Workers** (cluster-wide periodic maintenance):
|
||||
- `NeighborCleanupWorker` - Hourly cleanup of stale neighbors/ARP/MAC entries (runs every hour via cron)
|
||||
- `StaleAgentWorker` - Detects agents that haven't checked in for 10+ minutes (runs every minute)
|
||||
- `AgentLatencyEvaluator` - Latency-based agent reassignment (runs every 5 minutes)
|
||||
- `JobHealthCheckWorker` - Safety net to recover missing monitor/poller jobs (runs every 10 minutes)
|
||||
- Scheduled via `Oban.Plugins.Cron` in `config/dev.exs` and `config/runtime.exs`
|
||||
- Run cluster-wide (only one instance executes at a time across all pods)
|
||||
|
||||
**Queues:**
|
||||
- `default` (10 workers) - General background tasks
|
||||
- `discovery` (10 workers) - SNMP discovery operations
|
||||
- `pollers` (50 workers) - SNMP polling jobs (one per device)
|
||||
- `monitors` (50 workers) - Health check jobs (one per device)
|
||||
- `maintenance` (5 workers) - Periodic cleanup and health check workers
|
||||
|
||||
**Resilience:**
|
||||
- Oban Cron jobs run cluster-wide via PostgreSQL-based locking
|
||||
- If a pod dies, another pod immediately picks up scheduled cron jobs
|
||||
- Self-scheduling workers (monitor/poller) are recovered by `JobHealthCheckWorker` every 10 minutes
|
||||
- All jobs visible in Oban dashboard at `/dev/dashboard` → Oban tab
|
||||
- Failed jobs automatically retried with exponential backoff
|
||||
|
||||
**Configuration:**
|
||||
- Dev: `config/dev.exs` - Oban config with cron plugin
|
||||
- Prod: `config/runtime.exs` - Same cron schedule as dev
|
||||
- Application: `lib/towerops/application.ex` - Oban supervision, no GenServer workers for periodic tasks
|
||||
|
||||
### LiveDashboard and Telemetry
|
||||
|
||||
LiveDashboard is available at `/dashboard` (requires authentication in production, `/dev/dashboard` in development).
|
||||
|
|
|
|||
|
|
@ -38,6 +38,18 @@ config :towerops, Oban,
|
|||
maintenance: 5
|
||||
],
|
||||
plugins: [
|
||||
# Cron jobs for periodic maintenance tasks
|
||||
{Oban.Plugins.Cron,
|
||||
crontab: [
|
||||
# Run neighbor cleanup every hour
|
||||
{"0 * * * *", Towerops.Snmp.NeighborCleanupWorker},
|
||||
# Check for stale agents every minute
|
||||
{"* * * * *", Towerops.Workers.StaleAgentWorker},
|
||||
# Evaluate latency-based agent reassignment every 5 minutes
|
||||
{"*/5 * * * *", Towerops.Workers.AgentLatencyEvaluator},
|
||||
# Health check for missing jobs every 10 minutes
|
||||
{"*/10 * * * *", Towerops.Workers.JobHealthCheckWorker}
|
||||
]},
|
||||
# Automatically delete completed jobs after 60 seconds
|
||||
{Oban.Plugins.Pruner, max_age: 60},
|
||||
# Rescue orphaned jobs (when node crashes)
|
||||
|
|
|
|||
|
|
@ -84,12 +84,26 @@ if config_env() == :prod do
|
|||
repo: Towerops.Repo,
|
||||
queues: [
|
||||
default: 10,
|
||||
discovery: 10,
|
||||
# SNMP polling jobs - one per device
|
||||
pollers: 50,
|
||||
# Device monitoring jobs - health checks
|
||||
monitors: 50
|
||||
monitors: 50,
|
||||
maintenance: 5
|
||||
],
|
||||
plugins: [
|
||||
# Cron jobs for periodic maintenance tasks
|
||||
{Oban.Plugins.Cron,
|
||||
crontab: [
|
||||
# Run neighbor cleanup every hour
|
||||
{"0 * * * *", Towerops.Snmp.NeighborCleanupWorker},
|
||||
# Check for stale agents every minute
|
||||
{"* * * * *", Towerops.Workers.StaleAgentWorker},
|
||||
# Evaluate latency-based agent reassignment every 5 minutes
|
||||
{"*/5 * * * *", Towerops.Workers.AgentLatencyEvaluator},
|
||||
# Health check for missing jobs every 10 minutes
|
||||
{"*/10 * * * *", Towerops.Workers.JobHealthCheckWorker}
|
||||
]},
|
||||
# Automatically delete completed jobs after 60 seconds
|
||||
{Oban.Plugins.Pruner, max_age: 60},
|
||||
# Rescue orphaned jobs (when node crashes)
|
||||
|
|
|
|||
|
|
@ -61,13 +61,9 @@ defmodule Towerops.Application do
|
|||
else
|
||||
[
|
||||
# Start event logger (subscribes to PubSub)
|
||||
Towerops.Devices.EventLogger,
|
||||
# Start monitoring supervisor (includes NeighborCleanupWorker)
|
||||
Towerops.Monitoring.Supervisor,
|
||||
# Start stale agent detection worker
|
||||
Towerops.Workers.StaleAgentWorker,
|
||||
# Start latency-based agent reassignment worker
|
||||
Towerops.Workers.AgentLatencyEvaluator
|
||||
Towerops.Devices.EventLogger
|
||||
# Note: NeighborCleanupWorker, StaleAgentWorker, AgentLatencyEvaluator,
|
||||
# and JobHealthCheckWorker are now Oban Cron jobs (see config/runtime.exs)
|
||||
]
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -1,37 +0,0 @@
|
|||
defmodule Towerops.Monitoring.Supervisor do
|
||||
@moduledoc """
|
||||
Supervisor for monitoring-related background workers.
|
||||
|
||||
With the migration to Oban, device monitoring and polling are handled
|
||||
by Oban workers. This supervisor only manages auxiliary services like
|
||||
the neighbor cleanup worker.
|
||||
"""
|
||||
use Supervisor
|
||||
|
||||
alias Towerops.Snmp.NeighborCleanupWorker
|
||||
|
||||
require Logger
|
||||
|
||||
def start_link(init_arg) do
|
||||
Supervisor.start_link(__MODULE__, init_arg, name: __MODULE__)
|
||||
end
|
||||
|
||||
@impl true
|
||||
def init(_init_arg) do
|
||||
# Only start cleanup worker in production
|
||||
children =
|
||||
if test_mode?() do
|
||||
[]
|
||||
else
|
||||
[NeighborCleanupWorker]
|
||||
end
|
||||
|
||||
Supervisor.init(children, strategy: :one_for_one)
|
||||
end
|
||||
|
||||
# Check if running in test mode by inspecting the Repo pool configuration
|
||||
defp test_mode? do
|
||||
config = Application.get_env(:towerops, Towerops.Repo, [])
|
||||
Keyword.get(config, :pool) == Ecto.Adapters.SQL.Sandbox
|
||||
end
|
||||
end
|
||||
|
|
@ -1,40 +1,27 @@
|
|||
defmodule Towerops.Snmp.NeighborCleanupWorker do
|
||||
@moduledoc """
|
||||
GenServer that periodically cleans up stale neighbor, ARP, and MAC FDB records.
|
||||
Oban worker that periodically cleans up stale neighbor, ARP, and MAC FDB records.
|
||||
|
||||
Neighbors, ARP entries, and MAC forwarding database entries that haven't been
|
||||
seen in the last 24 hours are considered stale and will be automatically
|
||||
removed from the database.
|
||||
|
||||
Runs hourly via Oban.Plugins.Cron (configured in application config).
|
||||
"""
|
||||
use GenServer
|
||||
use Oban.Worker, queue: :maintenance
|
||||
|
||||
alias Towerops.Devices
|
||||
alias Towerops.Snmp
|
||||
|
||||
require Logger
|
||||
|
||||
# Check for stale neighbors every hour
|
||||
@cleanup_interval to_timeout(hour: 1)
|
||||
|
||||
# Consider neighbors stale if not seen in 24 hours
|
||||
@stale_threshold_hours 24
|
||||
|
||||
def start_link(opts \\ []) do
|
||||
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
|
||||
end
|
||||
|
||||
@impl true
|
||||
def init(_opts) do
|
||||
# Schedule first cleanup shortly after startup
|
||||
schedule_cleanup(10_000)
|
||||
{:ok, %{}}
|
||||
end
|
||||
|
||||
@impl true
|
||||
def handle_info(:cleanup, state) do
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
cleanup_stale_records()
|
||||
schedule_cleanup(@cleanup_interval)
|
||||
{:noreply, state}
|
||||
:ok
|
||||
end
|
||||
|
||||
defp cleanup_stale_records do
|
||||
|
|
@ -79,8 +66,4 @@ defmodule Towerops.Snmp.NeighborCleanupWorker do
|
|||
|
||||
:ok
|
||||
end
|
||||
|
||||
defp schedule_cleanup(delay) do
|
||||
Process.send_after(self(), :cleanup, delay)
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -1,8 +1,8 @@
|
|||
defmodule Towerops.Workers.AgentLatencyEvaluator do
|
||||
@moduledoc """
|
||||
GenServer that periodically evaluates device-to-agent latency and performs automatic reassignments.
|
||||
Oban worker that periodically evaluates device-to-agent latency and performs automatic reassignments.
|
||||
|
||||
Runs every 5 minutes to:
|
||||
Runs every 5 minutes via Oban.Plugins.Cron to:
|
||||
1. Find devices where an alternative agent has significantly better latency (20%+ improvement)
|
||||
2. Reassign devices based on their assignment source (site, organization, global, or device-level)
|
||||
3. Broadcast PubSub events for UI updates
|
||||
|
|
@ -11,7 +11,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluator do
|
|||
Only considers devices with automatic assignments (site, organization, global, or none).
|
||||
Manual device-level assignments are never automatically changed.
|
||||
"""
|
||||
use GenServer
|
||||
use Oban.Worker, queue: :maintenance
|
||||
|
||||
alias Towerops.Agents
|
||||
alias Towerops.Agents.Stats
|
||||
|
|
@ -22,49 +22,14 @@ defmodule Towerops.Workers.AgentLatencyEvaluator do
|
|||
|
||||
require Logger
|
||||
|
||||
# Run latency evaluation every 5 minutes
|
||||
@check_interval to_timeout(minute: 5)
|
||||
|
||||
# Minimum improvement required for reassignment (20%)
|
||||
@min_improvement_percent 20
|
||||
|
||||
# Minimum checks per agent for statistical validity
|
||||
@min_checks_per_agent 10
|
||||
|
||||
def start_link(opts \\ []) do
|
||||
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
|
||||
end
|
||||
|
||||
@doc """
|
||||
Manually trigger a latency evaluation.
|
||||
|
||||
Useful for testing and manual triggering from admin UI.
|
||||
"""
|
||||
def trigger_evaluation do
|
||||
GenServer.cast(__MODULE__, :evaluate_latency)
|
||||
end
|
||||
|
||||
@impl true
|
||||
def init(_opts) do
|
||||
# Schedule first evaluation shortly after startup (30 seconds)
|
||||
schedule_evaluation(30_000)
|
||||
{:ok, %{last_evaluation_at: nil, total_reassignments: 0}}
|
||||
end
|
||||
|
||||
@impl true
|
||||
def handle_cast(:evaluate_latency, state) do
|
||||
new_state = evaluate_and_reassign(state)
|
||||
{:noreply, new_state}
|
||||
end
|
||||
|
||||
@impl true
|
||||
def handle_info(:evaluate_latency, state) do
|
||||
new_state = evaluate_and_reassign(state)
|
||||
schedule_evaluation(@check_interval)
|
||||
{:noreply, new_state}
|
||||
end
|
||||
|
||||
defp evaluate_and_reassign(state) do
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
start_time = System.monotonic_time(:millisecond)
|
||||
|
||||
candidates =
|
||||
|
|
@ -85,11 +50,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluator do
|
|||
"#{reassignment_count} device(s) reassigned in #{duration_ms}ms"
|
||||
)
|
||||
|
||||
%{
|
||||
state
|
||||
| last_evaluation_at: DateTime.utc_now(),
|
||||
total_reassignments: state.total_reassignments + reassignment_count
|
||||
}
|
||||
:ok
|
||||
end
|
||||
|
||||
defp reassign_device(candidate) do
|
||||
|
|
@ -166,8 +127,4 @@ defmodule Towerops.Workers.AgentLatencyEvaluator do
|
|||
{:device_reassigned, candidate}
|
||||
)
|
||||
end
|
||||
|
||||
defp schedule_evaluation(delay) do
|
||||
Process.send_after(self(), :evaluate_latency, delay)
|
||||
end
|
||||
end
|
||||
|
|
|
|||
99
lib/towerops/workers/job_health_check_worker.ex
Normal file
99
lib/towerops/workers/job_health_check_worker.ex
Normal file
|
|
@ -0,0 +1,99 @@
|
|||
defmodule Towerops.Workers.JobHealthCheckWorker do
|
||||
@moduledoc """
|
||||
Oban worker that ensures all enabled devices have active monitoring and polling jobs.
|
||||
|
||||
Runs every 10 minutes via Oban.Plugins.Cron as a safety net to:
|
||||
1. Find devices with monitoring_enabled=true but no active DeviceMonitorWorker job
|
||||
2. Find devices with snmp_enabled=true but no active DevicePollerWorker job
|
||||
3. Automatically create missing jobs
|
||||
4. Log when missing jobs are found and recovered
|
||||
|
||||
This handles edge cases where jobs might get cancelled unexpectedly or database
|
||||
inconsistencies occur during deployments or pod failures.
|
||||
"""
|
||||
use Oban.Worker, queue: :maintenance
|
||||
|
||||
import Ecto.Query
|
||||
|
||||
alias Towerops.Devices
|
||||
alias Towerops.Repo
|
||||
alias Towerops.Workers.DeviceMonitorWorker
|
||||
alias Towerops.Workers.DevicePollerWorker
|
||||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
monitor_recoveries = recover_missing_monitor_jobs()
|
||||
poller_recoveries = recover_missing_poller_jobs()
|
||||
|
||||
total_recoveries = monitor_recoveries + poller_recoveries
|
||||
|
||||
if total_recoveries > 0 do
|
||||
Logger.warning(
|
||||
"Job health check recovered #{total_recoveries} missing job(s): " <>
|
||||
"#{monitor_recoveries} monitor(s), #{poller_recoveries} poller(s)"
|
||||
)
|
||||
else
|
||||
Logger.debug("Job health check completed: all jobs healthy")
|
||||
end
|
||||
|
||||
:ok
|
||||
end
|
||||
|
||||
defp recover_missing_monitor_jobs do
|
||||
devices_with_monitoring = Devices.list_monitored_devices()
|
||||
active_monitor_job_device_ids = get_active_monitor_job_device_ids()
|
||||
|
||||
missing_monitor_devices = Enum.reject(devices_with_monitoring, &MapSet.member?(active_monitor_job_device_ids, &1.id))
|
||||
|
||||
Enum.each(missing_monitor_devices, fn device ->
|
||||
Logger.warning(
|
||||
"Device '#{device.name}' has monitoring_enabled=true but no active monitor job, creating job",
|
||||
device_id: device.id,
|
||||
device_name: device.name
|
||||
)
|
||||
|
||||
DeviceMonitorWorker.start_monitoring(device.id)
|
||||
end)
|
||||
|
||||
length(missing_monitor_devices)
|
||||
end
|
||||
|
||||
defp recover_missing_poller_jobs do
|
||||
devices_with_snmp = Devices.list_snmp_enabled_devices()
|
||||
active_poller_job_device_ids = get_active_poller_job_device_ids()
|
||||
|
||||
missing_poller_devices = Enum.reject(devices_with_snmp, &MapSet.member?(active_poller_job_device_ids, &1.id))
|
||||
|
||||
Enum.each(missing_poller_devices, fn device ->
|
||||
Logger.warning(
|
||||
"Device '#{device.name}' has snmp_enabled=true but no active poller job, creating job",
|
||||
device_id: device.id,
|
||||
device_name: device.name
|
||||
)
|
||||
|
||||
DevicePollerWorker.start_polling(device.id)
|
||||
end)
|
||||
|
||||
length(missing_poller_devices)
|
||||
end
|
||||
|
||||
defp get_active_monitor_job_device_ids do
|
||||
Oban.Job
|
||||
|> where([j], j.worker == "Towerops.Workers.DeviceMonitorWorker")
|
||||
|> where([j], j.state in ["available", "scheduled", "executing", "retryable"])
|
||||
|> select([j], fragment("args->>'device_id'"))
|
||||
|> Repo.all()
|
||||
|> MapSet.new()
|
||||
end
|
||||
|
||||
defp get_active_poller_job_device_ids do
|
||||
Oban.Job
|
||||
|> where([j], j.worker == "Towerops.Workers.DevicePollerWorker")
|
||||
|> where([j], j.state in ["available", "scheduled", "executing", "retryable"])
|
||||
|> select([j], fragment("args->>'device_id'"))
|
||||
|> Repo.all()
|
||||
|> MapSet.new()
|
||||
end
|
||||
end
|
||||
|
|
@ -1,8 +1,8 @@
|
|||
defmodule Towerops.Workers.StaleAgentWorker do
|
||||
@moduledoc """
|
||||
GenServer that periodically checks for stale agents.
|
||||
Oban worker that periodically checks for stale agents.
|
||||
|
||||
Runs every minute to:
|
||||
Runs every minute via Oban.Plugins.Cron to:
|
||||
1. Find agents that haven't sent a heartbeat in over 10 minutes
|
||||
2. Log warnings for operators
|
||||
3. Broadcast alerts via PubSub for UI updates
|
||||
|
|
@ -12,7 +12,7 @@ defmodule Towerops.Workers.StaleAgentWorker do
|
|||
- They have checked in at least once (last_seen_at is not nil)
|
||||
- They haven't checked in for more than 10 minutes
|
||||
"""
|
||||
use GenServer
|
||||
use Oban.Worker, queue: :maintenance
|
||||
|
||||
import Ecto.Query
|
||||
|
||||
|
|
@ -21,55 +21,43 @@ defmodule Towerops.Workers.StaleAgentWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
# Check for stale agents every minute
|
||||
@check_interval to_timeout(minute: 1)
|
||||
|
||||
# Consider agents stale if not seen in 10 minutes
|
||||
@stale_threshold_minutes 10
|
||||
|
||||
def start_link(opts \\ []) do
|
||||
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
|
||||
end
|
||||
# Consider agents "newly stale" if not seen in 10-15 minutes (for warning logs)
|
||||
@newly_stale_threshold_minutes 15
|
||||
|
||||
@impl true
|
||||
def init(_opts) do
|
||||
# Schedule first check shortly after startup
|
||||
schedule_check(10_000)
|
||||
{:ok, %{last_stale_ids: MapSet.new()}}
|
||||
end
|
||||
|
||||
@impl true
|
||||
def handle_info(:check_stale, state) do
|
||||
new_state = check_for_stale_agents(state)
|
||||
schedule_check(@check_interval)
|
||||
{:noreply, new_state}
|
||||
end
|
||||
|
||||
defp check_for_stale_agents(state) do
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
stale_agents = find_stale_agents()
|
||||
current_stale_ids = MapSet.new(stale_agents, & &1.id)
|
||||
|
||||
# Log warnings for newly stale agents (not already reported)
|
||||
newly_stale = Enum.reject(stale_agents, &MapSet.member?(state.last_stale_ids, &1.id))
|
||||
|
||||
if Enum.any?(newly_stale) do
|
||||
Enum.each(newly_stale, &log_stale_agent/1)
|
||||
|
||||
Logger.warning("Detected #{length(newly_stale)} newly stale agent(s)")
|
||||
|
||||
# Broadcast for any UI listeners
|
||||
Phoenix.PubSub.broadcast(Towerops.PubSub, "agents:health", {:agents_stale, newly_stale})
|
||||
if !Enum.empty?(stale_agents) do
|
||||
process_stale_agents(stale_agents)
|
||||
end
|
||||
|
||||
# Track recovered agents (were stale, now healthy)
|
||||
recovered_ids = MapSet.difference(state.last_stale_ids, current_stale_ids)
|
||||
:ok
|
||||
end
|
||||
|
||||
if MapSet.size(recovered_ids) > 0 do
|
||||
Logger.info("#{MapSet.size(recovered_ids)} agent(s) recovered from stale state")
|
||||
Phoenix.PubSub.broadcast(Towerops.PubSub, "agents:health", {:agents_recovered, recovered_ids})
|
||||
end
|
||||
defp process_stale_agents(stale_agents) do
|
||||
# Separate newly stale (10-15 min) from long-term stale (>15 min)
|
||||
{newly_stale, long_term_stale} =
|
||||
Enum.split_with(stale_agents, &newly_stale?/1)
|
||||
|
||||
%{state | last_stale_ids: current_stale_ids}
|
||||
# Log warnings for newly stale agents
|
||||
log_newly_stale_agents(newly_stale)
|
||||
|
||||
# Log debug for long-term stale agents (reduce log noise)
|
||||
Enum.each(long_term_stale, fn agent -> log_stale_agent(agent, :debug) end)
|
||||
|
||||
# Broadcast all stale agents for UI updates
|
||||
Phoenix.PubSub.broadcast(Towerops.PubSub, "agents:health", {:agents_stale, stale_agents})
|
||||
end
|
||||
|
||||
defp log_newly_stale_agents([]), do: :ok
|
||||
|
||||
defp log_newly_stale_agents(newly_stale) do
|
||||
Enum.each(newly_stale, fn agent -> log_stale_agent(agent, :warning) end)
|
||||
Logger.warning("Detected #{length(newly_stale)} newly stale agent(s)")
|
||||
end
|
||||
|
||||
@doc """
|
||||
|
|
@ -88,24 +76,31 @@ defmodule Towerops.Workers.StaleAgentWorker do
|
|||
|> Repo.all()
|
||||
end
|
||||
|
||||
defp log_stale_agent(agent) do
|
||||
defp newly_stale?(agent) do
|
||||
newly_stale_cutoff = DateTime.add(DateTime.utc_now(), -@newly_stale_threshold_minutes, :minute)
|
||||
DateTime.after?(agent.last_seen_at, newly_stale_cutoff)
|
||||
end
|
||||
|
||||
defp log_stale_agent(agent, level) do
|
||||
minutes_since_seen =
|
||||
DateTime.utc_now()
|
||||
|> DateTime.diff(agent.last_seen_at, :second)
|
||||
|> div(60)
|
||||
|
||||
Logger.warning(
|
||||
"Agent '#{agent.name}' is stale - last seen #{minutes_since_seen} minutes ago",
|
||||
message = "Agent '#{agent.name}' is stale - last seen #{minutes_since_seen} minutes ago"
|
||||
|
||||
metadata = [
|
||||
agent_token_id: agent.id,
|
||||
agent_name: agent.name,
|
||||
organization_id: agent.organization_id,
|
||||
last_seen_at: agent.last_seen_at,
|
||||
last_ip: agent.last_ip,
|
||||
minutes_since_seen: minutes_since_seen
|
||||
)
|
||||
end
|
||||
]
|
||||
|
||||
defp schedule_check(delay) do
|
||||
Process.send_after(self(), :check_stale, delay)
|
||||
case level do
|
||||
:warning -> Logger.warning(message, metadata)
|
||||
:debug -> Logger.debug(message, metadata)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -95,12 +95,6 @@ defmodule Towerops.Snmp.NeighborCleanupWorkerTest do
|
|||
end
|
||||
|
||||
describe "cleanup process" do
|
||||
test "worker starts successfully" do
|
||||
{:ok, pid} = NeighborCleanupWorker.start_link()
|
||||
assert Process.alive?(pid)
|
||||
GenServer.stop(pid)
|
||||
end
|
||||
|
||||
test "cleanup removes stale neighbors across all devices", %{
|
||||
device_schema1: device_schema1,
|
||||
device_schema2: device_schema2,
|
||||
|
|
@ -146,12 +140,8 @@ defmodule Towerops.Snmp.NeighborCleanupWorkerTest do
|
|||
last_discovered_at: one_hour_ago
|
||||
})
|
||||
|
||||
# Start worker and trigger cleanup immediately
|
||||
{:ok, pid} = NeighborCleanupWorker.start_link()
|
||||
send(pid, :cleanup)
|
||||
|
||||
# Wait for cleanup to complete
|
||||
Process.sleep(100)
|
||||
# Trigger cleanup via Oban worker
|
||||
assert :ok = NeighborCleanupWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Old neighbors should be deleted
|
||||
assert Repo.get(Neighbor, old_neighbor1.id) == nil
|
||||
|
|
@ -159,22 +149,14 @@ defmodule Towerops.Snmp.NeighborCleanupWorkerTest do
|
|||
|
||||
# Recent neighbor should still exist
|
||||
assert Repo.get(Neighbor, recent_neighbor.id)
|
||||
|
||||
GenServer.stop(pid)
|
||||
end
|
||||
|
||||
test "cleanup handles device with no neighbors gracefully", %{device1: device1} do
|
||||
# Ensure no neighbors exist
|
||||
assert Snmp.list_neighbors(device1.id) == []
|
||||
|
||||
{:ok, pid} = NeighborCleanupWorker.start_link()
|
||||
send(pid, :cleanup)
|
||||
|
||||
# Should not crash
|
||||
Process.sleep(100)
|
||||
assert Process.alive?(pid)
|
||||
|
||||
GenServer.stop(pid)
|
||||
assert :ok = NeighborCleanupWorker.perform(%Oban.Job{args: %{}})
|
||||
end
|
||||
|
||||
test "cleanup preserves neighbors within 24 hour threshold", %{
|
||||
|
|
@ -209,18 +191,13 @@ defmodule Towerops.Snmp.NeighborCleanupWorkerTest do
|
|||
last_discovered_at: twenty_five_hours_ago
|
||||
})
|
||||
|
||||
{:ok, pid} = NeighborCleanupWorker.start_link()
|
||||
send(pid, :cleanup)
|
||||
|
||||
Process.sleep(100)
|
||||
assert :ok = NeighborCleanupWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Recent neighbor should still exist
|
||||
assert Repo.get(Neighbor, recent_neighbor.id)
|
||||
|
||||
# Stale neighbor should be deleted
|
||||
assert Repo.get(Neighbor, stale_neighbor.id) == nil
|
||||
|
||||
GenServer.stop(pid)
|
||||
end
|
||||
|
||||
test "cleanup does not affect device without SNMP enabled" do
|
||||
|
|
@ -274,16 +251,11 @@ defmodule Towerops.Snmp.NeighborCleanupWorkerTest do
|
|||
last_discovered_at: two_days_ago
|
||||
})
|
||||
|
||||
{:ok, pid} = NeighborCleanupWorker.start_link()
|
||||
send(pid, :cleanup)
|
||||
|
||||
Process.sleep(100)
|
||||
assert :ok = NeighborCleanupWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Neighbor should still exist (device not included in cleanup)
|
||||
# because cleanup only processes SNMP-enabled devices
|
||||
assert Repo.get(Neighbor, neighbor.id)
|
||||
|
||||
GenServer.stop(pid)
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -333,10 +305,7 @@ defmodule Towerops.Snmp.NeighborCleanupWorkerTest do
|
|||
last_seen_at: one_hour_ago
|
||||
})
|
||||
|
||||
{:ok, pid} = NeighborCleanupWorker.start_link()
|
||||
send(pid, :cleanup)
|
||||
|
||||
Process.sleep(100)
|
||||
assert :ok = NeighborCleanupWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Old ARP entries should be deleted
|
||||
assert Repo.get(ArpEntry, old_arp1.id) == nil
|
||||
|
|
@ -344,8 +313,6 @@ defmodule Towerops.Snmp.NeighborCleanupWorkerTest do
|
|||
|
||||
# Recent ARP entry should still exist
|
||||
assert Repo.get(ArpEntry, recent_arp.id)
|
||||
|
||||
GenServer.stop(pid)
|
||||
end
|
||||
|
||||
test "cleanup preserves ARP entries within 24 hour threshold", %{
|
||||
|
|
@ -380,36 +347,13 @@ defmodule Towerops.Snmp.NeighborCleanupWorkerTest do
|
|||
last_seen_at: twenty_five_hours_ago
|
||||
})
|
||||
|
||||
{:ok, pid} = NeighborCleanupWorker.start_link()
|
||||
send(pid, :cleanup)
|
||||
|
||||
Process.sleep(100)
|
||||
assert :ok = NeighborCleanupWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Recent ARP entry should still exist
|
||||
assert Repo.get(ArpEntry, recent_arp.id)
|
||||
|
||||
# Stale ARP entry should be deleted
|
||||
assert Repo.get(ArpEntry, stale_arp.id) == nil
|
||||
|
||||
GenServer.stop(pid)
|
||||
end
|
||||
end
|
||||
|
||||
describe "scheduled cleanup" do
|
||||
test "worker schedules periodic cleanup" do
|
||||
{:ok, pid} = NeighborCleanupWorker.start_link()
|
||||
|
||||
# Check that worker is alive and has scheduled a cleanup message
|
||||
assert Process.alive?(pid)
|
||||
|
||||
# Get process info
|
||||
{:messages, _messages} = Process.info(pid, :messages)
|
||||
|
||||
# Should eventually have a :cleanup message scheduled
|
||||
# (may take a moment for the initial schedule)
|
||||
Process.sleep(20)
|
||||
|
||||
GenServer.stop(pid)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -9,13 +9,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
alias Towerops.Sites
|
||||
alias Towerops.Workers.AgentLatencyEvaluator
|
||||
|
||||
setup do
|
||||
# Start the worker for testing
|
||||
{:ok, pid} = start_supervised(AgentLatencyEvaluator)
|
||||
%{worker_pid: pid}
|
||||
end
|
||||
|
||||
describe "evaluate_latency/0" do
|
||||
describe "perform/1" do
|
||||
test "reassigns device to faster agent when 20%+ improvement available" do
|
||||
user = user_fixture()
|
||||
org = organization_fixture(user.id)
|
||||
|
|
@ -59,10 +53,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
end
|
||||
|
||||
# Trigger evaluation
|
||||
AgentLatencyEvaluator.trigger_evaluation()
|
||||
|
||||
# Wait for async processing
|
||||
Process.sleep(100)
|
||||
assert :ok = AgentLatencyEvaluator.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Verify site was reassigned
|
||||
updated_site = Sites.get_site!(site.id)
|
||||
|
|
@ -112,10 +103,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
end
|
||||
|
||||
# Trigger evaluation
|
||||
AgentLatencyEvaluator.trigger_evaluation()
|
||||
|
||||
# Wait for async processing
|
||||
Process.sleep(100)
|
||||
assert :ok = AgentLatencyEvaluator.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Verify site was NOT reassigned
|
||||
updated_site = Sites.get_site!(site.id)
|
||||
|
|
@ -165,10 +153,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
end
|
||||
|
||||
# Trigger evaluation
|
||||
AgentLatencyEvaluator.trigger_evaluation()
|
||||
|
||||
# Wait for async processing
|
||||
Process.sleep(100)
|
||||
assert :ok = AgentLatencyEvaluator.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Verify organization was reassigned
|
||||
updated_org = Towerops.Organizations.get_organization!(org.id)
|
||||
|
|
@ -217,10 +202,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
end
|
||||
|
||||
# Trigger evaluation
|
||||
AgentLatencyEvaluator.trigger_evaluation()
|
||||
|
||||
# Wait for async processing
|
||||
Process.sleep(100)
|
||||
assert :ok = AgentLatencyEvaluator.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Verify device assignment was created
|
||||
assignment = Agents.get_device_assignment(device.id)
|
||||
|
|
@ -281,10 +263,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
Repo.insert_all(Towerops.Monitoring.Check, checks)
|
||||
|
||||
# Trigger evaluation
|
||||
AgentLatencyEvaluator.trigger_evaluation()
|
||||
|
||||
# Wait for async processing
|
||||
Process.sleep(100)
|
||||
assert :ok = AgentLatencyEvaluator.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Verify device assignment was created to the fast agent
|
||||
assignment = Agents.get_device_assignment(device.id)
|
||||
|
|
@ -333,10 +312,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
end
|
||||
|
||||
# Trigger evaluation
|
||||
AgentLatencyEvaluator.trigger_evaluation()
|
||||
|
||||
# Wait for async processing
|
||||
Process.sleep(100)
|
||||
assert :ok = AgentLatencyEvaluator.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Verify device assignment was NOT changed
|
||||
assignment = Agents.get_device_assignment(device.id)
|
||||
|
|
@ -386,10 +362,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
end
|
||||
|
||||
# Trigger evaluation
|
||||
AgentLatencyEvaluator.trigger_evaluation()
|
||||
|
||||
# Wait for async processing
|
||||
Process.sleep(100)
|
||||
assert :ok = AgentLatencyEvaluator.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Verify site was NOT reassigned (insufficient data)
|
||||
updated_site = Sites.get_site!(site.id)
|
||||
|
|
@ -458,10 +431,7 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
end
|
||||
|
||||
# Trigger evaluation
|
||||
AgentLatencyEvaluator.trigger_evaluation()
|
||||
|
||||
# Wait for async processing
|
||||
Process.sleep(100)
|
||||
assert :ok = AgentLatencyEvaluator.perform(%Oban.Job{args: %{}})
|
||||
|
||||
# Verify both sites were reassigned
|
||||
updated_site1 = Sites.get_site!(site1.id)
|
||||
|
|
@ -471,59 +441,4 @@ defmodule Towerops.Workers.AgentLatencyEvaluatorTest do
|
|||
assert updated_site2.agent_token_id == fast_agent.id
|
||||
end
|
||||
end
|
||||
|
||||
describe "handle_info(:evaluate_latency)" do
|
||||
test "worker processes evaluation on schedule" do
|
||||
user = user_fixture()
|
||||
org = organization_fixture(user.id)
|
||||
{:ok, slow_agent, _} = Agents.create_agent_token(org.id, "Slow Agent")
|
||||
{:ok, fast_agent, _} = Agents.create_agent_token(org.id, "Fast Agent")
|
||||
|
||||
{:ok, site} =
|
||||
Sites.create_site(%{
|
||||
name: "Test Site",
|
||||
organization_id: org.id,
|
||||
agent_token_id: slow_agent.id
|
||||
})
|
||||
|
||||
{:ok, device} =
|
||||
Towerops.Devices.create_device(%{
|
||||
name: "Router",
|
||||
ip_address: "192.168.1.1",
|
||||
site_id: site.id,
|
||||
snmp_enabled: true,
|
||||
snmp_community: "public",
|
||||
snmp_version: "2c"
|
||||
})
|
||||
|
||||
# Create checks
|
||||
for _ <- 1..15 do
|
||||
Monitoring.create_check(%{
|
||||
device_id: device.id,
|
||||
agent_token_id: slow_agent.id,
|
||||
status: :success,
|
||||
response_time_ms: 100,
|
||||
checked_at: DateTime.utc_now()
|
||||
})
|
||||
|
||||
Monitoring.create_check(%{
|
||||
device_id: device.id,
|
||||
agent_token_id: fast_agent.id,
|
||||
status: :success,
|
||||
response_time_ms: 50,
|
||||
checked_at: DateTime.utc_now()
|
||||
})
|
||||
end
|
||||
|
||||
# Send the message directly to the worker
|
||||
send(Process.whereis(AgentLatencyEvaluator), :evaluate_latency)
|
||||
|
||||
# Wait for async processing
|
||||
Process.sleep(100)
|
||||
|
||||
# Verify site was reassigned
|
||||
updated_site = Sites.get_site!(site.id)
|
||||
assert updated_site.agent_token_id == fast_agent.id
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue