prop/lib/microwaveprop/propagation.ex
Graham McIntire 32fdb2583d
Shrink PropagationGridWorker per-forecast-hour memory footprint
Two related optimizations in the propagation chain hot path. Both
land on the f01..f18 step that was previously OOM-killing prod
pods after the HRRR pressure-level footprint halved wasn't enough
to fit inside 4 Gi.

1. Skip the GridCache broadcast on forecast hours.

   /weather only ever renders the analysis hour (latest_grid_valid_time
   feeds the map). Building 92k rows, serializing them through PubSub,
   and rebuilding the {lat,lon}→row map on all three replicas was
   pure waste for f01..f18 — no consumer was reading that data. Only
   f00 now calls build_grid_cache_rows + broadcast_put. Point lookups
   for non-analysis hours still work through ProfilesFile on disk
   (weather_point_detail_from_profiles/3) exactly as before.

2. Fold replace_scores into a single streaming pass.

   The old path did `Enum.to_list/1` on the ~460k-entry score stream
   followed by `Enum.group_by/2`, holding two full copies of the grid
   before any file was written. A single `Enum.reduce/3` that folds
   each score into a per-band accumulator keeps only one copy and
   eliminates the group_by intermediate entirely. The public
   signature — an Enumerable in, {:ok, count} out — is unchanged.
2026-04-15 08:36:53 -05:00

451 lines
16 KiB
Elixir

