265 lines
6.6 KiB
Elixir
265 lines
6.6 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
|
|
# Oban jobs have ~149 params each; PG limit is 65535 params per query
|
|
@insert_batch_size 400
|
|
|
|
@doc """
|
|
Enqueue all enrichment jobs (weather, HRRR, terrain, IEMRE) for a single QSO.
|
|
Called directly from submission flow — no Oban indirection.
|
|
"""
|
|
def enqueue_for_qso(qso) do
|
|
weather_jobs = build_weather_jobs([qso])
|
|
hrrr_jobs = build_hrrr_jobs([qso])
|
|
terrain_jobs = build_terrain_jobs([qso])
|
|
iemre_jobs = build_iemre_jobs([qso])
|
|
|
|
all_jobs = weather_jobs ++ hrrr_jobs ++ terrain_jobs ++ iemre_jobs
|
|
|
|
if all_jobs != [] do
|
|
insert_all_chunked(all_jobs)
|
|
end
|
|
|
|
qso_ids = [qso.id]
|
|
if qso.pos1, do: Radio.mark_weather_queued!(qso_ids)
|
|
if qso.pos1, do: Radio.mark_hrrr_queued!(qso_ids)
|
|
if qso.pos1 && qso.pos2, do: Radio.mark_terrain_queued!(qso_ids)
|
|
if qso.pos1, do: Radio.mark_iemre_queued!(qso_ids)
|
|
|
|
:ok
|
|
end
|
|
|
|
@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
|
|
insert_all_chunked(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
|
|
insert_all_chunked(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
|
|
insert_all_chunked(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.filter(fn qso -> qso.pos1 && qso.pos2 end)
|
|
|> 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
|
|
# Group all QSO points by HRRR hour, then create one batch job per hour
|
|
# instead of one job per point. This deduplicates GRIB2 downloads.
|
|
qsos
|
|
|> Enum.flat_map(&hrrr_points_for_qso/1)
|
|
|> Enum.group_by(fn {_point, hour} -> hour end, fn {point, _hour} -> point end)
|
|
|> Enum.map(fn {hour, points} ->
|
|
unique_points =
|
|
points
|
|
|> Enum.uniq()
|
|
|> Enum.map(fn {lat, lon} -> %{"lat" => lat, "lon" => lon} end)
|
|
|
|
HrrrFetchWorker.new(%{
|
|
"points" => unique_points,
|
|
"valid_time" => DateTime.to_iso8601(hour)
|
|
})
|
|
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
|
|
qso
|
|
|> Radio.qso_path_points()
|
|
|> Enum.flat_map(fn {lat, lon} ->
|
|
build_asos_jobs(lat, lon, qso.qso_timestamp) ++
|
|
build_raob_jobs(lat, lon, qso.qso_timestamp)
|
|
end)
|
|
end
|
|
|
|
defp enqueue_iemre_jobs do
|
|
case Radio.unprocessed_iemre_qsos() do
|
|
[] ->
|
|
:ok
|
|
|
|
qsos ->
|
|
jobs = build_iemre_jobs(qsos)
|
|
|
|
if jobs != [] do
|
|
insert_all_chunked(jobs)
|
|
end
|
|
|
|
qso_ids = Enum.map(qsos, & &1.id)
|
|
Radio.mark_iemre_queued!(qso_ids)
|
|
|
|
enqueue_iemre_jobs()
|
|
end
|
|
end
|
|
|
|
defp hrrr_points_for_qso(%{pos1: nil}), do: []
|
|
|
|
defp hrrr_points_for_qso(qso) do
|
|
rounded_time = HrrrClient.nearest_hrrr_hour(qso.qso_timestamp)
|
|
|
|
qso
|
|
|> Radio.qso_path_points()
|
|
|> Enum.flat_map(fn {lat, lon} ->
|
|
{rlat, rlon} = Weather.round_to_hrrr_grid(lat, lon)
|
|
|
|
if Weather.has_hrrr_profile?(rlat, rlon, rounded_time) do
|
|
[]
|
|
else
|
|
[{{rlat, rlon}, rounded_time}]
|
|
end
|
|
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.reject(fn station ->
|
|
Weather.has_surface_observations?(station.id, start_dt, end_dt)
|
|
end)
|
|
|> 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
|
|
date = DateTime.to_date(qso.qso_timestamp)
|
|
|
|
qso
|
|
|> Radio.qso_path_points()
|
|
|> Enum.flat_map(fn {lat, lon} ->
|
|
{rlat, rlon} = Weather.round_to_iemre_grid(lat, lon)
|
|
|
|
if Weather.has_iemre_observation?(rlat, rlon, date) do
|
|
[]
|
|
else
|
|
[
|
|
IemreFetchWorker.new(%{
|
|
"lat" => rlat,
|
|
"lon" => rlon,
|
|
"date" => Date.to_iso8601(date)
|
|
})
|
|
]
|
|
end
|
|
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,
|
|
not Weather.has_sounding?(station.id, sounding_time) do
|
|
WeatherFetchWorker.new(%{
|
|
"fetch_type" => "raob",
|
|
"station_id" => station.id,
|
|
"station_code" => station.station_code,
|
|
"sounding_time" => DateTime.to_iso8601(sounding_time)
|
|
})
|
|
end
|
|
end
|
|
|
|
defp insert_all_chunked(jobs) do
|
|
jobs
|
|
|> Enum.chunk_every(@insert_batch_size)
|
|
|> Enum.each(&Oban.insert_all/1)
|
|
end
|
|
end
|