prop/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex
Graham McIntire 3228835636
Add IEM precipitation data to QSO weather pipeline
- Add precip_1h_in and wx_codes fields to ASOS surface observations
- Add IEMRE reanalysis schema for radar-derived hourly precipitation
- Add IemreFetchWorker with exponential backoff and idempotency
- Integrate IEMRE enqueue into cron weather backfill pipeline
- All existing QSOs marked iemre_queued=false for automatic backfill
2026-03-30 13:18:11 -05:00

216 lines
4.9 KiB
Elixir

defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do
@moduledoc false
use Oban.Worker, queue: :enqueue, max_attempts: 3
alias Microwaveprop.Radio
alias Microwaveprop.Weather
alias Microwaveprop.Weather.HrrrClient
alias Microwaveprop.Workers.HrrrFetchWorker
alias Microwaveprop.Workers.IemreFetchWorker
alias Microwaveprop.Workers.TerrainProfileWorker
alias Microwaveprop.Workers.WeatherFetchWorker
@asos_radius_km 150
@sounding_radius_km 300
@impl Oban.Worker
def perform(%Oban.Job{}) do
enqueue_weather_jobs()
enqueue_hrrr_jobs()
enqueue_terrain_jobs()
enqueue_iemre_jobs()
:ok
end
defp enqueue_weather_jobs do
case Radio.unprocessed_qsos() do
[] ->
:ok
qsos ->
Radio.backfill_distances(qsos)
jobs = build_weather_jobs(qsos)
if jobs != [] do
Oban.insert_all(jobs)
end
qso_ids = Enum.map(qsos, & &1.id)
Radio.mark_weather_queued!(qso_ids)
enqueue_weather_jobs()
end
end
defp enqueue_hrrr_jobs do
case Radio.unprocessed_hrrr_qsos() do
[] ->
:ok
qsos ->
jobs = build_hrrr_jobs(qsos)
if jobs != [] do
Oban.insert_all(jobs)
end
qso_ids = Enum.map(qsos, & &1.id)
Radio.mark_hrrr_queued!(qso_ids)
enqueue_hrrr_jobs()
end
end
defp enqueue_terrain_jobs do
case Radio.unprocessed_terrain_qsos() do
[] ->
:ok
qsos ->
jobs = build_terrain_jobs(qsos)
if jobs != [] do
Oban.insert_all(jobs)
end
qso_ids = Enum.map(qsos, & &1.id)
Radio.mark_terrain_queued!(qso_ids)
# Continue until all QSOs are enqueued
enqueue_terrain_jobs()
end
end
def build_terrain_jobs(qsos) do
qsos
|> Enum.map(fn qso ->
TerrainProfileWorker.new(%{"qso_id" => qso.id})
end)
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
end
def build_weather_jobs(qsos) do
qsos
|> Enum.flat_map(&jobs_for_qso/1)
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
end
def build_hrrr_jobs(qsos) do
qsos
|> Enum.flat_map(&hrrr_job_for_qso/1)
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
end
def build_iemre_jobs(qsos) do
qsos
|> Enum.flat_map(&iemre_job_for_qso/1)
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
end
defp jobs_for_qso(qso) do
lat = qso.pos1["lat"]
lon = qso.pos1["lon"] || qso.pos1["lng"]
asos_jobs = build_asos_jobs(lat, lon, qso.qso_timestamp)
raob_jobs = build_raob_jobs(lat, lon, qso.qso_timestamp)
asos_jobs ++ raob_jobs
end
defp enqueue_iemre_jobs do
case Radio.unprocessed_iemre_qsos() do
[] ->
:ok
qsos ->
jobs = build_iemre_jobs(qsos)
if jobs != [] do
Oban.insert_all(jobs)
end
qso_ids = Enum.map(qsos, & &1.id)
Radio.mark_iemre_queued!(qso_ids)
enqueue_iemre_jobs()
end
end
defp hrrr_job_for_qso(%{pos1: nil}), do: []
defp hrrr_job_for_qso(qso) do
lat = qso.pos1["lat"]
lon = qso.pos1["lon"] || qso.pos1["lng"]
if lat && lon do
{rlat, rlon} = Weather.round_to_hrrr_grid(lat, lon)
rounded_time = HrrrClient.nearest_hrrr_hour(qso.qso_timestamp)
[
HrrrFetchWorker.new(%{
"lat" => rlat,
"lon" => rlon,
"valid_time" => DateTime.to_iso8601(rounded_time)
})
]
else
[]
end
end
defp build_asos_jobs(lat, lon, timestamp) do
start_dt = DateTime.add(timestamp, -2 * 3600, :second)
end_dt = DateTime.add(timestamp, 2 * 3600, :second)
lat
|> Weather.nearby_stations(lon, "asos", @asos_radius_km)
|> Enum.map(fn station ->
WeatherFetchWorker.new(%{
"fetch_type" => "asos",
"station_id" => station.id,
"station_code" => station.station_code,
"start_dt" => DateTime.to_iso8601(start_dt),
"end_dt" => DateTime.to_iso8601(end_dt)
})
end)
end
defp iemre_job_for_qso(%{pos1: nil}), do: []
defp iemre_job_for_qso(qso) do
lat = qso.pos1["lat"]
lon = qso.pos1["lon"] || qso.pos1["lng"]
if lat && lon do
{rlat, rlon} = Weather.round_to_iemre_grid(lat, lon)
date = DateTime.to_date(qso.qso_timestamp)
[
IemreFetchWorker.new(%{
"lat" => rlat,
"lon" => rlon,
"date" => Date.to_iso8601(date)
})
]
else
[]
end
end
defp build_raob_jobs(lat, lon, timestamp) do
sounding_times = Weather.sounding_times_around(timestamp)
stations = Weather.nearby_stations(lat, lon, "sounding", @sounding_radius_km)
for station <- stations, sounding_time <- sounding_times do
WeatherFetchWorker.new(%{
"fetch_type" => "raob",
"station_id" => station.id,
"station_code" => station.station_code,
"sounding_time" => DateTime.to_iso8601(sounding_time)
})
end
end
end