prop/lib/microwaveprop/workers/propagation_grid_worker.ex
Graham McIntire 2405c5f169
Stop PropagationGridWorker crash-looping on forecast hours
warm_grid_cache_and_broadcast was re-SELECTing 92k hrrr_profiles rows
with JSONB profile/duct_characteristics columns on every forecast hour.
Row decoding exceeded the 15s default DB connection timeout, killing
the worker before f01-f18 could run. Oban retried from f00 and got
stuck in a loop: for the last several hours, no forecast hours were
being written and even f00 was stale.

Build the weather cache rows directly from the in-memory grid_data
the worker already has (Weather.build_grid_cache_rows/2) and
GridCache.broadcast_put directly. No DB round trip, no JSONB decode.

Also widen the cron from hourly to every 3 hours so a full f00-f18
sweep (~95 min) can actually complete before the next run starts,
and drop the redundant prune_old_scores() at worker start
(PropagationPruneWorker already runs every 15 min). Add a 60s query
timeout to load_weather_grid_from_db as defense-in-depth for the
cold-cache fill path that still uses it.
2026-04-13 08:20:28 -05:00

238 lines
7.5 KiB
Elixir
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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.BandConfig
alias Microwaveprop.Propagation.Grid
alias Microwaveprop.Propagation.ScoreCache
alias Microwaveprop.Weather
alias Microwaveprop.Weather.GridCache
alias Microwaveprop.Weather.HrrrClient
alias Microwaveprop.Weather.HrrrNativeClient
alias Microwaveprop.Weather.SoundingParams
require Logger
@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}")
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)
# Reclaim grid_data/profiles from this forecast hour before starting the next
:erlang.garbage_collect()
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)")
Weather.prune_old_grid_profiles()
Propagation.prune_old_scores()
:ok
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)
:erlang.garbage_collect()
# Fetch native duct metrics and merge into grid_data for scoring
grid_data = merge_native_duct_data(grid_data, run_time, forecast_hour)
:erlang.garbage_collect()
# Build weather cache rows from in-memory grid_data — avoids the ~20s
# JSONB round trip that was crashing the worker before f01f18 could run.
rows = Weather.build_grid_cache_rows(grid_data, valid_time)
GridCache.broadcast_put(valid_time, rows)
Phoenix.PubSub.broadcast(
Microwaveprop.PubSub,
"weather:updated",
{:weather_updated, 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}")
warm_cache(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 warm_cache(valid_time) do
Enum.each(BandConfig.all_bands(), fn band ->
Propagation.warm_cache_and_broadcast(band.freq_mhz, valid_time)
end)
ScoreCache.prune_older_than(DateTime.add(DateTime.utc_now(), -2, :hour))
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
count =
grid_data
|> Stream.flat_map(fn {{lat, lon}, profile} ->
build_profile_attrs_for_grid(lat, lon, valid_time, profile)
end)
|> Stream.chunk_every(500)
|> Enum.reduce(0, fn chunk, acc ->
Weather.upsert_hrrr_profiles_batch(chunk)
acc + length(chunk)
end)
Logger.info("PropagationGrid: stored #{count} HRRR profiles")
end
defp build_profile_attrs_for_grid(lat, lon, valid_time, profile) do
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
defp merge_native_duct_data(grid_data, run_time, forecast_hour) do
hour_dt = HrrrClient.nearest_hrrr_hour(run_time)
date = DateTime.to_date(hour_dt)
hour = hour_dt.hour
grid_spec = Grid.wgrib2_grid_spec()
case timed("native", fn ->
HrrrNativeClient.fetch_native_duct_grid(date, hour, grid_spec, forecast_hour)
end) do
{:ok, duct_grid} ->
Logger.info("PropagationGrid: merged #{map_size(duct_grid)} native duct cells")
apply_duct_grid(grid_data, duct_grid)
{:error, reason} ->
Logger.warning("PropagationGrid: native duct fetch failed (continuing without): #{inspect(reason)}")
grid_data
end
end
defp apply_duct_grid(grid_data, duct_grid) do
Map.new(grid_data, fn {point, profile} ->
case Map.get(duct_grid, point) do
nil -> {point, profile}
duct -> {point, Map.merge(profile, duct)}
end
end)
end
defp compute_scores(grid_data, valid_time) do
# Algorithm is the primary scorer. ML score stored in factors for comparison.
compute_scores_algorithm(grid_data, valid_time)
end
defp compute_scores_algorithm(grid_data, valid_time) do
grid_data
|> Task.async_stream(
fn {{lat, lon}, profile} ->
band_scores = Propagation.score_grid_point(profile, valid_time, lat, 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
)
|> Stream.flat_map(fn
{:ok, results} -> results
{:exit, _reason} -> []
end)
end
end