defmodule Microwaveprop.Weather do @moduledoc false import Ecto.Query require Logger alias Microwaveprop.Repo alias Microwaveprop.Weather.IemClient alias Microwaveprop.Weather.HrrrProfile alias Microwaveprop.Weather.IemreObservation alias Microwaveprop.Weather.SolarIndex alias Microwaveprop.Weather.Sounding alias Microwaveprop.Weather.Station alias Microwaveprop.Weather.SurfaceObservation # Approximate km per degree latitude @km_per_deg_lat 111.0 def find_or_create_station(attrs) do code = attrs[:station_code] || attrs["station_code"] type = attrs[:station_type] || attrs["station_type"] if code && type do case Repo.get_by(Station, station_code: code, station_type: type) do nil -> %Station{} |> Station.changeset(attrs) |> Repo.insert() station -> {:ok, station} end else %Station{} |> Station.changeset(attrs) |> Repo.insert() end end def upsert_surface_observation(%Station{} = station, attrs) do attrs = Map.put(attrs, :station_id, station.id) %SurfaceObservation{} |> SurfaceObservation.changeset(attrs) |> Repo.insert( on_conflict: {:replace_all_except, [:id, :station_id, :observed_at, :inserted_at]}, conflict_target: [:station_id, :observed_at], returning: true ) end def upsert_sounding(%Station{} = station, attrs) do attrs = Map.put(attrs, :station_id, station.id) %Sounding{} |> Sounding.changeset(attrs) |> Repo.insert( on_conflict: {:replace_all_except, [:id, :station_id, :observed_at, :inserted_at]}, conflict_target: [:station_id, :observed_at], returning: true ) end def upsert_solar_index(attrs) do %SolarIndex{} |> SolarIndex.changeset(attrs) |> Repo.insert( on_conflict: {:replace_all_except, [:id, :date, :inserted_at]}, conflict_target: [:date], returning: true ) end def has_surface_observations?(station_id, start_dt, end_dt) do SurfaceObservation |> where([o], o.station_id == ^station_id) |> where([o], o.observed_at >= ^start_dt and o.observed_at <= ^end_dt) |> Repo.exists?() end def has_sounding?(station_id, observed_at) do Sounding |> where([s], s.station_id == ^station_id and s.observed_at == ^observed_at) |> Repo.exists?() end def get_solar_index(date) do Repo.get_by(SolarIndex, date: date) end def existing_solar_dates do SolarIndex |> select([s], s.date) |> Repo.all() |> MapSet.new() end def sync_stations! do asos = for s <- ~w(AK AL AR AZ CA CO CT DE FL GA HI IA ID IL IN KS KY LA MA MD ME MI MN MO MS MT NC ND NE NH NJ NM NV NY OH OK OR PA RI SC SD TN TX UT VA VT WA WI WV WY), do: "#{s}_ASOS" for network <- asos ++ ["RAOB"] do type = if String.contains?(network, "ASOS"), do: "asos", else: "sounding" case IemClient.fetch_network(network) do {:ok, stations} -> for s <- stations do %Station{} |> Station.changeset(Map.put(s, :station_type, type)) |> Repo.insert(on_conflict: :nothing, conflict_target: [:station_code, :station_type]) end Logger.info("Synced #{length(stations)} stations from #{network}") {:error, e} -> Logger.warning("Failed to sync #{network}: #{inspect(e)}") end Process.sleep(200) end count = Repo.aggregate(Station, :count) Logger.info("Weather stations sync complete: #{count} total") :ok end def nearby_stations(lat, lon, station_type, radius_km) do dlat = radius_km / @km_per_deg_lat dlon = radius_km / (@km_per_deg_lat * :math.cos(lat * :math.pi() / 180)) Station |> where([s], s.station_type == ^station_type) |> where( [s], s.lat >= ^(lat - dlat) and s.lat <= ^(lat + dlat) and s.lon >= ^(lon - dlon) and s.lon <= ^(lon + dlon) ) |> Repo.all() end def sounding_times_around(dt) do date = DateTime.to_date(dt) times = if dt.hour < 12 do [ DateTime.new!(Date.add(date, -1), ~T[12:00:00], "Etc/UTC"), DateTime.new!(date, ~T[00:00:00], "Etc/UTC") ] else [ DateTime.new!(date, ~T[00:00:00], "Etc/UTC"), DateTime.new!(date, ~T[12:00:00], "Etc/UTC") ] end Enum.uniq(times) end def weather_for_contact(contact_params, opts \\ []) do lat = contact_params[:lat] || contact_params.lat lon = contact_params[:lon] || contact_params.lon timestamp = contact_params[:timestamp] || contact_params.timestamp radius_km = Keyword.get(opts, :radius_km, 150) time_window_hours = Keyword.get(opts, :time_window_hours, 6) # Bounding box in degrees dlat = radius_km / @km_per_deg_lat dlon = radius_km / (@km_per_deg_lat * :math.cos(lat * :math.pi() / 180)) time_start = DateTime.add(timestamp, -time_window_hours * 3600, :second) time_end = DateTime.add(timestamp, time_window_hours * 3600, :second) station_ids = Station |> where( [s], s.lat >= ^(lat - dlat) and s.lat <= ^(lat + dlat) and s.lon >= ^(lon - dlon) and s.lon <= ^(lon + dlon) ) |> select([s], s.id) surface_observations = SurfaceObservation |> where([o], o.station_id in subquery(station_ids)) |> where([o], o.observed_at >= ^time_start and o.observed_at <= ^time_end) |> preload(:station) |> Repo.all() soundings = Sounding |> where([s], s.station_id in subquery(station_ids)) |> where([s], s.observed_at >= ^time_start and s.observed_at <= ^time_end) |> preload(:station) |> Repo.all() %{surface_observations: surface_observations, soundings: soundings} end def upsert_hrrr_profile(attrs) do %HrrrProfile{} |> HrrrProfile.changeset(attrs) |> Repo.insert( on_conflict: {:replace_all_except, [:id, :inserted_at]}, conflict_target: [:lat, :lon, :valid_time], returning: true ) end def upsert_hrrr_profiles_batch(profiles) do now = DateTime.truncate(DateTime.utc_now(), :second) entries = Enum.map(profiles, fn attrs -> Map.merge(attrs, %{ id: Ecto.UUID.generate(), inserted_at: now, updated_at: now }) end) Repo.insert_all(HrrrProfile, entries, on_conflict: {:replace_all_except, [:id, :inserted_at]}, conflict_target: [:lat, :lon, :valid_time] ) end def has_hrrr_profile?(lat, lon, valid_time) do {rlat, rlon} = round_to_hrrr_grid(lat, lon) HrrrProfile |> where([h], h.lat == ^rlat and h.lon == ^rlon and h.valid_time == ^valid_time) |> Repo.exists?() end def hrrr_for_contact(%{pos1: nil}), do: nil def hrrr_for_contact(contact) do lat = contact.pos1["lat"] lon = contact.pos1["lon"] || contact.pos1["lng"] if lat && lon do find_nearest_hrrr(lat, lon, contact.qso_timestamp) end end def find_nearest_hrrr(lat, lon, timestamp) do dlat = 0.05 dlon = 0.05 time_start = DateTime.add(timestamp, -3600, :second) time_end = DateTime.add(timestamp, 3600, :second) HrrrProfile |> where( [h], h.lat >= ^(lat - dlat) and h.lat <= ^(lat + dlat) and h.lon >= ^(lon - dlon) and h.lon <= ^(lon + dlon) and h.valid_time >= ^time_start and h.valid_time <= ^time_end ) |> order_by([h], asc: fragment( "ABS(? - ?) + ABS(? - ?) + ABS(EXTRACT(EPOCH FROM ? - ?))", h.lat, ^lat, h.lon, ^lon, h.valid_time, ^timestamp ) ) |> limit(1) |> Repo.one() end def hrrr_profiles_for_path(%{pos1: nil}), do: [] def hrrr_profiles_for_path(contact) do contact |> Microwaveprop.Radio.contact_path_points() |> Enum.map(fn {lat, lon} -> find_nearest_hrrr(lat, lon, contact.qso_timestamp) end) |> Enum.reject(&is_nil/1) end def round_to_hrrr_grid(lat, lon) do {Float.round(lat / 1.0, 2), Float.round(lon / 1.0, 2)} end @doc """ Delete HRRR grid profiles older than 24 hours that sit on the 0.125° grid. Preserves QSO-linked profiles which are at arbitrary lat/lon positions. With 19 forecast hours per run, this keeps ~2 runs worth of grid profiles. """ def prune_old_grid_profiles do cutoff = DateTime.add(DateTime.utc_now(), -24, :hour) # Bound the scan to recent partitions only (grid profiles are at most ~48h old). # Without a lower bound, Postgres scans every partition back to 2016. floor = DateTime.add(DateTime.utc_now(), -72, :hour) # Only delete profiles on the 0.125° propagation grid, preserving QSO-linked # profiles which sit at arbitrary lat/lon positions (HRRR's native ~3km grid). %{num_rows: deleted} = Repo.query!( """ DELETE FROM hrrr_profiles WHERE valid_time >= $2 AND valid_time < $1 AND (lat * 1000)::int % 125 = 0 AND (lon * 1000)::int % 125 = 0 """, [cutoff, floor], timeout: 300_000 ) if deleted > 0 do require Logger Logger.info("Pruned #{deleted} old HRRR grid profiles (valid_time < #{cutoff})") end deleted end def round_to_iemre_grid(lat, lon) do {Float.round(lat * 8) / 8, Float.round(lon * 8) / 8} end def upsert_iemre_observation(attrs) do %IemreObservation{} |> IemreObservation.changeset(attrs) |> Repo.insert( on_conflict: {:replace_all_except, [:id, :inserted_at]}, conflict_target: [:lat, :lon, :date], returning: true ) end def has_iemre_observation?(lat, lon, date) do IemreObservation |> where([i], i.lat == ^lat and i.lon == ^lon and i.date == ^date) |> Repo.exists?() end def iemre_for_contact(%{pos1: nil}), do: nil def iemre_for_contact(contact) do lat = contact.pos1["lat"] lon = contact.pos1["lon"] || contact.pos1["lng"] if lat && lon do find_nearest_iemre(lat, lon, contact.qso_timestamp) end end def find_nearest_iemre(lat, lon, timestamp) do {rlat, rlon} = round_to_iemre_grid(lat, lon) date = DateTime.to_date(timestamp) IemreObservation |> where([i], i.lat == ^rlat and i.lon == ^rlon and i.date == ^date) |> Repo.one() end def iemre_for_path(%{pos1: nil}), do: [] def iemre_for_path(contact) do contact |> Microwaveprop.Radio.contact_path_points() |> Enum.map(fn {lat, lon} -> find_nearest_iemre(lat, lon, contact.qso_timestamp) end) |> Enum.reject(&is_nil/1) end end