feat(hrdps): activate Canadian propagation chain end-to-end

Stages 6, 7, 10 of the HRDPS plan, plus the Elixir read-side merge
that makes Rust's `.hrdps.prop` writes show up on /map.

Read side:

- `ScoresFile.path_for_hrdps/2` + `read_hrdps/2` — companions to the
  HRRR equivalents. Same decode path; different filename.
- `ScoresFile.read_bounds/3` now reads HRRR + HRDPS files and
  concatenates the cells. Files are disjoint by construction
  (Grid.hrdps_only_points excludes CONUS) so no de-dup needed.
- `ScoresFile.read_point/4` checks `.prop`, then `.hrdps.prop`, then
  legacy `.ntms` — first hit wins. Cells outside their owning region
  read as no-data in the other file's grid, so the order doesn't
  matter for correctness.
- `Propagation.warm_cache_and_broadcast/2` reads both files and
  merges before broadcasting to the cluster ScoreCache. A single
  missing file (HRDPS pre-cycle, HRRR briefly absent) is now OK —
  the cache warms with whichever side is available.

Activation:

- runtime.exs cron: HRDPS seed worker fires at HH:35 of each cycle
  hour (05:35Z, 11:35Z, 17:35Z, 23:35Z) — 5h after cycle, 30 min
  staggered from HRRR's */15 schedule.
- HrdpsGridWorker moduledoc updated: no longer dormant; Rust
  prop-grid-rs handles source='hrdps' rows.
- map_live North America bbox extended to (24.5, 60.0, -141.0, -52.0)
  for the out-of-coverage banner check; visitors anywhere in CONUS
  or the Canadian extent now stay banner-free.

Cleanup:

- mix hrdps.probe deleted — its risks are retired and the production
  HrdpsClient covers the same surface.

What's still in motion (not blocking the Canadian /map):

- Per-valid_time profile + scalar artifacts for HRDPS (would let
  /weather show Canadian temperature/refractivity layers). Rust
  pipeline currently skips those so the HRRR-written companion files
  aren't clobbered — separate stage when /weather UX is decided.
- FreshnessMonitor HRDPS staleness tracking — HRDPS has its own cron
  and a different cadence than HRRR; HRRR-stale logic doesn't apply
  cleanly. Defer until prod tells us what alarms we actually want.
This commit is contained in:
Graham McIntire 2026-04-29 17:29:37 -05:00
parent 46c2d84e5c
commit 1da78a80e5
No known key found for this signature in database
GPG key ID: F4ABF488E6029E59
7 changed files with 123 additions and 196 deletions

View file

@ -247,6 +247,13 @@ if config_env() == :prod do
# so duplicate fires for an already-seeded run_time are a single # so duplicate fires for an already-seeded run_time are a single
# idempotent INSERT batch. # idempotent INSERT batch.
{"5,20,35,50 * * * *", Microwaveprop.Workers.PropagationGridWorker}, {"5,20,35,50 * * * *", Microwaveprop.Workers.PropagationGridWorker},
# HRDPS Canadian propagation chain. ECCC publishes HRDPS f000
# ~3-4h after each 00/06/12/18Z cycle, so seeding at HH:35 of
# each cycle hour (5h after cycle, 30 min offset from HRRR's
# */15 schedule) gives a comfortable margin before the f000
# files land. The Rust prop-grid-rs HRDPS branch (run_chain_step_hrdps)
# drains the resulting source='hrdps' rows.
{"35 5,11,17,23 * * *", Microwaveprop.Workers.HrdpsGridWorker},
# GEFS runs publish ~3-4h after the 00/06/12/18Z cycle; seed the # GEFS runs publish ~3-4h after the 00/06/12/18Z cycle; seed the
# Day 2-7 outlook 5h after each run to stay safely past NOMADS' # Day 2-7 outlook 5h after each run to stay safely past NOMADS'
# publication lag. Seeder enqueues f024-f168 at 6-hour cadence # publication lag. Seeder enqueues f024-f168 at 6-hour cadence

View file

