feat: add failed jobs and recent completions queries
Add list_failed_jobs/0 to return retryable and cancelled jobs and list_recent_completions/1 to return completed, cancelled, or discarded jobs with configurable limit. Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
This commit is contained in:
parent
a3cbf7f412
commit
adbaeef83d
2 changed files with 92 additions and 1 deletions
|
|
@ -7,8 +7,8 @@ defmodule Towerops.JobMonitoring do
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import Ecto.Query
|
import Ecto.Query
|
||||||
alias Towerops.Repo
|
|
||||||
alias Oban.Job
|
alias Oban.Job
|
||||||
|
alias Towerops.Repo
|
||||||
|
|
||||||
@worker_names [
|
@worker_names [
|
||||||
"Towerops.Workers.DevicePollerWorker",
|
"Towerops.Workers.DevicePollerWorker",
|
||||||
|
|
@ -53,4 +53,37 @@ defmodule Towerops.JobMonitoring do
|
||||||
)
|
)
|
||||||
|> Repo.all()
|
|> Repo.all()
|
||||||
end
|
end
|
||||||
|
|
||||||
|
@doc """
|
||||||
|
Lists jobs that have failed and are retrying or have been cancelled.
|
||||||
|
|
||||||
|
Limited to last 100 jobs.
|
||||||
|
"""
|
||||||
|
@spec list_failed_jobs() :: [Job.t()]
|
||||||
|
def list_failed_jobs do
|
||||||
|
from(j in Job,
|
||||||
|
where: j.state in ["retryable", "cancelled"],
|
||||||
|
where: j.worker in ^@worker_names,
|
||||||
|
order_by: [desc: j.inserted_at],
|
||||||
|
limit: 100
|
||||||
|
)
|
||||||
|
|> Repo.all()
|
||||||
|
end
|
||||||
|
|
||||||
|
@doc """
|
||||||
|
Lists recently completed jobs (completed, cancelled, or discarded).
|
||||||
|
|
||||||
|
## Options
|
||||||
|
- `:limit` - Maximum number of jobs to return (default: 100)
|
||||||
|
"""
|
||||||
|
@spec list_recent_completions(integer()) :: [Job.t()]
|
||||||
|
def list_recent_completions(limit \\ 100) do
|
||||||
|
from(j in Job,
|
||||||
|
where: j.state in ["completed", "cancelled", "discarded"],
|
||||||
|
where: j.worker in ^@worker_names,
|
||||||
|
order_by: [desc: j.completed_at],
|
||||||
|
limit: ^limit
|
||||||
|
)
|
||||||
|
|> Repo.all()
|
||||||
|
end
|
||||||
end
|
end
|
||||||
|
|
|
||||||
|
|
@ -202,4 +202,62 @@ defmodule Towerops.JobMonitoringTest do
|
||||||
assert stuck |> Enum.map(& &1.id) |> Enum.sort() == [stuck_discovery.id, stuck_poller.id] |> Enum.sort()
|
assert stuck |> Enum.map(& &1.id) |> Enum.sort() == [stuck_discovery.id, stuck_poller.id] |> Enum.sort()
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
|
describe "list_failed_jobs/0" do
|
||||||
|
test "returns retryable and cancelled jobs" do
|
||||||
|
device = device_fixture()
|
||||||
|
|
||||||
|
oban_job_fixture(%{
|
||||||
|
worker: "Towerops.Workers.DevicePollerWorker",
|
||||||
|
state: "retryable",
|
||||||
|
args: %{"device_id" => device.id}
|
||||||
|
})
|
||||||
|
|
||||||
|
oban_job_fixture(%{
|
||||||
|
worker: "Towerops.Workers.DiscoveryWorker",
|
||||||
|
state: "cancelled",
|
||||||
|
args: %{"device_id" => device.id}
|
||||||
|
})
|
||||||
|
|
||||||
|
# Should not include completed
|
||||||
|
oban_job_fixture(%{
|
||||||
|
worker: "Towerops.Workers.DevicePollerWorker",
|
||||||
|
state: "completed",
|
||||||
|
args: %{"device_id" => device.id}
|
||||||
|
})
|
||||||
|
|
||||||
|
failed = JobMonitoring.list_failed_jobs()
|
||||||
|
|
||||||
|
assert length(failed) == 2
|
||||||
|
assert Enum.all?(failed, fn j -> j.state in ["retryable", "cancelled"] end)
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
|
describe "list_recent_completions/1" do
|
||||||
|
test "returns completed, cancelled, and discarded jobs up to limit" do
|
||||||
|
device = device_fixture()
|
||||||
|
|
||||||
|
# Create 5 completed jobs
|
||||||
|
for _ <- 1..5 do
|
||||||
|
oban_job_fixture(%{
|
||||||
|
worker: "Towerops.Workers.DevicePollerWorker",
|
||||||
|
state: "completed",
|
||||||
|
args: %{"device_id" => device.id},
|
||||||
|
completed_at: DateTime.utc_now()
|
||||||
|
})
|
||||||
|
end
|
||||||
|
|
||||||
|
# Should not include executing
|
||||||
|
oban_job_fixture(%{
|
||||||
|
worker: "Towerops.Workers.DevicePollerWorker",
|
||||||
|
state: "executing",
|
||||||
|
args: %{"device_id" => device.id}
|
||||||
|
})
|
||||||
|
|
||||||
|
recent = JobMonitoring.list_recent_completions(3)
|
||||||
|
|
||||||
|
assert length(recent) == 3
|
||||||
|
assert Enum.all?(recent, fn j -> j.state == "completed" end)
|
||||||
|
end
|
||||||
|
end
|
||||||
end
|
end
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue