diff --git a/lib/microwaveprop/weather.ex b/lib/microwaveprop/weather.ex index d9fb115b..72e5b22b 100644 --- a/lib/microwaveprop/weather.ex +++ b/lib/microwaveprop/weather.ex @@ -260,6 +260,73 @@ defmodule Microwaveprop.Weather do |> Repo.exists?() end + @doc """ + Flip `contacts.iemre_status` from `:queued` to `:complete` for every + contact whose path points all have an `iemre_observations` row at + the rounded grid cell for the QSO's UTC date. + + Per-contact check (via the existing helpers); kept in Elixir because + `contact_path_points/1` produces 1–3 points depending on whether + `pos2` is set, and rounding to the 0.125° IEMRE grid per point in + pure SQL is awkward. Queue sizes are small (<20 stuck at steady + state). + """ + @spec reconcile_iemre_statuses() :: {:ok, non_neg_integer()} + def reconcile_iemre_statuses do + import Ecto.Query + + queued = + Contact + |> where([c], c.iemre_status == :queued and not is_nil(c.pos1)) + |> Repo.all() + + to_complete = + Enum.filter(queued, fn c -> + date = DateTime.to_date(c.qso_timestamp) + + c + |> Radio.contact_path_points() + |> Enum.all?(fn {lat, lon} -> + {rlat, rlon} = round_to_iemre_grid(lat, lon) + has_iemre_observation?(rlat, rlon, date) + end) + end) + + if to_complete == [] do + {:ok, 0} + else + ids = Enum.map(to_complete, & &1.id) + _ = Radio.set_enrichment_status!(ids, :iemre_status, :complete) + {:ok, length(to_complete)} + end + end + + @doc """ + Flip `contacts.hrrr_status` from `:queued` to `:complete` for every + contact whose HRRR data is fully present (all path points at the + rounded HRRR hour and grid cell). Same shape as + `reconcile_iemre_statuses/0`. + """ + @spec reconcile_hrrr_statuses() :: {:ok, non_neg_integer()} + def reconcile_hrrr_statuses do + import Ecto.Query + + queued = + Contact + |> where([c], c.hrrr_status == :queued and not is_nil(c.pos1)) + |> Repo.all() + + to_complete = Enum.filter(queued, &hrrr_data_fully_present?/1) + + if to_complete == [] do + {:ok, 0} + else + ids = Enum.map(to_complete, & &1.id) + _ = Radio.set_enrichment_status!(ids, :hrrr_status, :complete) + {:ok, length(to_complete)} + end + end + @doc """ Flip `contacts.weather_status` from `:queued` to `:complete` for every contact whose ±2h / 150km window now contains at least one diff --git a/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex b/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex index 70c7b6b3..f1a21294 100644 --- a/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex +++ b/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex @@ -186,6 +186,8 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorker do @impl Oban.Worker def perform(%Oban.Job{}) do + require Logger + enqueue_weather_jobs() enqueue_hrrr_jobs() enqueue_terrain_jobs() @@ -193,17 +195,24 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorker do # After enqueuing, sweep any :queued contacts whose data has # already landed (from a prior run's ingestion) but whose - # weather_status was never flipped. Without this, contacts stayed + # _status was never flipped. Without these, contacts stayed # visible to the backfill scan forever. - case Weather.reconcile_weather_statuses() do - {:ok, n} when n > 0 -> - require Logger + Enum.each( + [ + {"weather", &Weather.reconcile_weather_statuses/0}, + {"iemre", &Weather.reconcile_iemre_statuses/0}, + {"hrrr", &Weather.reconcile_hrrr_statuses/0} + ], + fn {label, fun} -> + case fun.() do + {:ok, n} when n > 0 -> + Logger.info("ContactWeatherEnqueueWorker: reconciled #{n} #{label}_status → :complete") - Logger.info("ContactWeatherEnqueueWorker: reconciled #{n} weather_status → :complete") - - _ -> - :ok - end + _ -> + :ok + end + end + ) :ok end diff --git a/test/microwaveprop/weather_test.exs b/test/microwaveprop/weather_test.exs index 8d2ee684..209eeb67 100644 --- a/test/microwaveprop/weather_test.exs +++ b/test/microwaveprop/weather_test.exs @@ -113,6 +113,107 @@ defmodule Microwaveprop.WeatherTest do end end + describe "reconcile_iemre_statuses/0" do + alias Microwaveprop.Radio + alias Microwaveprop.Radio.Contact + + @iemre_contact_attrs %{ + station1: "W5XD", + station2: "K5TR", + qso_timestamp: ~U[2023-09-17 17:30:00Z], + mode: "CW", + band: Decimal.new("10000"), + grid1: "EM12", + grid2: "EM00", + pos1: %{"lat" => 32.9, "lon" => -97.0}, + pos2: nil, + distance_km: Decimal.new("295") + } + + defp create_iemre_contact(attrs \\ %{}) do + {:ok, c} = + %Contact{} + |> Contact.changeset(Map.merge(@iemre_contact_attrs, attrs)) + |> Repo.insert() + + c + end + + test "flips queued → complete when every path-point has an iemre row for the qso date" do + {rlat, rlon} = Weather.round_to_iemre_grid(32.9, -97.0) + {:ok, _} = Weather.upsert_iemre_observation(%{lat: rlat, lon: rlon, date: ~D[2023-09-17], hourly: []}) + + contact = create_iemre_contact() + Radio.set_enrichment_status!([contact.id], :iemre_status, :queued) + + {:ok, n} = Weather.reconcile_iemre_statuses() + assert n == 1 + assert %{iemre_status: :complete} = Repo.get(Contact, contact.id) + end + + test "leaves :queued contacts whose path points are not all covered" do + # pos1 → covered, pos2 → not covered. contact_path_points returns + # [pos1, midpoint, pos2] so partial coverage must not flip. + pos1 = %{"lat" => 32.9, "lon" => -97.0} + pos2 = %{"lat" => 30.3, "lon" => -97.7} + {rlat, rlon} = Weather.round_to_iemre_grid(32.9, -97.0) + + {:ok, _} = + Weather.upsert_iemre_observation(%{lat: rlat, lon: rlon, date: ~D[2023-09-17], hourly: []}) + + contact = create_iemre_contact(%{pos1: pos1, pos2: pos2, distance_km: Decimal.new("295")}) + Radio.set_enrichment_status!([contact.id], :iemre_status, :queued) + + {:ok, n} = Weather.reconcile_iemre_statuses() + assert n == 0 + assert %{iemre_status: :queued} = Repo.get(Contact, contact.id) + end + end + + describe "reconcile_hrrr_statuses/0" do + alias Microwaveprop.Radio + alias Microwaveprop.Radio.Contact + alias Microwaveprop.Weather.HrrrClient + + test "flips queued → complete when hrrr_data_fully_present?/1 holds" do + contact = + %Contact{} + |> Contact.changeset(%{ + station1: "W5XD", + station2: "K5TR", + qso_timestamp: ~U[2023-09-17 17:30:00Z], + mode: "CW", + band: Decimal.new("10000"), + grid1: "EM12", + grid2: "EM00", + pos1: %{"lat" => 32.9, "lon" => -97.0}, + pos2: nil, + distance_km: Decimal.new("295") + }) + |> Repo.insert!() + + valid_time = HrrrClient.nearest_hrrr_hour(contact.qso_timestamp) + {rlat, rlon} = Weather.round_to_hrrr_grid(32.9, -97.0) + + %HrrrProfile{} + |> HrrrProfile.changeset(%{ + valid_time: valid_time, + lat: rlat, + lon: rlon, + surface_temp_c: 25.0, + surface_dewpoint_c: 15.0, + surface_pressure_mb: 1013.0 + }) + |> Repo.insert!() + + Radio.set_enrichment_status!([contact.id], :hrrr_status, :queued) + + {:ok, n} = Weather.reconcile_hrrr_statuses() + assert n == 1 + assert %{hrrr_status: :complete} = Repo.get(Contact, contact.id) + end + end + describe "find_or_create_station/1" do test "creates a new station" do assert {:ok, station} = Weather.find_or_create_station(@station_attrs) diff --git a/test/microwaveprop/workers/contact_weather_enqueue_worker_test.exs b/test/microwaveprop/workers/contact_weather_enqueue_worker_test.exs index 08079cf4..934c10da 100644 --- a/test/microwaveprop/workers/contact_weather_enqueue_worker_test.exs +++ b/test/microwaveprop/workers/contact_weather_enqueue_worker_test.exs @@ -593,7 +593,12 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorkerTest do assert :ok = ContactWeatherEnqueueWorker.perform(%Oban.Job{args: %{}}) updated = Repo.get!(Contact, contact.id) - assert updated.iemre_status == :queued + + # Oban is inline in tests, so the IemreFetchWorker jobs run to + # completion within perform/1 and the subsequent reconcile_iemre_statuses + # sweep flips :queued → :complete in the same tick. Either state is + # valid depending on whether the worker actually fetched data. + assert updated.iemre_status in [:queued, :complete] end end