diff --git a/lib/microwaveprop/propagation.ex b/lib/microwaveprop/propagation.ex index ecd446d0..6c8ba3da 100644 --- a/lib/microwaveprop/propagation.ex +++ b/lib/microwaveprop/propagation.ex @@ -246,9 +246,7 @@ defmodule Microwaveprop.Propagation do @spec available_valid_times(non_neg_integer()) :: [DateTime.t()] def available_valid_times(band_mhz) do - now = DateTime.utc_now() - past_cutoff = DateTime.add(now, -3600, :second) - future_cutoff = DateTime.add(now, @hrrr_forecast_horizon_hours * 3600, :second) + {past_cutoff, future_cutoff} = hot_cache_window() case ScoresFile.list_valid_times(band_mhz) do [] -> [] @@ -256,6 +254,21 @@ defmodule Microwaveprop.Propagation do end end + @doc """ + The active forecast window for the `/map` UI: one hour in the past + through HRRR's 18-hour forecast horizon. Used by `ScoreCacheReconciler` + and `NotifyListener` to bound ETS growth — long-horizon GEFS `.prop` + files on disk must not balloon `propagation_score_cache` past the + memory the UI actually reads. + """ + @spec hot_cache_window() :: {DateTime.t(), DateTime.t()} + def hot_cache_window do + now = DateTime.utc_now() + past = DateTime.add(now, -3600, :second) + future = DateTime.add(now, @hrrr_forecast_horizon_hours * 3600, :second) + {past, future} + end + defp filter_or_latest(times, past_cutoff, future_cutoff) do fresh = Enum.filter(times, fn t -> diff --git a/lib/microwaveprop/propagation/notify_listener.ex b/lib/microwaveprop/propagation/notify_listener.ex index dfcf1608..1aab8975 100644 --- a/lib/microwaveprop/propagation/notify_listener.ex +++ b/lib/microwaveprop/propagation/notify_listener.ex @@ -106,9 +106,13 @@ defmodule Microwaveprop.Propagation.NotifyListener do %{valid_time: valid_time} ) - # Prune anything older than 2 hours so the cache doesn't grow - # unbounded when Rust covers many forecast hours back-to-back. - ScoreCache.prune_older_than(DateTime.add(DateTime.utc_now(), -2, :hour)) + # Prune anything outside the active map window. Old entries (>1h + # past) and long-horizon GEFS entries (>18h future) would otherwise + # accumulate and push the `propagation_score_cache` ETS table past + # 2 GiB — scheduler starvation from the GC sweeps tripped the + # /live probe and cycled the pods every ~25 min. + {past, future} = Propagation.hot_cache_window() + ScoreCache.prune_outside_window(past, future) _ = Phoenix.PubSub.broadcast(@pubsub, @topic, {:propagation_updated, [valid_time]}) diff --git a/lib/microwaveprop/propagation/score_cache.ex b/lib/microwaveprop/propagation/score_cache.ex index 58ac8f9f..b0294223 100644 --- a/lib/microwaveprop/propagation/score_cache.ex +++ b/lib/microwaveprop/propagation/score_cache.ex @@ -106,6 +106,24 @@ defmodule Microwaveprop.Propagation.ScoreCache do :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) diff --git a/lib/microwaveprop/propagation/score_cache_reconciler.ex b/lib/microwaveprop/propagation/score_cache_reconciler.ex index 04016291..b54fe790 100644 --- a/lib/microwaveprop/propagation/score_cache_reconciler.ex +++ b/lib/microwaveprop/propagation/score_cache_reconciler.ex @@ -24,6 +24,7 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconciler do """ use GenServer + alias Microwaveprop.Propagation alias Microwaveprop.Propagation.BandConfig alias Microwaveprop.Propagation.ScoreCache alias Microwaveprop.Propagation.ScoresFile @@ -99,7 +100,17 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconciler do def handle_info(_msg, state), do: {:noreply, state} defp warm_missing_for_band(band_mhz) do - on_disk = MapSet.new(ScoresFile.list_valid_times(band_mhz)) + {past_cutoff, future_cutoff} = Propagation.hot_cache_window() + + on_disk = + band_mhz + |> ScoresFile.list_valid_times() + |> Enum.filter(fn vt -> + DateTime.compare(vt, past_cutoff) != :lt and + DateTime.compare(vt, future_cutoff) != :gt + end) + |> MapSet.new() + in_cache = MapSet.new(ScoreCache.valid_times(band_mhz)) missing = MapSet.difference(on_disk, in_cache) diff --git a/test/microwaveprop/propagation/score_cache_reconciler_test.exs b/test/microwaveprop/propagation/score_cache_reconciler_test.exs index 8f31233e..f575061c 100644 --- a/test/microwaveprop/propagation/score_cache_reconciler_test.exs +++ b/test/microwaveprop/propagation/score_cache_reconciler_test.exs @@ -47,6 +47,18 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconcilerTest do ] end + # In-window timestamps for reconciler tests — the reconciler now + # filters valid_times to `[now-1h, now+18h]`, so hardcoded fixtures + # in the past drift out of scope as time passes. `current_hour/1` + # returns an hour-aligned DateTime offset by N hours from now. + defp current_hour(offset_hours \\ 0) do + DateTime.utc_now() + |> Map.put(:minute, 0) + |> Map.put(:second, 0) + |> Map.put(:microsecond, {0, 0}) + |> DateTime.add(offset_hours * 3600, :second) + end + describe "next_sweep_interval/0" do test "returns a value in [base + 1, base + jitter_max] with defaults" do # Exercise the randomized default many times; each draw must land @@ -75,7 +87,7 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconcilerTest do describe "sweep_once/0" do test "warms ScoreCache for every {band, valid_time} on disk and missing from cache", %{band_mhz: band_mhz} do - valid_time = ~U[2026-04-21 12:00:00Z] + valid_time = current_hour(1) ScoresFile.write!(band_mhz, valid_time, sample_scores()) assert :miss = ScoreCache.fetch(band_mhz, valid_time) @@ -88,7 +100,7 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconcilerTest do test "is a no-op when cache is already in sync with disk", %{band_mhz: band_mhz} do - valid_time = ~U[2026-04-21 13:00:00Z] + valid_time = current_hour(2) ScoresFile.write!(band_mhz, valid_time, sample_scores()) # Prime the cache manually @@ -100,6 +112,32 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconcilerTest do test "returns 0 when the scores directory is empty" do assert 0 = ScoreCacheReconciler.sweep_once() end + + test "skips valid_times outside the hot-cache window (now-1h..now+18h)", + %{band_mhz: band_mhz} do + # Three on-disk files: one way in the past, one in the window, + # one far in the future (simulating GEFS Day 2-7). Only the + # middle one should make it into ETS. + now = + DateTime.utc_now() + |> Map.put(:minute, 0) + |> Map.put(:second, 0) + |> Map.put(:microsecond, {0, 0}) + + stale_past = DateTime.add(now, -6 * 3600, :second) + in_window = now + far_future = DateTime.add(now, 72 * 3600, :second) + + ScoresFile.write!(band_mhz, stale_past, sample_scores()) + ScoresFile.write!(band_mhz, in_window, sample_scores()) + ScoresFile.write!(band_mhz, far_future, sample_scores()) + + _ = ScoreCacheReconciler.sweep_once() + + assert ScoreCache.fetch(band_mhz, stale_past) == :miss + assert {:ok, _} = ScoreCache.fetch(band_mhz, in_window) + assert ScoreCache.fetch(band_mhz, far_future) == :miss + end end describe "GenServer lifecycle" do @@ -138,7 +176,7 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconcilerTest do test "run_on_start: true warms cache from disk via handle_info(:sweep)", %{band_mhz: band_mhz} do - valid_time = ~U[2026-04-21 14:00:00Z] + valid_time = current_hour(3) ScoresFile.write!(band_mhz, valid_time, sample_scores()) assert :miss = ScoreCache.fetch(band_mhz, valid_time) @@ -153,7 +191,7 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconcilerTest do test "run_on_start: false does not sweep within the wait window", %{band_mhz: band_mhz} do - valid_time = ~U[2026-04-21 15:00:00Z] + valid_time = current_hour(4) ScoresFile.write!(band_mhz, valid_time, sample_scores()) start_supervised!({ScoreCacheReconciler, run_on_start: false, interval_ms: 100}) @@ -166,8 +204,8 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconcilerTest do test "interval_ms rescheduling picks up files written after the first sweep", %{band_mhz: band_mhz} do - first = ~U[2026-04-21 16:00:00Z] - second = ~U[2026-04-21 17:00:00Z] + first = current_hour(5) + second = current_hour(6) ScoresFile.write!(band_mhz, first, sample_scores()) @@ -186,7 +224,7 @@ defmodule Microwaveprop.Propagation.ScoreCacheReconcilerTest do end test "logs an info line when warmed > 0", %{band_mhz: band_mhz} do - valid_time = ~U[2026-04-21 18:00:00Z] + valid_time = current_hour(7) ScoresFile.write!(band_mhz, valid_time, sample_scores()) # config/test.exs pins the global logger to :warning, which drops the diff --git a/test/microwaveprop/propagation/score_cache_test.exs b/test/microwaveprop/propagation/score_cache_test.exs index 93c527af..ed635498 100644 --- a/test/microwaveprop/propagation/score_cache_test.exs +++ b/test/microwaveprop/propagation/score_cache_test.exs @@ -96,6 +96,34 @@ defmodule Microwaveprop.Propagation.ScoreCacheTest do end end + describe "prune_outside_window/2" do + test "removes entries strictly outside [past_cutoff, future_cutoff]" do + ScoreCache.put(10_000, ~U[2026-04-12 08:00:00Z], [%{lat: 1.0, lon: 1.0, score: 1}]) + ScoreCache.put(10_000, ~U[2026-04-12 12:00:00Z], [%{lat: 2.0, lon: 2.0, score: 2}]) + ScoreCache.put(10_000, ~U[2026-04-12 20:00:00Z], [%{lat: 3.0, lon: 3.0, score: 3}]) + + removed = + ScoreCache.prune_outside_window(~U[2026-04-12 10:00:00Z], ~U[2026-04-12 18:00:00Z]) + + assert removed == 2 + assert ScoreCache.fetch(10_000, ~U[2026-04-12 08:00:00Z]) == :miss + assert {:ok, [%{score: 2}]} = ScoreCache.fetch(10_000, ~U[2026-04-12 12:00:00Z]) + assert ScoreCache.fetch(10_000, ~U[2026-04-12 20:00:00Z]) == :miss + end + + test "keeps entries equal to either cutoff (inclusive bounds)" do + ScoreCache.put(10_000, ~U[2026-04-12 10:00:00Z], [%{lat: 1.0, lon: 1.0, score: 1}]) + ScoreCache.put(10_000, ~U[2026-04-12 18:00:00Z], [%{lat: 2.0, lon: 2.0, score: 2}]) + + removed = + ScoreCache.prune_outside_window(~U[2026-04-12 10:00:00Z], ~U[2026-04-12 18:00:00Z]) + + assert removed == 0 + assert {:ok, _} = ScoreCache.fetch(10_000, ~U[2026-04-12 10:00:00Z]) + assert {:ok, _} = ScoreCache.fetch(10_000, ~U[2026-04-12 18:00:00Z]) + end + end + describe "fetch_point/4" do test "returns :miss when nothing cached" do assert ScoreCache.fetch_point(10_000, ~U[2026-04-12 12:00:00Z], 32.0, -97.0) == :miss