New AsosAdjustmentWorker runs every 10 minutes: - Fetches latest ASOS observations from all ~2900 US stations via IEM bulk currents API (parallel fetch across 51 state networks) - For each grid point within 75km of a reporting station, re-scores using fresh ASOS data (temp, dewpoint, wind, sky, pressure, precip) with HRRR refractivity gradient from the last hourly computation - Pushes updated scores to the map via PubSub Also stores HRRR profiles in the database during grid computation so the data persists for reference and ASOS blending.
133 lines
4.1 KiB
Elixir
133 lines
4.1 KiB
Elixir
defmodule Microwaveprop.Workers.PropagationGridWorker do
|
|
@moduledoc """
|
|
Hourly Oban worker that downloads the latest HRRR data and computes
|
|
propagation scores across the CONUS grid for all bands.
|
|
"""
|
|
|
|
use Oban.Worker,
|
|
queue: :propagation,
|
|
max_attempts: 3
|
|
|
|
alias Microwaveprop.Propagation
|
|
alias Microwaveprop.Propagation.Grid
|
|
alias Microwaveprop.Weather
|
|
alias Microwaveprop.Weather.HrrrClient
|
|
alias Microwaveprop.Weather.SoundingParams
|
|
|
|
require Logger
|
|
|
|
# Pause these queues while the grid fetch runs so they don't compete for bandwidth
|
|
@pause_queues [:hrrr, :weather, :iemre, :terrain]
|
|
|
|
@impl Oban.Worker
|
|
def perform(%Oban.Job{}) do
|
|
# HRRR takes ~45min to publish after the hour. Use 2 hours ago to ensure availability.
|
|
two_hours_ago = DateTime.add(DateTime.utc_now(), -2, :hour)
|
|
valid_time = two_hours_ago |> HrrrClient.nearest_hrrr_hour() |> DateTime.truncate(:second)
|
|
points = Grid.conus_points()
|
|
|
|
Logger.info("PropagationGrid: pausing backfill queues, fetching HRRR for #{valid_time}, #{length(points)} points")
|
|
Enum.each(@pause_queues, &Oban.pause_queue(queue: &1))
|
|
|
|
result =
|
|
with {:ok, grid_data} <- HrrrClient.fetch_grid(points, valid_time) do
|
|
# Store HRRR profiles in the database for future reference
|
|
store_hrrr_profiles(grid_data, valid_time)
|
|
|
|
scores = compute_scores(grid_data, valid_time)
|
|
|
|
Logger.info("PropagationGrid: computed #{length(scores)} scores, upserting")
|
|
|
|
case Propagation.upsert_scores(scores) do
|
|
{:ok, count} ->
|
|
Logger.info("PropagationGrid: upserted #{count} scores for #{valid_time}")
|
|
Phoenix.PubSub.broadcast(Microwaveprop.PubSub, "propagation:updated", {:propagation_updated, valid_time})
|
|
:ok
|
|
|
|
error ->
|
|
Logger.error("PropagationGrid: upsert failed: #{inspect(error)}")
|
|
error
|
|
end
|
|
end
|
|
|
|
Logger.info("PropagationGrid: resuming backfill queues")
|
|
Enum.each(@pause_queues, &Oban.resume_queue(queue: &1))
|
|
|
|
result
|
|
end
|
|
|
|
defp store_hrrr_profiles(grid_data, valid_time) do
|
|
stored =
|
|
Enum.flat_map(grid_data, fn {{lat, lon}, profile} ->
|
|
if profile.surface_temp_c do
|
|
params =
|
|
if is_list(profile.profile) and length(profile.profile) >= 3,
|
|
do: SoundingParams.derive(profile.profile)
|
|
|
|
attrs = %{
|
|
valid_time: valid_time,
|
|
lat: lat,
|
|
lon: lon,
|
|
run_time: profile.run_time,
|
|
profile: profile.profile || [],
|
|
hpbl_m: profile.hpbl_m,
|
|
pwat_mm: profile.pwat_mm,
|
|
surface_temp_c: profile.surface_temp_c,
|
|
surface_dewpoint_c: profile.surface_dewpoint_c,
|
|
surface_pressure_mb: profile.surface_pressure_mb
|
|
}
|
|
|
|
attrs =
|
|
if params do
|
|
Map.merge(attrs, %{
|
|
surface_refractivity: params.surface_refractivity,
|
|
min_refractivity_gradient: params.min_refractivity_gradient,
|
|
ducting_detected: params.ducting_detected,
|
|
duct_characteristics: params.duct_characteristics
|
|
})
|
|
else
|
|
attrs
|
|
end
|
|
|
|
[attrs]
|
|
else
|
|
[]
|
|
end
|
|
end)
|
|
|
|
# Batch upsert — skip duplicates
|
|
stored
|
|
|> Enum.chunk_every(500)
|
|
|> Enum.each(fn chunk ->
|
|
Weather.upsert_hrrr_profiles_batch(chunk)
|
|
end)
|
|
|
|
Logger.info("PropagationGrid: stored #{length(stored)} HRRR profiles")
|
|
end
|
|
|
|
defp compute_scores(grid_data, valid_time) do
|
|
grid_data
|
|
|> Task.async_stream(
|
|
fn {{lat, lon}, profile} ->
|
|
band_scores = Propagation.score_grid_point(profile, valid_time)
|
|
|
|
Enum.map(band_scores, fn r ->
|
|
%{
|
|
lat: lat,
|
|
lon: lon,
|
|
valid_time: valid_time,
|
|
band_mhz: r.band_mhz,
|
|
score: r.score,
|
|
factors: r.factors
|
|
}
|
|
end)
|
|
end,
|
|
max_concurrency: System.schedulers_online() * 2,
|
|
timeout: 30_000
|
|
)
|
|
|> Enum.flat_map(fn
|
|
{:ok, results} -> results
|
|
{:exit, _reason} -> []
|
|
end)
|
|
end
|
|
end
|