prop/lib/microwaveprop/workers/rtma_fetch_worker.ex
Graham McIntire dea407c8ca Add ERA5 reanalysis and RTMA data sources
ERA5 (Copernicus CDS API):
- Era5Profile schema matching HRRR profile structure for interop
- Era5Client with async job submission, polling, GRIB2 download
- Era5FetchWorker (Oban queue: era5, max_attempts: 5)
- Unified lookup: Weather.best_profile_for_contact/1 tries HRRR
  first, falls back to ERA5 for pre-2014 contacts

RTMA (NOAA S3, 2.5km/15-min):
- RtmaObservation schema for surface-only fields
- RtmaClient with byte-range GRIB2 requests (same pattern as HRRR)
- RtmaFetchWorker (Oban queue: rtma, max_attempts: 10)
- Weather.find_nearest_rtma/3 for surface condition lookup

Both queues added to production Oban config (2 workers each).
2026-04-07 12:04:16 -05:00

68 lines
2.1 KiB
Elixir

defmodule Microwaveprop.Workers.RtmaFetchWorker do
@moduledoc """
Fetches RTMA surface observations for real-time propagation scoring.
15-minute resolution provides finer temporal detail than HRRR's hourly cycle.
"""
use Oban.Worker, queue: :rtma, max_attempts: 10
alias Microwaveprop.Repo
alias Microwaveprop.Weather.RtmaClient
alias Microwaveprop.Weather.RtmaObservation
require Logger
@impl Oban.Worker
def backoff(%Oban.Job{attempt: attempt}) do
min(60 * Integer.pow(2, attempt - 1), _six_hours = 21_600)
end
@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 * 40) / 40
rlon = Float.round(lon * 40) / 40
if has_rtma_observation?(rlat, rlon, valid_time) do
Logger.debug("RTMA: observation exists for #{rlat},#{rlon} @ #{valid_time_str}")
:ok
else
Logger.info("RTMA: fetching for #{rlat},#{rlon} @ #{valid_time_str}")
case RtmaClient.fetch_observation(lat, lon, valid_time) do
{:ok, attrs} ->
%RtmaObservation{}
|> RtmaObservation.changeset(attrs)
|> Repo.insert(
on_conflict: :nothing,
conflict_target: [:lat, :lon, :valid_time]
)
Logger.info("RTMA: stored observation for #{rlat},#{rlon} @ #{valid_time_str}")
:ok
{:error, reason} ->
Logger.warning("RTMA: failed for #{rlat},#{rlon} @ #{valid_time_str}: #{inspect(reason)}")
{:error, reason}
end
end
end
defp has_rtma_observation?(lat, lon, valid_time) do
import Ecto.Query
dlat = 0.03
dlon = 0.03
time_start = DateTime.add(valid_time, -450, :second)
time_end = DateTime.add(valid_time, 450, :second)
RtmaObservation
|> where(
[o],
o.lat >= ^(lat - dlat) and o.lat <= ^(lat + dlat) and
o.lon >= ^(lon - dlon) and o.lon <= ^(lon + dlon) and
o.valid_time >= ^time_start and o.valid_time <= ^time_end
)
|> Repo.exists?()
end
end