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
This commit is contained in:
parent
10e8feb486
commit
3228835636
18 changed files with 734 additions and 22 deletions
|
|
@ -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)},
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
28
lib/microwaveprop/weather/iemre_observation.ex
Normal file
28
lib/microwaveprop/weather/iemre_observation.ex
Normal file
|
|
@ -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
|
||||
|
|
@ -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
|
||||
|
|
|
|||
64
lib/microwaveprop/workers/iemre_fetch_worker.ex
Normal file
64
lib/microwaveprop/workers/iemre_fetch_worker.ex
Normal file
|
|
@ -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
|
||||
|
|
@ -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)
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -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
|
||||
|
|
@ -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
|
||||
|
|
|
|||
61
test/microwaveprop/weather/iemre_observation_test.exs
Normal file
61
test/microwaveprop/weather/iemre_observation_test.exs
Normal file
|
|
@ -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
|
||||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
121
test/microwaveprop/workers/iemre_fetch_worker_test.exs
Normal file
121
test/microwaveprop/workers/iemre_fetch_worker_test.exs
Normal file
|
|
@ -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
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue