206 lines
6.8 KiB
Elixir
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
|