prop/lib/microwaveprop/workers/propagation_grid_worker.ex
Graham McIntire 65693ed415
feat(prop): Phase 2 cutover — Rust owns f01–f18
Rust `prop-grid-rs` has been validated end-to-end on talos5 on the
streamed-bands image (main-1776635096-6d91461): real chain steps
complete in ~7 s each and the peak-heap refactor landed.

Changes:
- `PropagationGridWorker.seed_chain` now enqueues a single f00 Oban
  job instead of f00–f18. The f00 analysis-hour keeps its enrichment
  (native-level duct, NEXRAD composite, commercial-link boost,
  ProfilesFile write) and the GridCache broadcast.
- `GridTaskEnqueuer.seed/1` is called again so grid_tasks gets the
  18 forecast-hour rows the Rust worker claims.
- Test asserts only one Oban job is enqueued and the matching 18
  grid_tasks rows land in the DB.
- `deployment-grid-rs.yaml` flips `PROP_SCORES_DIR` from
  `/data/scores_shadow` to `/data/scores` so Rust writes to the live
  score tree. No file collision — Elixir only writes the f00 files.

The hot pod memory shrink (6 Gi → 1 Gi) is deferred until one full
hourly cycle runs clean under the new split.
2026-04-19 16:51:14 -05:00

