- Schema: Era5Profile → NarrProfile; table renamed via migration rename_era5_profiles_to_narr_profiles (ALTER TABLE + rename indexes, atomic and instant, no row rewrite) - Weather helpers: find_nearest_era5/era5_for_contact/era5_profiles_for_path → find_nearest_narr/narr_for_contact/narr_profiles_for_path - BackfillEnqueueWorker: :era5 → :narr type (virtual, no narr_status column); reconcile_stale_queued_to_unavailable and status_priority_order now skip virtual types via @virtual_types. Fixes the prod crash where the cron tried to update a non-existent era5_status column. - ContactWeatherEnqueueWorker: build_era5_jobs → build_narr_jobs, era5_jobs_for_contact → narr_jobs_for_contact - ContactLive: @era5 → @narr, maybe_enqueue_era5 → maybe_enqueue_narr; UI label "ERA5 (0.25°)" → "NARR (32 km)" - Cron (runtime.exs): args types list now "narr" instead of "era5" - /status: progress row, status-by-type card, totals.narr_profiles, and table-count lookup all target narr_profiles - Drop obsolete Era5CdsJob schema + era5_cds_jobs table (inspection artifact from the retired CDS pipeline; 34 orphan rows in prod) - Misc docstring/comment cleanups (skew_t, about, wgrib2, propagation_train) Includes a regression test for the virtual-type crash.
99 lines
2.9 KiB
Elixir
99 lines
2.9 KiB
Elixir
defmodule Microwaveprop.Workers.NarrFetchWorker do
|
|
@moduledoc """
|
|
Fetches one NCEP NARR profile from NCEI via `NarrClient.fetch_profile_at/2`
|
|
and inserts it as a `NarrProfile` row.
|
|
|
|
Covers pre-2014 contacts where the HRRR archive doesn't reach. Replaces
|
|
the retired ERA5/CDS pipeline. NARR fetches are per-point (one analysis
|
|
hour, one location), so this worker does the fetch + insert directly —
|
|
no downstream batch worker, no routing.
|
|
|
|
Jobs are unique on `{lat, lon, valid_time}` so duplicate enrichment
|
|
requests collapse. `valid_time` MUST already be snapped to a NARR
|
|
3-hourly analysis slot (00/03/06/09/12/15/18/21 UTC on the hour) by the
|
|
caller — the worker raises `ArgumentError` otherwise via
|
|
`NarrClient.url_for/1`.
|
|
|
|
See `docs/plans/2026-04-15-merra2-historical-backfill.md` for the full
|
|
architecture.
|
|
"""
|
|
use Oban.Worker,
|
|
queue: :narr,
|
|
max_attempts: 3,
|
|
unique: [
|
|
period: :infinity,
|
|
states: [:available, :scheduled, :executing, :retryable]
|
|
]
|
|
|
|
import Ecto.Query
|
|
|
|
alias Microwaveprop.Repo
|
|
alias Microwaveprop.Weather.NarrClient
|
|
alias Microwaveprop.Weather.NarrProfile
|
|
|
|
require Logger
|
|
|
|
@impl Oban.Worker
|
|
def perform(%Oban.Job{args: %{"lat" => lat, "lon" => lon, "valid_time" => valid_time_str}}) do
|
|
{:ok, valid_time, _} = DateTime.from_iso8601(valid_time_str)
|
|
|
|
rlat = Float.round(lat * 4) / 4
|
|
rlon = Float.round(lon * 4) / 4
|
|
|
|
if has_narr_profile?(rlat, rlon, valid_time) do
|
|
:ok
|
|
else
|
|
fetch_and_insert(rlat, rlon, valid_time)
|
|
end
|
|
end
|
|
|
|
defp fetch_and_insert(rlat, rlon, valid_time) do
|
|
case NarrClient.fetch_profile_at(valid_time, {rlat, rlon}) do
|
|
{:ok, attrs} ->
|
|
insert_profile(attrs, rlat, rlon, valid_time)
|
|
:ok
|
|
|
|
{:error, reason} = error ->
|
|
Logger.warning(
|
|
"NARR: fetch_profile_at failed for #{rlat},#{rlon} @ #{DateTime.to_iso8601(valid_time)}: #{inspect(reason)}"
|
|
)
|
|
|
|
error
|
|
end
|
|
end
|
|
|
|
defp insert_profile(attrs, rlat, rlon, valid_time) do
|
|
now = DateTime.truncate(DateTime.utc_now(), :second)
|
|
|
|
row =
|
|
attrs
|
|
|> Map.put(:lat, rlat)
|
|
|> Map.put(:lon, rlon)
|
|
|> Map.put(:valid_time, valid_time)
|
|
|> Map.put(:inserted_at, now)
|
|
|> Map.put(:updated_at, now)
|
|
|
|
Repo.insert_all(NarrProfile, [row],
|
|
on_conflict: :nothing,
|
|
conflict_target: [:lat, :lon, :valid_time]
|
|
)
|
|
end
|
|
|
|
defp has_narr_profile?(lat, lon, valid_time) do
|
|
# NARR is exact 3-hourly, so the time window can be tight (±15 min).
|
|
# Lat/lon tolerance matches the 0.25° dedup granularity.
|
|
dlat = 0.13
|
|
dlon = 0.13
|
|
time_start = DateTime.add(valid_time, -900, :second)
|
|
time_end = DateTime.add(valid_time, 900, :second)
|
|
|
|
NarrProfile
|
|
|> where(
|
|
[p],
|
|
p.lat >= ^(lat - dlat) and p.lat <= ^(lat + dlat) and
|
|
p.lon >= ^(lon - dlon) and p.lon <= ^(lon + dlon) and
|
|
p.valid_time >= ^time_start and p.valid_time <= ^time_end
|
|
)
|
|
|> Repo.exists?()
|
|
end
|
|
end
|