diff --git a/lib/microwaveprop/weather/grid_cache.ex b/lib/microwaveprop/weather/grid_cache.ex index 87d1d5dd..e87ed74d 100644 --- a/lib/microwaveprop/weather/grid_cache.ex +++ b/lib/microwaveprop/weather/grid_cache.ex @@ -91,27 +91,28 @@ defmodule Microwaveprop.Weather.GridCache do @spec clear() :: :ok def clear do - :ets.delete_all_objects(@table) - :ok + GenServer.call(__MODULE__, :clear) end @doc """ Atomically claim the right to fill the cache for `valid_time`. Returns `true` if this caller won the claim and should run the fill; `false` if another caller is already filling. Prevents N concurrent /weather mounts - after a pod restart from each firing the 15-second `load_weather_grid_from_db` - query and starving the Postgres connection pool. + after a pod restart from each firing the 15-second cold-fill read and + starving the Postgres connection pool. + + The GenServer `Process.monitor`s the caller: if the caller crashes + before calling `release_fill/1`, the lock is released automatically. """ @spec claim_fill(DateTime.t()) :: boolean() def claim_fill(valid_time) do - :ets.insert_new(@lock_table, {valid_time, :in_progress}) + GenServer.call(__MODULE__, {:claim_fill, valid_time, self()}) end @doc "Release a fill lock claimed via `claim_fill/1`." @spec release_fill(DateTime.t()) :: :ok def release_fill(valid_time) do - :ets.delete(@lock_table, valid_time) - :ok + GenServer.call(__MODULE__, {:release_fill, valid_time}) end @spec sync() :: :ok @@ -122,15 +123,51 @@ defmodule Microwaveprop.Weather.GridCache do @impl true def init(_opts) do :ets.new(@table, [:set, :named_table, :public, :compressed, read_concurrency: true]) - :ets.new(@lock_table, [:set, :named_table, :public]) + :ets.new(@lock_table, [:set, :named_table, :protected]) PubSub.subscribe(@pubsub, @topic) - {:ok, %{}} + {:ok, %{monitors: %{}}} end @impl true def handle_call(:sync, _from, state), do: {:reply, :ok, state} + def handle_call(:clear, _from, state) do + :ets.delete_all_objects(@table) + + # Demonitor all tracked callers and drop every lock so tests start clean. + for {ref, _vt} <- state.monitors, do: Process.demonitor(ref, [:flush]) + :ets.delete_all_objects(@lock_table) + + {:reply, :ok, %{state | monitors: %{}}} + end + + def handle_call({:claim_fill, valid_time, caller}, _from, state) do + if :ets.insert_new(@lock_table, {valid_time, caller}) do + ref = Process.monitor(caller) + {:reply, true, %{state | monitors: Map.put(state.monitors, ref, valid_time)}} + else + {:reply, false, state} + end + end + + def handle_call({:release_fill, valid_time}, _from, state) do + {monitors, _matched} = pop_monitor_for(state.monitors, valid_time) + :ets.delete(@lock_table, valid_time) + {:reply, :ok, %{state | monitors: monitors}} + end + @impl true + def handle_info({:DOWN, ref, :process, _pid, _reason}, state) do + case Map.pop(state.monitors, ref) do + {nil, _} -> + {:noreply, state} + + {valid_time, monitors} -> + :ets.delete(@lock_table, valid_time) + {:noreply, %{state | monitors: monitors}} + end + end + def handle_info({:weather_cache_refresh, valid_time, rows}, state) do put(valid_time, rows) {:noreply, state} @@ -138,6 +175,17 @@ defmodule Microwaveprop.Weather.GridCache do def handle_info(_msg, state), do: {:noreply, state} + defp pop_monitor_for(monitors, valid_time) do + case Enum.find(monitors, fn {_ref, vt} -> vt == valid_time end) do + nil -> + {monitors, nil} + + {ref, ^valid_time} -> + Process.demonitor(ref, [:flush]) + {Map.delete(monitors, ref), ref} + end + end + # ---------- Internal ---------- defp list_to_grid(rows) do diff --git a/lib/microwaveprop/workers/admin_task_worker.ex b/lib/microwaveprop/workers/admin_task_worker.ex index 40e1725e..e9cfcba2 100644 --- a/lib/microwaveprop/workers/admin_task_worker.ex +++ b/lib/microwaveprop/workers/admin_task_worker.ex @@ -135,43 +135,8 @@ defmodule Microwaveprop.Workers.AdminTaskWorker do Logger.info("AdminTask: deriving fields for #{length(profiles)} native profiles") - count = - Enum.count(profiles, fn profile -> - p = %{ - heights_m: profile.heights_m, - temp_k: profile.temp_k, - spfh: profile.spfh, - pressure_pa: profile.pressure_pa, - u_wind_ms: profile.u_wind_ms, - v_wind_ms: profile.v_wind_ms, - tke_m2s2: profile.tke_m2s2, - level_count: profile.level_count - } - - inversion_fields = - case Inversion.find_inversion_top(p) do - {:ok, %{height_m: top_h, level_idx: top_idx, base_idx: base_idx}} -> - [ - inversion_top_m: top_h, - bulk_richardson: Inversion.bulk_richardson(p, base_idx, top_idx), - shear_at_top_ms: Inversion.shear_magnitude(p, base_idx, top_idx), - theta_e_jump_k: ThetaE.theta_e_jump(p, base_idx, top_idx) - ] - - :none -> - [inversion_top_m: nil, bulk_richardson: nil, shear_at_top_ms: nil, theta_e_jump_k: nil] - end - - duct_result = Duct.analyze(p) - fields = inversion_fields ++ [ducts: duct_result.ducts, best_duct_band_ghz: duct_result.best_duct_band_ghz] - - {1, _} = - HrrrNativeProfile - |> where([pr], pr.id == ^profile.id) - |> Repo.update_all(set: fields) - - true - end) + rows = Enum.map(profiles, &derive_native_row/1) + count = bulk_update_native_derivations(rows) Logger.info("AdminTask: derived fields for #{count} profiles") :ok @@ -231,4 +196,99 @@ defmodule Microwaveprop.Workers.AdminTaskWorker do Logger.info("AdminTask: backtest for #{feature_name} complete, wrote #{path}") :ok end + + # Bulk-update via one `UPDATE ... FROM unnest(...)` statement per batch: + # one Postgres round-trip and one fsync instead of N, which was the + # dominant cost for runs of 10k profiles against the remote DB host. + @bulk_update_chunk 2000 + + defp derive_native_row(profile) do + p = %{ + heights_m: profile.heights_m, + temp_k: profile.temp_k, + spfh: profile.spfh, + pressure_pa: profile.pressure_pa, + u_wind_ms: profile.u_wind_ms, + v_wind_ms: profile.v_wind_ms, + tke_m2s2: profile.tke_m2s2, + level_count: profile.level_count + } + + {top_m, bulk, shear, theta_jump} = + case Inversion.find_inversion_top(p) do + {:ok, %{height_m: top_h, level_idx: top_idx, base_idx: base_idx}} -> + { + top_h, + Inversion.bulk_richardson(p, base_idx, top_idx), + Inversion.shear_magnitude(p, base_idx, top_idx), + ThetaE.theta_e_jump(p, base_idx, top_idx) + } + + :none -> + {nil, nil, nil, nil} + end + + duct_result = Duct.analyze(p) + + %{ + id: profile.id, + inversion_top_m: top_m, + bulk_richardson: bulk, + shear_at_top_ms: shear, + theta_e_jump_k: theta_jump, + ducts: duct_result.ducts, + best_duct_band_ghz: duct_result.best_duct_band_ghz + } + end + + defp bulk_update_native_derivations([]), do: 0 + + defp bulk_update_native_derivations(rows) do + rows + |> Enum.chunk_every(@bulk_update_chunk) + |> Enum.reduce(0, fn chunk, acc -> acc + execute_bulk_update(chunk) end) + end + + defp execute_bulk_update(chunk) do + # Ecto loads binary_id as the dashed-string form; Postgrex expects raw + # 16-byte binaries for `uuid[]` params. Dump once per id before send. + ids = Enum.map(chunk, fn %{id: id} -> Ecto.UUID.dump!(id) end) + tops = Enum.map(chunk, & &1.inversion_top_m) + bulks = Enum.map(chunk, & &1.bulk_richardson) + shears = Enum.map(chunk, & &1.shear_at_top_ms) + theta_jumps = Enum.map(chunk, & &1.theta_e_jump_k) + # Ducts are sent as a `text[]` of encoded JSON; each element is cast + # to `jsonb` per-row via the v.ducts::jsonb reference in the UPDATE + # so Postgres parses the JSON payload on ingest rather than storing + # it as a literal string. + ducts_json = Enum.map(chunk, fn row -> Jason.encode!(row.ducts || []) end) + bands = Enum.map(chunk, & &1.best_duct_band_ghz) + + sql = """ + UPDATE hrrr_native_profiles p + SET inversion_top_m = v.inversion_top_m, + bulk_richardson = v.bulk_richardson, + shear_at_top_ms = v.shear_at_top_ms, + theta_e_jump_k = v.theta_e_jump_k, + ducts = v.ducts::jsonb, + best_duct_band_ghz = v.best_duct_band_ghz, + updated_at = now() + FROM ( + SELECT + unnest($1::uuid[]) AS id, + unnest($2::float8[]) AS inversion_top_m, + unnest($3::float8[]) AS bulk_richardson, + unnest($4::float8[]) AS shear_at_top_ms, + unnest($5::float8[]) AS theta_e_jump_k, + unnest($6::text[]) AS ducts, + unnest($7::float8[]) AS best_duct_band_ghz + ) v + WHERE p.id = v.id + """ + + %{num_rows: n} = + Repo.query!(sql, [ids, tops, bulks, shears, theta_jumps, ducts_json, bands]) + + n + end end diff --git a/lib/microwaveprop/workers/radar_frame_worker.ex b/lib/microwaveprop/workers/radar_frame_worker.ex index c20d430c..f8c23a39 100644 --- a/lib/microwaveprop/workers/radar_frame_worker.ex +++ b/lib/microwaveprop/workers/radar_frame_worker.ex @@ -68,18 +68,42 @@ defmodule Microwaveprop.Workers.RadarFrameWorker do case NexradClient.fetch_decoded_frame(rounded) do {:ok, pixels, width} -> process_frame(contacts, rounded, pixels, width) + :ok {:error, reason} -> - Enum.each(contacts, &mark_status(&1, :unavailable)) + if permanent_error?(reason) do + Enum.each(contacts, &mark_status(&1, :unavailable)) - Logger.info( - "RadarFrameWorker: no frame for #{DateTime.to_iso8601(rounded)} (#{length(contacts)} contacts): #{inspect(reason)}" - ) + Logger.info( + "RadarFrameWorker: no frame for #{DateTime.to_iso8601(rounded)} (#{length(contacts)} contacts): #{inspect(reason)}" + ) + + :ok + else + # Transient (5xx, timeout, connrefused) — let Oban retry so the + # contacts stay :queued and land on the next attempt rather than + # getting pinned :unavailable after a temporary NEXRAD outage. + Logger.warning( + "RadarFrameWorker: transient error for #{DateTime.to_iso8601(rounded)} (#{length(contacts)} contacts): #{inspect(reason)}" + ) + + {:error, reason} + end end - - :ok end + # Permanent: 4xx means the IEM archive does not have this frame (the + # archive genuinely has gaps), so the frame will never exist and + # retrying is pointless. Any other shape is treated as transient. + defp permanent_error?("NEXRAD n0q HTTP " <> code) do + case Integer.parse(code) do + {status, _} when status in 400..499 -> true + _ -> false + end + end + + defp permanent_error?(_), do: false + @doc """ Process a batch of contacts against an already-decoded frame. Exposed for tests and for any caller that has the pixel buffer in hand. diff --git a/test/microwaveprop/weather/grid_cache_test.exs b/test/microwaveprop/weather/grid_cache_test.exs index 31e45fe4..d99d952e 100644 --- a/test/microwaveprop/weather/grid_cache_test.exs +++ b/test/microwaveprop/weather/grid_cache_test.exs @@ -70,6 +70,57 @@ defmodule Microwaveprop.Weather.GridCacheTest do end end + describe "claim_fill/1 and release_fill/1" do + @valid_time ~U[2026-04-21 12:00:00Z] + + test "first claimer wins, subsequent claimers see :in_progress" do + assert GridCache.claim_fill(@valid_time) == true + assert GridCache.claim_fill(@valid_time) == false + end + + test "release_fill allows a re-claim" do + assert GridCache.claim_fill(@valid_time) == true + :ok = GridCache.release_fill(@valid_time) + assert GridCache.claim_fill(@valid_time) == true + end + + test "lock is auto-released when the claimer process crashes" do + parent = self() + + {:ok, claimer} = + Task.start(fn -> + true = GridCache.claim_fill(@valid_time) + send(parent, :claimed) + # Block until we're killed + Process.sleep(:infinity) + end) + + receive do + :claimed -> :ok + after + 1000 -> flunk("claimer never signaled") + end + + # Lock is held by the crashed process + assert GridCache.claim_fill(@valid_time) == false + + # Kill the claimer + ref = Process.monitor(claimer) + Process.exit(claimer, :kill) + + receive do + {:DOWN, ^ref, :process, ^claimer, _} -> :ok + after + 1000 -> flunk("claimer didn't die") + end + + # Flush the GenServer so the :DOWN handler runs before we re-probe ETS. + :ok = GridCache.sync() + + assert GridCache.claim_fill(@valid_time) == true + end + end + describe "latest_valid_time/0" do test "returns the most recent cached valid_time" do GridCache.put(~U[2026-04-12 10:00:00Z], []) diff --git a/test/microwaveprop/workers/admin_task_worker_test.exs b/test/microwaveprop/workers/admin_task_worker_test.exs index ed9fcb63..6b41c4bb 100644 --- a/test/microwaveprop/workers/admin_task_worker_test.exs +++ b/test/microwaveprop/workers/admin_task_worker_test.exs @@ -23,8 +23,6 @@ defmodule Microwaveprop.Workers.AdminTaskWorkerTest do * `recalibrate` — delegates to `Recalibrator.fit/1`. With no contacts in the sandbox it returns the insufficient-data shape. - * `scorer_diff` — documented no-op. - * unknown task — catch-all returns `{:error, _}`. The tests characterize current behavior so downstream refactors @@ -346,14 +344,6 @@ defmodule Microwaveprop.Workers.AdminTaskWorkerTest do end end - describe "perform/1 — task=scorer_diff" do - test "is a documented no-op" do - # The propagation_scores table was dropped; this branch only - # exists to drain in-flight Oban rows from the pre-cutover era. - assert :ok = AdminTaskWorker.perform(%Oban.Job{args: %{"task" => "scorer_diff"}}) - end - end - describe "perform/1 — unknown task" do test "returns {:error, reason} for an unrecognized task name" do assert {:error, "unknown task: frobnicate"} = @@ -374,7 +364,7 @@ defmodule Microwaveprop.Workers.AdminTaskWorkerTest do # Worker config is load-bearing: admin tasks are long-running and # must not run concurrently (unique) or retry on failure # (max_attempts=1). Pin it here so a config regression fails loudly. - changeset = AdminTaskWorker.new(%{"task" => "scorer_diff"}) + changeset = AdminTaskWorker.new(%{"task" => "recalibrate"}) assert Ecto.Changeset.get_field(changeset, :queue) == "admin" assert Ecto.Changeset.get_field(changeset, :max_attempts) == 1 diff --git a/test/microwaveprop/workers/radar_frame_worker_test.exs b/test/microwaveprop/workers/radar_frame_worker_test.exs index 313bf3ad..e1719a98 100644 --- a/test/microwaveprop/workers/radar_frame_worker_test.exs +++ b/test/microwaveprop/workers/radar_frame_worker_test.exs @@ -83,6 +83,44 @@ defmodule Microwaveprop.Workers.RadarFrameWorkerTest do assert %Contact{radar_status: :unavailable} = Repo.get!(Contact, c2.id) end + test "returns {:error, _} and leaves contacts :queued on a 5xx server error (so Oban retries)" do + c1 = insert_contact(%{station1: "S1"}) + c2 = insert_contact(%{station1: "S2"}) + + Req.Test.stub(NexradClient, fn conn -> + Plug.Conn.send_resp(conn, 503, "service unavailable") + end) + + assert {:error, _reason} = + perform_job(RadarFrameWorker, %{ + "frame_ts" => "2024-09-15T18:30:00Z", + "contact_ids" => [c1.id, c2.id] + }) + + # Neither contact should be pinned :unavailable — a retry might succeed. + assert %Contact{radar_status: status1} = Repo.get!(Contact, c1.id) + assert %Contact{radar_status: status2} = Repo.get!(Contact, c2.id) + refute status1 == :unavailable + refute status2 == :unavailable + end + + test "returns {:error, _} on a transport error without pinning contacts :unavailable" do + c1 = insert_contact(%{station1: "T1"}) + + Req.Test.stub(NexradClient, fn conn -> + Req.Test.transport_error(conn, :timeout) + end) + + assert {:error, _reason} = + perform_job(RadarFrameWorker, %{ + "frame_ts" => "2024-09-15T18:30:00Z", + "contact_ids" => [c1.id] + }) + + assert %Contact{radar_status: status} = Repo.get!(Contact, c1.id) + refute status == :unavailable + end + test "skips contacts already at a terminal status (idempotent replay)" do c_done = insert_contact(%{station1: "DONE"}) {:ok, _} = c_done |> Ecto.Changeset.change(%{radar_status: :complete}) |> Repo.update()