diff --git a/config/runtime.exs b/config/runtime.exs index 91d5e189..123007a8 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -371,6 +371,5 @@ if config_env() == :prod do config :microwaveprop, :prom_ex_include_oban_plugin, false config :microwaveprop, hrrr_base_url: System.get_env("HRRR_BASE_URL", "http://skippy.w5isp.com:8080") config :microwaveprop, srtm_tiles_dir: "/data/srtm" - config :microwaveprop, start_freshness_monitor: false end diff --git a/lib/microwaveprop/weather/iem_client.ex b/lib/microwaveprop/weather/iem_client.ex index 560e9401..a398dad1 100644 --- a/lib/microwaveprop/weather/iem_client.ex +++ b/lib/microwaveprop/weather/iem_client.ex @@ -12,10 +12,13 @@ defmodule Microwaveprop.Weather.IemClient do "#{@iem_base}/json/network.py?network=#{network}" end - @spec asos_url(String.t(), DateTime.t(), DateTime.t()) :: String.t() - def asos_url(station_id, start_dt, end_dt) do + @spec asos_url(String.t() | [String.t()], DateTime.t(), DateTime.t()) :: String.t() + def asos_url(station_id_or_codes, start_dt, end_dt) do + codes = List.wrap(station_id_or_codes) + station_params = Enum.map_join(codes, "&", &"station=#{&1}") + "#{@iem_base}/cgi-bin/request/asos.py" <> - "?station=#{station_id}" <> + "?#{station_params}" <> "&data=tmpf&data=dwpf&data=relh&data=sknt&data=drct" <> "&data=mslp&data=alti&data=skyc1&data=p01i&data=wxcodes" <> "&year1=#{start_dt.year}&month1=#{start_dt.month}&day1=#{start_dt.day}" <> @@ -75,6 +78,43 @@ defmodule Microwaveprop.Weather.IemClient do end) end + @doc """ + Batched variant of `fetch_asos/3` — one HTTP call for all `codes`, + returning `{:ok, %{code => [rows]}}`. Stations absent from the + response appear in the map with `[]` so callers can distinguish + "IEM had no data" from "we never asked about this code". + + IEM's ASOS CSV endpoint accepts multiple `station=` params and + keys the response rows on the first column. A ~15-station batch + replaces ~15 individual requests, each of which pays the + per-pod `IemRateLimiter` gap + 429 retry tail. + """ + @spec fetch_asos_batch([String.t()], DateTime.t(), DateTime.t()) :: + {:ok, %{String.t() => [map()]}} | {:error, term()} + def fetch_asos_batch(codes, start_dt, end_dt) when is_list(codes) and codes != [] do + url = asos_url(codes, start_dt, end_dt) + + Microwaveprop.Instrument.span([:iem, :fetch_asos_batch], %{count: length(codes)}, fn -> + IemRateLimiter.acquire() + + case Req.get(url, req_options()) do + {:ok, %{status: 200, body: body}} -> + rows = parse_asos_csv(body) + grouped = Enum.group_by(rows, & &1.station_code) + # Ensure every requested code has at least an empty list in the + # map so downstream iteration doesn't miss "no data" cases. + filled = Enum.reduce(codes, grouped, &Map.put_new(&2, &1, [])) + {:ok, filled} + + {:ok, %{status: status}} -> + {:error, "IEM ASOS HTTP #{status}"} + + {:error, reason} -> + {:error, reason} + end + end) + end + @spec fetch_raob(String.t(), DateTime.t()) :: {:ok, [map()]} | {:error, term()} def fetch_raob(station_id, dt) do url = raob_url(station_id, dt) @@ -252,6 +292,7 @@ defmodule Microwaveprop.Weather.IemClient do parts = String.split(line, ",") %{ + station_code: parse_nullable_string(Enum.at(parts, 0)), observed_at: parse_asos_timestamp(Enum.at(parts, 1)), temp_f: parse_float(Enum.at(parts, 2)), dewpoint_f: parse_float(Enum.at(parts, 3)), diff --git a/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex b/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex index beefee93..5b6f5e1d 100644 --- a/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex +++ b/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex @@ -392,17 +392,32 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorker do station_ids = Enum.map(stations, & &1.id) covered = Weather.station_ids_with_surface_observations(station_ids, start_dt, end_dt) - stations - |> Enum.reject(fn station -> MapSet.member?(covered, station.id) end) - |> Enum.map(fn station -> - WeatherFetchWorker.new(%{ - "fetch_type" => "asos", - "station_id" => station.id, - "station_code" => station.station_code, - "start_dt" => DateTime.to_iso8601(start_dt), - "end_dt" => DateTime.to_iso8601(end_dt) - }) - end) + # Group every uncovered station for this (lat, lon, window) into one + # batch job. IEM's ASOS CSV endpoint accepts multiple station= params + # in a single request and we group the result rows by station_code in + # the worker. Collapsing ~N→1 requests per endpoint was the largest + # lever against the IemRateLimiter gap + 429 retry tail. + uncovered = + stations + |> Enum.reject(fn s -> MapSet.member?(covered, s.id) end) + # Sort for deterministic args → Oban unique dedup works across enqueues. + |> Enum.sort_by(& &1.station_code) + + case uncovered do + [] -> + [] + + stations -> + [ + WeatherFetchWorker.new(%{ + "fetch_type" => "asos_batch", + "station_ids" => Enum.map(stations, & &1.id), + "station_codes" => Enum.map(stations, & &1.station_code), + "start_dt" => DateTime.to_iso8601(start_dt), + "end_dt" => DateTime.to_iso8601(end_dt) + }) + ] + end end defp iemre_job_for_contact(%{pos1: nil}), do: [] diff --git a/lib/microwaveprop/workers/weather_fetch_worker.ex b/lib/microwaveprop/workers/weather_fetch_worker.ex index 34999515..28bd1b14 100644 --- a/lib/microwaveprop/workers/weather_fetch_worker.ex +++ b/lib/microwaveprop/workers/weather_fetch_worker.ex @@ -41,6 +41,28 @@ defmodule Microwaveprop.Workers.WeatherFetchWorker do end end + def perform(%Oban.Job{args: %{"fetch_type" => "asos_batch"} = args}) do + %{ + "station_ids" => station_ids, + "station_codes" => station_codes, + "start_dt" => start_dt_str, + "end_dt" => end_dt_str + } = args + + {:ok, start_dt, _} = DateTime.from_iso8601(start_dt_str) + {:ok, end_dt, _} = DateTime.from_iso8601(end_dt_str) + + Logger.info("WeatherFetch ASOS batch: fetching #{length(station_codes)} stations for #{start_dt_str}..#{end_dt_str}") + + case IemClient.fetch_asos_batch(station_codes, start_dt, end_dt) do + {:ok, grouped} -> + ingest_asos_batch(station_ids, station_codes, grouped, start_dt, end_dt) + + {:error, reason} -> + handle_iem_error("ASOS batch", "#{length(station_codes)} stations", reason) + end + end + def perform(%Oban.Job{args: %{"fetch_type" => "raob"} = args}) do %{ "station_id" => station_id, @@ -141,6 +163,33 @@ defmodule Microwaveprop.Workers.WeatherFetchWorker do |> Repo.insert(on_conflict: :nothing, conflict_target: [:station_id, :observed_at]) end + # Iterate each requested (station_id, station_code) pair, upsert the + # rows the batch response contained for that code, and fall back to + # stub rows when a code was absent from the response. Stations that + # have been deleted between enqueue and perform are silently skipped. + defp ingest_asos_batch(station_ids, station_codes, grouped, start_dt, end_dt) do + stations_by_id = + Station + |> Repo.all() + |> Enum.filter(&(&1.id in station_ids)) + |> Map.new(&{&1.id, &1}) + + station_ids + |> Enum.zip(station_codes) + |> Enum.each(fn {station_id, code} -> + case Map.get(stations_by_id, station_id) do + nil -> + Logger.warning("WeatherFetch ASOS batch: station #{station_id} (#{code}) missing, skipping") + + station -> + rows = Map.get(grouped, code, []) + ingest_asos_rows(station, station_id, code, rows, start_dt, end_dt) + end + end) + + :ok + end + # -- RAOB helpers -- defp fetch_raob(station, station_id, station_code, sounding_time, sounding_time_str, contact_id) do diff --git a/test/microwaveprop/weather/iem_client_test.exs b/test/microwaveprop/weather/iem_client_test.exs index 290e8c1a..f98459e3 100644 --- a/test/microwaveprop/weather/iem_client_test.exs +++ b/test/microwaveprop/weather/iem_client_test.exs @@ -118,6 +118,67 @@ defmodule Microwaveprop.Weather.IemClientTest do assert url =~ "data=p01i" assert url =~ "data=wxcodes" end + + test "accepts a list of station codes and emits one station= param each" do + url = + IemClient.asos_url( + ["KDFW", "KFTW", "KAFW"], + ~U[2026-03-28 12:00:00Z], + ~U[2026-03-28 18:00:00Z] + ) + + assert url =~ "station=KDFW" + assert url =~ "station=KFTW" + assert url =~ "station=KAFW" + end + end + + describe "parse_asos_csv/1 — batched response" do + test "tags each row with its station_code" do + csv = """ + 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,null,null + KFTW,2026-03-28 18:55,73.0,53.0,49.00,10,190,1013.0,29.91,FEW,null,null + """ + + rows = IemClient.parse_asos_csv(csv) + assert length(rows) == 2 + + codes = Enum.map(rows, & &1.station_code) + assert "KDFW" in codes + assert "KFTW" in codes + end + end + + describe "fetch_asos_batch/3" do + test "returns rows grouped by station_code" do + Req.Test.stub(IemClient, fn conn -> + body = """ + 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,null,null + KDFW,2026-03-28 19:53,76.0,54.0,46.00,10,190,1013.0,29.91,FEW,null,null + KFTW,2026-03-28 18:55,73.0,53.0,49.00,10,190,1013.0,29.91,FEW,null,null + """ + + conn + |> Plug.Conn.put_resp_content_type("text/plain") + |> Plug.Conn.resp(200, body) + end) + + assert {:ok, grouped} = + IemClient.fetch_asos_batch( + ["KDFW", "KFTW", "KSPS"], + ~U[2026-03-28 18:00:00Z], + ~U[2026-03-28 20:00:00Z] + ) + + assert is_map(grouped) + assert length(Map.fetch!(grouped, "KDFW")) == 2 + assert length(Map.fetch!(grouped, "KFTW")) == 1 + # Stations in the request but absent from the response get an empty list + # so callers can distinguish "no data" from "not asked about this code". + assert Map.fetch!(grouped, "KSPS") == [] + end end describe "raob_url/2" do diff --git a/test/microwaveprop/workers/contact_weather_enqueue_worker_test.exs b/test/microwaveprop/workers/contact_weather_enqueue_worker_test.exs index 14e44d6e..e86db057 100644 --- a/test/microwaveprop/workers/contact_weather_enqueue_worker_test.exs +++ b/test/microwaveprop/workers/contact_weather_enqueue_worker_test.exs @@ -62,7 +62,7 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorkerTest do end describe "build_weather_jobs/1" do - test "builds ASOS jobs for nearby stations" do + test "builds ASOS batch jobs grouping nearby stations per time window" do station = create_asos_station("KDFW", 32.90, -97.04) contact = create_contact() @@ -70,18 +70,42 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorkerTest do asos_jobs = Enum.filter(jobs, fn j -> - j.changes.args["fetch_type"] == "asos" + j.changes.args["fetch_type"] == "asos_batch" end) assert asos_jobs != [] - job = hd(asos_jobs) - assert job.changes.args["station_id"] == station.id - assert job.changes.args["station_code"] == "KDFW" + # Every uncovered station near this endpoint rides in one batch job, + # not one job per station. + job = Enum.find(asos_jobs, &(station.id in &1.changes.args["station_ids"])) + assert job + assert "KDFW" in job.changes.args["station_codes"] assert job.changes.args["start_dt"] assert job.changes.args["end_dt"] end + test "builds several batch jobs for nearby stations — one per endpoint window" do + create_asos_station("KDFW", 32.90, -97.04) + create_asos_station("KFTW", 32.82, -97.36) + create_asos_station("KAFW", 32.98, -97.32) + contact = create_contact() + + jobs = ContactWeatherEnqueueWorker.build_weather_jobs([contact]) + + asos_jobs = + Enum.filter(jobs, fn j -> + j.changes.args["fetch_type"] == "asos_batch" + end) + + # At least one batch covers pos1's cluster; all three nearby stations + # should land in the same batch (shared 4-hour window, same endpoint). + pos1_batch = Enum.find(asos_jobs, &(length(&1.changes.args["station_codes"]) == 3)) + assert pos1_batch + + assert MapSet.new(pos1_batch.changes.args["station_codes"]) == + MapSet.new(["KDFW", "KFTW", "KAFW"]) + end + test "builds RAOB jobs for nearby sounding stations" do station = create_sounding_station("FWD", 32.83, -97.30) contact = create_contact() @@ -104,7 +128,7 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorkerTest do test "deduplicates jobs with identical args" do _station = create_asos_station("KDFW", 32.90, -97.04) - # Two QSOs at the exact same time produce identical args → deduplicated + # Two QSOs at the exact same time produce identical batch args → deduped q1 = create_contact(%{station1: "A1", qso_timestamp: ~U[2026-03-28 18:00:00Z]}) q2 = create_contact(%{station1: "A2", qso_timestamp: ~U[2026-03-28 18:00:00Z]}) @@ -112,7 +136,7 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorkerTest do asos_jobs = Enum.filter(jobs, fn j -> - j.changes.args["fetch_type"] == "asos" + j.changes.args["fetch_type"] == "asos_batch" end) assert length(asos_jobs) == 1 @@ -125,8 +149,10 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorkerTest do jobs = ContactWeatherEnqueueWorker.build_weather_jobs([contact]) - station_ids = Enum.map(jobs, & &1.changes.args["station_id"]) - assert station_near_pos2.id in station_ids + all_station_ids = + Enum.flat_map(jobs, fn j -> j.changes.args["station_ids"] || [] end) + + assert station_near_pos2.id in all_station_ids end test "returns empty list when no stations nearby" do @@ -149,7 +175,7 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorkerTest do asos_jobs = Enum.filter(jobs, fn j -> - j.changes.args["fetch_type"] == "asos" + j.changes.args["fetch_type"] == "asos_batch" end) assert asos_jobs == [] diff --git a/test/microwaveprop/workers/weather_fetch_worker_test.exs b/test/microwaveprop/workers/weather_fetch_worker_test.exs index 87adb021..1ce653f3 100644 --- a/test/microwaveprop/workers/weather_fetch_worker_test.exs +++ b/test/microwaveprop/workers/weather_fetch_worker_test.exs @@ -59,6 +59,62 @@ defmodule Microwaveprop.Workers.WeatherFetchWorkerTest do } end + describe "ASOS batch perform/1" do + test "fetches one URL covering all stations and ingests per-station rows" do + {:ok, s1} = + Weather.find_or_create_station(%{ + station_code: "KDFW", + station_type: "asos", + name: "DFW", + lat: 32.9, + lon: -97.0 + }) + + {:ok, s2} = + Weather.find_or_create_station(%{ + station_code: "KFTW", + station_type: "asos", + name: "FTW", + lat: 32.82, + lon: -97.36 + }) + + {:ok, s3} = + Weather.find_or_create_station(%{ + station_code: "KAFW", + station_type: "asos", + name: "AFW", + lat: 32.98, + lon: -97.32 + }) + + # One HTTP call returns rows for s1 + s2; s3 is absent → should stub. + Req.Test.stub(IemClient, fn conn -> + body = """ + #DEBUG + 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,null,null + KDFW,2026-03-28 19:53,76.0,54.0,46.0,10,190,1013.0,29.91,FEW,null,null + KFTW,2026-03-28 18:55,73.0,53.0,49.0,10,190,1013.0,29.91,FEW,null,null + """ + + Req.Test.text(conn, body) + end) + + args = %{ + "fetch_type" => "asos_batch", + "station_ids" => [s1.id, s2.id, s3.id], + "station_codes" => ["KDFW", "KFTW", "KAFW"], + "start_dt" => "2026-03-28T16:00:00Z", + "end_dt" => "2026-03-28T20:00:00Z" + } + + assert :ok = WeatherFetchWorker.perform(%Oban.Job{args: args}) + # 2 real rows for KDFW + 1 for KFTW + 1 stub for KAFW = 4 + assert Repo.aggregate(SurfaceObservation, :count) == 4 + end + end + describe "ASOS perform/1" do setup :create_asos_station