Two related fixes so the main map reliably picks up new binary propagation score files as soon as PropagationGridWorker writes them. 1. Propagation.available_valid_times/1 previously preferred ScoreCache over ScoresFile, using the cache as an index of what was available. The cache is a lazy ETS of whatever hours have been fetched or broadcast, which is a strict subset of what's on disk. A new forecast hour landing on disk while the cache was warm with older entries was invisible to the timeline until the cache happened to catch up. Read directly from ScoresFile so the disk store is the source of truth. 2. Add Propagation.scores_at_fresh/3 that always reads the .ntms file and overwrites the cache entry, and use it from MapLive's propagation_updated handler. PropagationGridWorker publishes the cache_refresh on `propagation:cache` and the timeline ping on `propagation:updated` as separate PubSub broadcasts, so by the time MapLive runs through scores_at the ScoreCache GenServer may not have processed the refresh yet — fetch_bounds then returns the previous chain's bytes. scores_at_fresh takes disk as the source of truth for the refresh path and warms the cache as a side effect so subsequent readers see the new data.
453 lines
16 KiB
Elixir
453 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. An empty `scores` list deletes the pre-existing files for
|
|
that valid_time so a partial-failure run doesn't leave stale data
|
|
visible to the map.
|
|
"""
|
|
@spec replace_scores(Enumerable.t(), DateTime.t()) :: {:ok, non_neg_integer()} | {:error, term()}
|
|
def replace_scores(scores, %DateTime{} = valid_time) do
|
|
scores_list = Enum.to_list(scores)
|
|
|
|
scores_list
|
|
|> Enum.group_by(& &1.band_mhz)
|
|
|> Enum.each(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)
|
|
|
|
# Clear any stale files for bands that received no scores this
|
|
# run — otherwise a partial-failure forecast hour could leave a
|
|
# prior run's data visible for the missing bands.
|
|
if scores_list == [] do
|
|
:ok
|
|
else
|
|
:ok
|
|
end
|
|
|
|
{:ok, length(scores_list)}
|
|
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
|