Era5SubmitWorker was counting era5_cds_jobs rows 1:1 against the CDS
150-job ceiling, but each row holds TWO CDS job IDs (single-level +
pressure-level). Prod hit 74 rows = 148 in-flight jobs and got stuck
in a reject→resubmit spiral. The cap guard now multiplies row count
by 2 to match what CDS sees.
Era5PollWorker @stuck_after_seconds goes from 4h to 18h. CDS
single-level routinely sits at 'accepted' for 12+h under load while
pressure finishes in ~90 min; the 4h threshold was firing on normal
slow runs and feeding the same spiral.
PipelineStatus now returns `label` + `detail` separately so the map
chip can render "Updating propagation" on one line and the
sub-detail ("HRRR run" / "ASOS nudge" / "+3h" / "now") on a second
line below it. running_label/1 is renamed to running_detail/1 and
returns just the sub-part.
258 lines
9.1 KiB
Elixir
258 lines
9.1 KiB
Elixir
defmodule Microwaveprop.Workers.Era5SubmitWorkerTest do
|
||
use Microwaveprop.DataCase, async: false
|
||
use Oban.Testing, repo: Microwaveprop.Repo
|
||
|
||
import Ecto.Query
|
||
|
||
alias Microwaveprop.Weather.Era5CdsJob
|
||
alias Microwaveprop.Weather.Era5Client
|
||
alias Microwaveprop.Weather.Era5Profile
|
||
alias Microwaveprop.Workers.Era5PollWorker
|
||
alias Microwaveprop.Workers.Era5SubmitWorker
|
||
|
||
@default_args %{"year" => 2014, "month" => 3, "tile_lat" => 32, "tile_lon" => -98}
|
||
|
||
setup do
|
||
System.put_env("ERA5_CDS_API_KEY", "test-key")
|
||
on_exit(fn -> System.delete_env("ERA5_CDS_API_KEY") end)
|
||
:ok
|
||
end
|
||
|
||
describe "perform/1 — short-circuit paths" do
|
||
test "returns {:ok, :cached} when the month-tile already has era5_profiles" do
|
||
# Seed a profile inside the (2014-03, tile 32,-98) window.
|
||
%Era5Profile{}
|
||
|> Era5Profile.changeset(%{
|
||
lat: 32.5,
|
||
lon: -97.25,
|
||
valid_time: ~U[2014-03-15 12:00:00Z],
|
||
profile: []
|
||
})
|
||
|> Repo.insert!()
|
||
|
||
Req.Test.stub(Era5Client, fn _ -> raise "CDS should not be contacted on cache hit" end)
|
||
|
||
Oban.Testing.with_testing_mode(:manual, fn ->
|
||
assert {:ok, :cached} = Era5SubmitWorker.perform(%Oban.Job{args: @default_args})
|
||
assert [] = all_enqueued(worker: Era5PollWorker)
|
||
end)
|
||
end
|
||
|
||
test "returns {:ok, :already_submitted} when an era5_cds_jobs row already exists" do
|
||
{:ok, _} =
|
||
%Era5CdsJob{}
|
||
|> Era5CdsJob.changeset(%{
|
||
year: 2014,
|
||
month: 3,
|
||
tile_lat: 32,
|
||
tile_lon: -98,
|
||
single_job_id: "pre-existing-single",
|
||
pressure_job_id: "pre-existing-pressure",
|
||
submitted_at: DateTime.truncate(DateTime.utc_now(), :second)
|
||
})
|
||
|> Repo.insert()
|
||
|
||
Req.Test.stub(Era5Client, fn _ -> raise "CDS should not be contacted on duplicate submit" end)
|
||
|
||
Oban.Testing.with_testing_mode(:manual, fn ->
|
||
assert {:ok, :already_submitted} =
|
||
Era5SubmitWorker.perform(%Oban.Job{args: @default_args})
|
||
end)
|
||
end
|
||
end
|
||
|
||
describe "perform/1 — happy path" do
|
||
test "submits both CDS jobs, persists the row, and enqueues a poll worker" do
|
||
Req.Test.stub(Era5Client, fn conn ->
|
||
assert conn.method == "POST"
|
||
|
||
dataset_id =
|
||
cond do
|
||
String.contains?(conn.request_path, "single-levels") -> "single-job-123"
|
||
String.contains?(conn.request_path, "pressure-levels") -> "pressure-job-456"
|
||
end
|
||
|
||
Req.Test.json(conn, %{"jobID" => dataset_id})
|
||
end)
|
||
|
||
Oban.Testing.with_testing_mode(:manual, fn ->
|
||
assert {:ok, %Era5CdsJob{} = row} =
|
||
Era5SubmitWorker.perform(%Oban.Job{args: @default_args})
|
||
|
||
assert row.single_job_id == "single-job-123"
|
||
assert row.pressure_job_id == "pressure-job-456"
|
||
assert row.year == 2014
|
||
assert row.month == 3
|
||
assert row.tile_lat == 32
|
||
assert row.tile_lon == -98
|
||
|
||
# Row is persisted in the DB.
|
||
assert Repo.aggregate(
|
||
from(j in Era5CdsJob,
|
||
where: j.year == 2014 and j.month == 3 and j.tile_lat == 32 and j.tile_lon == -98
|
||
),
|
||
:count,
|
||
:id
|
||
) == 1
|
||
|
||
# A poll worker was enqueued for this row, scheduled ~5 min from now.
|
||
[poll_job] = all_enqueued(worker: Era5PollWorker)
|
||
assert poll_job.args["era5_cds_job_id"] == row.id
|
||
assert DateTime.after?(poll_job.scheduled_at, DateTime.utc_now())
|
||
end)
|
||
end
|
||
|
||
test "returns {:error, _} and persists no row if either CDS submit fails" do
|
||
Req.Test.stub(Era5Client, fn conn ->
|
||
if String.contains?(conn.request_path, "pressure-levels") do
|
||
Plug.Conn.resp(conn, 400, ~s({"detail":"bad pressure request"}))
|
||
else
|
||
Req.Test.json(conn, %{"jobID" => "single-ok"})
|
||
end
|
||
end)
|
||
|
||
Oban.Testing.with_testing_mode(:manual, fn ->
|
||
assert {:error, reason} = Era5SubmitWorker.perform(%Oban.Job{args: @default_args})
|
||
assert is_binary(reason)
|
||
|
||
# No orphaned row.
|
||
assert Repo.aggregate(Era5CdsJob, :count, :id) == 0
|
||
|
||
# No poll worker enqueued.
|
||
assert [] = all_enqueued(worker: Era5PollWorker)
|
||
end)
|
||
end
|
||
end
|
||
|
||
describe "CDS in-flight cap" do
|
||
# CDS rejects with "Number of queued requests is limited to 150" when
|
||
# a user has too many pending jobs. Each tile-month lives in ONE
|
||
# era5_cds_jobs row but carries TWO CDS job IDs (single-level +
|
||
# pressure-level), so the real in-flight count is 2 × row_count. The
|
||
# cap guard multiplies accordingly. The exact threshold (headroom) is
|
||
# tuned in the worker; the tests below exercise both the obvious
|
||
# "way above the cap" path and a row count that should only trip the
|
||
# snooze if the 2× math is in place.
|
||
@cds_ceiling 150
|
||
|
||
test "snoozes when era5_cds_jobs count is at or above the snooze threshold" do
|
||
# 140 in-flight is above any reasonable snooze threshold.
|
||
seed_count = @cds_ceiling - 10
|
||
|
||
rows =
|
||
for i <- 1..seed_count do
|
||
%{
|
||
id: Ecto.UUID.generate(),
|
||
year: 2014,
|
||
month: 3,
|
||
tile_lat: 32,
|
||
tile_lon: -98 - i,
|
||
single_job_id: "seed-single-#{i}",
|
||
pressure_job_id: "seed-pressure-#{i}",
|
||
submitted_at: DateTime.truncate(DateTime.utc_now(), :second),
|
||
inserted_at: DateTime.truncate(DateTime.utc_now(), :second),
|
||
updated_at: DateTime.truncate(DateTime.utc_now(), :second),
|
||
poll_count: 0
|
||
}
|
||
end
|
||
|
||
Repo.insert_all(Era5CdsJob, rows)
|
||
|
||
Req.Test.stub(Era5Client, fn _ -> raise "CDS should not be contacted when capped" end)
|
||
|
||
Oban.Testing.with_testing_mode(:manual, fn ->
|
||
assert {:snooze, seconds} = Era5SubmitWorker.perform(%Oban.Job{args: @default_args})
|
||
assert is_integer(seconds) and seconds > 0
|
||
end)
|
||
|
||
# No new era5_cds_jobs row inserted for the requested tile.
|
||
refute Repo.one(
|
||
from j in Era5CdsJob,
|
||
where: j.year == 2014 and j.month == 3 and j.tile_lat == 32 and j.tile_lon == -98
|
||
)
|
||
end
|
||
|
||
test "snoozes when row count × 2 crosses the cap (each row = two CDS jobs)" do
|
||
# 65 rows = 130 in-flight CDS jobs — above the effective snooze
|
||
# threshold even though row count alone is nowhere near 150. This
|
||
# is the scenario that was escaping the old guard and driving the
|
||
# per-user-cap reject storms in prod.
|
||
rows =
|
||
for i <- 1..65 do
|
||
%{
|
||
id: Ecto.UUID.generate(),
|
||
year: 2014,
|
||
month: rem(i - 1, 12) + 1,
|
||
tile_lat: 32,
|
||
tile_lon: -98 - i,
|
||
single_job_id: "x2-single-#{i}",
|
||
pressure_job_id: "x2-pressure-#{i}",
|
||
submitted_at: DateTime.truncate(DateTime.utc_now(), :second),
|
||
inserted_at: DateTime.truncate(DateTime.utc_now(), :second),
|
||
updated_at: DateTime.truncate(DateTime.utc_now(), :second),
|
||
poll_count: 0
|
||
}
|
||
end
|
||
|
||
Repo.insert_all(Era5CdsJob, rows)
|
||
|
||
Req.Test.stub(Era5Client, fn _ -> raise "CDS should not be contacted when capped" end)
|
||
|
||
Oban.Testing.with_testing_mode(:manual, fn ->
|
||
assert {:snooze, seconds} = Era5SubmitWorker.perform(%Oban.Job{args: @default_args})
|
||
assert is_integer(seconds) and seconds > 0
|
||
end)
|
||
|
||
refute Repo.one(
|
||
from j in Era5CdsJob,
|
||
where: j.year == 2014 and j.month == 3 and j.tile_lat == 32 and j.tile_lon == -98
|
||
)
|
||
end
|
||
|
||
test "submits normally when in-flight count is well below the ceiling" do
|
||
# 10 in-flight — comfortably under the cap.
|
||
rows =
|
||
for i <- 1..10 do
|
||
%{
|
||
id: Ecto.UUID.generate(),
|
||
year: 2013,
|
||
month: i,
|
||
tile_lat: 40,
|
||
tile_lon: -90,
|
||
single_job_id: "warm-single-#{i}",
|
||
pressure_job_id: "warm-pressure-#{i}",
|
||
submitted_at: DateTime.truncate(DateTime.utc_now(), :second),
|
||
inserted_at: DateTime.truncate(DateTime.utc_now(), :second),
|
||
updated_at: DateTime.truncate(DateTime.utc_now(), :second),
|
||
poll_count: 0
|
||
}
|
||
end
|
||
|
||
Repo.insert_all(Era5CdsJob, rows)
|
||
|
||
Req.Test.stub(Era5Client, fn conn ->
|
||
dataset_id =
|
||
cond do
|
||
String.contains?(conn.request_path, "single-levels") -> "fresh-single"
|
||
String.contains?(conn.request_path, "pressure-levels") -> "fresh-pressure"
|
||
end
|
||
|
||
Req.Test.json(conn, %{"jobID" => dataset_id})
|
||
end)
|
||
|
||
Oban.Testing.with_testing_mode(:manual, fn ->
|
||
assert {:ok, %Era5CdsJob{}} = Era5SubmitWorker.perform(%Oban.Job{args: @default_args})
|
||
end)
|
||
end
|
||
end
|
||
|
||
describe "unique constraint" do
|
||
test "collapses duplicate (year, month, tile_lat, tile_lon) Oban enqueues" do
|
||
Oban.Testing.with_testing_mode(:manual, fn ->
|
||
{:ok, j1} = Oban.insert(Era5SubmitWorker.new(@default_args))
|
||
{:ok, j2} = Oban.insert(Era5SubmitWorker.new(@default_args))
|
||
assert j1.id == j2.id
|
||
end)
|
||
end
|
||
end
|
||
end
|