Score time-of-day per grid point using longitude/15 solar offset instead of hardcoded CST/CDT. Add PWAT as 10th scoring factor. Refine pressure thresholds. Update ML model and training pipeline to use local solar time.
185 lines
5.6 KiB
Elixir
185 lines
5.6 KiB
Elixir
defmodule Microwaveprop.Workers.PropagationGridWorker do
|
|
@moduledoc """
|
|
Hourly Oban worker that downloads HRRR data (analysis + forecast hours)
|
|
and computes propagation scores across the CONUS grid for all bands.
|
|
Fetches f00-f18 to provide 18-hour forecast timeline.
|
|
"""
|
|
|
|
use Oban.Worker,
|
|
queue: :propagation,
|
|
max_attempts: 3,
|
|
unique: [period: 3600, states: [:available, :scheduled, :executing, :retryable]]
|
|
|
|
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]
|
|
@max_forecast_hour 18
|
|
|
|
@impl Oban.Worker
|
|
def perform(%Oban.Job{}) do
|
|
t_start = System.monotonic_time(:millisecond)
|
|
|
|
# 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)
|
|
run_time = two_hours_ago |> HrrrClient.nearest_hrrr_hour() |> DateTime.truncate(:second)
|
|
points = Grid.conus_points()
|
|
|
|
Logger.info("PropagationGrid: run_time=#{run_time}, #{length(points)} points, f00-f#{@max_forecast_hour}")
|
|
|
|
Enum.each(@pause_queues, &Oban.pause_queue(queue: &1))
|
|
|
|
try do
|
|
for_result =
|
|
for fh <- 0..@max_forecast_hour do
|
|
valid_time = DateTime.add(run_time, fh * 3600, :second)
|
|
result = process_forecast_hour(points, run_time, fh, valid_time)
|
|
|
|
case result do
|
|
:ok -> valid_time
|
|
_ -> nil
|
|
end
|
|
end
|
|
|
|
valid_times = Enum.reject(for_result, &is_nil/1)
|
|
|
|
if valid_times != [] do
|
|
Phoenix.PubSub.broadcast(
|
|
Microwaveprop.PubSub,
|
|
"propagation:updated",
|
|
{:propagation_updated, valid_times}
|
|
)
|
|
end
|
|
|
|
total_ms = System.monotonic_time(:millisecond) - t_start
|
|
Logger.info("PropagationGrid: total time #{format_duration(total_ms)} (#{length(valid_times)} forecast hours)")
|
|
|
|
:ok
|
|
after
|
|
Logger.info("PropagationGrid: resuming backfill queues")
|
|
Enum.each(@pause_queues, &Oban.resume_queue(queue: &1))
|
|
|
|
Weather.prune_old_grid_profiles()
|
|
Propagation.prune_old_scores()
|
|
end
|
|
end
|
|
|
|
defp process_forecast_hour(points, run_time, forecast_hour, valid_time) do
|
|
label = "f#{String.pad_leading(Integer.to_string(forecast_hour), 2, "0")}"
|
|
|
|
case timed(label, fn ->
|
|
HrrrClient.fetch_grid(points, run_time, forecast_hour: forecast_hour)
|
|
end) do
|
|
{:ok, grid_data} ->
|
|
store_hrrr_profiles(grid_data, valid_time)
|
|
scores = compute_scores(grid_data, valid_time)
|
|
|
|
case Propagation.upsert_scores(scores, prune: false) do
|
|
{:ok, count} ->
|
|
Logger.info("PropagationGrid: #{label} → #{count} scores for #{valid_time}")
|
|
:ok
|
|
|
|
error ->
|
|
Logger.error("PropagationGrid: #{label} upsert failed: #{inspect(error)}")
|
|
error
|
|
end
|
|
|
|
error ->
|
|
Logger.warning("PropagationGrid: #{label} fetch failed: #{inspect(error)}")
|
|
error
|
|
end
|
|
end
|
|
|
|
defp timed(label, fun) do
|
|
t0 = System.monotonic_time(:millisecond)
|
|
result = fun.()
|
|
elapsed = System.monotonic_time(:millisecond) - t0
|
|
Logger.info("PropagationGrid: #{label} took #{format_duration(elapsed)}")
|
|
result
|
|
end
|
|
|
|
defp format_duration(ms) when ms < 1000, do: "#{ms}ms"
|
|
defp format_duration(ms), do: "#{Float.round(ms / 1000, 1)}s"
|
|
|
|
defp store_hrrr_profiles(grid_data, valid_time) do
|
|
stored =
|
|
Enum.flat_map(grid_data, fn {{lat, lon}, profile} ->
|
|
temp = profile.surface_temp_c
|
|
|
|
if temp && temp > -80 && temp < 60 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)
|
|
|
|
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, lon)
|
|
|
|
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
|