defmodule Microwaveprop.Propagation do
@moduledoc false
alias Microwaveprop.Propagation.BandConfig
alias Microwaveprop.Propagation.Grid
alias Microwaveprop.Propagation.ProfilesFile
alias Microwaveprop.Propagation.ScoreCache
alias Microwaveprop.Propagation.Scorer
alias Microwaveprop.Propagation.ScoresFile
alias Microwaveprop.Weather.SoundingParams
require Logger
@ml_key :propagation_ml
@ml_module Microwaveprop.Propagation.Model
@doc """
Loads the ML model from disk, compiles the predict function, and caches both
in persistent_term. No-op if the model file doesn't exist or ML deps unavailable.
"""
@spec load_ml_model() :: :ok
def load_ml_model do
if Code.ensure_loaded?(@ml_module) do
case @ml_module.load() do
{:ok, params} ->
predict_fn = @ml_module.compile_predict()
:persistent_term.put(@ml_key, {predict_fn, params})
Logger.info("PropagationML: model loaded and compiled")
:ok
:error ->
Logger.info("PropagationML: no model file found, using algorithm scorer only")
:ok
end
else
Logger.info("PropagationML: ML dependencies not available")
:ok
end
end
@doc "Returns cached {predict_fn, params} tuple, or nil if not loaded."
@spec ml_model() :: {function(), term()} | nil
def ml_model do
:persistent_term.get(@ml_key, nil)
end
@doc """
Score a single grid point across all bands using HRRR profile data.
Uses ML model if loaded, falls back to algorithm scorer.
Returns a list of %{band_mhz, score, factors} maps.
"""
@spec score_grid_point(map(), DateTime.t(), float(), float()) ::
[%{band_mhz: non_neg_integer(), score: non_neg_integer(), factors: map()}]
def score_grid_point(hrrr_profile, valid_time, latitude, longitude) do
derived = derive_from_hrrr(hrrr_profile)
temp_c = hrrr_profile.surface_temp_c
dewpoint_c = hrrr_profile.surface_dewpoint_c
# Skip points with missing or physically impossible surface data
if is_nil(temp_c) or is_nil(dewpoint_c) or temp_c < -80 or temp_c > 60 or
dewpoint_c < -80 or dewpoint_c > 50 do
[]
else
score_grid_point_with_data(hrrr_profile, valid_time, temp_c, dewpoint_c, derived, latitude, longitude)
end
end
defp score_grid_point_with_data(hrrr_profile, valid_time, temp_c, dewpoint_c, derived, latitude, longitude) do
# Algorithm is the primary scorer — always used for the map score.
# ML score stored in factors as :ml_score for comparison/analysis.
score_with_algorithm(hrrr_profile, valid_time, temp_c, dewpoint_c, derived, latitude, longitude)
end
defp score_with_algorithm(hrrr_profile, valid_time, temp_c, dewpoint_c, derived, _latitude, longitude) do
temp_f = Scorer.c_to_f(temp_c)
dewpoint_f = Scorer.c_to_f(dewpoint_c)
conditions = %{
abs_humidity: Scorer.absolute_humidity(temp_c, dewpoint_c),
temp_f: temp_f,
dewpoint_f: dewpoint_f,
wind_speed_kts: Scorer.wind_speed_kts(hrrr_profile[:wind_u], hrrr_profile[:wind_v]),
sky_cover_pct: hrrr_profile[:cloud_cover_pct],
utc_hour: valid_time.hour,
utc_minute: valid_time.minute,
month: valid_time.month,
longitude: longitude,
pressure_mb: hrrr_profile.surface_pressure_mb,
prev_pressure_mb: nil,
rain_rate_mmhr: merged_rain_rate(hrrr_profile),
min_refractivity_gradient: hrrr_profile[:native_min_gradient] || derived[:min_refractivity_gradient],
bl_depth_m: hrrr_profile[:hpbl_m],
pwat_mm: hrrr_profile[:pwat_mm],
best_duct_band_ghz: hrrr_profile[:best_duct_freq_ghz] || hrrr_profile[:best_duct_band_ghz]
}
duct_info =
if hrrr_profile[:duct_count] && hrrr_profile[:duct_count] > 0 do
%{
duct_count: hrrr_profile[:duct_count],
best_duct_freq_ghz: hrrr_profile[:best_duct_freq_ghz],
max_duct_thickness_m: hrrr_profile[:max_duct_thickness_m],
ducts: hrrr_profile[:ducts] || []
}
end
link_degradation = hrrr_profile[:commercial_link_degradation]
Enum.map(BandConfig.all_bands(), fn band_config ->
result = Scorer.composite_score(conditions, band_config)
boosted_score =
if link_degradation do
Scorer.commercial_link_boost(result.score, link_degradation)
else
result.score
end
result =
result
|> Map.put(:score, boosted_score)
|> Map.put(:band_mhz, band_config.freq_mhz)
result =
if link_degradation do
put_in(result, [:factors, :commercial_link_degradation], link_degradation)
else
result
end
if duct_info, do: put_in(result, [:factors, :duct_info], duct_info), else: result
end)
end
# Pick the heavier of HRRR's hourly accumulation-derived rate and NEXRAD's
# reflectivity-derived rate. NEXRAD catches fast-moving convective cells that
# fall between HRRR hourly analyses; HRRR catches broad stratiform rain that
# NEXRAD reports as low dBZ. Taking max lets either source trigger the rain
# penalty without double-counting.
defp merged_rain_rate(hrrr_profile) do
hrrr_rate = Scorer.precip_to_rate_mmhr(hrrr_profile[:precip_mm])
nexrad_rate = Scorer.dbz_to_rain_rate_mmhr(hrrr_profile[:nexrad_max_reflectivity_dbz])
max(hrrr_rate, nexrad_rate)
end
@doc """
Replace every propagation score for `valid_time` with `scores`.
Used by `PropagationGridWorker` on the hot path. Scores are written
as binary files on disk via `ScoresFile.write!/3`, one file per
band.
Consumes `scores` in a single streaming pass that folds each score
straight into a per-band accumulator. Previously this function ran
`Enum.to_list/1` followed by `Enum.group_by/2`, which held two full
copies of the ~460k-entry grid (list + grouped list) in memory at
once — the hot path's largest transient spike after native-duct
merge. The single-pass reduce keeps only one copy and buys back
~100 MB of headroom per forecast-hour step.
"""
@spec replace_scores(Enumerable.t(), DateTime.t()) :: {:ok, non_neg_integer()} | {:error, term()}
def replace_scores(scores, %DateTime{} = valid_time) do
{per_band, total} =
Enum.reduce(scores, {%{}, 0}, fn score, {acc, count} ->
{Map.update(acc, score.band_mhz, [score], &[score | &1]), count + 1}
end)
Enum.each(per_band, fn {band_mhz, band_scores} ->
try do
ScoresFile.write!(band_mhz, valid_time, band_scores)
rescue
e ->
Logger.warning("Propagation: ScoresFile write failed for band=#{band_mhz} vt=#{valid_time}: #{inspect(e)}")
end
end)
{:ok, total}
end
@doc """
Remove score files with valid_times older than 2 hours. Called on
a cron by `Microwaveprop.Workers.PropagationPruneWorker`.
"""
@spec prune_old_scores() :: :ok
def prune_old_scores do
cutoff = DateTime.add(DateTime.utc_now(), -2, :hour)
file_deleted = ScoresFile.prune_older_than(cutoff)
profiles_deleted = ProfilesFile.prune_older_than(cutoff)
if file_deleted + profiles_deleted > 0 do
Logger.info(
"PropagationScores: pruned #{file_deleted} old score files + " <>
"#{profiles_deleted} profile files (before #{cutoff})"
)
end
:ok
end
@doc """
Returns distinct valid_times for a band, ordered ascending. Always
reads from the on-disk `ScoresFile` store — the `ScoreCache` only
holds whatever hours have been fetched or broadcast, which can be a
partial view of what's actually on disk, so using it as the source
of truth for the timeline makes new forecast hours invisible until
the cache happens to catch up. Filters out times more than 1 hour in
the past, but always includes the most recent valid_time so there's
always data to display.
"""
@spec available_valid_times(non_neg_integer()) :: [DateTime.t()]
def available_valid_times(band_mhz) do
cutoff = DateTime.add(DateTime.utc_now(), -3600, :second)
case ScoresFile.list_valid_times(band_mhz) do
[] -> []
times -> filter_or_latest(times, cutoff)
end
end
defp filter_or_latest(times, cutoff) do
fresh = Enum.filter(times, fn t -> DateTime.compare(t, cutoff) != :lt end)
if fresh == [] do
[Enum.max(times, DateTime)]
else
fresh
end
end
@doc """
Get scores for a band at a specific valid_time, optionally within a bounding box.
If valid_time is nil, uses the earliest available (current analysis hour).
Excludes factors for performance.
"""
@spec scores_at(non_neg_integer(), DateTime.t() | nil, map() | nil) ::
[%{lat: float(), lon: float(), score: non_neg_integer(), valid_time: DateTime.t()}]
def scores_at(band_mhz, valid_time, bounds \\ nil) do
time = valid_time || earliest_valid_time(band_mhz)
case time do
nil ->
[]
_ ->
scores_at_fetch(band_mhz, time, bounds)
end
end
defp scores_at_fetch(band_mhz, time, bounds) do
case ScoreCache.fetch_bounds(band_mhz, time, bounds) do
{:ok, scores} ->
Enum.map(scores, &Map.put(&1, :valid_time, time))
:miss ->
read_from_disk_and_cache(band_mhz, time, bounds)
end
end
@doc """
Variant of `scores_at/3` that always reads from the `.ntms` file on
disk and overwrites the cache entry, rather than returning whatever
the cache happens to hold. Use from update paths (the map's
`propagation_updated` handler) where the underlying file has just
been rewritten but the cache may still contain the previous chain's
scores because of the race between `propagation:cache` fan-out and
`propagation:updated` delivery.
"""
@spec scores_at_fresh(non_neg_integer(), DateTime.t(), map() | nil) ::
[%{lat: float(), lon: float(), score: non_neg_integer(), valid_time: DateTime.t()}]
def scores_at_fresh(band_mhz, %DateTime{} = valid_time, bounds \\ nil) do
read_from_disk_and_cache(band_mhz, valid_time, bounds)
end
defp read_from_disk_and_cache(band_mhz, time, bounds) do
full = ScoresFile.read_bounds(band_mhz, time)
ScoreCache.put(band_mhz, time, full)
full
|> filter_bounds(bounds)
|> Enum.map(&Map.put(&1, :valid_time, time))
end
@doc """
Load the full CONUS score set for `{band_mhz, valid_time}` from the
on-disk binary file and broadcast it to every `ScoreCache` in the
cluster. Called from `PropagationGridWorker` after each forecast
hour so all pods have a warm cache by the time clients begin
requesting the new hour.
"""
@spec warm_cache_and_broadcast(non_neg_integer(), DateTime.t()) :: :ok
def warm_cache_and_broadcast(band_mhz, valid_time) do
scores = ScoresFile.read_bounds(band_mhz, valid_time)
ScoreCache.broadcast_put(band_mhz, valid_time, scores)
:ok
end
defp filter_bounds(scores, nil), do: scores
defp filter_bounds(scores, %{"south" => s, "north" => n, "west" => w, "east" => e}) do
Enum.filter(scores, fn %{lat: lat, lon: lon} ->
lat >= s and lat <= n and lon >= w and lon <= e
end)
end
@doc "Get the latest scores for a band (alias for scores_at with earliest valid_time)."
@spec latest_scores(non_neg_integer(), map() | nil) ::
[%{lat: float(), lon: float(), score: non_neg_integer(), valid_time: DateTime.t()}]
def latest_scores(band_mhz, bounds \\ nil) do
scores_at(band_mhz, nil, bounds)
end
defp earliest_valid_time(band_mhz) do
case ScoresFile.list_valid_times(band_mhz) do
[earliest | _] -> earliest
[] -> nil
end
end
@doc "Get scores across all forecast hours for a single grid point (for sparkline)."
@spec point_forecast(non_neg_integer(), float(), float()) ::
[%{valid_time: DateTime.t(), score: non_neg_integer()}]
def point_forecast(band_mhz, lat, lon) do
{snapped_lat, snapped_lon} = snap_to_grid(lat, lon)
now = DateTime.utc_now()
case ScoreCache.valid_times(band_mhz) do
[] -> point_forecast_from_store(band_mhz, snapped_lat, snapped_lon, now)
cached -> point_forecast_from_cache(band_mhz, snapped_lat, snapped_lon, now, cached)
end
end
defp point_forecast_from_store(band_mhz, lat, lon, now) do
band_mhz
|> ScoresFile.list_valid_times()
|> Enum.filter(fn t -> DateTime.compare(t, now) != :lt end)
|> Enum.map(&point_forecast_entry(&1, band_mhz, lat, lon))
|> Enum.reject(&is_nil/1)
end
defp point_forecast_entry(valid_time, band_mhz, lat, lon) do
case ScoresFile.read_point(band_mhz, valid_time, lat, lon) do
nil -> nil
score -> %{valid_time: valid_time, score: score}
end
end
defp point_forecast_from_cache(band_mhz, lat, lon, now, cached_times) do
cached_times
|> Enum.filter(fn t -> DateTime.compare(t, now) != :lt end)
|> Enum.map(fn t ->
case ScoreCache.fetch_point(band_mhz, t, lat, lon) do
{:ok, score} -> %{valid_time: t, score: score}
:miss -> nil
end
end)
|> Enum.reject(&is_nil/1)
end
defp snap_to_grid(lat, lon) do
step = Grid.step()
{Float.round(Float.round(lat / step) * step, 3), Float.round(Float.round(lon / step) * step, 3)}
end
@doc "Get the full score and factors for a specific grid point, snapped to nearest grid."
@spec point_detail(non_neg_integer(), float(), float(), DateTime.t() | nil) ::
%{
lat: float(),
lon: float(),
score: non_neg_integer(),
factors: map(),
valid_time: DateTime.t()
}
| nil
def point_detail(band_mhz, lat, lon, valid_time \\ nil) do
{snapped_lat, snapped_lon} = snap_to_grid(lat, lon)
time = valid_time || latest_valid_time(band_mhz)
case time do
nil ->
nil
_ ->
case ScoresFile.read_point(band_mhz, time, snapped_lat, snapped_lon) do
nil ->
nil
score ->
%{
lat: snapped_lat,
lon: snapped_lon,
score: score,
factors: factors_for(band_mhz, time, snapped_lat, snapped_lon),
valid_time: time
}
end
end
end
# Rebuild the factor breakdown for a clicked grid cell by rescoring
# the persisted HRRR profile. Forecast hours don't persist profiles
# (see PropagationGridWorker.process_forecast_hour) so they just get
# an empty map, which the JS popup tolerates by omitting the
# breakdown block.
defp factors_for(band_mhz, valid_time, lat, lon) do
case ProfilesFile.read_point(valid_time, lat, lon) do
nil ->
%{}
profile ->
profile
|> score_grid_point(valid_time, lat, lon)
|> Enum.find(fn r -> r.band_mhz == band_mhz end)
|> case do
%{factors: factors} when is_map(factors) -> factors
_ -> %{}
end
end
end
@doc "Get the latest valid_time across all bands."
@spec latest_valid_time() :: DateTime.t() | nil
def latest_valid_time do
ScoresFile.latest_valid_time()
end
@doc "Get the latest valid_time for a specific band."
@spec latest_valid_time(non_neg_integer()) :: DateTime.t() | nil
def latest_valid_time(band_mhz) do
case ScoresFile.list_valid_times(band_mhz) do
[] -> nil
times -> Enum.max(times, DateTime)
end
end
# Prefer the persisted scalar — `hrrr_profiles` already stored this at
# ingestion time and AsosAdjustmentWorker loads 92k rows per tick without
# the JSONB `profile` column to avoid a Jason.decode! storm on the DB pool.
defp derive_from_hrrr(%{min_refractivity_gradient: grad}) when is_number(grad) do
%{min_refractivity_gradient: grad * 1.0}
end
defp derive_from_hrrr(%{profile: profile}) when is_list(profile) and length(profile) >= 3 do
case SoundingParams.derive(profile) do
nil -> %{}
derived -> %{min_refractivity_gradient: derived.min_refractivity_gradient}
end
end
defp derive_from_hrrr(_), do: %{}
end