towerops/test/mix/tasks/oban.cancel_stuck_discovery_test.exs
2026-02-03 10:19:56 -06:00

206 lines
6.8 KiB
Elixir

defmodule Mix.Tasks.Oban.CancelStuckDiscoveryTest do
use Towerops.DataCase, async: true
import ExUnit.CaptureIO
import Mix.Tasks.Oban.CancelStuckDiscovery, only: [run: 1]
alias Towerops.Repo
describe "run/1" do
test "cancels stuck discovery jobs in scheduled state" do
# Insert a scheduled discovery job
{:ok, job} =
%{device_id: Ecto.UUID.generate()}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
assert job.state == "available"
# Change state to scheduled
job = Repo.update!(Ecto.Changeset.change(job, state: "scheduled"))
output = capture_io(fn -> run([]) end)
# Verify job was cancelled
updated_job = Repo.reload(job)
assert updated_job.state == "cancelled"
assert output =~ "Cancelled 1 stuck discovery jobs"
end
test "cancels stuck discovery jobs in retryable state" do
{:ok, job} =
%{device_id: Ecto.UUID.generate()}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
# Change state to retryable
job = Repo.update!(Ecto.Changeset.change(job, state: "retryable"))
output = capture_io(fn -> run([]) end)
updated_job = Repo.reload(job)
assert updated_job.state == "cancelled"
assert output =~ "Cancelled 1 stuck discovery jobs"
end
test "cancels stuck discovery jobs in executing state" do
{:ok, job} =
%{device_id: Ecto.UUID.generate()}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
# Change state to executing
job =
Repo.update!(
Ecto.Changeset.change(job,
state: "executing",
attempted_at: DateTime.utc_now()
)
)
output = capture_io(fn -> run([]) end)
updated_job = Repo.reload(job)
assert updated_job.state == "cancelled"
assert output =~ "Cancelled 1 stuck discovery jobs"
end
test "cancels multiple stuck discovery jobs" do
device_id1 = Ecto.UUID.generate()
device_id2 = Ecto.UUID.generate()
device_id3 = Ecto.UUID.generate()
# Create multiple stuck jobs
{:ok, job1} =
%{device_id: device_id1}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
{:ok, job2} =
%{device_id: device_id2}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
{:ok, job3} =
%{device_id: device_id3}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
# Make them stuck
Repo.update!(Ecto.Changeset.change(job1, state: "scheduled"))
Repo.update!(Ecto.Changeset.change(job2, state: "retryable"))
Repo.update!(Ecto.Changeset.change(job3, state: "executing", attempted_at: DateTime.utc_now()))
output = capture_io(fn -> run([]) end)
# Verify all were cancelled
assert Repo.reload(job1).state == "cancelled"
assert Repo.reload(job2).state == "cancelled"
assert Repo.reload(job3).state == "cancelled"
assert output =~ "Cancelled 3 stuck discovery jobs"
end
test "does not cancel completed discovery jobs" do
{:ok, job} =
%{device_id: Ecto.UUID.generate()}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
# Change state to completed
job = Repo.update!(Ecto.Changeset.change(job, state: "completed"))
output = capture_io(fn -> run([]) end)
# Job should remain completed
updated_job = Repo.reload(job)
assert updated_job.state == "completed"
assert output =~ "Cancelled 0 stuck discovery jobs"
end
test "does not cancel cancelled discovery jobs" do
{:ok, job} =
%{device_id: Ecto.UUID.generate()}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
# Change state to cancelled
job = Repo.update!(Ecto.Changeset.change(job, state: "cancelled"))
output = capture_io(fn -> run([]) end)
# Job should remain cancelled
updated_job = Repo.reload(job)
assert updated_job.state == "cancelled"
assert output =~ "Cancelled 0 stuck discovery jobs"
end
test "does not cancel jobs from other workers" do
{:ok, discovery_job} =
%{device_id: Ecto.UUID.generate()}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
{:ok, poller_job} =
%{device_id: Ecto.UUID.generate()}
|> Oban.Job.new(worker: "Towerops.Workers.DevicePollerWorker", queue: :pollers)
|> Repo.insert()
# Make both stuck
discovery_job = Repo.update!(Ecto.Changeset.change(discovery_job, state: "scheduled"))
poller_job = Repo.update!(Ecto.Changeset.change(poller_job, state: "scheduled"))
output = capture_io(fn -> run([]) end)
# Only discovery job should be cancelled
assert Repo.reload(discovery_job).state == "cancelled"
assert Repo.reload(poller_job).state == "scheduled"
assert output =~ "Cancelled 1 stuck discovery jobs"
end
test "reports zero cancellations when no stuck jobs exist" do
output = capture_io(fn -> run([]) end)
assert output =~ "Cancelled 0 stuck discovery jobs"
end
test "handles jobs with no device_id in args" do
{:ok, job} =
%{some_other_field: "value"}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
job = Repo.update!(Ecto.Changeset.change(job, state: "scheduled"))
output = capture_io(fn -> run([]) end)
updated_job = Repo.reload(job)
assert updated_job.state == "cancelled"
assert output =~ "Cancelled 1 stuck discovery jobs"
end
test "continues on cancellation failure for one job" do
# Create two stuck jobs
{:ok, job1} =
%{device_id: Ecto.UUID.generate()}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
{:ok, job2} =
%{device_id: Ecto.UUID.generate()}
|> Oban.Job.new(worker: "Towerops.Workers.DiscoveryWorker", queue: :discovery)
|> Repo.insert()
job1 = Repo.update!(Ecto.Changeset.change(job1, state: "scheduled"))
job2 = Repo.update!(Ecto.Changeset.change(job2, state: "scheduled"))
# Delete job1 to simulate a cancellation failure (job not found)
Repo.delete!(job1)
output = capture_io(fn -> run([]) end)
# job2 should still be cancelled
assert Repo.reload(job2).state == "cancelled"
# Should report 1 successful cancellation (job2 only)
assert output =~ "Cancelled 1 stuck discovery jobs"
end
end
end