423 lines
16 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 """
Oban worker that downloads HRRR data and computes propagation scores
across the CONUS grid for all bands, one forecast hour at a time.
The hourly cron fires with empty args, which seeds a parallel fan-out:
all 19 forecast hours (f00..f18) are enqueued as independent jobs
against the `:propagation` queue. Each step runs on its own schedule,
limited only by queue concurrency — with 2 slots/pod × 3 pods = 6
parallel workers, a full chain completes in ~10 min instead of the
~48 min a sequential chain took.
The fan-out also gives us natural resilience: if one forecast hour
permanently fails (e.g. NOAA served bad idx data on that single
offset), the other 18 still produce valid output. Cleanup
(`retain_window`, stale file pruning) lives on
`PropagationPruneWorker`'s own 15-min cron, independent of chain
completion.
"""
use Oban.Worker,
queue: :propagation,
# Highest priority on the shared :propagation queue. MrmsFetchWorker
# and PropagationPruneWorker also live here; a backlog of those
# (e.g. after a deploy, or MRMS firing every 2 minutes) would
# otherwise starve the hourly chain because Oban dispatches
# priority-then-FIFO. Explicit so the invariant is visible.
priority: 0,
# Higher than the default 3 so a few DynamicLifeline rescues
# (e.g., a rolling deploy that kills a mid-flight chain step)
# don't exhaust the chain's retry budget and discard the whole
# run. Legitimate scoring errors still give up after 5 attempts.
max_attempts: 5,
# Deduplicate identical jobs across a 1-hour window. Uniqueness
# is over the full args set, so the seed (`%{}`) collapses with
# itself — FreshnessMonitor's 5-minute ticks during a long outage
# no longer stack 24 redundant f00-f18 chains — while chain steps
# with distinct `run_time` + `forecast_hour` args remain distinct
# from each other and from the seed. The period is the same as
# the hourly cron interval so successive hourly cron fires move
# past the window naturally.
unique: [period: 3600, states: [:available, :scheduled, :executing, :retryable]]
alias Microwaveprop.Commercial
alias Microwaveprop.Propagation
alias Microwaveprop.Propagation.BandConfig
alias Microwaveprop.Propagation.Grid
alias Microwaveprop.Propagation.GridTaskEnqueuer
alias Microwaveprop.Propagation.ProfilesFile
alias Microwaveprop.Propagation.ScoreCache
alias Microwaveprop.Weather
alias Microwaveprop.Weather.GridCache
alias Microwaveprop.Weather.HrrrClient
alias Microwaveprop.Weather.HrrrNativeClient
alias Microwaveprop.Weather.NexradClient
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
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 (+f01-f#{@max_forecast_hour} via grid_tasks)")
# Post-cutover: Elixir only runs the f00 analysis-hour step. f00 carries
# the expensive enrichment that Rust doesn't cover yet (native-level
# duct merge, NEXRAD composite, commercial-link degradation) and writes
# the ProfilesFile that /weather reads. f01..f18 go through the
# `grid_tasks` handoff queue to the Rust `prop-grid-rs` worker, which
# fetches + decodes + scores each forecast hour independently.
Oban.insert_all([new(%{"run_time" => DateTime.to_iso8601(run_time), "forecast_hour" => 0})])
_ = GridTaskEnqueuer.seed(run_time)
:ok
end
defp run_chain_step(run_time, fh) do
t_start = System.monotonic_time(:millisecond)
started_at = DateTime.utc_now()
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()
total_ms = System.monotonic_time(:millisecond) - t_start
record_timing(run_time, fh, valid_time, started_at, total_ms, result)
case result do
:ok ->
Phoenix.PubSub.broadcast(
Microwaveprop.PubSub,
"propagation:updated",
{:propagation_updated, [valid_time]}
)
Logger.info("PropagationGrid: fh=#{fh} step finished in #{format_duration(total_ms)}")
:ok
other ->
other
end
end
# Persist one row to propagation_run_timings. Swallow DB failures so
# the instrumentation can never brick the chain — if the table is
# missing or the connection is flaky, we log and carry on.
defp record_timing(run_time, fh, valid_time, started_at, duration_ms, result) do
{status, error} =
case result do
:ok -> {:ok, nil}
other -> {:failed, inspect(other)}
end
attrs = %{
run_time: DateTime.truncate(run_time, :second),
forecast_hour: fh,
valid_time: DateTime.truncate(valid_time, :second),
started_at: started_at,
finished_at: DateTime.add(started_at, duration_ms, :millisecond),
duration_ms: duration_ms,
status: status,
error: error
}
case Propagation.record_run_timing(attrs) do
{:ok, _} ->
:ok
{:error, changeset} ->
Logger.warning("PropagationGrid: timing insert failed: #{inspect(changeset.errors)}")
end
rescue
e -> Logger.warning("PropagationGrid: timing insert raised: #{inspect(e)}")
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} ->
# HRRR profiles used to be persisted to the hrrr_profiles table
# here for AsosAdjustmentWorker to re-score from. That was ~12
# minutes of JSONB inserts per chain (92k rows × 19 forecast
# hours) with no user-visible benefit — the scores live in the
# /data/scores files now, and AsosAdjustmentWorker is disabled.
# Per-contact HRRR enrichment still uses HrrrFetchWorker, which
# writes its own `is_grid_point: false` rows.
# Native-level duct fetch + wgrib2 pass costs ~7-11 min/hour and
# is only run on f00. Forecast hours fall back to
# derived[:min_refractivity_gradient] from the pressure-level
# profile — coarser (~250 m vs ~10-50 m) but good enough for
# forecast-hour ducting, which is inherently lower-confidence.
grid_data =
if forecast_hour == 0 do
merge_native_duct_data(grid_data, run_time, forecast_hour)
else
grid_data
end
: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()
# Persist the fully-enriched grid_data for every forecast hour
# so (a) /weather can show forecast-hour data after a pod
# restart and (b) point_detail can rebuild the factor
# breakdown for a clicked cell at any forecast hour by
# re-running the scorer against the stored profile.
persist_profiles(grid_data, valid_time)
# Weather map only shows the analysis hour — f01..f18 are
# forecast data that /weather doesn't render. Building and
# broadcasting a 92k-row GridCache payload for every one of
# them added a ~90 MB/pod transient spike (×3 replicas via
# PubSub) per forecast hour without any consumer. Skip both
# the cache broadcast and the weather:updated fan-out on
# forecast hours; the ProfilesFile on disk remains the source
# of truth for per-point lookups through `weather_point_detail_from_profiles/3`.
if forecast_hour == 0 do
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}
)
end
scores = compute_scores(grid_data, valid_time, forecast_hour)
case Propagation.replace_scores(scores, valid_time) do
{:ok, count} ->
Logger.info("PropagationGrid: #{label}#{count} scores for #{valid_time}")
warm_cache(valid_time)
# Broadcast progress *after* persistence so the map's
# pipeline chip only advances to "through +Nh" once that
# hour is actually readable from the scores file. Emitting
# this before the fetch would push the chip ahead of the
# map by the full forecast-hour wall time (~10 minutes).
Phoenix.PubSub.broadcast(
Microwaveprop.PubSub,
"propagation:pipeline",
{:propagation_pipeline_progress, %{forecast_hour: forecast_hour, valid_time: valid_time}}
)
:ok
error ->
Logger.error("PropagationGrid: #{label} replace failed: #{inspect(error)}")
error
end
error ->
Logger.warning("PropagationGrid: #{label} fetch failed: #{inspect(error)}")
error
end
end
defp persist_profiles(grid_data, valid_time) do
timed("profiles", fn ->
try do
ProfilesFile.write!(valid_time, grid_data)
rescue
e ->
Logger.warning("PropagationGrid: profiles write failed: #{inspect(e)}")
end
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 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
# Precompute per-link degradation once (≤10 SQL queries total).
# Commercial links cluster around DFW so most grid cells see nil —
# the per-cell path is now a pure haversine check, not a DB query.
lookup = Commercial.build_link_lookup(valid_time)
{merged, boosted} =
Enum.reduce(grid_data, {%{}, 0}, fn {{lat, lon} = point, profile}, {acc, count} ->
case Commercial.link_degradation_from_lookup({lat, lon}, lookup) do
nil ->
{Map.put(acc, point, profile), count}
degradation ->
{Map.put(acc, point, Map.put(profile, :commercial_link_degradation, degradation)), count + 1}
end
end)
if boosted > 0 do
Logger.info("PropagationGrid: commercial-link degradation available for #{boosted} grid cells")
end
merged
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, forecast_hour) do
# Algorithm is the primary scorer. `factors` is only populated for
# f00 (the analysis hour) — forecast hours skip the JSONB write so
# the scoring+upsert phase can land in under a minute instead of
# ~4-5 minutes. point_detail on forecast hours returns a nil
# breakdown, which the UI tolerates.
compute_scores_algorithm(grid_data, valid_time, forecast_hour == 0)
end
@doc false
# Public for testing. grid_data is a %{{lat, lon} => profile} map.
def compute_scores_algorithm(grid_data, valid_time, include_factors?) do
Microwaveprop.Instrument.span(
[:propagation_grid, :score_band],
%{point_count: map_size(grid_data)},
fn ->
grid_data
|> Task.async_stream(
&score_one_point(&1, valid_time, include_factors?),
max_concurrency: System.schedulers_online() * 2,
timeout: 30_000
)
|> Stream.flat_map(fn
{:ok, results} -> results
{:exit, _reason} -> []
end)
|> Enum.to_list()
end
)
end
defp score_one_point({{lat, lon}, profile}, valid_time, include_factors?) do
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: if(include_factors?, do: r.factors)
}
end)
end
end