prop/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex
Graham McIntire ce4253f412
Loop enqueue workers until all unprocessed QSOs are covered
Previously each run only enqueued 500 QSOs (the query limit).
Now weather, HRRR, and terrain enqueue functions recurse until
no unprocessed QSOs remain.
2026-03-30 10:11:34 -05:00

167 lines
3.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.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()
: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
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 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 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