diff --git a/lib/towerops/job_monitoring.ex b/lib/towerops/job_monitoring.ex index 456e4d2e..bf0fa861 100644 --- a/lib/towerops/job_monitoring.ex +++ b/lib/towerops/job_monitoring.ex @@ -15,6 +15,9 @@ defmodule Towerops.JobMonitoring do "Towerops.Workers.DiscoveryWorker" ] + @polling_threshold_seconds 120 # 2 minutes + @discovery_threshold_seconds 300 # 5 minutes + @doc """ Lists all currently executing polling and discovery jobs. """ @@ -27,4 +30,27 @@ defmodule Towerops.JobMonitoring do ) |> Repo.all() end + + @doc """ + Lists jobs that are executing longer than expected thresholds. + + Thresholds: + - Polling: 2 minutes + - Discovery: 5 minutes + """ + @spec list_stuck_jobs() :: [Job.t()] + def list_stuck_jobs do + now = DateTime.utc_now() + polling_threshold = DateTime.add(now, -@polling_threshold_seconds, :second) + discovery_threshold = DateTime.add(now, -@discovery_threshold_seconds, :second) + + from(j in Job, + where: j.state == "executing", + where: + (j.worker == "Towerops.Workers.DevicePollerWorker" and j.attempted_at < ^polling_threshold) or + (j.worker == "Towerops.Workers.DiscoveryWorker" and j.attempted_at < ^discovery_threshold), + order_by: [asc: j.attempted_at] + ) + |> Repo.all() + end end diff --git a/test/towerops/job_monitoring_test.exs b/test/towerops/job_monitoring_test.exs index 5824c2bb..c2ad6169 100644 --- a/test/towerops/job_monitoring_test.exs +++ b/test/towerops/job_monitoring_test.exs @@ -68,4 +68,56 @@ defmodule Towerops.JobMonitoringTest do assert hd(jobs).worker == "Towerops.Workers.DevicePollerWorker" end end + + describe "list_stuck_jobs/0" do + test "returns polling jobs executing longer than 2 minutes" do + device = device_fixture() + + # Stuck polling job (3 minutes ago) + stuck_job = oban_job_fixture(%{ + worker: "Towerops.Workers.DevicePollerWorker", + state: "executing", + args: %{"device_id" => device.id}, + attempted_at: DateTime.add(DateTime.utc_now(), -180, :second) + }) + + # Recent polling job (30 seconds ago) - not stuck + oban_job_fixture(%{ + worker: "Towerops.Workers.DevicePollerWorker", + state: "executing", + args: %{"device_id" => device.id}, + attempted_at: DateTime.add(DateTime.utc_now(), -30, :second) + }) + + stuck = JobMonitoring.list_stuck_jobs() + + assert length(stuck) == 1 + assert hd(stuck).id == stuck_job.id + end + + test "returns discovery jobs executing longer than 5 minutes" do + device = device_fixture() + + # Stuck discovery job (6 minutes ago) + stuck_job = oban_job_fixture(%{ + worker: "Towerops.Workers.DiscoveryWorker", + state: "executing", + args: %{"device_id" => device.id}, + attempted_at: DateTime.add(DateTime.utc_now(), -360, :second) + }) + + # Recent discovery job (3 minutes ago) - not stuck + oban_job_fixture(%{ + worker: "Towerops.Workers.DiscoveryWorker", + state: "executing", + args: %{"device_id" => device.id}, + attempted_at: DateTime.add(DateTime.utc_now(), -180, :second) + }) + + stuck = JobMonitoring.list_stuck_jobs() + + assert length(stuck) == 1 + assert hd(stuck).id == stuck_job.id + end + end end