71 lines
2.1 KiB
Elixir
71 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,
|
|
unique: [period: 300, fields: [:args], states: [:scheduled, :available]]
|
|
|
|
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
|