diff --git a/CLAUDE.md b/CLAUDE.md index 59077d3f..56efd46d 100644 --- a/CLAUDE.md +++ b/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). diff --git a/config/dev.exs b/config/dev.exs index a08035ed..63874a0c 100644 --- a/config/dev.exs +++ b/config/dev.exs @@ -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) diff --git a/config/runtime.exs b/config/runtime.exs index a0205059..cc5e5ae3 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -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) diff --git a/lib/towerops/application.ex b/lib/towerops/application.ex index f1db35ec..78508179 100644 --- a/lib/towerops/application.ex +++ b/lib/towerops/application.ex @@ -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 diff --git a/lib/towerops/monitoring/supervisor.ex b/lib/towerops/monitoring/supervisor.ex deleted file mode 100644 index 48f76391..00000000 --- a/lib/towerops/monitoring/supervisor.ex +++ /dev/null @@ -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 diff --git a/lib/towerops/snmp/neighbor_cleanup_worker.ex b/lib/towerops/snmp/neighbor_cleanup_worker.ex index e147ac4c..c731014f 100644 --- a/lib/towerops/snmp/neighbor_cleanup_worker.ex +++ b/lib/towerops/snmp/neighbor_cleanup_worker.ex @@ -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 diff --git a/lib/towerops/workers/agent_latency_evaluator.ex b/lib/towerops/workers/agent_latency_evaluator.ex index 40001bc1..0d5510df 100644 --- a/lib/towerops/workers/agent_latency_evaluator.ex +++ b/lib/towerops/workers/agent_latency_evaluator.ex @@ -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 diff --git a/lib/towerops/workers/job_health_check_worker.ex b/lib/towerops/workers/job_health_check_worker.ex new file mode 100644 index 00000000..a86a136e --- /dev/null +++ b/lib/towerops/workers/job_health_check_worker.ex @@ -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 diff --git a/lib/towerops/workers/stale_agent_worker.ex b/lib/towerops/workers/stale_agent_worker.ex index 9d9a7105..e70458c4 100644 --- a/lib/towerops/workers/stale_agent_worker.ex +++ b/lib/towerops/workers/stale_agent_worker.ex @@ -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 diff --git a/test/towerops/snmp/neighbor_cleanup_worker_test.exs b/test/towerops/snmp/neighbor_cleanup_worker_test.exs index 923c6bc7..1d84c2c8 100644 --- a/test/towerops/snmp/neighbor_cleanup_worker_test.exs +++ b/test/towerops/snmp/neighbor_cleanup_worker_test.exs @@ -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 diff --git a/test/towerops/workers/agent_latency_evaluator_test.exs b/test/towerops/workers/agent_latency_evaluator_test.exs index c698fe47..072a0698 100644 --- a/test/towerops/workers/agent_latency_evaluator_test.exs +++ b/test/towerops/workers/agent_latency_evaluator_test.exs @@ -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