defmodule Microwaveprop.Workers.Era5SubmitWorker do @moduledoc """ First half of the split ERA5 backfill pipeline. Submits the two CDS jobs for a month-tile (single-level + pressure-level) in parallel, persists the returned CDS job IDs to `era5_cds_jobs`, and enqueues an `Era5PollWorker` job scheduled a few minutes out to check for completion. The fast submit → enqueue-poll split means an Oban slot is held for ~1 second (the time to POST and write a row), not the 30-45 minutes CDS takes to queue, assemble, and return a month-tile. A pod restart between submit and poll does not lose the CDS work: the row in `era5_cds_jobs` persists the job IDs and the poll worker picks up where we left off. Idempotent: if the month-tile already has profiles (`era5_profiles`) or an in-flight `era5_cds_jobs` row, no submit is made. """ use Oban.Worker, queue: :era5_submit, max_attempts: 5, unique: [ period: :infinity, states: [:available, :scheduled, :executing, :retryable], keys: [:year, :month, :tile_lat, :tile_lon] ] import Ecto.Query alias Microwaveprop.Repo alias Microwaveprop.Weather.Era5BatchClient alias Microwaveprop.Weather.Era5CdsJob alias Microwaveprop.Weather.Era5Client alias Microwaveprop.Weather.Era5Profile alias Microwaveprop.Workers.Era5PollWorker require Logger @poll_delay_seconds 5 * 60 # CDS enforces a server-side per-user cap of 150 in-flight jobs. Each # tile-month submit lives in ONE era5_cds_jobs row but carries TWO CDS # job IDs (single-level + pressure-level), so the real in-flight count # CDS sees is 2 × row_count. The cap guard multiplies accordingly. # Snoozing (instead of failing or sleeping) releases the Oban slot # immediately and lets the poll workers drain in-flight jobs before we # push more in. # # Headroom is deliberately generous (30 job slots, effective ceiling # 120 jobs = 60 rows) because the cap check races between concurrent # workers: two workers that both observe row_count=59 can both submit, # landing us at row_count=61 (122 jobs) after both finish. With 30 # slots of headroom the cluster would have to race 15 workers wide to # clip CDS's hard 150 ceiling. Observed in prod: the old 1× math was # allowing 74 rows (148 jobs) which triggered CDS reject storms. @cds_ceiling 150 @cap_headroom 30 @snooze_seconds 5 * 60 @impl Oban.Worker def backoff(%Oban.Job{attempt: attempt}) do # The submit itself is cheap; retry fast so transient CDS/HTTP blips # don't stall the tile for an hour waiting on exponential backoff. min(30 * Integer.pow(2, attempt - 1), _one_hour = 3_600) end @impl Oban.Worker def perform(%Oban.Job{args: args}) do year = fetch_int!(args, "year") month = fetch_int!(args, "month") tile_lat = fetch_int!(args, "tile_lat") tile_lon = fetch_int!(args, "tile_lon") cond do month_tile_cached?(year, month, tile_lat, tile_lon) -> Logger.info("Era5Submit: #{year}-#{pad(month)} tile #{tile_lat},#{tile_lon} already cached — skip") {:ok, :cached} existing = find_in_flight(year, month, tile_lat, tile_lon) -> Logger.info( "Era5Submit: #{year}-#{pad(month)} tile #{tile_lat},#{tile_lon} already submitted (cds_job_id=#{existing.id}) — skip" ) {:ok, :already_submitted} cds_queue_full?() -> row_count = Repo.aggregate(Era5CdsJob, :count, :id) in_flight_jobs = row_count * 2 Logger.info( "Era5Submit: CDS in-flight cap reached (#{in_flight_jobs}/#{@cds_ceiling} jobs across #{row_count} rows) — snoozing #{year}-#{pad(month)} tile #{tile_lat},#{tile_lon}" ) {:snooze, @snooze_seconds} true -> submit_and_persist(year, month, tile_lat, tile_lon) end end defp cds_queue_full? do in_flight_jobs_after_submit = Repo.aggregate(Era5CdsJob, :count, :id) * 2 + 2 in_flight_jobs_after_submit > @cds_ceiling - @cap_headroom end defp submit_and_persist(year, month, tile_lat, tile_lon) do days = Calendar.ISO.days_in_month(year, month) area = tile_area(tile_lat, tile_lon) # Run both submits in parallel — they're independent HTTP POSTs. single_task = Task.async(fn -> Era5Client.submit_job("reanalysis-era5-single-levels", single_request(year, month, days, area)) end) pressure_task = Task.async(fn -> Era5Client.submit_job("reanalysis-era5-pressure-levels", pressure_request(year, month, days, area)) end) try do with {:ok, single_id} <- Task.await(single_task, :infinity), {:ok, pressure_id} <- Task.await(pressure_task, :infinity), {:ok, row} <- insert_cds_job_row(year, month, tile_lat, tile_lon, single_id, pressure_id) do enqueue_poll(row) Logger.info( "Era5Submit: submitted #{year}-#{pad(month)} tile #{tile_lat},#{tile_lon} (single=#{single_id}, pressure=#{pressure_id})" ) {:ok, row} end after Task.shutdown(single_task, :brutal_kill) Task.shutdown(pressure_task, :brutal_kill) end end defp insert_cds_job_row(year, month, tile_lat, tile_lon, single_id, pressure_id) do %Era5CdsJob{} |> Era5CdsJob.changeset(%{ year: year, month: month, tile_lat: tile_lat, tile_lon: tile_lon, single_job_id: single_id, pressure_job_id: pressure_id, submitted_at: DateTime.truncate(DateTime.utc_now(), :second) }) |> Repo.insert() |> case do {:ok, row} -> {:ok, row} {:error, changeset} -> {:error, "era5_cds_jobs insert failed: #{inspect(changeset.errors)}"} end end defp enqueue_poll(%Era5CdsJob{id: id}) do %{"era5_cds_job_id" => id} |> Era5PollWorker.new(schedule_in: @poll_delay_seconds) |> Oban.insert() end defp find_in_flight(year, month, tile_lat, tile_lon) do Repo.one( from j in Era5CdsJob, where: j.year == ^year and j.month == ^month and j.tile_lat == ^tile_lat and j.tile_lon == ^tile_lon, limit: 1 ) end defp month_tile_cached?(year, month, tile_lat, tile_lon) do {first_dt, last_dt} = month_bounds(year, month) tile_degrees = Era5BatchClient.tile_degrees() Era5Profile |> where( [p], p.valid_time >= ^first_dt and p.valid_time < ^last_dt and p.lat >= ^(tile_lat * 1.0) and p.lat <= ^(tile_lat * 1.0 + tile_degrees) and p.lon >= ^(tile_lon * 1.0) and p.lon <= ^(tile_lon * 1.0 + tile_degrees) ) |> Repo.exists?() end defp month_bounds(year, month) do {:ok, first} = Date.new(year, month, 1) {:ok, first_dt} = DateTime.new(first, ~T[00:00:00], "Etc/UTC") days = Calendar.ISO.days_in_month(year, month) last_dt = DateTime.add(first_dt, days * 86_400, :second) {first_dt, last_dt} end defp tile_area(tile_lat, tile_lon) do tile_degrees = Era5BatchClient.tile_degrees() Enum.map( [tile_lat + tile_degrees, tile_lon, tile_lat, tile_lon + tile_degrees], &(&1 * 1.0) ) end defp single_request(year, month, days, area) do %{ "product_type" => ["reanalysis"], "variable" => Era5Client.single_level_vars(), "year" => [Integer.to_string(year)], "month" => [pad(month)], "day" => Enum.map(1..days, &pad/1), "time" => Enum.map(0..23, &"#{pad(&1)}:00"), "area" => area, "data_format" => "grib" } end defp pressure_request(year, month, days, area) do %{ "product_type" => ["reanalysis"], "variable" => Era5Client.pressure_level_vars(), "pressure_level" => Era5Client.pressure_levels(), "year" => [Integer.to_string(year)], "month" => [pad(month)], "day" => Enum.map(1..days, &pad/1), "time" => Enum.map(0..23, &"#{pad(&1)}:00"), "area" => area, "data_format" => "grib" } end defp pad(n), do: String.pad_leading(Integer.to_string(n), 2, "0") defp fetch_int!(args, key) do case Map.fetch!(args, key) do v when is_integer(v) -> v v when is_binary(v) -> String.to_integer(v) end end end