From 35c73288315caa8f02b417dfdf30c9cd4bb88d8a Mon Sep 17 00:00:00 2001 From: mayor Date: Fri, 6 Feb 2026 16:15:18 -0600 Subject: [PATCH] feat: add JobMonitoring context with list_active_jobs query Co-Authored-By: Claude Sonnet 4.5 --- lib/towerops/job_monitoring.ex | 29 ++++++++++++++++++ test/support/fixtures/jobs_fixtures.ex | 35 +++++++++++++++++++++ test/towerops/job_monitoring_test.exs | 42 ++++++++++++++++++++++++++ 3 files changed, 106 insertions(+) create mode 100644 lib/towerops/job_monitoring.ex create mode 100644 test/support/fixtures/jobs_fixtures.ex create mode 100644 test/towerops/job_monitoring_test.exs diff --git a/lib/towerops/job_monitoring.ex b/lib/towerops/job_monitoring.ex new file mode 100644 index 00000000..5d02f945 --- /dev/null +++ b/lib/towerops/job_monitoring.ex @@ -0,0 +1,29 @@ +defmodule Towerops.JobMonitoring do + @moduledoc """ + Context for monitoring Oban jobs, specifically polling and discovery operations. + + Provides queries for active, stuck, failed, and completed jobs with device + and agent context for operational monitoring. + """ + + import Ecto.Query + alias Towerops.Repo + alias Oban.Job + + @worker_names [ + "Towerops.Workers.DevicePollerWorker", + "Towerops.Workers.DiscoveryWorker" + ] + + @doc """ + Lists all currently executing polling and discovery jobs. + """ + def list_active_jobs do + from(j in Job, + where: j.state == "executing", + where: j.worker in ^@worker_names, + order_by: [asc: j.attempted_at] + ) + |> Repo.all() + end +end diff --git a/test/support/fixtures/jobs_fixtures.ex b/test/support/fixtures/jobs_fixtures.ex new file mode 100644 index 00000000..d7d3e281 --- /dev/null +++ b/test/support/fixtures/jobs_fixtures.ex @@ -0,0 +1,35 @@ +defmodule Towerops.JobsFixtures do + @moduledoc """ + This module defines test helpers for creating Oban jobs for testing. + """ + + alias Towerops.Repo + + @doc """ + Generate an Oban job with all required fields. + """ + def oban_job_fixture(attrs \\ %{}) do + default_attrs = %{ + worker: "Towerops.Workers.DevicePollerWorker", + queue: "pollers", + args: %{"device_id" => 1}, + state: "available", + attempt: 0, + max_attempts: 3, + inserted_at: DateTime.utc_now(), + scheduled_at: DateTime.utc_now() + } + + merged_attrs = + Map.merge( + default_attrs, + attrs |> Map.new(fn {k, v} -> {to_atom_key(k), v} end) + ) + + job = struct!(Oban.Job, merged_attrs) + Repo.insert!(job) + end + + defp to_atom_key(key) when is_atom(key), do: key + defp to_atom_key(key) when is_binary(key), do: String.to_existing_atom(key) +end diff --git a/test/towerops/job_monitoring_test.exs b/test/towerops/job_monitoring_test.exs new file mode 100644 index 00000000..5a38ee51 --- /dev/null +++ b/test/towerops/job_monitoring_test.exs @@ -0,0 +1,42 @@ +defmodule Towerops.JobMonitoringTest do + use Towerops.DataCase, async: true + + import Towerops.DevicesFixtures + import Towerops.JobsFixtures + + alias Towerops.JobMonitoring + + describe "list_active_jobs/0" do + test "returns jobs in executing state for polling and discovery workers" do + device = device_fixture() + + # Create executing polling job + oban_job_fixture(%{ + worker: "Towerops.Workers.DevicePollerWorker", + state: "executing", + args: %{"device_id" => device.id}, + attempted_at: DateTime.utc_now() + }) + + # Create executing discovery job + oban_job_fixture(%{ + worker: "Towerops.Workers.DiscoveryWorker", + state: "executing", + args: %{"device_id" => device.id}, + attempted_at: DateTime.utc_now() + }) + + # Create completed job (should not be returned) + oban_job_fixture(%{ + worker: "Towerops.Workers.DevicePollerWorker", + state: "completed", + args: %{"device_id" => device.id} + }) + + jobs = JobMonitoring.list_active_jobs() + + assert length(jobs) == 2 + assert Enum.all?(jobs, fn j -> j.state == "executing" end) + end + end +end