@ -345,27 +345,38 @@ defmodule Microwaveprop.Propagation do
end end
@doc """ @doc """
Load the full CONUS score set for `{band_mhz, valid_time}` from the Load the full North America score set for `{band_mhz, valid_time}` from
on-disk binary file and broadcast it to every `ScoreCache` in the the on-disk binary files (HRRR `.prop` + HRDPS `.hrdps.prop`, merged) and
cluster. Called from `PropagationGridWorker` after each forecast broadcast it to every `ScoreCache` in the cluster. Called from
hour so all pods have a warm cache by the time clients begin `PropagationGridWorker` after each forecast hour so all pods have a warm
requesting the new hour. cache by the time clients begin requesting the new hour.
Returns `{:error, reason}` when the score file is missing or Returns `{:error, :enoent}` only when neither file exists. A single
corrupt callers distinguish those from the successful empty-grid missing file (HRDPS pre-cycle, HRRR briefly absent) is OK the cache
case to avoid poisoning the cache with `[]` on a bad read. warms with whichever side is available.
""" """
@spec warm_cache_and_broadcast(non_neg_integer(), DateTime.t()) :: @spec warm_cache_and_broadcast(non_neg_integer(), DateTime.t()) ::
:ok | {:error, :enoent | :invalid_format} :ok | {:error, :enoent | :invalid_format}
def warm_cache_and_broadcast(band_mhz, valid_time) do def warm_cache_and_broadcast(band_mhz, valid_time) do
case ScoresFile.read(band_mhz, valid_time) do hrrr =
{:ok, payload} -> case ScoresFile.read(band_mhz, valid_time) do
scores = ScoresFile.extract_points(payload, nil) {:ok, payload} -> ScoresFile.extract_points(payload, nil)
ScoreCache.broadcast_put(band_mhz, valid_time, scores) _ -> nil
:ok end
{:error, reason} -> hrdps =
{:error, reason} case ScoresFile.read_hrdps(band_mhz, valid_time) do
{:ok, payload} -> ScoresFile.extract_points(payload, nil)
_ -> nil
end
case {hrrr, hrdps} do
{nil, nil} ->
{:error, :enoent}
{h, c} ->
ScoreCache.broadcast_put(band_mhz, valid_time, (h || []) ++ (c || []))
:ok
end end
end end

View file

@ -50,13 +50,25 @@ defmodule Microwaveprop.Propagation.ScoresFile do
Application.get_env(:microwaveprop, :propagation_scores_dir, "/data/scores") Application.get_env(:microwaveprop, :propagation_scores_dir, "/data/scores")
end end
@doc "Returns the full path for a `(band_mhz, valid_time)` score file." @doc "Returns the full path for a `(band_mhz, valid_time)` HRRR score file."
@spec path_for(non_neg_integer(), DateTime.t()) :: String.t() @spec path_for(non_neg_integer(), DateTime.t()) :: String.t()
def path_for(band_mhz, %DateTime{} = valid_time) when is_integer(band_mhz) do def path_for(band_mhz, %DateTime{} = valid_time) when is_integer(band_mhz) do
iso = valid_time |> DateTime.truncate(:second) |> DateTime.to_iso8601() iso = valid_time |> DateTime.truncate(:second) |> DateTime.to_iso8601()
Path.join([base_dir(), Integer.to_string(band_mhz), "#{iso}.prop"]) Path.join([base_dir(), Integer.to_string(band_mhz), "#{iso}.prop"])
end end
@doc """
Returns the path for the HRDPS-source companion score file. Sibling of
`path_for/2` both files coexist for the same `(band, valid_time)`.
Read-side functions transparently merge the two so callers get the
union of CONUS + Canadian scored cells.
"""
@spec path_for_hrdps(non_neg_integer(), DateTime.t()) :: String.t()
def path_for_hrdps(band_mhz, %DateTime{} = valid_time) when is_integer(band_mhz) do
iso = valid_time |> DateTime.truncate(:second) |> DateTime.to_iso8601()
Path.join([base_dir(), Integer.to_string(band_mhz), "#{iso}.hrdps.prop"])
end
@doc false @doc false
@spec legacy_path_for(non_neg_integer(), DateTime.t()) :: String.t() @spec legacy_path_for(non_neg_integer(), DateTime.t()) :: String.t()
def legacy_path_for(band_mhz, %DateTime{} = valid_time) when is_integer(band_mhz) do def legacy_path_for(band_mhz, %DateTime{} = valid_time) when is_integer(band_mhz) do
@ -241,9 +253,35 @@ defmodule Microwaveprop.Propagation.ScoresFile do
@spec read_bounds(non_neg_integer(), DateTime.t(), map() | nil) :: @spec read_bounds(non_neg_integer(), DateTime.t(), map() | nil) ::
[%{lat: float(), lon: float(), score: non_neg_integer()}] [%{lat: float(), lon: float(), score: non_neg_integer()}]
def read_bounds(band_mhz, %DateTime{} = valid_time, bounds \\ nil) do def read_bounds(band_mhz, %DateTime{} = valid_time, bounds \\ nil) do
case read(band_mhz, valid_time) do hrrr_cells =
{:ok, payload} -> extract_points(payload, bounds) case read(band_mhz, valid_time) do
_ -> [] {:ok, payload} -> extract_points(payload, bounds)
_ -> []
end
hrdps_cells =
case read_hrdps(band_mhz, valid_time) do
{:ok, payload} -> extract_points(payload, bounds)
_ -> []
end
# The two grids are disjoint by construction (Grid.hrdps_only_points
# excludes CONUS), so concat without dedup. Order isn't load-bearing —
# callers downstream of /scores/cells don't depend on ordering.
hrrr_cells ++ hrdps_cells
end
@doc """
Read the HRDPS-source score file for `(band_mhz, valid_time)`. Same
decode shape as `read/2`; uses the `.hrdps.prop` extension.
"""
@spec read_hrdps(non_neg_integer(), DateTime.t()) ::
{:ok, map()} | {:error, :enoent | :invalid_format}
def read_hrdps(band_mhz, %DateTime{} = valid_time) when is_integer(band_mhz) do
case File.read(path_for_hrdps(band_mhz, valid_time)) do
{:ok, binary} -> decode(binary)
{:error, :enoent} -> {:error, :enoent}
{:error, other} -> {:error, other}
end end
end end
@ -261,7 +299,11 @@ defmodule Microwaveprop.Propagation.ScoresFile do
@spec read_point(non_neg_integer(), DateTime.t(), float(), float()) :: non_neg_integer() | nil @spec read_point(non_neg_integer(), DateTime.t(), float(), float()) :: non_neg_integer() | nil
def read_point(band_mhz, %DateTime{} = valid_time, lat, lon) do def read_point(band_mhz, %DateTime{} = valid_time, lat, lon) do
Enum.find_value( Enum.find_value(
[path_for(band_mhz, valid_time), legacy_path_for(band_mhz, valid_time)], [
path_for(band_mhz, valid_time),
path_for_hrdps(band_mhz, valid_time),
legacy_path_for(band_mhz, valid_time)
],
&open_and_read_point(&1, lat, lon) &open_and_read_point(&1, lat, lon)
) )
end end

