Show on /weather where the detected duct geometry traps each microwave band. Overview mode bins cells by their lowest trapped frequency (Bean & Dutton cutoff); band-picker mode masks cells whose duct doesn't reach the selected band. Cutoff is pre-computed per cell in both writers (Elixir derive + Rust derive_row reads best_duct_freq_ghz from CellValues), persisted to ScalarFile, and shipped to the JS hook via the existing binary cell pack. Also cleans up six pre-existing length/1 credo warnings in unrelated test files.
532 lines
17 KiB
Elixir
532 lines
17 KiB
Elixir
defmodule Microwaveprop.Weather.ScalarFile do
|
||
@moduledoc """
|
||
On-disk store for the per-cell *derived* weather rows that feed `/weather`.
|
||
One directory per `valid_time`, with rows bucketed into 5°×5° chunk files
|
||
so viewport reads only decode the chunks that overlap the requested bounds.
|
||
|
||
This is the cheap-read sibling of `Microwaveprop.Propagation.ProfilesFile`:
|
||
|
||
* `ProfilesFile` keeps the raw HRRR profile per cell (~10 MB decoded) and
|
||
is the source of truth for advanced diagnostics, terrain analysis, and
|
||
point-detail breakdowns.
|
||
* `ScalarFile` keeps only the scalar fields the weather map renders
|
||
(temperature, dewpoint depression, refractivity, lapse rates, duct
|
||
summary, etc.), already pushed through `WeatherLayers.derive/1`.
|
||
|
||
Because the scalar shape is small and the derivation has already happened,
|
||
a tile request can read just the chunks it needs and skip the
|
||
`SoundingParams.derive/1` + `WeatherLayers.derive/1` work entirely.
|
||
|
||
## Wire format
|
||
|
||
Each chunk is a gzipped MessagePack `[row_map, ...]`. MessagePack was
|
||
picked over ETF so the Rust `prop_grid_rs` worker can produce these
|
||
files inline in the propagation pipeline using `rmp-serde` and Elixir
|
||
reads them with `Msgpax`. Both producers write identical bytes — there
|
||
is no Elixir-only or Rust-only branch.
|
||
|
||
Map keys are strings on the wire (MessagePack has no atom type).
|
||
Whitelisted keys are atomized on read so callers see the same
|
||
Elixir-shaped row regardless of producer.
|
||
|
||
## Layout
|
||
|
||
<base_dir>/weather_scalars/
|
||
<iso>/ # e.g. 2026-04-28T12:00:00Z/
|
||
<lat_band>_<lon_band>.mp.gz
|
||
|
||
`lat_band = floor(lat / 5)`, `lon_band = floor(lon / 5)`. Writes go
|
||
through `rename(2)` → atomic on NFS.
|
||
|
||
## Chunk size
|
||
|
||
5°×5° was picked so the typical /weather viewport (a US state at z=6-7)
|
||
overlaps 1-4 chunks. Smaller chunks would balloon directory entries on
|
||
NFS without measurable read benefit.
|
||
"""
|
||
|
||
require Logger
|
||
|
||
@chunk_step 5
|
||
@subdir "weather_scalars"
|
||
|
||
# Keys produced by the writer that should land as atoms in Elixir.
|
||
# Anything outside this set stays a string and is ignored by callers
|
||
# (weather.ex / tile_renderer.ex use atom access throughout).
|
||
@atom_keys MapSet.new([
|
||
"lat",
|
||
"lon",
|
||
"valid_time",
|
||
"temperature",
|
||
"dewpoint_depression",
|
||
"surface_rh",
|
||
"surface_pressure_mb",
|
||
"surface_refractivity",
|
||
"refractivity_gradient",
|
||
"bl_height",
|
||
"pwat",
|
||
"temp_850mb",
|
||
"dewpoint_850mb",
|
||
"temp_700mb",
|
||
"dewpoint_700mb",
|
||
"lapse_rate",
|
||
"mid_lapse_rate",
|
||
"inversion_strength",
|
||
"inversion_base_m",
|
||
"ducting",
|
||
"duct_base_m",
|
||
"duct_strength",
|
||
"duct_cutoff_ghz"
|
||
])
|
||
|
||
@type row :: %{required(:lat) => float(), required(:lon) => float(), optional(atom()) => term()}
|
||
@type bounds :: %{optional(String.t()) => number()}
|
||
|
||
@spec base_dir() :: String.t()
|
||
def base_dir do
|
||
Path.join(
|
||
Application.get_env(:microwaveprop, :propagation_scores_dir, "/data/scores"),
|
||
@subdir
|
||
)
|
||
end
|
||
|
||
@spec dir_for(DateTime.t()) :: String.t()
|
||
def dir_for(%DateTime{} = valid_time) do
|
||
Path.join(base_dir(), iso_key(valid_time))
|
||
end
|
||
|
||
@doc """
|
||
Sibling dir for HRDPS-derived scalar chunks. The HRRR `dir_for/1`
|
||
writer wipes its destination dir of stale chunks before each write;
|
||
routing HRDPS to a parallel directory means the two sources don't
|
||
clobber each other. Read paths read both dirs and merge.
|
||
"""
|
||
@spec dir_for_hrdps(DateTime.t()) :: String.t()
|
||
def dir_for_hrdps(%DateTime{} = valid_time) do
|
||
Path.join(base_dir(), iso_key(valid_time) <> ".hrdps")
|
||
end
|
||
|
||
@spec exists?(DateTime.t()) :: boolean()
|
||
def exists?(%DateTime{} = valid_time) do
|
||
has_chunks?(dir_for(valid_time)) or has_chunks?(dir_for_hrdps(valid_time))
|
||
end
|
||
|
||
defp has_chunks?(dir) do
|
||
case File.ls(dir) do
|
||
{:ok, [_ | _]} -> true
|
||
_ -> false
|
||
end
|
||
end
|
||
|
||
@doc """
|
||
Persist `rows` for `valid_time`. Rows are bucketed by 5°×5° chunk; each
|
||
chunk is written atomically via tmp + rename. Existing chunks for the
|
||
same `valid_time` are removed first so a smaller follow-up write doesn't
|
||
leave stale chunks behind.
|
||
"""
|
||
@spec write!(DateTime.t(), [row()]) :: :ok
|
||
def write!(%DateTime{} = valid_time, rows) when is_list(rows) do
|
||
dir = dir_for(valid_time)
|
||
File.mkdir_p!(dir)
|
||
|
||
# Drop any pre-existing chunks for this valid_time so this write is
|
||
# the canonical state — otherwise a smaller follow-up write would
|
||
# leave stale chunks behind.
|
||
case File.ls(dir) do
|
||
{:ok, entries} ->
|
||
for entry <- entries do
|
||
_ = File.rm(Path.join(dir, entry))
|
||
end
|
||
|
||
_ ->
|
||
:ok
|
||
end
|
||
|
||
rows
|
||
|> Enum.group_by(&chunk_key/1)
|
||
|> Enum.each(fn {{lat_band, lon_band}, chunk_rows} ->
|
||
path = Path.join(dir, "#{lat_band}_#{lon_band}.mp.gz")
|
||
binary = encode_chunk(chunk_rows)
|
||
|
||
tmp = path <> ".tmp." <> unique_suffix()
|
||
File.write!(tmp, binary, [:binary])
|
||
File.rename!(tmp, path)
|
||
end)
|
||
|
||
:ok
|
||
end
|
||
|
||
@doc """
|
||
Read every persisted row for `valid_time` whose lat/lon falls within
|
||
`bounds`. Pass `nil` to read every chunk. Returns `[]` if no scalar
|
||
file exists for `valid_time`.
|
||
"""
|
||
@spec read_bounds(DateTime.t(), bounds() | nil) :: [row()]
|
||
def read_bounds(%DateTime{} = valid_time, bounds) do
|
||
# Read HRRR + HRDPS chunks, then concat. HRRR wins where both
|
||
# sources cover (shouldn't happen in practice — Rust's HRDPS
|
||
# pipeline drops cells in HRRR's CONUS bbox before writing — but
|
||
# de-dup defensively in case mask-disagreement ever produces
|
||
# overlap). The user-visible intent: show HRRR-derived weather data
|
||
# everywhere CONUS, then fill the rest of North America with HRDPS.
|
||
hrrr = read_chunks(list_chunk_files(valid_time), bounds)
|
||
hrdps = read_chunks(list_chunk_files_hrdps(valid_time), bounds)
|
||
|
||
case {hrrr, hrdps} do
|
||
{[], []} -> []
|
||
{h, []} -> h
|
||
{[], c} -> c
|
||
{h, c} -> merge_prefer_hrrr(h, c)
|
||
end
|
||
end
|
||
|
||
@doc """
|
||
Read only the HRDPS sibling chunks for `valid_time` whose lat/lon falls
|
||
within `bounds`. Used by the `/weather-ca` endpoint, which intentionally
|
||
excludes HRRR-derived rows so the user can see Canadian-source data
|
||
without being eclipsed by the higher-resolution HRRR scalars over the
|
||
US/Canada border.
|
||
"""
|
||
@spec read_bounds_hrdps(DateTime.t(), bounds() | nil) :: [row()]
|
||
def read_bounds_hrdps(%DateTime{} = valid_time, bounds) do
|
||
read_chunks(list_chunk_files_hrdps(valid_time), bounds)
|
||
end
|
||
|
||
@doc """
|
||
Read HRDPS rows for the closest available valid_time within
|
||
`window_seconds` of `target`. Returns rows from the snapped time, or
|
||
`[]` if no HRDPS dir falls inside the window.
|
||
|
||
HRDPS publishes 4×/day with ~3-4 h latency, so the /weather timeline
|
||
(driven by HRRR's hourly cadence) will frequently request a time that
|
||
has no HRDPS dir on disk. Snapping to the nearest available time lets
|
||
the Canadian overlay render slightly stale rather than disappear
|
||
entirely whenever the user's selected forecast hour misses HRDPS.
|
||
"""
|
||
@spec read_bounds_hrdps_nearest(DateTime.t(), bounds() | nil, non_neg_integer()) :: [row()]
|
||
def read_bounds_hrdps_nearest(%DateTime{} = target, bounds, window_seconds)
|
||
when is_integer(window_seconds) and window_seconds >= 0 do
|
||
case nearest_hrdps_valid_time(target, window_seconds) do
|
||
nil -> []
|
||
vt -> read_chunks(list_chunk_files_hrdps(vt), bounds)
|
||
end
|
||
end
|
||
|
||
defp nearest_hrdps_valid_time(target, window_seconds) do
|
||
target_unix = DateTime.to_unix(target)
|
||
|
||
list_valid_times_hrdps()
|
||
|> Enum.map(fn vt -> {vt, abs(DateTime.to_unix(vt) - target_unix)} end)
|
||
|> Enum.filter(fn {_vt, delta} -> delta <= window_seconds end)
|
||
|> Enum.min_by(fn {_vt, delta} -> delta end, fn -> nil end)
|
||
|> case do
|
||
nil -> nil
|
||
{vt, _delta} -> vt
|
||
end
|
||
end
|
||
|
||
defp read_chunks([], _bounds), do: []
|
||
|
||
defp read_chunks(files, bounds) do
|
||
files
|
||
|> Enum.filter(fn {key, _path} -> chunk_intersects_bounds?(key, bounds) end)
|
||
|> Enum.flat_map(fn {_key, path} -> decode_chunk(path) end)
|
||
|> filter_bounds(bounds)
|
||
end
|
||
|
||
defp merge_prefer_hrrr(hrrr_rows, hrdps_rows) do
|
||
hrrr_keys =
|
||
MapSet.new(hrrr_rows, fn %{lat: lat, lon: lon} -> {lat, lon} end)
|
||
|
||
extras =
|
||
Enum.reject(hrdps_rows, fn %{lat: lat, lon: lon} -> MapSet.member?(hrrr_keys, {lat, lon}) end)
|
||
|
||
hrrr_rows ++ extras
|
||
end
|
||
|
||
@doc """
|
||
Look up a single derived row at `(valid_time, lat, lon)`. Returns
|
||
`{:ok, row}` or `:miss` when either the chunk file is absent or the
|
||
cell isn't present.
|
||
"""
|
||
@spec read_point(DateTime.t(), float(), float()) :: {:ok, row()} | :miss
|
||
def read_point(%DateTime{} = valid_time, lat, lon) when is_number(lat) and is_number(lon) do
|
||
key = {chunk_band(lat * 1.0), chunk_band(lon * 1.0)}
|
||
chunk_name = chunk_filename(key)
|
||
snapped_lat = snap(lat)
|
||
snapped_lon = snap(lon)
|
||
|
||
# Try HRRR first (where most points live), then HRDPS for Canadian
|
||
# cells outside HRRR coverage.
|
||
Enum.find_value(
|
||
[
|
||
Path.join(dir_for(valid_time), chunk_name),
|
||
Path.join(dir_for_hrdps(valid_time), chunk_name)
|
||
],
|
||
:miss,
|
||
&lookup_in_chunk(&1, snapped_lat, snapped_lon)
|
||
)
|
||
end
|
||
|
||
defp lookup_in_chunk(path, snapped_lat, snapped_lon) do
|
||
if File.exists?(path) do
|
||
hit_or_false(find_point_in_chunk(path, snapped_lat, snapped_lon))
|
||
else
|
||
false
|
||
end
|
||
end
|
||
|
||
defp hit_or_false({:ok, _} = hit), do: hit
|
||
defp hit_or_false(:miss), do: false
|
||
|
||
defp find_point_in_chunk(path, snapped_lat, snapped_lon) do
|
||
case Enum.find(decode_chunk(path), fn r ->
|
||
r.lat == snapped_lat and r.lon == snapped_lon
|
||
end) do
|
||
nil -> :miss
|
||
row -> {:ok, row}
|
||
end
|
||
end
|
||
|
||
@doc """
|
||
Every persisted valid_time, sorted ascending. Used to back the /weather
|
||
forecast timeline once scalars take over from raw profiles as the
|
||
primary read.
|
||
"""
|
||
@spec list_valid_times() :: [DateTime.t()]
|
||
def list_valid_times do
|
||
case File.ls(base_dir()) do
|
||
{:ok, entries} ->
|
||
times =
|
||
for entry <- entries,
|
||
{:ok, dt, _} <- [DateTime.from_iso8601(entry)],
|
||
do: dt
|
||
|
||
Enum.sort(times, DateTime)
|
||
|
||
_ ->
|
||
[]
|
||
end
|
||
end
|
||
|
||
@doc """
|
||
Like `list_valid_times/0` but returns only the valid_times that have a
|
||
`<iso>.hrdps` directory present. Backs the `/weather-ca` timeline so
|
||
the scrubber surfaces HRDPS-only hours that wouldn't otherwise appear
|
||
on the standard `/weather` timeline.
|
||
"""
|
||
@spec list_valid_times_hrdps() :: [DateTime.t()]
|
||
def list_valid_times_hrdps do
|
||
case File.ls(base_dir()) do
|
||
{:ok, entries} ->
|
||
times =
|
||
for entry <- entries,
|
||
String.ends_with?(entry, ".hrdps"),
|
||
iso = String.trim_trailing(entry, ".hrdps"),
|
||
{:ok, dt, _} <- [DateTime.from_iso8601(iso)],
|
||
do: dt
|
||
|
||
Enum.sort(times, DateTime)
|
||
|
||
_ ->
|
||
[]
|
||
end
|
||
end
|
||
|
||
@doc "Delete scalar dirs whose valid_time is strictly before `cutoff`. Returns count removed."
|
||
@spec prune_older_than(DateTime.t()) :: non_neg_integer()
|
||
def prune_older_than(%DateTime{} = cutoff) do
|
||
cutoff_unix = DateTime.to_unix(cutoff)
|
||
|
||
base_dir()
|
||
|> list_valid_time_dirs()
|
||
|> Enum.reduce(0, fn {path, dt}, acc ->
|
||
if DateTime.to_unix(dt) < cutoff_unix do
|
||
_ = File.rm_rf(path)
|
||
acc + 1
|
||
else
|
||
acc
|
||
end
|
||
end)
|
||
end
|
||
|
||
@doc """
|
||
Keep only scalar dirs whose valid_time is inside the closed window
|
||
`[run_time, run_time + max_forecast_hour * 3600]`. Mirrors
|
||
`ProfilesFile.retain_window/2`.
|
||
"""
|
||
@spec retain_window(DateTime.t(), non_neg_integer()) :: non_neg_integer()
|
||
def retain_window(%DateTime{} = run_time, max_forecast_hour) when max_forecast_hour >= 0 do
|
||
lo = DateTime.to_unix(run_time)
|
||
hi = lo + max_forecast_hour * 3600
|
||
|
||
base_dir()
|
||
|> list_valid_time_dirs()
|
||
|> Enum.reduce(0, fn {path, dt}, acc ->
|
||
unix = DateTime.to_unix(dt)
|
||
|
||
if unix < lo or unix > hi do
|
||
_ = File.rm_rf(path)
|
||
acc + 1
|
||
else
|
||
acc
|
||
end
|
||
end)
|
||
end
|
||
|
||
# ---------- Internal ----------
|
||
|
||
defp iso_key(%DateTime{} = valid_time) do
|
||
valid_time |> DateTime.truncate(:second) |> DateTime.to_iso8601()
|
||
end
|
||
|
||
defp chunk_key(%{lat: lat, lon: lon}) do
|
||
{chunk_band(lat * 1.0), chunk_band(lon * 1.0)}
|
||
end
|
||
|
||
defp chunk_band(value) when is_float(value) do
|
||
(value / @chunk_step) |> Float.floor() |> trunc()
|
||
end
|
||
|
||
defp chunk_filename({lat_band, lon_band}), do: "#{lat_band}_#{lon_band}.mp.gz"
|
||
|
||
defp chunk_intersects_bounds?(_key, nil), do: true
|
||
|
||
defp chunk_intersects_bounds?({lat_band, lon_band}, %{"south" => s, "north" => n, "west" => w, "east" => e}) do
|
||
chunk_south = lat_band * @chunk_step
|
||
chunk_north = chunk_south + @chunk_step
|
||
chunk_west = lon_band * @chunk_step
|
||
chunk_east = chunk_west + @chunk_step
|
||
|
||
chunk_north >= s and chunk_south <= n and chunk_east >= w and chunk_west <= e
|
||
end
|
||
|
||
defp filter_bounds(rows, nil), do: rows
|
||
|
||
defp filter_bounds(rows, %{"south" => s, "north" => n, "west" => w, "east" => e}) do
|
||
Enum.filter(rows, fn %{lat: lat, lon: lon} ->
|
||
lat >= s and lat <= n and lon >= w and lon <= e
|
||
end)
|
||
end
|
||
|
||
# Encode rows as gzipped MessagePack. We strip atoms to strings on the
|
||
# wire (Msgpax doesn't speak atom keys natively, and Rust can't read
|
||
# them anyway) and re-atomize on the read side.
|
||
defp encode_chunk(rows) do
|
||
rows
|
||
|> Enum.map(&prepare_row_for_encode/1)
|
||
|> Msgpax.pack!(iodata: false)
|
||
|> :zlib.gzip()
|
||
end
|
||
|
||
# Convert atom keys to strings and DateTime values to ISO8601 strings
|
||
# so the row is representable in MessagePack.
|
||
defp prepare_row_for_encode(row) when is_map(row) do
|
||
Map.new(row, fn {k, v} ->
|
||
key = if is_atom(k), do: Atom.to_string(k), else: k
|
||
{key, encode_value(v)}
|
||
end)
|
||
end
|
||
|
||
defp encode_value(%DateTime{} = dt), do: DateTime.to_iso8601(DateTime.truncate(dt, :second))
|
||
defp encode_value(value), do: value
|
||
|
||
defp decode_chunk(path) do
|
||
with {:ok, gz} <- File.read(path),
|
||
{:ok, binary} <- safe_gunzip(gz),
|
||
{:ok, decoded} <- Msgpax.unpack(binary) do
|
||
decoded
|
||
|> List.wrap()
|
||
|> Enum.map(&normalize_row/1)
|
||
else
|
||
error ->
|
||
Logger.warning("ScalarFile chunk read failed path=#{path} reason=#{inspect(error)}")
|
||
[]
|
||
end
|
||
end
|
||
|
||
defp safe_gunzip(binary) do
|
||
{:ok, :zlib.gunzip(binary)}
|
||
rescue
|
||
e -> {:error, {:gunzip, e}}
|
||
end
|
||
|
||
# Atomize whitelisted keys and re-parse the `valid_time` ISO8601 back
|
||
# into a DateTime so callers see the same shape as the in-memory rows
|
||
# produced by `Microwaveprop.Weather.build_grid_cache_rows/3`.
|
||
defp normalize_row(row) when is_map(row) do
|
||
Map.new(row, fn {k, v} ->
|
||
key = atomize_key(k)
|
||
{key, normalize_value(key, v)}
|
||
end)
|
||
end
|
||
|
||
defp atomize_key(k) when is_binary(k) do
|
||
if MapSet.member?(@atom_keys, k), do: String.to_atom(k), else: k
|
||
end
|
||
|
||
defp atomize_key(k), do: k
|
||
|
||
defp normalize_value(:valid_time, v) when is_binary(v) do
|
||
case DateTime.from_iso8601(v) do
|
||
{:ok, dt, _offset} -> dt
|
||
_ -> v
|
||
end
|
||
end
|
||
|
||
defp normalize_value(_key, v), do: v
|
||
|
||
defp list_chunk_files(%DateTime{} = valid_time) do
|
||
list_chunk_files_in(dir_for(valid_time))
|
||
end
|
||
|
||
defp list_chunk_files_hrdps(%DateTime{} = valid_time) do
|
||
list_chunk_files_in(dir_for_hrdps(valid_time))
|
||
end
|
||
|
||
defp list_chunk_files_in(dir) do
|
||
case File.ls(dir) do
|
||
{:ok, entries} ->
|
||
for entry <- entries,
|
||
key = parse_chunk_filename(entry),
|
||
key != nil,
|
||
do: {key, Path.join(dir, entry)}
|
||
|
||
_ ->
|
||
[]
|
||
end
|
||
end
|
||
|
||
defp parse_chunk_filename(filename) do
|
||
with [_, lat_str, lon_str] <- Regex.run(~r/^(-?\d+)_(-?\d+)\.mp\.gz$/, filename),
|
||
{lat_band, ""} <- Integer.parse(lat_str),
|
||
{lon_band, ""} <- Integer.parse(lon_str) do
|
||
{lat_band, lon_band}
|
||
else
|
||
_ -> nil
|
||
end
|
||
end
|
||
|
||
defp list_valid_time_dirs(base) do
|
||
case File.ls(base) do
|
||
{:ok, entries} ->
|
||
for entry <- entries,
|
||
{:ok, dt, _} <- [DateTime.from_iso8601(entry)],
|
||
do: {Path.join(base, entry), dt}
|
||
|
||
_ ->
|
||
[]
|
||
end
|
||
end
|
||
|
||
# Match ProfilesFile's snap step (0.125°) so a click at any lat/lon
|
||
# rounds to the same key the writer used.
|
||
defp snap(value) do
|
||
step = 0.125
|
||
Float.round(Float.round(value / step) * step, 3)
|
||
end
|
||
|
||
defp unique_suffix do
|
||
"#{System.system_time(:nanosecond)}.#{:erlang.unique_integer([:positive])}"
|
||
end
|
||
end
|