A full f00-f18 sweep takes ~170 min of wall time, longer than the ~90 min gap between deploys on this cluster. Over the last 7 days every run was killed mid-sweep and zero runs completed. The propagation_scores table only ever held the scraps from partial runs that landed before the pod died, which is why the map looked like it was showing "the current hour" (or nothing). The worker now processes exactly one forecast hour per perform/1 (~8-10 min) and enqueues the next hour as a fresh Oban job. A cron fire with empty args seeds the chain at f00; subsequent runs carry run_time + forecast_hour args. At f18 the chain stops and the pruner runs. Each step broadcasts propagation:updated immediately so new hours appear on the map as they land. Wall time per perform drops from ~170 min to ~10 min, so Lifeline's 2h rescue window is no longer a factor and a pod restart loses at most one forecast hour instead of the whole sweep. Retries and max_attempts now describe a single fh, not the whole chain. Also: bump the :propagation queue from 1 to 2 slots so PropagationPruneWorker can run alongside the chain job, drop the worker timeout/1 from 90 min to 20 min to match single-fh runs, and pull Lifeline rescue_after from 120 min to 45 min to keep the safety net above the step timeout. Adds a "Data from HH:MM UTC · Nh ago/now/+Nh" indicator above the pipeline chip in both the mobile and desktop sidebars so the user can tell which valid_time the map is showing when the bottom timeline isn't visible.
377 lines
12 KiB
Elixir
377 lines
12 KiB
Elixir
defmodule Microwaveprop.Workers.PropagationGridWorker do
|
||
@moduledoc """
|
||
Oban worker that downloads HRRR data and computes propagation scores
|
||
across the CONUS grid for all bands, one forecast hour at a time.
|
||
|
||
Each `perform/1` processes a single forecast hour (~8–10 min of wall
|
||
time) and enqueues the next hour in the chain. A cron fire with
|
||
empty args seeds the chain at f00; subsequent runs carry
|
||
`forecast_hour` + `run_time` args. At f18 the chain stops and the
|
||
pruner cleans up old scores.
|
||
|
||
Splitting by forecast hour is a resilience play: a full sweep takes
|
||
~3 hours of wall time, longer than the typical pod-restart interval
|
||
on this deployment. Under the old "one big perform" design, any
|
||
deploy mid-sweep killed the whole run, and max_attempts would
|
||
exhaust without recording an error. Per-hour jobs survive deploys
|
||
because Lifeline only needs to rescue a single 10-minute step, and
|
||
retries re-fetch just that forecast hour.
|
||
"""
|
||
|
||
use Oban.Worker,
|
||
queue: :propagation,
|
||
max_attempts: 3
|
||
|
||
alias Microwaveprop.Commercial
|
||
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.NexradClient
|
||
alias Microwaveprop.Weather.SoundingParams
|
||
|
||
require Logger
|
||
|
||
# Hard ceiling for one forecast hour. A healthy step is ~8-10 min.
|
||
# 20 min gives 2× headroom for a slow HRRR fetch or scoring batch.
|
||
# Oban kills the executing process on timeout, which closes linked
|
||
# ports and cascades SIGKILL to any child wgrib2 subprocess.
|
||
@run_timeout_ms 20 * 60 * 1000
|
||
|
||
@impl Oban.Worker
|
||
def timeout(_job), do: @run_timeout_ms
|
||
|
||
@max_forecast_hour 18
|
||
|
||
@impl Oban.Worker
|
||
def perform(%Oban.Job{args: %{"forecast_hour" => fh, "run_time" => run_time_iso}}) do
|
||
{:ok, run_time, _} = DateTime.from_iso8601(run_time_iso)
|
||
run_chain_step(run_time, fh)
|
||
end
|
||
|
||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||
seed_chain()
|
||
end
|
||
|
||
@doc """
|
||
Enqueue the next chain step after a successful forecast-hour run.
|
||
|
||
Returns `{:ok, job}` when a new step is enqueued, or `:final` when
|
||
`fh` is already `@max_forecast_hour` so the chain has no more work.
|
||
Public so the chain entry point and tests can both exercise the
|
||
same enqueue path.
|
||
"""
|
||
@spec enqueue_next_step(DateTime.t(), non_neg_integer()) :: {:ok, Oban.Job.t()} | :final
|
||
def enqueue_next_step(_run_time, fh) when fh >= @max_forecast_hour, do: :final
|
||
|
||
def enqueue_next_step(%DateTime{} = run_time, fh) when fh >= 0 do
|
||
{:ok, _job} =
|
||
%{"run_time" => DateTime.to_iso8601(run_time), "forecast_hour" => fh + 1}
|
||
|> new()
|
||
|> Oban.insert()
|
||
end
|
||
|
||
defp seed_chain 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)
|
||
run_time = two_hours_ago |> HrrrClient.nearest_hrrr_hour() |> DateTime.truncate(:second)
|
||
|
||
Logger.info("PropagationGrid: seeding chain run_time=#{run_time}, f00-f#{@max_forecast_hour}")
|
||
|
||
{:ok, _job} =
|
||
%{"run_time" => DateTime.to_iso8601(run_time), "forecast_hour" => 0}
|
||
|> new()
|
||
|> Oban.insert()
|
||
|
||
:ok
|
||
end
|
||
|
||
defp run_chain_step(run_time, fh) do
|
||
t_start = System.monotonic_time(:millisecond)
|
||
points = Grid.conus_points()
|
||
valid_time = DateTime.add(run_time, fh * 3600, :second)
|
||
|
||
Logger.info("PropagationGrid: chain step run_time=#{run_time} fh=#{fh} (#{length(points)} points)")
|
||
|
||
result = process_forecast_hour(points, run_time, fh, valid_time)
|
||
:erlang.garbage_collect()
|
||
|
||
case result do
|
||
:ok ->
|
||
Phoenix.PubSub.broadcast(
|
||
Microwaveprop.PubSub,
|
||
"propagation:updated",
|
||
{:propagation_updated, [valid_time]}
|
||
)
|
||
|
||
total_ms = System.monotonic_time(:millisecond) - t_start
|
||
Logger.info("PropagationGrid: fh=#{fh} step finished in #{format_duration(total_ms)}")
|
||
|
||
case enqueue_next_step(run_time, fh) do
|
||
:final ->
|
||
Weather.prune_old_grid_profiles()
|
||
Propagation.prune_old_scores()
|
||
Logger.info("PropagationGrid: chain complete for run_time=#{run_time}")
|
||
|
||
{:ok, _job} ->
|
||
:ok
|
||
end
|
||
|
||
:ok
|
||
|
||
other ->
|
||
other
|
||
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")}"
|
||
|
||
# Emit progress so the map page's pipeline status chip can show
|
||
# which forecast hour is currently being scored ("now", "+1h", …).
|
||
Phoenix.PubSub.broadcast(
|
||
Microwaveprop.PubSub,
|
||
"propagation:pipeline",
|
||
{:propagation_pipeline_progress, %{forecast_hour: forecast_hour, valid_time: valid_time}}
|
||
)
|
||
|
||
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()
|
||
|
||
# NEXRAD current-hour composite reflectivity catches fast-moving
|
||
# convective cells between HRRR hourly analyses. Only useful for
|
||
# f00 — forecast hours can't see the future radar image.
|
||
grid_data =
|
||
if forecast_hour == 0 do
|
||
merge_nexrad_data(grid_data, valid_time)
|
||
else
|
||
grid_data
|
||
end
|
||
|
||
# Commercial-link inverse sensor — only meaningful for f00 because
|
||
# the measurement is of the current atmospheric state, not a forecast.
|
||
grid_data =
|
||
if forecast_hour == 0 do
|
||
merge_commercial_link_data(grid_data, valid_time)
|
||
else
|
||
grid_data
|
||
end
|
||
|
||
:erlang.garbage_collect()
|
||
|
||
# Build weather cache rows from in-memory grid_data — avoids the ~20s
|
||
# JSONB round trip that was crashing the worker before f01–f18 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 merge_commercial_link_data(grid_data, valid_time) do
|
||
# Compute link degradation once per distinct link cluster and cache by
|
||
# point. Commercial links only cluster around DFW so most grid points see
|
||
# nil — cheap no-op path dominates.
|
||
boosted =
|
||
Enum.count(grid_data, fn {{lat, lon}, _profile} ->
|
||
degradation = Commercial.link_degradation_at({lat, lon}, valid_time)
|
||
not is_nil(degradation)
|
||
end)
|
||
|
||
if boosted > 0 do
|
||
Logger.info("PropagationGrid: commercial-link degradation available for #{boosted} grid cells")
|
||
end
|
||
|
||
Map.new(grid_data, fn {{lat, lon} = point, profile} ->
|
||
case Commercial.link_degradation_at({lat, lon}, valid_time) do
|
||
nil -> {point, profile}
|
||
degradation -> {point, Map.put(profile, :commercial_link_degradation, degradation)}
|
||
end
|
||
end)
|
||
end
|
||
|
||
defp merge_nexrad_data(grid_data, valid_time) do
|
||
points = Map.keys(grid_data)
|
||
|
||
case timed("nexrad", fn -> NexradClient.fetch_frame(valid_time, points) end) do
|
||
{:ok, observations} ->
|
||
apply_nexrad_observations(grid_data, observations)
|
||
|
||
{:error, reason} ->
|
||
Logger.warning("PropagationGrid: NEXRAD fetch failed (continuing without): #{inspect(reason)}")
|
||
grid_data
|
||
end
|
||
end
|
||
|
||
defp apply_nexrad_observations(grid_data, observations) do
|
||
index = Map.new(observations, fn obs -> {{obs.lat, obs.lon}, obs.max_reflectivity_dbz} end)
|
||
non_zero = Enum.count(index, fn {_pt, dbz} -> dbz > 0 end)
|
||
Logger.info("PropagationGrid: NEXRAD merged (#{non_zero} cells with precip)")
|
||
|
||
Map.new(grid_data, fn {point, profile} ->
|
||
case Map.get(index, point) do
|
||
nil -> {point, profile}
|
||
dbz -> {point, Map.put(profile, :nexrad_max_reflectivity_dbz, dbz)}
|
||
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
|