From 322883563606492c42cfaba27d1df1673f95657f Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Mon, 30 Mar 2026 13:18:11 -0500 Subject: [PATCH] Add IEM precipitation data to QSO weather pipeline - Add precip_1h_in and wx_codes fields to ASOS surface observations - Add IEMRE reanalysis schema for radar-derived hourly precipitation - Add IemreFetchWorker with exponential backoff and idempotency - Integrate IEMRE enqueue into cron weather backfill pipeline - All existing QSOs marked iemre_queued=false for automatic backfill --- config/config.exs | 2 +- lib/microwaveprop/radio.ex | 14 ++ lib/microwaveprop/radio/qso.ex | 1 + lib/microwaveprop/weather.ex | 37 ++++++ lib/microwaveprop/weather/iem_client.ex | 29 ++++- .../weather/iemre_observation.ex | 28 ++++ .../weather/surface_observation.ex | 4 +- .../workers/iemre_fetch_worker.ex | 64 +++++++++ .../workers/qso_weather_enqueue_worker.ex | 49 +++++++ ..._precipitation_to_surface_observations.exs | 10 ++ ...emre_observations_and_add_iemre_queued.exs | 25 ++++ .../microwaveprop/weather/iem_client_test.exs | 95 ++++++++++++-- .../weather/iemre_observation_test.exs | 61 +++++++++ .../weather/surface_observation_test.exs | 6 +- test/microwaveprop/weather_test.exs | 86 +++++++++++++ .../workers/iemre_fetch_worker_test.exs | 121 ++++++++++++++++++ .../qso_weather_enqueue_worker_test.exs | 118 ++++++++++++++++- .../workers/weather_fetch_worker_test.exs | 6 +- 18 files changed, 734 insertions(+), 22 deletions(-) create mode 100644 lib/microwaveprop/weather/iemre_observation.ex create mode 100644 lib/microwaveprop/workers/iemre_fetch_worker.ex create mode 100644 priv/repo/migrations/20260330180421_add_precipitation_to_surface_observations.exs create mode 100644 priv/repo/migrations/20260330180721_create_iemre_observations_and_add_iemre_queued.exs create mode 100644 test/microwaveprop/weather/iemre_observation_test.exs create mode 100644 test/microwaveprop/workers/iemre_fetch_worker_test.exs diff --git a/config/config.exs b/config/config.exs index 62ca5075..7863d6fb 100644 --- a/config/config.exs +++ b/config/config.exs @@ -44,7 +44,7 @@ config :microwaveprop, MicrowavepropWeb.Endpoint, config :microwaveprop, Oban, repo: Microwaveprop.Repo, - queues: [solar: 1, weather: 3, enqueue: 1, hrrr: 20, terrain: 4, commercial: 2], + queues: [solar: 1, weather: 3, enqueue: 1, hrrr: 20, terrain: 4, commercial: 2, iemre: 5], plugins: [ {Oban.Plugins.Pruner, max_age: 3600 * 24}, {Oban.Plugins.Lifeline, rescue_after: to_timeout(minute: 30)}, diff --git a/lib/microwaveprop/radio.ex b/lib/microwaveprop/radio.ex index 608986ec..e2bdc392 100644 --- a/lib/microwaveprop/radio.ex +++ b/lib/microwaveprop/radio.ex @@ -86,6 +86,20 @@ defmodule Microwaveprop.Radio do |> Repo.update_all(set: [terrain_queued: true]) end + def unprocessed_iemre_qsos(limit \\ 500) do + Qso + |> where([q], q.iemre_queued == false and not is_nil(q.pos1)) + |> order_by([q], asc: q.qso_timestamp) + |> limit(^limit) + |> Repo.all() + end + + def mark_iemre_queued!(qso_ids) do + Qso + |> where([q], q.id in ^qso_ids) + |> Repo.update_all(set: [iemre_queued: true]) + end + @earth_radius_km 6371.0 def haversine_km(lat1, lon1, lat2, lon2) do diff --git a/lib/microwaveprop/radio/qso.ex b/lib/microwaveprop/radio/qso.ex index 4d772b43..ad9b8d13 100644 --- a/lib/microwaveprop/radio/qso.ex +++ b/lib/microwaveprop/radio/qso.ex @@ -23,6 +23,7 @@ defmodule Microwaveprop.Radio.Qso do field :weather_queued, :boolean, default: false field :hrrr_queued, :boolean, default: false field :terrain_queued, :boolean, default: false + field :iemre_queued, :boolean, default: false field :user_submitted, :boolean, default: false field :submitter_email, :string diff --git a/lib/microwaveprop/weather.ex b/lib/microwaveprop/weather.ex index 31b81b06..3c698e89 100644 --- a/lib/microwaveprop/weather.ex +++ b/lib/microwaveprop/weather.ex @@ -5,6 +5,7 @@ defmodule Microwaveprop.Weather do alias Microwaveprop.Repo alias Microwaveprop.Weather.HrrrProfile + alias Microwaveprop.Weather.IemreObservation alias Microwaveprop.Weather.SolarIndex alias Microwaveprop.Weather.Sounding alias Microwaveprop.Weather.Station @@ -224,4 +225,40 @@ defmodule Microwaveprop.Weather do def round_to_hrrr_grid(lat, lon) do {Float.round(lat / 1.0, 2), Float.round(lon / 1.0, 2)} 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_qso(%{pos1: nil}), do: nil + + def iemre_for_qso(qso) do + lat = qso.pos1["lat"] + lon = qso.pos1["lon"] || qso.pos1["lng"] + + if lat && lon do + {rlat, rlon} = round_to_iemre_grid(lat, lon) + date = DateTime.to_date(qso.qso_timestamp) + + IemreObservation + |> where([i], i.lat == ^rlat and i.lon == ^rlon and i.date == ^date) + |> Repo.one() + end + end end diff --git a/lib/microwaveprop/weather/iem_client.ex b/lib/microwaveprop/weather/iem_client.ex index 06c86e99..3a508d83 100644 --- a/lib/microwaveprop/weather/iem_client.ex +++ b/lib/microwaveprop/weather/iem_client.ex @@ -13,7 +13,7 @@ defmodule Microwaveprop.Weather.IemClient do "#{@iem_base}/cgi-bin/request/asos.py" <> "?station=#{station_id}" <> "&data=tmpf&data=dwpf&data=relh&data=sknt&data=drct" <> - "&data=mslp&data=alti&data=skyc1" <> + "&data=mslp&data=alti&data=skyc1&data=p01i&data=wxcodes" <> "&year1=#{start_dt.year}&month1=#{start_dt.month}&day1=#{start_dt.day}" <> "&hour1=#{start_dt.hour}&minute1=0" <> "&year2=#{end_dt.year}&month2=#{end_dt.month}&day2=#{end_dt.day}" <> @@ -26,6 +26,10 @@ defmodule Microwaveprop.Weather.IemClient do "#{@iem_base}/json/raob.py?ts=#{ts}&station=#{station_id}" end + def iemre_url(lat, lon, date) do + "#{@iem_base}/iemre/hourly/#{date}/#{lat}/#{lon}/json" + end + # --- HTTP fetchers --- def fetch_network(network) do @@ -73,6 +77,21 @@ defmodule Microwaveprop.Weather.IemClient do end end + def fetch_iemre(lat, lon, date) do + url = iemre_url(lat, lon, date) + + case Req.get(url, req_options()) do + {:ok, %{status: 200, body: body}} -> + {:ok, parse_iemre_json(body)} + + {:ok, %{status: status}} -> + {:error, "IEM IEMRE HTTP #{status}"} + + {:error, reason} -> + {:error, reason} + end + end + defp req_options do defaults = [retry: &retry?/2, max_retries: 5, retry_delay: &retry_delay/1] overrides = Application.get_env(:microwaveprop, :iem_req_options, []) @@ -95,6 +114,10 @@ defmodule Microwaveprop.Weather.IemClient do # --- Parsers --- + def parse_iemre_json(json) when is_map(json) do + Map.get(json, "data", []) + end + def parse_network_json(json) when is_map(json) do json |> Map.get("stations", []) @@ -142,7 +165,9 @@ defmodule Microwaveprop.Weather.IemClient do wind_direction_deg: parse_int(Enum.at(parts, 6)), sea_level_pressure_mb: parse_float(Enum.at(parts, 7)), altimeter_setting: parse_float(Enum.at(parts, 8)), - sky_condition: parse_nullable_string(Enum.at(parts, 9)) + sky_condition: parse_nullable_string(Enum.at(parts, 9)), + precip_1h_in: parse_float(Enum.at(parts, 10)), + wx_codes: parse_nullable_string(Enum.at(parts, 11)) } end diff --git a/lib/microwaveprop/weather/iemre_observation.ex b/lib/microwaveprop/weather/iemre_observation.ex new file mode 100644 index 00000000..8e9737ce --- /dev/null +++ b/lib/microwaveprop/weather/iemre_observation.ex @@ -0,0 +1,28 @@ +defmodule Microwaveprop.Weather.IemreObservation do + @moduledoc false + use Ecto.Schema + + import Ecto.Changeset + + @primary_key {:id, :binary_id, autogenerate: true} + @foreign_key_type :binary_id + + schema "iemre_observations" do + field :lat, :float + field :lon, :float + field :date, :date + field :hourly, {:array, :map}, default: [] + + timestamps(type: :utc_datetime) + end + + @required_fields ~w(lat lon date)a + @optional_fields ~w(hourly)a + + def changeset(observation, attrs) do + observation + |> cast(attrs, @required_fields ++ @optional_fields) + |> validate_required(@required_fields) + |> unique_constraint([:lat, :lon, :date]) + end +end diff --git a/lib/microwaveprop/weather/surface_observation.ex b/lib/microwaveprop/weather/surface_observation.ex index bcb48e31..ce5e3025 100644 --- a/lib/microwaveprop/weather/surface_observation.ex +++ b/lib/microwaveprop/weather/surface_observation.ex @@ -18,12 +18,14 @@ defmodule Microwaveprop.Weather.SurfaceObservation do field :sea_level_pressure_mb, :float field :altimeter_setting, :float field :sky_condition, :string + field :precip_1h_in, :float + field :wx_codes, :string timestamps(type: :utc_datetime) end @required_fields ~w(station_id observed_at)a - @optional_fields ~w(temp_f dewpoint_f relative_humidity wind_speed_kts wind_direction_deg sea_level_pressure_mb altimeter_setting sky_condition)a + @optional_fields ~w(temp_f dewpoint_f relative_humidity wind_speed_kts wind_direction_deg sea_level_pressure_mb altimeter_setting sky_condition precip_1h_in wx_codes)a def changeset(observation, attrs) do observation diff --git a/lib/microwaveprop/workers/iemre_fetch_worker.ex b/lib/microwaveprop/workers/iemre_fetch_worker.ex new file mode 100644 index 00000000..98b34274 --- /dev/null +++ b/lib/microwaveprop/workers/iemre_fetch_worker.ex @@ -0,0 +1,64 @@ +defmodule Microwaveprop.Workers.IemreFetchWorker do + @moduledoc false + use Oban.Worker, queue: :iemre, max_attempts: 20 + + alias Microwaveprop.Weather + alias Microwaveprop.Weather.IemClient + + require Logger + + @impl Oban.Worker + def backoff(%Oban.Job{attempt: attempt}) do + min(120 * Integer.pow(2, attempt - 1), _six_hours = 21_600) + end + + @impl Oban.Worker + def perform(%Oban.Job{args: args}) do + %{"lat" => lat, "lon" => lon, "date" => date_str} = args + date = Date.from_iso8601!(date_str) + + if Weather.has_iemre_observation?(lat, lon, date) do + Logger.info("IEMRE observation already exists for #{lat},#{lon} @ #{date_str}") + :ok + else + Logger.info("Fetching IEMRE data for #{lat},#{lon} @ #{date_str}") + + case IemClient.fetch_iemre(lat, lon, date) do + {:ok, []} -> + Logger.info("IEMRE returned empty data for #{lat},#{lon} @ #{date_str}") + :ok + + {:ok, data} -> + Weather.upsert_iemre_observation(%{ + lat: lat, + lon: lon, + date: date, + hourly: data + }) + + Logger.info("IEMRE observation saved for #{lat},#{lon} @ #{date_str} (#{length(data)} hours)") + :ok + + {:error, reason} -> + if transient_failure?(reason) do + Logger.error("IEMRE transient error for #{lat},#{lon} @ #{date_str}: #{inspect(reason)}") + {:error, reason} + else + Logger.warning("IEMRE permanent failure for #{lat},#{lon} @ #{date_str}: #{inspect(reason)}") + {:cancel, reason} + end + end + end + end + + defp transient_failure?(%{__exception__: true}), do: true + defp transient_failure?("IEM IEMRE HTTP " <> status), do: server_error?(status) + defp transient_failure?(_), do: false + + defp server_error?(status) do + case Integer.parse(status) do + {code, _} when code in [429, 500, 502, 503, 504] -> true + _ -> false + end + end +end diff --git a/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex b/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex index ed2bbaaf..925e94ae 100644 --- a/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex +++ b/lib/microwaveprop/workers/qso_weather_enqueue_worker.ex @@ -6,6 +6,7 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do alias Microwaveprop.Weather alias Microwaveprop.Weather.HrrrClient alias Microwaveprop.Workers.HrrrFetchWorker + alias Microwaveprop.Workers.IemreFetchWorker alias Microwaveprop.Workers.TerrainProfileWorker alias Microwaveprop.Workers.WeatherFetchWorker @@ -17,6 +18,7 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do enqueue_weather_jobs() enqueue_hrrr_jobs() enqueue_terrain_jobs() + enqueue_iemre_jobs() :ok end @@ -101,6 +103,12 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do |> Enum.uniq_by(fn changeset -> changeset.changes.args end) end + def build_iemre_jobs(qsos) do + qsos + |> Enum.flat_map(&iemre_job_for_qso/1) + |> Enum.uniq_by(fn changeset -> changeset.changes.args end) + end + defp jobs_for_qso(qso) do lat = qso.pos1["lat"] lon = qso.pos1["lon"] || qso.pos1["lng"] @@ -111,6 +119,25 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do asos_jobs ++ raob_jobs end + defp enqueue_iemre_jobs do + case Radio.unprocessed_iemre_qsos() do + [] -> + :ok + + qsos -> + jobs = build_iemre_jobs(qsos) + + if jobs != [] do + Oban.insert_all(jobs) + end + + qso_ids = Enum.map(qsos, & &1.id) + Radio.mark_iemre_queued!(qso_ids) + + enqueue_iemre_jobs() + end + end + defp hrrr_job_for_qso(%{pos1: nil}), do: [] defp hrrr_job_for_qso(qso) do @@ -150,6 +177,28 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorker do end) end + defp iemre_job_for_qso(%{pos1: nil}), do: [] + + defp iemre_job_for_qso(qso) do + lat = qso.pos1["lat"] + lon = qso.pos1["lon"] || qso.pos1["lng"] + + if lat && lon do + {rlat, rlon} = Weather.round_to_iemre_grid(lat, lon) + date = DateTime.to_date(qso.qso_timestamp) + + [ + IemreFetchWorker.new(%{ + "lat" => rlat, + "lon" => rlon, + "date" => Date.to_iso8601(date) + }) + ] + else + [] + end + end + defp build_raob_jobs(lat, lon, timestamp) do sounding_times = Weather.sounding_times_around(timestamp) diff --git a/priv/repo/migrations/20260330180421_add_precipitation_to_surface_observations.exs b/priv/repo/migrations/20260330180421_add_precipitation_to_surface_observations.exs new file mode 100644 index 00000000..486a747a --- /dev/null +++ b/priv/repo/migrations/20260330180421_add_precipitation_to_surface_observations.exs @@ -0,0 +1,10 @@ +defmodule Microwaveprop.Repo.Migrations.AddPrecipitationToSurfaceObservations do + use Ecto.Migration + + def change do + alter table(:surface_observations) do + add :precip_1h_in, :float + add :wx_codes, :string + end + end +end diff --git a/priv/repo/migrations/20260330180721_create_iemre_observations_and_add_iemre_queued.exs b/priv/repo/migrations/20260330180721_create_iemre_observations_and_add_iemre_queued.exs new file mode 100644 index 00000000..7b13a310 --- /dev/null +++ b/priv/repo/migrations/20260330180721_create_iemre_observations_and_add_iemre_queued.exs @@ -0,0 +1,25 @@ +defmodule Microwaveprop.Repo.Migrations.CreateIemreObservationsAndAddIemreQueued do + use Ecto.Migration + + def change do + create table(:iemre_observations, primary_key: false) do + add :id, :binary_id, primary_key: true + add :lat, :float, null: false + add :lon, :float, null: false + add :date, :date, null: false + add :hourly, {:array, :map}, default: [] + + timestamps(type: :utc_datetime) + end + + create unique_index(:iemre_observations, [:lat, :lon, :date]) + + alter table(:qsos) do + add :iemre_queued, :boolean, default: false, null: false + end + + execute "UPDATE qsos SET iemre_queued = false", "SELECT 1" + + create index(:qsos, [:iemre_queued], where: "iemre_queued = false") + end +end diff --git a/test/microwaveprop/weather/iem_client_test.exs b/test/microwaveprop/weather/iem_client_test.exs index 55918384..0e68cdfe 100644 --- a/test/microwaveprop/weather/iem_client_test.exs +++ b/test/microwaveprop/weather/iem_client_test.exs @@ -4,11 +4,11 @@ defmodule Microwaveprop.Weather.IemClientTest do alias Microwaveprop.Weather.IemClient describe "parse_asos_csv/1" do - test "parses valid CSV rows" do + test "parses valid CSV rows including precipitation" do csv = """ - station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1 - KDFW,2026-03-28 18:53,75.0,55.0,49.12,12,180,1013.2,29.92,SCT - KDFW,2026-03-28 19:53,76.0,54.0,46.00,10,190,1013.0,29.91,FEW + station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1,p01i,wxcodes + KDFW,2026-03-28 18:53,75.0,55.0,49.12,12,180,1013.2,29.92,SCT,0.03,RA + KDFW,2026-03-28 19:53,76.0,54.0,46.00,10,190,1013.0,29.91,FEW,null,null """ rows = IemClient.parse_asos_csv(csv) @@ -24,12 +24,18 @@ defmodule Microwaveprop.Weather.IemClientTest do assert first.altimeter_setting == 29.92 assert first.sky_condition == "SCT" assert first.observed_at == ~U[2026-03-28 18:53:00Z] + assert first.precip_1h_in == 0.03 + assert first.wx_codes == "RA" + + second = Enum.at(rows, 1) + assert second.precip_1h_in == nil + assert second.wx_codes == nil end test "handles null values in CSV" do csv = """ - station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1 - KDFW,2026-03-28 18:53,null,null,null,null,null,null,null,null + station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1,p01i,wxcodes + KDFW,2026-03-28 18:53,null,null,null,null,null,null,null,null,null,null """ rows = IemClient.parse_asos_csv(csv) @@ -39,13 +45,15 @@ defmodule Microwaveprop.Weather.IemClientTest do assert first.temp_f == nil assert first.dewpoint_f == nil assert first.sky_condition == nil + assert first.precip_1h_in == nil + assert first.wx_codes == nil end test "skips comment lines and empty lines" do csv = """ #DEBUG something - station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1 - KDFW,2026-03-28 18:53,75.0,55.0,49.0,12,180,1013.2,29.92,SCT + station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1,p01i,wxcodes + KDFW,2026-03-28 18:53,75.0,55.0,49.0,12,180,1013.2,29.92,SCT,0.01, """ @@ -106,6 +114,8 @@ defmodule Microwaveprop.Weather.IemClientTest do assert url =~ "year2=2026" assert url =~ "hour2=18" assert url =~ "format=onlycomma" + assert url =~ "data=p01i" + assert url =~ "data=wxcodes" end end @@ -156,8 +166,8 @@ defmodule Microwaveprop.Weather.IemClientTest do describe "fetch_asos/3" do test "returns parsed observations on success" do csv = """ - station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1 - KDFW,2026-03-28 18:53,75.0,55.0,49.12,12,180,1013.2,29.92,SCT + station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1,p01i,wxcodes + KDFW,2026-03-28 18:53,75.0,55.0,49.12,12,180,1013.2,29.92,SCT,0.05,RA BR """ Req.Test.stub(IemClient, fn conn -> @@ -169,6 +179,8 @@ defmodule Microwaveprop.Weather.IemClientTest do assert obs.temp_f == 75.0 assert obs.observed_at == ~U[2026-03-28 18:53:00Z] + assert obs.precip_1h_in == 0.05 + assert obs.wx_codes == "RA BR" end test "returns error on non-200 status" do @@ -273,4 +285,67 @@ defmodule Microwaveprop.Weather.IemClientTest do assert station.lon == -97.038 end end + + describe "iemre_url/3" do + test "builds correct URL for lat/lon and date" do + url = IemClient.iemre_url(32.875, -97.0, ~D[2026-03-28]) + assert url == "https://mesonet.agron.iastate.edu/iemre/hourly/2026-03-28/32.875/-97.0/json" + end + end + + describe "parse_iemre_json/1" do + test "extracts data array from JSON response" do + json = %{ + "data" => [ + %{"utc_hour" => 0, "p01m_mm" => 0.0, "skyc_pct" => 25.0}, + %{"utc_hour" => 1, "p01m_mm" => 0.5, "skyc_pct" => 50.0} + ] + } + + result = IemClient.parse_iemre_json(json) + assert length(result) == 2 + assert hd(result)["utc_hour"] == 0 + end + + test "returns empty list for missing data key" do + assert IemClient.parse_iemre_json(%{}) == [] + end + + test "returns empty list for empty data array" do + assert IemClient.parse_iemre_json(%{"data" => []}) == [] + end + end + + describe "fetch_iemre/3" do + test "returns parsed hourly data on success" do + Req.Test.stub(IemClient, fn conn -> + Req.Test.json(conn, %{ + "data" => [ + %{"utc_hour" => 0, "p01m_mm" => 0.0}, + %{"utc_hour" => 1, "p01m_mm" => 1.2} + ] + }) + end) + + assert {:ok, data} = IemClient.fetch_iemre(32.875, -97.0, ~D[2026-03-28]) + assert length(data) == 2 + assert Enum.at(data, 1)["p01m_mm"] == 1.2 + end + + test "returns error on non-200 status" do + Req.Test.stub(IemClient, fn conn -> + Plug.Conn.send_resp(conn, 503, "Service Unavailable") + end) + + assert {:error, "IEM IEMRE HTTP 503"} = IemClient.fetch_iemre(32.875, -97.0, ~D[2026-03-28]) + end + + test "returns error on 404" do + Req.Test.stub(IemClient, fn conn -> + Plug.Conn.send_resp(conn, 404, "not found") + end) + + assert {:error, "IEM IEMRE HTTP 404"} = IemClient.fetch_iemre(32.875, -97.0, ~D[2026-03-28]) + end + end end diff --git a/test/microwaveprop/weather/iemre_observation_test.exs b/test/microwaveprop/weather/iemre_observation_test.exs new file mode 100644 index 00000000..aa658e12 --- /dev/null +++ b/test/microwaveprop/weather/iemre_observation_test.exs @@ -0,0 +1,61 @@ +defmodule Microwaveprop.Weather.IemreObservationTest do + use Microwaveprop.DataCase, async: true + + alias Microwaveprop.Weather.IemreObservation + + @valid_attrs %{ + lat: 32.875, + lon: -97.0, + date: ~D[2026-03-28], + hourly: [ + %{"hour" => 0, "p01m_mm" => 0.0}, + %{"hour" => 1, "p01m_mm" => 0.5} + ] + } + + describe "changeset/2" do + test "valid attributes produce a valid changeset" do + changeset = IemreObservation.changeset(%IemreObservation{}, @valid_attrs) + assert changeset.valid? + end + + test "requires lat, lon, and date" do + changeset = IemreObservation.changeset(%IemreObservation{}, %{}) + + assert %{ + lat: ["can't be blank"], + lon: ["can't be blank"], + date: ["can't be blank"] + } = errors_on(changeset) + end + + test "hourly is optional" do + attrs = Map.delete(@valid_attrs, :hourly) + changeset = IemreObservation.changeset(%IemreObservation{}, attrs) + assert changeset.valid? + end + + test "persists to the database" do + changeset = IemreObservation.changeset(%IemreObservation{}, @valid_attrs) + assert {:ok, obs} = Repo.insert(changeset) + assert obs.lat == 32.875 + assert obs.lon == -97.0 + assert obs.date == ~D[2026-03-28] + assert length(obs.hourly) == 2 + end + + test "enforces unique constraint on lat, lon, date" do + assert {:ok, _} = + %IemreObservation{} + |> IemreObservation.changeset(@valid_attrs) + |> Repo.insert() + + assert {:error, changeset} = + %IemreObservation{} + |> IemreObservation.changeset(@valid_attrs) + |> Repo.insert() + + assert %{lat: ["has already been taken"]} = errors_on(changeset) + end + end +end diff --git a/test/microwaveprop/weather/surface_observation_test.exs b/test/microwaveprop/weather/surface_observation_test.exs index 6e449fc9..ca11e59b 100644 --- a/test/microwaveprop/weather/surface_observation_test.exs +++ b/test/microwaveprop/weather/surface_observation_test.exs @@ -25,7 +25,9 @@ defmodule Microwaveprop.Weather.SurfaceObservationTest do wind_direction_deg: 180, sea_level_pressure_mb: 1013.25, altimeter_setting: 29.92, - sky_condition: "SCT" + sky_condition: "SCT", + precip_1h_in: 0.03, + wx_codes: "RA" } describe "changeset/2" do @@ -64,6 +66,8 @@ defmodule Microwaveprop.Weather.SurfaceObservationTest do assert {:ok, obs} = Repo.insert(changeset) assert obs.temp_f == 75.0 assert obs.sky_condition == "SCT" + assert obs.precip_1h_in == 0.03 + assert obs.wx_codes == "RA" assert obs.station_id == station.id end diff --git a/test/microwaveprop/weather_test.exs b/test/microwaveprop/weather_test.exs index 4e34e1db..0325dc5d 100644 --- a/test/microwaveprop/weather_test.exs +++ b/test/microwaveprop/weather_test.exs @@ -519,4 +519,90 @@ defmodule Microwaveprop.WeatherTest do assert Weather.round_to_hrrr_grid(-33.456, -97.999) == {-33.46, -98.0} end end + + describe "round_to_iemre_grid/2" do + test "rounds to nearest 0.125 degree" do + assert Weather.round_to_iemre_grid(32.9047, -97.0382) == {32.875, -97.0} + end + + test "handles exact grid points" do + assert Weather.round_to_iemre_grid(33.0, -97.125) == {33.0, -97.125} + end + + test "handles negative values correctly" do + assert Weather.round_to_iemre_grid(-33.44, -97.99) == {-33.5, -98.0} + end + end + + @iemre_attrs %{ + lat: 32.875, + lon: -97.0, + date: ~D[2026-03-28], + hourly: [%{"hour" => 0, "p01m_mm" => 0.0}] + } + + describe "upsert_iemre_observation/1" do + test "inserts a new IEMRE observation" do + assert {:ok, obs} = Weather.upsert_iemre_observation(@iemre_attrs) + assert obs.lat == 32.875 + assert obs.date == ~D[2026-03-28] + end + + test "updates existing observation on conflict" do + {:ok, first} = Weather.upsert_iemre_observation(@iemre_attrs) + + {:ok, second} = + Weather.upsert_iemre_observation(%{ + @iemre_attrs + | hourly: [%{"hour" => 0, "p01m_mm" => 1.5}] + }) + + assert first.id == second.id + assert hd(second.hourly)["p01m_mm"] == 1.5 + end + end + + describe "has_iemre_observation?/3" do + test "returns true when observation exists" do + Weather.upsert_iemre_observation(@iemre_attrs) + assert Weather.has_iemre_observation?(32.875, -97.0, ~D[2026-03-28]) + end + + test "returns false when no observation exists" do + refute Weather.has_iemre_observation?(32.875, -97.0, ~D[2026-03-28]) + end + + test "returns false for different location" do + Weather.upsert_iemre_observation(@iemre_attrs) + refute Weather.has_iemre_observation?(40.0, -80.0, ~D[2026-03-28]) + end + end + + describe "iemre_for_qso/1" do + test "returns IEMRE observation matching QSO pos1 and date" do + Weather.upsert_iemre_observation(@iemre_attrs) + + qso = %{ + pos1: %{"lat" => 32.9, "lon" => -97.0}, + qso_timestamp: ~U[2026-03-28 18:00:00Z] + } + + obs = Weather.iemre_for_qso(qso) + assert obs + assert obs.lat == 32.875 + end + + test "returns nil when no IEMRE observation exists" do + qso = %{ + pos1: %{"lat" => 40.0, "lon" => -80.0}, + qso_timestamp: ~U[2026-03-28 18:00:00Z] + } + + assert Weather.iemre_for_qso(qso) == nil + end + + test "returns nil when QSO has no pos1" do + assert Weather.iemre_for_qso(%{pos1: nil, qso_timestamp: ~U[2026-03-28 18:00:00Z]}) == nil + end + end end diff --git a/test/microwaveprop/workers/iemre_fetch_worker_test.exs b/test/microwaveprop/workers/iemre_fetch_worker_test.exs new file mode 100644 index 00000000..2f65d35c --- /dev/null +++ b/test/microwaveprop/workers/iemre_fetch_worker_test.exs @@ -0,0 +1,121 @@ +defmodule Microwaveprop.Workers.IemreFetchWorkerTest do + use Microwaveprop.DataCase, async: true + + alias Microwaveprop.Weather + alias Microwaveprop.Weather.IemClient + alias Microwaveprop.Weather.IemreObservation + alias Microwaveprop.Workers.IemreFetchWorker + + @sample_iemre_data [ + %{"utc_hour" => 0, "p01m_mm" => 0.0, "skyc_pct" => 25.0}, + %{"utc_hour" => 1, "p01m_mm" => 0.5, "skyc_pct" => 50.0} + ] + + describe "perform/1" do + test "skips if IEMRE observation already exists" do + Weather.upsert_iemre_observation(%{ + lat: 32.875, + lon: -97.0, + date: ~D[2026-03-28], + hourly: @sample_iemre_data + }) + + # Stub should NOT be called + Req.Test.stub(IemClient, fn conn -> + Plug.Conn.send_resp(conn, 500, "Should not be called") + end) + + job = %Oban.Job{ + args: %{ + "lat" => 32.875, + "lon" => -97.0, + "date" => "2026-03-28" + } + } + + assert :ok = IemreFetchWorker.perform(job) + assert Repo.aggregate(IemreObservation, :count) == 1 + end + + test "fetches and stores IEMRE data on success" do + Req.Test.stub(IemClient, fn conn -> + Req.Test.json(conn, %{"data" => @sample_iemre_data}) + end) + + job = %Oban.Job{ + args: %{ + "lat" => 32.875, + "lon" => -97.0, + "date" => "2026-03-28" + } + } + + assert :ok = IemreFetchWorker.perform(job) + assert Repo.aggregate(IemreObservation, :count) == 1 + + obs = Repo.one(IemreObservation) + assert obs.lat == 32.875 + assert obs.lon == -97.0 + assert obs.date == ~D[2026-03-28] + assert length(obs.hourly) == 2 + end + + test "returns :ok on empty data" do + Req.Test.stub(IemClient, fn conn -> + Req.Test.json(conn, %{"data" => []}) + end) + + job = %Oban.Job{ + args: %{ + "lat" => 32.875, + "lon" => -97.0, + "date" => "2026-03-28" + } + } + + assert :ok = IemreFetchWorker.perform(job) + assert Repo.aggregate(IemreObservation, :count) == 0 + end + + test "retries on transient server error" do + Req.Test.stub(IemClient, fn conn -> + Plug.Conn.send_resp(conn, 503, "Service Unavailable") + end) + + job = %Oban.Job{ + args: %{ + "lat" => 32.875, + "lon" => -97.0, + "date" => "2026-03-28" + } + } + + assert {:error, "IEM IEMRE HTTP 503"} = IemreFetchWorker.perform(job) + end + + test "cancels on permanent failure (404)" do + Req.Test.stub(IemClient, fn conn -> + Plug.Conn.send_resp(conn, 404, "not found") + end) + + job = %Oban.Job{ + args: %{ + "lat" => 32.875, + "lon" => -97.0, + "date" => "2026-03-28" + } + } + + assert {:cancel, "IEM IEMRE HTTP 404"} = IemreFetchWorker.perform(job) + end + end + + describe "backoff/1" do + test "uses exponential backoff capped at 6 hours" do + assert IemreFetchWorker.backoff(%Oban.Job{attempt: 1}) == 120 + assert IemreFetchWorker.backoff(%Oban.Job{attempt: 2}) == 240 + assert IemreFetchWorker.backoff(%Oban.Job{attempt: 10}) == 21_600 + assert IemreFetchWorker.backoff(%Oban.Job{attempt: 20}) == 21_600 + end + end +end diff --git a/test/microwaveprop/workers/qso_weather_enqueue_worker_test.exs b/test/microwaveprop/workers/qso_weather_enqueue_worker_test.exs index 975ee4e9..53f880db 100644 --- a/test/microwaveprop/workers/qso_weather_enqueue_worker_test.exs +++ b/test/microwaveprop/workers/qso_weather_enqueue_worker_test.exs @@ -125,8 +125,14 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorkerTest do # Stub IEM client so inline Oban doesn't blow up when executing enqueued jobs Req.Test.stub(IemClient, fn conn -> case conn.request_path do - "/cgi-bin/request/asos.py" -> Req.Test.text(conn, "#DEBUG,\nstation,valid,tmpf\n") - _ -> Req.Test.json(conn, %{"profiles" => []}) + "/cgi-bin/request/asos.py" -> + Req.Test.text(conn, "#DEBUG,\nstation,valid,tmpf\n") + + "/iemre/" <> _ -> + Req.Test.json(conn, %{"data" => [%{"utc_hour" => 0, "p01m_mm" => 0.0}]}) + + _ -> + Req.Test.json(conn, %{"profiles" => []}) end end) @@ -266,8 +272,14 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorkerTest do setup do Req.Test.stub(IemClient, fn conn -> case conn.request_path do - "/cgi-bin/request/asos.py" -> Req.Test.text(conn, "#DEBUG,\nstation,valid,tmpf\n") - _ -> Req.Test.json(conn, %{"profiles" => []}) + "/cgi-bin/request/asos.py" -> + Req.Test.text(conn, "#DEBUG,\nstation,valid,tmpf\n") + + "/iemre/" <> _ -> + Req.Test.json(conn, %{"data" => [%{"utc_hour" => 0, "p01m_mm" => 0.0}]}) + + _ -> + Req.Test.json(conn, %{"profiles" => []}) end end) @@ -303,4 +315,102 @@ defmodule Microwaveprop.Workers.QsoWeatherEnqueueWorkerTest do refute Enum.any?(all_qsos, & &1.terrain_queued) end end + + describe "build_iemre_jobs/1" do + test "builds one IEMRE job per QSO with rounded lat/lon" do + qso = create_qso(%{pos1: %{"lat" => 32.907, "lon" => -97.038}}) + + jobs = QsoWeatherEnqueueWorker.build_iemre_jobs([qso]) + + assert length(jobs) == 1 + job = hd(jobs) + assert job.changes.args["lat"] == 32.875 + assert job.changes.args["lon"] == -97.0 + assert job.changes.args["date"] == "2026-03-28" + end + + test "deduplicates jobs with identical args" do + q1 = + create_qso(%{ + station1: "A1", + qso_timestamp: ~U[2026-03-28 18:10:00Z], + pos1: %{"lat" => 32.90, "lon" => -97.04} + }) + + q2 = + create_qso(%{ + station1: "A2", + qso_timestamp: ~U[2026-03-28 20:20:00Z], + pos1: %{"lat" => 32.90, "lon" => -97.04} + }) + + jobs = QsoWeatherEnqueueWorker.build_iemre_jobs([q1, q2]) + + # Same location and same date → 1 job + assert length(jobs) == 1 + end + + test "creates separate jobs for different dates" do + q1 = + create_qso(%{ + station1: "A1", + qso_timestamp: ~U[2026-03-28 23:00:00Z], + pos1: %{"lat" => 32.90, "lon" => -97.04} + }) + + q2 = + create_qso(%{ + station1: "A2", + qso_timestamp: ~U[2026-03-29 01:00:00Z], + pos1: %{"lat" => 32.90, "lon" => -97.04} + }) + + jobs = QsoWeatherEnqueueWorker.build_iemre_jobs([q1, q2]) + assert length(jobs) == 2 + end + + test "skips QSOs without pos1" do + qso = create_qso(%{pos1: nil}) + assert QsoWeatherEnqueueWorker.build_iemre_jobs([qso]) == [] + end + end + + describe "perform/1 iemre enqueue" do + setup do + Req.Test.stub(IemClient, fn conn -> + case conn.request_path do + "/cgi-bin/request/asos.py" -> + Req.Test.text(conn, "#DEBUG,\nstation,valid,tmpf\n") + + "/iemre/" <> _ -> + Req.Test.json(conn, %{"data" => [%{"utc_hour" => 0, "p01m_mm" => 0.0}]}) + + _ -> + Req.Test.json(conn, %{"profiles" => []}) + end + end) + + Req.Test.stub(HrrrClient, fn conn -> + Plug.Conn.send_resp(conn, 404, "not found") + end) + + Req.Test.stub(ElevationClient, fn conn -> + params = Plug.Conn.fetch_query_params(conn).query_params + lat_count = params["latitude"] |> String.split(",") |> length() + Req.Test.json(conn, %{"elevation" => List.duplicate(200.0, lat_count)}) + end) + + :ok + end + + test "enqueues IEMRE jobs and marks QSOs as iemre_queued" do + qso = create_qso() + assert qso.iemre_queued == false + + assert :ok = QsoWeatherEnqueueWorker.perform(%Oban.Job{args: %{}}) + + updated = Repo.get!(Qso, qso.id) + assert updated.iemre_queued == true + end + end end diff --git a/test/microwaveprop/workers/weather_fetch_worker_test.exs b/test/microwaveprop/workers/weather_fetch_worker_test.exs index db80bb6d..f6d77b9b 100644 --- a/test/microwaveprop/workers/weather_fetch_worker_test.exs +++ b/test/microwaveprop/workers/weather_fetch_worker_test.exs @@ -36,9 +36,9 @@ defmodule Microwaveprop.Workers.WeatherFetchWorkerTest do defp asos_csv_body do """ #DEBUG, - station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1 - KORD,2026-03-28 18:53,45.0,30.0,55.0,10.0,270,1013.5,29.92,CLR - KORD,2026-03-28 19:53,46.0,31.0,54.0,12.0,280,1013.2,29.91,FEW + station,valid,tmpf,dwpf,relh,sknt,drct,mslp,alti,skyc1,p01i,wxcodes + KORD,2026-03-28 18:53,45.0,30.0,55.0,10.0,270,1013.5,29.92,CLR,null,null + KORD,2026-03-28 19:53,46.0,31.0,54.0,12.0,280,1013.2,29.91,FEW,0.02,RA """ end