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