prop/lib/microwaveprop/workers/propagation_grid_worker.ex
Graham McIntire b80878056d
feat(grid-rs): Rust worker for HRRR f01..f18 propagation chain
Extracts the memory-hostile HRRR fetch → decode → score pipeline for
forecast hours 1–18 into a separate Rust service (`prop-grid-rs`).
Elixir retains f00 with its native-duct + NEXRAD + commercial-link
enrichment; Rust handles the other 18 steps per hourly chain.

Hand-off via a new `grid_tasks` table (`FOR UPDATE SKIP LOCKED` claim
from Rust, Postgrex `NOTIFY propagation_ready` back to Elixir). Rust
writes the same on-disk score-grid format Elixir already uses, to the
same `/data/scores` NFS tree. Phase 1 ships in shadow mode with
PROP_SCORES_DIR=/data/scores_shadow.

Rust crate layout at `rust/prop_grid_rs/` (1:1 module parity with the
Elixir source it ports):
  - grid, region, band_config, scores_file, sounding_params
  - scorer: all 10 factors + composite. Matches the Elixir scorer
    byte-for-byte across 115 golden-fixture samples (5 scenarios
    × 23 bands).
  - decoder: wgrib2 subprocess + Fortran-record lola-binary parser
  - fetcher: HRRR URL derivation + idx cache (1h TTL) + 8-way
    parallel byte-range downloads with 429/5xx retry
  - pipeline: end-to-end chain step, f00 rejected at the boundary
  - db: sqlx grid_tasks claim/complete + propagation_ready NOTIFY
  - bin/worker: tokio main loop, JSON logs, SIGTERM-safe

91 Rust tests + clippy -D warnings clean. 2,158 Elixir tests green.

Elixir additions:
  - `GridTaskEnqueuer` seeds fh=1..18 rows from
    `PropagationGridWorker.seed_chain/0`
  - `NotifyListener` LISTEN → warm `ScoreCache` → PubSub fan-out
  - `ShadowComparator` diffs prod vs shadow `.ntms` bodies
  - `mix rust.golden` task writes the Rust-side golden fixture

k8s: new `deployment-grid-rs.yaml` pinned to talos5 (32 GB NUC) via
`prop-grid-rs=primary` nodeSelector + `workload=grid-rs:NoSchedule`
toleration, 512 Mi limit, sharing the existing NFS `/data` mount.

The plan document is at `plans/vivid-hatching-quail.md` (local to my
workstation); phases 2 (cutover) and 3 (talos5 concurrency tuning)
follow after 72h of shadow-mode parity.
2026-04-19 15:42:49 -05:00

434 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-f#{@max_forecast_hour} (parallel)")
# Fan out all 19 forecast hours at once. Each step is independent
# (different HRRR URL, different valid_time file output) so there's
# no dependency that requires the old sequential chain pattern.
# With 2 slots/pod × 3 pods = 6 concurrent workers the chain wall
# time drops from ~48 min to ~10 min. Cleanup (retain_window,
# prune) runs on PropagationPruneWorker's own 15-min cron.
jobs =
for fh <- 0..@max_forecast_hour do
new(%{"run_time" => DateTime.to_iso8601(run_time), "forecast_hour" => fh})
end
Oban.insert_all(jobs)
# Seed the grid_tasks table in parallel so the Rust `prop-grid-rs`
# worker can pick up f01..f18 from a shared NFS mount. During Phase
# 1 shadow mode this runs alongside the Elixir chain; the Rust
# writes land in /data/scores_shadow and are diffed by
# `ShadowComparator`. At cutover the Elixir fan-out above is
# trimmed to just fh=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