Adds Microwaveprop.Propagation.AsosNudge: a pure IDW bias-field module that takes ASOS observations + HRRR profiles and returns re-scored grid rows for every cell within 250km of a reporting station. Upper-air fields (min_refractivity_gradient, pwat_mm, hpbl_m, profile, duct metadata) pass through unchanged so HRRR's signal isn't clobbered. The old AsosAdjustmentWorker was unwired and buggy — nil'd out ~22% of the scoring weight and wrote orphan timestamps. Replaced with a slim worker that queries the latest HRRR valid_time, fetches live ASOS currents, calls AsosNudge.compute/3, and upserts onto (lat, lon, valid_time, band_mhz) so nudged values overwrite the HRRR hour cleanly instead of polluting available_valid_times. After each upsert it warms ScoreCache and broadcasts propagation:updated so live /map clients refresh. Cron hooked up every 10 minutes in config.exs and dev.exs. Also cleaned up the stale "dev has propagation disabled" note in CLAUDE.md. 13 new AsosNudge unit tests cover: residual computation (co-located, out-of-grid, nil fields), IDW weighting (single station, far station, two equidistant stations, nil component handling), upper-air preservation, and the compute/3 entry point's shape and radius filter. Drive-by Styler formatting touched a handful of unrelated files from `mix format`.
417 lines
14 KiB
Elixir
417 lines
14 KiB
Elixir
defmodule Microwaveprop.Propagation do
|
|
@moduledoc false
|
|
|
|
import Ecto.Query
|
|
|
|
alias Microwaveprop.Propagation.BandConfig
|
|
alias Microwaveprop.Propagation.Grid
|
|
alias Microwaveprop.Propagation.GridScore
|
|
alias Microwaveprop.Propagation.ScoreCache
|
|
alias Microwaveprop.Propagation.Scorer
|
|
alias Microwaveprop.Repo
|
|
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: Scorer.precip_to_rate_mmhr(hrrr_profile[:precip_mm]),
|
|
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]
|
|
}
|
|
|
|
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
|
|
|
|
Enum.map(BandConfig.all_bands(), fn band_config ->
|
|
result = Scorer.composite_score(conditions, band_config)
|
|
result = Map.put(result, :band_mhz, band_config.freq_mhz)
|
|
if duct_info, do: put_in(result, [:factors, :duct_info], duct_info), else: result
|
|
end)
|
|
end
|
|
|
|
@doc """
|
|
Upsert propagation scores in batches within a transaction so readers see all-or-nothing.
|
|
|
|
Options:
|
|
* `:prune` - whether to prune old scores after upsert (default true)
|
|
"""
|
|
@spec upsert_scores(Enumerable.t(), keyword()) :: {:ok, non_neg_integer()} | {:error, term()}
|
|
def upsert_scores(scores, opts \\ []) do
|
|
now = DateTime.truncate(DateTime.utc_now(), :second)
|
|
|
|
result =
|
|
Repo.transaction(
|
|
fn ->
|
|
scores
|
|
|> Stream.map(fn s ->
|
|
%{
|
|
id: Ecto.UUID.generate(),
|
|
lat: s.lat,
|
|
lon: s.lon,
|
|
valid_time: s.valid_time,
|
|
band_mhz: s.band_mhz,
|
|
score: s.score,
|
|
factors: s.factors,
|
|
inserted_at: now,
|
|
updated_at: now
|
|
}
|
|
end)
|
|
|> Stream.chunk_every(500)
|
|
|> Enum.reduce(0, fn chunk, acc ->
|
|
{count, _} =
|
|
Repo.insert_all(GridScore, chunk,
|
|
on_conflict:
|
|
from(g in GridScore,
|
|
update: [
|
|
set: [
|
|
score: fragment("EXCLUDED.score"),
|
|
factors: fragment("EXCLUDED.factors"),
|
|
updated_at: fragment("EXCLUDED.updated_at")
|
|
]
|
|
],
|
|
where: g.score != fragment("EXCLUDED.score")
|
|
),
|
|
conflict_target: [:lat, :lon, :valid_time, :band_mhz]
|
|
)
|
|
|
|
acc + count
|
|
end)
|
|
end,
|
|
timeout: 600_000
|
|
)
|
|
|
|
if Keyword.get(opts, :prune, true) do
|
|
case result do
|
|
{:ok, _count} -> prune_old_scores()
|
|
_ -> :ok
|
|
end
|
|
end
|
|
|
|
result
|
|
end
|
|
|
|
@doc """
|
|
Remove scores with valid_times older than 2 hours. Called on a cron by
|
|
`Microwaveprop.Workers.PropagationPruneWorker` and also at the start of
|
|
each `PropagationGridWorker.perform/1` as a safety net. The 5-minute
|
|
timeout gives enough headroom for a catch-up run (potentially millions of
|
|
rows across four indexes) after a stretch of failed compute jobs.
|
|
"""
|
|
@spec prune_old_scores() :: :ok
|
|
def prune_old_scores do
|
|
cutoff = DateTime.add(DateTime.utc_now(), -2, :hour)
|
|
|
|
{deleted, _} =
|
|
Repo.delete_all(from(gs in GridScore, where: gs.valid_time < ^cutoff), timeout: 300_000)
|
|
|
|
if deleted > 0 do
|
|
Logger.info("PropagationScores: pruned #{deleted} old scores (before #{cutoff})")
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Returns distinct valid_times for a band, ordered ascending.
|
|
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 ScoreCache.valid_times(band_mhz) do
|
|
[] -> available_valid_times_from_db(band_mhz, cutoff)
|
|
cached -> filter_fresh(cached, cutoff)
|
|
end
|
|
end
|
|
|
|
defp filter_fresh(times, cutoff) do
|
|
Enum.filter(times, fn t -> DateTime.compare(t, cutoff) != :lt end)
|
|
end
|
|
|
|
defp available_valid_times_from_db(band_mhz, cutoff) do
|
|
times =
|
|
Repo.all(
|
|
from(gs in GridScore,
|
|
where: gs.band_mhz == ^band_mhz and gs.valid_time >= ^cutoff,
|
|
select: gs.valid_time,
|
|
distinct: gs.valid_time,
|
|
order_by: [asc: gs.valid_time]
|
|
)
|
|
)
|
|
|
|
if times == [] do
|
|
# No future times — return the single most recent as fallback
|
|
case latest_valid_time(band_mhz) do
|
|
nil -> []
|
|
t -> [t]
|
|
end
|
|
else
|
|
times
|
|
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 ->
|
|
full = load_scores_from_db(band_mhz, time)
|
|
ScoreCache.put(band_mhz, time, full)
|
|
|
|
full
|
|
|> filter_bounds(bounds)
|
|
|> Enum.map(&Map.put(&1, :valid_time, time))
|
|
end
|
|
end
|
|
|
|
defp load_scores_from_db(band_mhz, time) do
|
|
Repo.all(
|
|
from(gs in GridScore,
|
|
where: gs.band_mhz == ^band_mhz and gs.valid_time == ^time,
|
|
select: %{lat: gs.lat, lon: gs.lon, score: gs.score}
|
|
)
|
|
)
|
|
end
|
|
|
|
@doc """
|
|
Load the full CONUS score set for `{band_mhz, valid_time}` from the DB and
|
|
broadcast it to every `ScoreCache` in the cluster. Intended to be called from
|
|
`PropagationGridWorker` after each upsert 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 = load_scores_from_db(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
|
|
Repo.one(from(gs in GridScore, where: gs.band_mhz == ^band_mhz, select: min(gs.valid_time)))
|
|
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_db(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_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 point_forecast_from_db(band_mhz, lat, lon, now) do
|
|
Repo.all(
|
|
from(gs in GridScore,
|
|
where:
|
|
gs.band_mhz == ^band_mhz and
|
|
gs.lat == ^lat and
|
|
gs.lon == ^lon and
|
|
gs.valid_time >= ^now,
|
|
select: %{valid_time: gs.valid_time, score: gs.score},
|
|
order_by: [asc: gs.valid_time]
|
|
)
|
|
)
|
|
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
|
|
step = Grid.step()
|
|
snapped_lat = Float.round(Float.round(lat / step) * step, 3)
|
|
snapped_lon = Float.round(Float.round(lon / step) * step, 3)
|
|
|
|
time = valid_time || latest_valid_time(band_mhz)
|
|
|
|
case time do
|
|
nil ->
|
|
nil
|
|
|
|
_ ->
|
|
Repo.one(
|
|
from(gs in GridScore,
|
|
where:
|
|
gs.band_mhz == ^band_mhz and
|
|
gs.valid_time == ^time and
|
|
gs.lat == ^snapped_lat and
|
|
gs.lon == ^snapped_lon,
|
|
select: %{
|
|
lat: gs.lat,
|
|
lon: gs.lon,
|
|
score: gs.score,
|
|
factors: gs.factors,
|
|
valid_time: gs.valid_time
|
|
}
|
|
)
|
|
)
|
|
end
|
|
end
|
|
|
|
@doc "Get the latest valid_time across all scores."
|
|
@spec latest_valid_time() :: DateTime.t() | nil
|
|
def latest_valid_time do
|
|
Repo.one(from(gs in GridScore, select: max(gs.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
|
|
Repo.one(from(gs in GridScore, where: gs.band_mhz == ^band_mhz, select: max(gs.valid_time)))
|
|
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
|