View file

@ -13,11 +13,10 @@ defmodule Microwaveprop.Workers.HrdpsGridWorker do
## Activation ## Activation
**Not yet wired into runtime.exs cron.** The Rust prop-grid-rs worker Active. The Rust prop-grid-rs worker handles `source: "hrdps"` rows via
needs an HRDPS branch (plan stage 4) before it can drain `pipeline::run_chain_step_hrdps`, fetching from MSC Datamart, writing
`source: "hrdps"` rows correctly. Until then this module is dormant scores to `<base>/<band>/<iso>.hrdps.prop` (sibling of the HRRR
scaffolding building blocks ready for the cron wire-up once the `.prop` file for the same cycle). Cron entry lives in `runtime.exs`.
Rust counterpart ships.
""" """
use Oban.Worker, use Oban.Worker,

View file

@ -15,13 +15,14 @@ defmodule MicrowavepropWeb.MapLive do
@default_center %{lat: 32.897, lon: -97.038} @default_center %{lat: 32.897, lon: -97.038}
@default_zoom 7 @default_zoom 7
# Continental-US bounding box (excludes Alaska, Hawaii, Puerto Rico). Used # North American coverage bbox: HRRR (CONUS) + HRDPS (Canadian extent
# to decide whether to show the "data is CONUS-only" banner to visitors # to 60°N — Arctic deferred until SRTM gap is filled). Used to decide
# whose Cloudflare geolocation falls outside the dataset's coverage. # whether to show the "out of coverage area" banner to visitors whose
# Cloudflare geolocation falls outside the combined dataset.
@conus_lat_min 24.5 @conus_lat_min 24.5
@conus_lat_max 49.5 @conus_lat_max 60.0
@conus_lon_min -125.0 @conus_lon_min -141.0
@conus_lon_max -66.5 @conus_lon_max -52.0
# How often to re-query oban_jobs for pipeline state while the map is # How often to re-query oban_jobs for pipeline state while the map is
# open. Short enough that the "Updating…" chip shows up within a page- # open. Short enough that the "Updating…" chip shows up within a page-

