Contact detail now searches for soundings in progressively wider radii (150 → 300 → 600 → 1000 km) before declaring "No soundings found". When even 1000 km turns up nothing, enqueue WeatherFetchWorker RAOB jobs for every station at the widest non-empty radius. The worker passes the contact_id through so it can PubSub back to contact_enrichment:<id> as each sounding lands, and the LiveView rehydrates without a refresh. Also adds Weather.soundings_with_widening_radius/1 and ContactWeatherEnqueueWorker.enqueue_raob_fetch_with_widening/1 so the same logic can be reused from other callers.
474 lines
14 KiB
Elixir
474 lines
14 KiB
Elixir
defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorker do
|
|
@moduledoc false
|
|
use Oban.Worker, queue: :enqueue, max_attempts: 3
|
|
|
|
alias Microwaveprop.Radio
|
|
alias Microwaveprop.Weather
|
|
alias Microwaveprop.Weather.HrrrClient
|
|
alias Microwaveprop.Weather.NarrClient
|
|
alias Microwaveprop.Workers.CommonVolumeRadarWorker
|
|
alias Microwaveprop.Workers.HrrrFetchWorker
|
|
alias Microwaveprop.Workers.IemreFetchWorker
|
|
alias Microwaveprop.Workers.MechanismClassifyWorker
|
|
alias Microwaveprop.Workers.NarrFetchWorker
|
|
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 contact.
|
|
Called directly from submission flow — no Oban indirection.
|
|
Optionally pass a list of types to limit which enrichments run.
|
|
"""
|
|
def enqueue_for_contact(contact, types \\ [:hrrr, :weather, :terrain, :iemre, :radar, :mechanism]) do
|
|
contact = Radio.ensure_positions!(contact)
|
|
jobs_by_type = build_jobs_by_type(contact, types)
|
|
|
|
bulk_jobs =
|
|
jobs_by_type[:weather] ++
|
|
jobs_by_type[:hrrr] ++
|
|
jobs_by_type[:terrain] ++
|
|
jobs_by_type[:iemre] ++
|
|
jobs_by_type[:radar] ++
|
|
jobs_by_type[:mechanism]
|
|
|
|
if bulk_jobs != [] do
|
|
insert_all_chunked(bulk_jobs)
|
|
end
|
|
|
|
# NARR jobs must go through Oban.insert/1 so the worker's unique
|
|
# constraint is honored; Oban OSS insert_all does not check uniqueness.
|
|
insert_unique(jobs_by_type[:narr])
|
|
|
|
mark_enrichment_statuses(contact, types, jobs_by_type)
|
|
|
|
:ok
|
|
end
|
|
|
|
defp build_jobs_by_type(contact, types) do
|
|
%{
|
|
weather: if(:weather in types, do: build_weather_jobs([contact]), else: []),
|
|
hrrr: if(:hrrr in types, do: build_hrrr_jobs([contact]), else: []),
|
|
terrain: if(:terrain in types, do: build_terrain_jobs([contact]), else: []),
|
|
iemre: if(:iemre in types, do: build_iemre_jobs([contact]), else: []),
|
|
narr: if(:narr in types, do: build_narr_jobs([contact]), else: []),
|
|
radar: if(:radar in types, do: build_radar_jobs([contact]), else: []),
|
|
mechanism: if(:mechanism in types, do: build_mechanism_jobs([contact]), else: [])
|
|
}
|
|
end
|
|
|
|
defp mark_enrichment_statuses(contact, types, jobs_by_type) do
|
|
ids = [contact.id]
|
|
|
|
if contact.pos1 do
|
|
mark_pos1_statuses(contact, ids, types, jobs_by_type)
|
|
end
|
|
|
|
if contact.pos1 && contact.pos2 do
|
|
mark_path_statuses(contact, ids, types, jobs_by_type)
|
|
end
|
|
end
|
|
|
|
defp mark_pos1_statuses(contact, ids, types, jobs_by_type) do
|
|
if :weather in types, do: mark_status!(ids, :weather_status, jobs_by_type[:weather])
|
|
if :hrrr in types, do: mark_hrrr_status!(contact, ids, jobs_by_type[:hrrr])
|
|
if :iemre in types, do: mark_status!(ids, :iemre_status, jobs_by_type[:iemre])
|
|
end
|
|
|
|
defp mark_path_statuses(contact, ids, types, jobs_by_type) do
|
|
if :terrain in types, do: mark_status!(ids, :terrain_status, jobs_by_type[:terrain])
|
|
if :radar in types, do: mark_radar_status!(contact, ids, jobs_by_type[:radar])
|
|
|
|
if :mechanism in types,
|
|
do: mark_status!(ids, :propagation_mechanism_status, jobs_by_type[:mechanism])
|
|
end
|
|
|
|
# IEM n0q composite reflectivity archive only has meaningful CONUS
|
|
# coverage post-2014 (pre that, the archive is spotty mosaics). For
|
|
# pre-2014 contacts we can't get radar, so pin to :unavailable rather
|
|
# than leaving it :pending forever.
|
|
defp mark_radar_status!(contact, ids, []) do
|
|
if NarrClient.in_coverage?(contact.qso_timestamp) do
|
|
Radio.set_enrichment_status!(ids, :radar_status, :unavailable)
|
|
else
|
|
Radio.set_enrichment_status!(ids, :radar_status, :queued)
|
|
end
|
|
end
|
|
|
|
defp mark_radar_status!(_contact, ids, [_ | _]), do: Radio.set_enrichment_status!(ids, :radar_status, :queued)
|
|
|
|
# Pre-2014 contacts have no HRRR data — NarrClient.in_coverage?/1 is the
|
|
# authoritative "HRRR can never serve this" check. Mark :unavailable so
|
|
# BackfillEnqueueWorker's :narr filter picks them up and the reconcile
|
|
# loop stops flipping the row back to :queued each cron cycle.
|
|
defp mark_hrrr_status!(contact, ids, []) do
|
|
if NarrClient.in_coverage?(contact.qso_timestamp) do
|
|
Radio.set_enrichment_status!(ids, :hrrr_status, :unavailable)
|
|
else
|
|
Radio.set_enrichment_status!(ids, :hrrr_status, :complete)
|
|
end
|
|
end
|
|
|
|
defp mark_hrrr_status!(_contact, ids, [_ | _]), do: Radio.set_enrichment_status!(ids, :hrrr_status, :queued)
|
|
|
|
@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_contacts() do
|
|
[] ->
|
|
:ok
|
|
|
|
contacts ->
|
|
Radio.backfill_distances(contacts)
|
|
|
|
jobs = build_weather_jobs(contacts)
|
|
|
|
if jobs != [] do
|
|
insert_all_chunked(jobs)
|
|
end
|
|
|
|
ids = Enum.map(contacts, & &1.id)
|
|
Radio.mark_weather_queued!(ids)
|
|
|
|
enqueue_weather_jobs()
|
|
end
|
|
end
|
|
|
|
defp enqueue_hrrr_jobs do
|
|
case Radio.unprocessed_hrrr_contacts() do
|
|
[] ->
|
|
:ok
|
|
|
|
contacts ->
|
|
jobs = build_hrrr_jobs(contacts)
|
|
|
|
if jobs != [] do
|
|
insert_all_chunked(jobs)
|
|
end
|
|
|
|
ids = Enum.map(contacts, & &1.id)
|
|
Radio.mark_hrrr_queued!(ids)
|
|
|
|
enqueue_hrrr_jobs()
|
|
end
|
|
end
|
|
|
|
defp enqueue_terrain_jobs do
|
|
case Radio.unprocessed_terrain_contacts() do
|
|
[] ->
|
|
:ok
|
|
|
|
contacts ->
|
|
jobs = build_terrain_jobs(contacts)
|
|
|
|
if jobs != [] do
|
|
insert_all_chunked(jobs)
|
|
end
|
|
|
|
ids = Enum.map(contacts, & &1.id)
|
|
Radio.mark_terrain_queued!(ids)
|
|
|
|
enqueue_terrain_jobs()
|
|
end
|
|
end
|
|
|
|
def build_terrain_jobs(contacts) do
|
|
contacts
|
|
|> Enum.filter(fn contact -> contact.pos1 && contact.pos2 end)
|
|
|> Enum.map(fn contact ->
|
|
TerrainProfileWorker.new(%{"contact_id" => contact.id})
|
|
end)
|
|
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
|
|
end
|
|
|
|
def build_radar_jobs(contacts) do
|
|
contacts
|
|
|> Enum.filter(fn contact ->
|
|
contact.pos1 && contact.pos2 && not NarrClient.in_coverage?(contact.qso_timestamp)
|
|
end)
|
|
|> Enum.map(fn contact ->
|
|
CommonVolumeRadarWorker.new(%{"contact_id" => contact.id})
|
|
end)
|
|
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
|
|
end
|
|
|
|
def build_mechanism_jobs(contacts) do
|
|
contacts
|
|
|> Enum.filter(fn contact -> contact.pos1 && contact.pos2 end)
|
|
|> Enum.map(fn contact ->
|
|
MechanismClassifyWorker.new(%{"contact_id" => contact.id})
|
|
end)
|
|
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
|
|
end
|
|
|
|
def build_weather_jobs(contacts) do
|
|
contacts
|
|
|> Enum.flat_map(&jobs_for_contact/1)
|
|
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
|
|
end
|
|
|
|
def build_hrrr_jobs(contacts) do
|
|
contacts
|
|
|> Enum.flat_map(&hrrr_points_for_contact/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(contacts) do
|
|
contacts
|
|
|> Enum.flat_map(&iemre_job_for_contact/1)
|
|
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
|
|
end
|
|
|
|
def build_narr_jobs(contacts) do
|
|
contacts
|
|
|> Enum.flat_map(&narr_jobs_for_contact/1)
|
|
|> Enum.uniq_by(fn changeset -> changeset.changes.args end)
|
|
end
|
|
|
|
defp narr_jobs_for_contact(%{pos1: nil}), do: []
|
|
|
|
defp narr_jobs_for_contact(contact) do
|
|
# NARR analyses are 3-hourly (00/03/06/09/12/15/18/21 UTC). NarrFetchWorker
|
|
# rejects anything else, so snap here in the caller.
|
|
valid_time = NarrClient.snap_to_analysis_hour(contact.qso_timestamp)
|
|
|
|
# NCEI's NARR archive covers 1979-01-01 → 2014-10-02. Post-coverage contacts
|
|
# (hrrr_status = :unavailable but qso_timestamp >= 2014-10-02) are HRRR's
|
|
# responsibility — NARR would just 404.
|
|
if NarrClient.in_coverage?(valid_time) do
|
|
narr_jobs_for_path(contact, valid_time)
|
|
else
|
|
[]
|
|
end
|
|
end
|
|
|
|
defp narr_jobs_for_path(contact, valid_time) do
|
|
contact
|
|
|> Radio.contact_path_points()
|
|
|> Enum.flat_map(fn {lat, lon} ->
|
|
rlat = Float.round(lat * 4) / 4
|
|
rlon = Float.round(lon * 4) / 4
|
|
|
|
if Weather.find_nearest_narr(rlat, rlon, valid_time) do
|
|
[]
|
|
else
|
|
[
|
|
NarrFetchWorker.new(%{
|
|
"lat" => rlat,
|
|
"lon" => rlon,
|
|
"valid_time" => DateTime.to_iso8601(valid_time)
|
|
})
|
|
]
|
|
end
|
|
end)
|
|
end
|
|
|
|
defp jobs_for_contact(contact) do
|
|
contact
|
|
|> Radio.contact_path_points()
|
|
|> Enum.flat_map(fn {lat, lon} ->
|
|
build_asos_jobs(lat, lon, contact.qso_timestamp) ++
|
|
build_raob_jobs(lat, lon, contact.qso_timestamp)
|
|
end)
|
|
end
|
|
|
|
defp enqueue_iemre_jobs do
|
|
case Radio.unprocessed_iemre_contacts() do
|
|
[] ->
|
|
:ok
|
|
|
|
contacts ->
|
|
jobs = build_iemre_jobs(contacts)
|
|
|
|
if jobs != [] do
|
|
insert_all_chunked(jobs)
|
|
end
|
|
|
|
ids = Enum.map(contacts, & &1.id)
|
|
Radio.mark_iemre_queued!(ids)
|
|
|
|
enqueue_iemre_jobs()
|
|
end
|
|
end
|
|
|
|
defp hrrr_points_for_contact(%{pos1: nil}), do: []
|
|
|
|
defp hrrr_points_for_contact(contact) do
|
|
# HRRR archive starts mid-2014; NARR covers the 1979 → 2014-10-02 window.
|
|
# Short-circuit here so build_hrrr_jobs returns nothing for pre-coverage
|
|
# contacts and mark_hrrr_status! can pin them to :unavailable.
|
|
if NarrClient.in_coverage?(contact.qso_timestamp) do
|
|
[]
|
|
else
|
|
hrrr_points_for_path(contact)
|
|
end
|
|
end
|
|
|
|
defp hrrr_points_for_path(contact) do
|
|
rounded_time = HrrrClient.nearest_hrrr_hour(contact.qso_timestamp)
|
|
|
|
contact
|
|
|> Radio.contact_path_points()
|
|
|> Enum.flat_map(&hrrr_point_or_skip(&1, rounded_time))
|
|
end
|
|
|
|
defp hrrr_point_or_skip({lat, lon}, rounded_time) do
|
|
{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
|
|
|
|
defp build_asos_jobs(lat, lon, timestamp) do
|
|
start_dt = DateTime.add(timestamp, -2 * 3600, :second)
|
|
end_dt = DateTime.add(timestamp, 2 * 3600, :second)
|
|
|
|
stations = Weather.nearby_stations(lat, lon, "asos", @asos_radius_km)
|
|
station_ids = Enum.map(stations, & &1.id)
|
|
covered = Weather.station_ids_with_surface_observations(station_ids, start_dt, end_dt)
|
|
|
|
stations
|
|
|> Enum.reject(fn station -> MapSet.member?(covered, station.id) 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_contact(%{pos1: nil}), do: []
|
|
|
|
defp iemre_job_for_contact(contact) do
|
|
date = DateTime.to_date(contact.qso_timestamp)
|
|
|
|
contact
|
|
|> Radio.contact_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
|
|
build_raob_jobs(lat, lon, timestamp, @sounding_radius_km, nil)
|
|
end
|
|
|
|
defp build_raob_jobs(lat, lon, timestamp, radius_km, contact_id) do
|
|
sounding_times = Weather.sounding_times_around(timestamp)
|
|
stations = Weather.nearby_stations(lat, lon, "sounding", radius_km)
|
|
station_ids = Enum.map(stations, & &1.id)
|
|
covered = Weather.station_ids_with_soundings(station_ids, sounding_times)
|
|
|
|
for station <- stations,
|
|
sounding_time <- sounding_times,
|
|
not MapSet.member?(covered, {station.id, sounding_time}) do
|
|
args =
|
|
Map.merge(
|
|
%{
|
|
"fetch_type" => "raob",
|
|
"station_id" => station.id,
|
|
"station_code" => station.station_code,
|
|
"sounding_time" => DateTime.to_iso8601(sounding_time)
|
|
},
|
|
if(contact_id, do: %{"contact_id" => contact_id}, else: %{})
|
|
)
|
|
|
|
WeatherFetchWorker.new(args)
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Enqueue RAOB fetch jobs for the given contact at a progressively wider
|
|
radius until the radius yields at least one sounding station. Passes
|
|
the contact id through so `WeatherFetchWorker` can PubSub the contact
|
|
detail page when each sounding lands.
|
|
|
|
Returns `{:ok, %{radius_km, jobs_enqueued}}`, or `{:error, :no_stations}`
|
|
when even 1000 km finds nothing.
|
|
"""
|
|
@spec enqueue_raob_fetch_with_widening(Microwaveprop.Radio.Contact.t()) ::
|
|
{:ok, %{radius_km: pos_integer(), jobs_enqueued: non_neg_integer()}}
|
|
| {:error, :no_stations}
|
|
| {:error, :no_coords}
|
|
def enqueue_raob_fetch_with_widening(contact) do
|
|
case Radio.contact_path_points(contact) do
|
|
[] ->
|
|
{:error, :no_coords}
|
|
|
|
[{lat, lon} | _] ->
|
|
do_widening_raob(lat, lon, contact)
|
|
end
|
|
end
|
|
|
|
defp do_widening_raob(lat, lon, contact) do
|
|
Enum.reduce_while([150, 300, 600, 1000], {:error, :no_stations}, fn radius, _acc ->
|
|
stations = Weather.nearby_stations(lat, lon, "sounding", radius)
|
|
|
|
if stations == [] do
|
|
{:cont, {:error, :no_stations}}
|
|
else
|
|
{:halt, enqueue_raob_at_radius(lat, lon, contact, radius)}
|
|
end
|
|
end)
|
|
end
|
|
|
|
defp enqueue_raob_at_radius(lat, lon, contact, radius) do
|
|
jobs = build_raob_jobs(lat, lon, contact.qso_timestamp, radius, contact.id)
|
|
if jobs != [], do: insert_all_chunked(jobs)
|
|
{:ok, %{radius_km: radius, jobs_enqueued: length(jobs)}}
|
|
end
|
|
|
|
defp insert_all_chunked(jobs) do
|
|
jobs
|
|
|> Enum.chunk_every(@insert_batch_size)
|
|
|> Enum.each(&Oban.insert_all/1)
|
|
end
|
|
|
|
defp insert_unique(jobs) do
|
|
Enum.each(jobs, &Oban.insert/1)
|
|
end
|
|
|
|
defp mark_status!(ids, field, [_ | _]), do: Radio.set_enrichment_status!(ids, field, :queued)
|
|
defp mark_status!(ids, field, []), do: Radio.set_enrichment_status!(ids, field, :complete)
|
|
end
|