Hot pods were restart-looping every ~25 minutes on liveness probe timeouts. Root cause: ScoreCacheReconciler mirrored every .prop file from NFS into ETS, including 7 days of GEFS Day 2-7 forecasts. With 23 bands × 43 valid_times × ~2 MB per grid = 1.95 GB in the propagation_score_cache table alone; GC sweeps starved the scheduler enough that /live dropped its 3 s budget. The /map UI only ever requests the [now-1h, now+18h] window. Share that bound as Propagation.hot_cache_window/0 and apply it in the reconciler's disk-scan path plus NotifyListener's post-warm prune. Long-horizon GEFS files stay on disk and are still served via the lazy read_from_disk_and_cache path when requested directly. Adds ScoreCache.prune_outside_window/2 (inclusive bounds) and updates the reconciler tests to use relative-to-now timestamps since hardcoded fixture dates now drift out of window.
184 lines
5.8 KiB
Elixir
184 lines
5.8 KiB
Elixir
defmodule Microwaveprop.Propagation.ScoreCache do
|
|
@moduledoc """
|
|
Node-local ETS cache of propagation scores keyed by `{band_mhz, valid_time}`.
|
|
|
|
Populated eagerly by `Microwaveprop.Workers.PropagationGridWorker` after each
|
|
hourly compute (via `broadcast_put/3`) and lazily by `Microwaveprop.Propagation.scores_at/3`
|
|
on a cache miss. `broadcast_put/3` fans the payload out across the cluster on
|
|
the `"propagation:cache"` PubSub topic so every BEAM node stays in sync with
|
|
a single compute.
|
|
|
|
Internally each cache entry stores a `%{{lat, lon} => score}` map so that
|
|
per-point lookups (used by `point_forecast/3` and `point_detail/4`) are O(1).
|
|
"""
|
|
use GenServer
|
|
|
|
alias Phoenix.PubSub
|
|
|
|
@table :propagation_score_cache
|
|
@topic "propagation:cache"
|
|
@pubsub Microwaveprop.PubSub
|
|
|
|
@type score :: %{lat: float(), lon: float(), score: non_neg_integer()}
|
|
@type bounds :: %{optional(String.t()) => float()}
|
|
|
|
# ---------- Public API ----------
|
|
|
|
@spec start_link(keyword()) :: GenServer.on_start()
|
|
def start_link(opts) do
|
|
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
|
|
end
|
|
|
|
@spec fetch(non_neg_integer(), DateTime.t()) :: {:ok, [score()]} | :miss
|
|
def fetch(band_mhz, valid_time) do
|
|
case :ets.lookup(@table, {band_mhz, valid_time}) do
|
|
[{_, grid}] -> {:ok, grid_to_list(grid)}
|
|
[] -> :miss
|
|
end
|
|
end
|
|
|
|
@spec fetch_bounds(non_neg_integer(), DateTime.t(), bounds() | nil) ::
|
|
{:ok, [score()]} | :miss
|
|
def fetch_bounds(band_mhz, valid_time, bounds) do
|
|
case :ets.lookup(@table, {band_mhz, valid_time}) do
|
|
[{_, grid}] -> {:ok, grid_to_filtered_list(grid, bounds)}
|
|
[] -> :miss
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Look up the score for a single `{lat, lon}` point at `{band_mhz, valid_time}`.
|
|
Returns `{:ok, score}` on hit or `:miss` on cache miss. The coordinates are
|
|
used as-is (no snap-to-grid) — callers that need the nearest grid cell should
|
|
snap first.
|
|
"""
|
|
@spec fetch_point(non_neg_integer(), DateTime.t(), float(), float()) ::
|
|
{:ok, non_neg_integer()} | :miss
|
|
def fetch_point(band_mhz, valid_time, lat, lon) do
|
|
case :ets.lookup(@table, {band_mhz, valid_time}) do
|
|
[{_, grid}] ->
|
|
case Map.get(grid, {lat, lon}) do
|
|
nil -> :miss
|
|
score -> {:ok, score}
|
|
end
|
|
|
|
[] ->
|
|
:miss
|
|
end
|
|
end
|
|
|
|
@spec put(non_neg_integer(), DateTime.t(), [score()]) :: :ok
|
|
def put(band_mhz, valid_time, scores) do
|
|
grid = list_to_grid(scores)
|
|
:ets.insert(@table, {{band_mhz, valid_time}, grid})
|
|
:ok
|
|
end
|
|
|
|
@doc """
|
|
Insert locally AND broadcast to peer nodes via PubSub.
|
|
|
|
Call this from the single pod that computed the scores; every cluster member
|
|
(including the caller) will receive the PubSub message and populate its local
|
|
ETS. Use `put/3` directly if you already have a cluster-wide fan-out.
|
|
"""
|
|
@spec broadcast_put(non_neg_integer(), DateTime.t(), [score()]) :: :ok
|
|
def broadcast_put(band_mhz, valid_time, scores) do
|
|
_ = PubSub.broadcast(@pubsub, @topic, {:cache_refresh, band_mhz, valid_time, scores})
|
|
:ok
|
|
end
|
|
|
|
@doc "Returns the sorted list of valid_times cached for `band_mhz`."
|
|
@spec valid_times(non_neg_integer()) :: [DateTime.t()]
|
|
def valid_times(band_mhz) do
|
|
match_spec = [{{{band_mhz, :"$1"}, :_}, [], [:"$1"]}]
|
|
|
|
@table
|
|
|> :ets.select(match_spec)
|
|
|> Enum.sort({:asc, DateTime})
|
|
end
|
|
|
|
@spec prune_older_than(DateTime.t()) :: non_neg_integer()
|
|
def prune_older_than(cutoff) do
|
|
match_spec = [
|
|
{{{:_, :"$1"}, :_}, [{:<, :"$1", {:const, cutoff}}], [true]}
|
|
]
|
|
|
|
:ets.select_delete(@table, match_spec)
|
|
end
|
|
|
|
@doc """
|
|
Delete every entry whose `valid_time` falls outside the inclusive
|
|
window `[past_cutoff, future_cutoff]`. Used to bound the cache to
|
|
the active forecast horizon (past ~1 h + HRRR's 18 h forward) so
|
|
long-horizon GEFS valid_times never accumulate in ETS.
|
|
"""
|
|
@spec prune_outside_window(DateTime.t(), DateTime.t()) :: non_neg_integer()
|
|
def prune_outside_window(past_cutoff, future_cutoff) do
|
|
match_spec = [
|
|
{{{:_, :"$1"}, :_},
|
|
[
|
|
{:orelse, {:<, :"$1", {:const, past_cutoff}}, {:>, :"$1", {:const, future_cutoff}}}
|
|
], [true]}
|
|
]
|
|
|
|
:ets.select_delete(@table, match_spec)
|
|
end
|
|
|
|
@spec clear() :: :ok
|
|
def clear do
|
|
:ets.delete_all_objects(@table)
|
|
:ok
|
|
end
|
|
|
|
@doc "Flush the GenServer mailbox so any in-flight cache_refresh messages are applied."
|
|
@spec sync() :: :ok
|
|
def sync do
|
|
GenServer.call(__MODULE__, :sync)
|
|
end
|
|
|
|
# ---------- GenServer callbacks ----------
|
|
|
|
@impl true
|
|
def init(_opts) do
|
|
_ = :ets.new(@table, [:set, :named_table, :public, :compressed, read_concurrency: true])
|
|
:ok = PubSub.subscribe(@pubsub, @topic)
|
|
{:ok, %{}}
|
|
end
|
|
|
|
@impl true
|
|
def handle_call(:sync, _from, state), do: {:reply, :ok, state}
|
|
|
|
@impl true
|
|
def handle_info({:cache_refresh, band_mhz, valid_time, scores}, state) do
|
|
put(band_mhz, valid_time, scores)
|
|
{:noreply, state}
|
|
end
|
|
|
|
def handle_info(_msg, state), do: {:noreply, state}
|
|
|
|
# ---------- Internal ----------
|
|
|
|
defp list_to_grid(scores) do
|
|
Map.new(scores, fn %{lat: lat, lon: lon, score: score} -> {{lat, lon}, score} end)
|
|
end
|
|
|
|
defp grid_to_list(grid) do
|
|
Enum.map(grid, fn {{lat, lon}, score} -> %{lat: lat, lon: lon, score: score} end)
|
|
end
|
|
|
|
defp grid_to_filtered_list(grid, nil), do: grid_to_list(grid)
|
|
|
|
defp grid_to_filtered_list(grid, %{"south" => s, "north" => n, "west" => w, "east" => e}) do
|
|
# Single pass over 92k cells: the former Enum.filter |> Enum.map
|
|
# walked the map twice and allocated an intermediate list of
|
|
# {key, value} tuples between them. Enum.reduce emits only the
|
|
# in-bounds result maps directly.
|
|
Enum.reduce(grid, [], fn {{lat, lon}, score}, acc ->
|
|
if lat >= s and lat <= n and lon >= w and lon <= e do
|
|
[%{lat: lat, lon: lon, score: score} | acc]
|
|
else
|
|
acc
|
|
end
|
|
end)
|
|
end
|
|
end
|