View file

@ -1,165 +0,0 @@
defmodule Mix.Tasks.Hrdps.Probe do
@shortdoc "Throwaway probe: fetch one HRDPS variable, decode at sample cities"
@moduledoc """
One-shot probe for the Canadian HRDPS model. Fetches a single GRIB2
file from MSC Datamart and uses `wgrib2 -lon` to extract values at
five Canadian cities, printing a table for visual sanity-checking.
Designed as a wedge to retire risk before committing to a real
HrdpsClient: validates URL pattern, GRIB2 decode, and rotated-lat/lon
reprojection without touching the DB, Oban, or the score grid.
Throwaway code. Will be deleted once the real client lands.
mix hrdps.probe # latest cycle, f000, 2m temperature
mix hrdps.probe --hour 6 # f006 from latest cycle
mix hrdps.probe --cycle 20260429T12Z
URL pattern (current as of 2026-04-29; ECCC reorganized away from the
old `/model_hrdps/` root to a date-prefixed structure):
https://dd.weather.gc.ca/{YYYYMMDD}/WXO-DD/model_hrdps/continental/2.5km/
{HH}/{FFF}/{YYYYMMDD}T{HH}Z_MSC_HRDPS_{VAR}_{LEVEL}_RLatLon0.0225_PT{FFF}H.grib2
"""
use Mix.Task
@datamart_base "https://dd.weather.gc.ca"
# Variable / level slug for 2m temperature in HRDPS naming. HRRR-side
# equivalent is `TMP:2 m above ground`; HRDPS uses MSC's per-variable
# filename convention.
@variable_slug "TMP_AGL-2m"
@cities [
{"YYZ Toronto", 43.68, -79.63},
{"YVR Vancouver", 49.19, -123.18},
{"YYC Calgary", 51.11, -114.02},
{"YOW Ottawa", 45.32, -75.67},
{"YHZ Halifax", 44.88, -63.51}
]
@impl Mix.Task
def run(argv) do
{opts, _, _} =
OptionParser.parse(argv,
strict: [cycle: :string, hour: :integer],
aliases: [c: :cycle, h: :hour]
)
{:ok, _} = Application.ensure_all_started(:req)
forecast_hour = Keyword.get(opts, :hour, 0)
cycle = opts |> Keyword.get(:cycle) |> resolve_cycle()
url = build_url(cycle, forecast_hour)
IO.puts("Fetching: #{url}")
{:ok, path} = download(url)
IO.puts("Downloaded #{path} (#{File.stat!(path).size} bytes)")
print_inventory(path)
print_city_table(path, cycle, forecast_hour)
end
# ----- Cycle resolution -----
defp resolve_cycle(nil), do: latest_published_cycle()
defp resolve_cycle(<<year::binary-size(4), month::binary-size(2), day::binary-size(2), "T", hour::binary-size(2), "Z">>) do
%{date: "#{year}#{month}#{day}", hour: hour}
end
defp resolve_cycle(other) do
Mix.raise("Invalid --cycle #{inspect(other)} — expected YYYYMMDDTHHZ")
end
# HRDPS publishes ~3-4h after each cycle (00/06/12/18Z). Walk back from
# "now - 4h" until we find a cycle whose root directory exists.
defp latest_published_cycle do
now = DateTime.utc_now()
candidates = for h <- 0..3, do: DateTime.add(now, -((4 + h * 6) * 3600), :second)
candidates
|> Enum.map(&snap_to_cycle/1)
|> Enum.uniq()
|> Enum.find(&cycle_published?/1)
|> case do
nil -> Mix.raise("No published HRDPS cycle found in the last 24h")
cycle -> cycle
end
end
defp snap_to_cycle(%DateTime{} = dt) do
cycle_hour = div(dt.hour, 6) * 6
%{date: Calendar.strftime(dt, "%Y%m%d"), hour: String.pad_leading("#{cycle_hour}", 2, "0")}
end
defp cycle_published?(%{date: date, hour: hour}) do
url = "#{@datamart_base}/#{date}/WXO-DD/model_hrdps/continental/2.5km/#{hour}/000/"
case Req.head(url, retry: false, connect_options: [timeout: 5_000]) do
{:ok, %{status: 200}} -> true
_ -> false
end
end
# ----- URL + download -----
defp build_url(%{date: date, hour: hour}, forecast_hour) do
fff = String.pad_leading("#{forecast_hour}", 3, "0")
filename = "#{date}T#{hour}Z_MSC_HRDPS_#{@variable_slug}_RLatLon0.0225_PT#{fff}H.grib2"
"#{@datamart_base}/#{date}/WXO-DD/model_hrdps/continental/2.5km/#{hour}/#{fff}/#{filename}"
end
defp download(url) do
path = Path.join(System.tmp_dir!(), "hrdps_probe_#{System.unique_integer([:positive])}.grib2")
case Req.get(url, into: File.stream!(path), retry: false) do
{:ok, %{status: 200}} -> {:ok, path}
{:ok, %{status: status}} -> Mix.raise("Datamart returned HTTP #{status} for #{url}")
{:error, reason} -> Mix.raise("Fetch failed: #{inspect(reason)}")
end
end
# ----- wgrib2 inspection -----
defp print_inventory(path) do
{output, 0} = System.cmd("wgrib2", [path], stderr_to_stdout: true)
IO.puts("\nGRIB2 inventory:")
IO.puts(" " <> String.trim(output))
end
defp print_city_table(path, cycle, forecast_hour) do
args = [path] ++ Enum.flat_map(@cities, fn {_, lat, lon} -> ["-lon", "#{lon}", "#{lat}"] end)
{output, 0} = System.cmd("wgrib2", args, stderr_to_stdout: true)
values = parse_wgrib2_lon_output(output)
IO.puts("\nHRDPS 2m temperature, cycle #{cycle.date}T#{cycle.hour}Z f#{pad3(forecast_hour)}:\n")
IO.puts(" city lat,lon temp")
IO.puts(" -------------------- -------------------- ------")
@cities
|> Enum.zip(values)
|> Enum.each(fn {{name, lat, lon}, kelvin} ->
celsius = kelvin - 273.15
IO.puts(" #{String.pad_trailing(name, 20)} #{format_coord(lat, lon)} #{format_temp(celsius)}")
end)
end
# wgrib2 -lon emits one line for the GRIB record with all `lon=…lat=…val=…`
# entries colon-separated. For five cities we get five `val=` matches.
defp parse_wgrib2_lon_output(output) do
~r/val=([\-\d\.]+)/
|> Regex.scan(output)
|> Enum.map(fn [_, v] -> String.to_float(v) end)
end
defp format_coord(lat, lon), do: "#{:io_lib.format("~6.2f", [lat])},#{:io_lib.format("~7.2f", [lon])}"
defp format_temp(celsius), do: "#{:io_lib.format("~5.1f", [celsius])}°C"
defp pad3(n), do: String.pad_leading("#{n}", 3, "0")
end

View file

@ -253,6 +253,38 @@ defmodule Microwaveprop.Propagation.ScoresFileTest do
test "returns an empty list when the file doesn't exist" do test "returns an empty list when the file doesn't exist" do
assert ScoresFile.read_bounds(10_000, ~U[2026-04-14 18:00:00Z]) == [] assert ScoresFile.read_bounds(10_000, ~U[2026-04-14 18:00:00Z]) == []
end end
test "merges cells from a sibling .hrdps.prop file when present" do
# Write the HRRR file via the public writer, then copy it to the
# HRDPS path. Both files now contain the same cell — read_bounds
# should return both, doubling the cell list. (Production HRDPS
# files have a different bbox; this test only verifies the merge
# plumbing routes through both paths.)
vt = ~U[2026-04-14 18:00:00Z]
ScoresFile.write!(10_000, vt, [%{lat: 25.0, lon: -125.0, score: 42}])
hrrr_path = ScoresFile.path_for(10_000, vt)
hrdps_path = ScoresFile.path_for_hrdps(10_000, vt)
File.cp!(hrrr_path, hrdps_path)
scores = ScoresFile.read_bounds(10_000, vt)
assert length(scores) == 2
assert Enum.all?(scores, &(&1.score == 42))
end
test "returns HRDPS cells alone when only the .hrdps.prop file exists" do
vt = ~U[2026-04-14 18:00:00Z]
# Use write! to land a file at the HRRR path, then move it to
# the HRDPS path so only the HRDPS file remains.
ScoresFile.write!(10_000, vt, [%{lat: 25.0, lon: -125.0, score: 80}])
File.rename!(
ScoresFile.path_for(10_000, vt),
ScoresFile.path_for_hrdps(10_000, vt)
)
assert [%{score: 80}] = ScoresFile.read_bounds(10_000, vt)
end
end end
describe "read_point/4" do describe "read_point/4" do