prop/lib/microwaveprop/propagation/grid_task_enqueuer.ex
Graham McIntire 020f05f8eb
feat(hrdps): grid_tasks source column, Canadian mask, HrdpsGridWorker
Stage 3 (Elixir foundation) of the HRDPS plan. Lands the scaffolding
that lets the Rust prop-grid-rs HRDPS branch be wired up cleanly when
it ships, without further DB or Elixir-side changes.

What ships:

- `grid_tasks.source` column (default 'hrrr', backfill is a no-op).
  Replaces the (run_time, forecast_hour, kind) unique index with a
  4-column key (..., source) so HRRR and HRDPS analysis/forecast rows
  for the same cycle coexist.
- `Grid.hrdps_only_points/0` — Canadian cells inside HRDPS bbox
  (49-60°N, -141 to -52°W) but outside HRRR's CONUS bbox. Disjoint
  from `conus_points/0` by construction so the two grids never
  double-write the same (lat, lon).
- `Grid.point_source/1` — `:hrrr | :hrdps | :neither` classification
  for a (lat, lon) pair. Routes per-QSO enrichment to the right NWP
  source.
- `GridTaskEnqueuer.seed_with_analysis/2` and `seed/2` accept a
  `:source` option; defaults to "hrrr" so existing callers are
  unchanged.
- `HrdpsGridWorker` cycle seeder mirroring PropagationGridWorker's
  shape but for HRDPS's 4×/day cadence. `:hrdps` queue added to base
  Oban config.

What's deferred:

- The Rust `prop-grid-rs` worker doesn't yet branch on
  `task.source`. Until that ships, `HrdpsGridWorker` stays out of
  runtime.exs cron — it would just produce rows that the Rust worker
  mishandles as HRRR. Module doc explicitly notes the dormant state.
- Map overlay extension and FreshnessMonitor HRDPS staleness deferred
  to follow-up sessions once real data flows.

Activation path documented in the updated plan file.

Coverage cap recorded as Stage 8 decision: 60°N for v1 (SRTM stops
there). Arctic CDEM deferred to a follow-up plan.
2026-04-29 17:13:00 -05:00

186 lines
6.2 KiB
Elixir

defmodule Microwaveprop.Propagation.GridTaskEnqueuer do
@moduledoc """
Seeds `grid_tasks` rows for the Rust `prop-grid-rs` worker.
Called from `PropagationGridWorker.seed_chain/0`. Rust claims
kind='forecast' and kind='analysis' rows, with analysis lanes
taking priority (see `claim_next_analysis` in the Rust db module).
Inserts are idempotent via the `(run_time, forecast_hour, kind)`
unique index — re-seeding the same cycle is a no-op.
"""
import Ecto.Query
alias Microwaveprop.Repo
require Logger
@max_forecast_hour 18
# A `running` row older than this is an orphan: the Rust worker that
# claimed it was SIGKILLed / OOMed / evicted before it could emit a
# terminal status. `claim_next` uses `FOR UPDATE SKIP LOCKED` so the
# row stays `running` indefinitely until something reclaims it.
# 15 min comfortably exceeds the ~3-4 min wall time of a healthy
# forecast-hour chain step; any row past that is stuck.
@stale_running_cutoff_seconds 15 * 60
# After this many claim/reclaim cycles, a row is permanently flipped
# to `failed` so the status-page spinner stops and the next hourly
# reseed can replace it. Rust's `claim_next` increments `attempt` on
# each pickup, so this is a total-attempts ceiling.
@max_reclaim_attempts 5
@doc """
Seed one kind='analysis' row (f00) plus 18 kind='forecast' rows
(f01..f18) for `run_time`. This is the Phase 3 Stream A cutover
shape: Rust owns the entire chain end-to-end.
## Options
* `:source` — `"hrrr"` (default) or `"hrdps"`. The Rust worker branches
on this column to pick the right URL pattern + decoder. HRDPS rows
are queued only once the Rust HRDPS branch ships (deferred plan
stage 4); until then no caller passes `:source`.
"""
@spec seed_with_analysis(DateTime.t(), keyword()) :: {:ok, non_neg_integer()} | {:error, term()}
def seed_with_analysis(%DateTime{} = run_time, opts \\ []) do
source = Keyword.get(opts, :source, "hrrr")
run_time = DateTime.truncate(run_time, :second)
now = DateTime.truncate(DateTime.utc_now(), :microsecond)
analysis_row = %{
id: Ecto.UUID.bingenerate(),
run_time: run_time,
forecast_hour: 0,
valid_time: run_time,
status: "queued",
attempt: 0,
kind: "analysis",
source: source,
claimed_at: nil,
completed_at: nil,
error: nil,
inserted_at: now,
updated_at: now
}
forecast_rows =
for fh <- 1..@max_forecast_hour do
%{
id: Ecto.UUID.bingenerate(),
run_time: run_time,
forecast_hour: fh,
valid_time: DateTime.add(run_time, fh * 3600, :second),
status: "queued",
attempt: 0,
kind: "forecast",
source: source,
claimed_at: nil,
completed_at: nil,
error: nil,
inserted_at: now,
updated_at: now
}
end
do_insert([analysis_row | forecast_rows], run_time, source)
end
@doc """
Legacy: seed only kind='forecast' rows (f01..f18). Retained for
tooling that needs to re-seed the forecast lane without touching
the analysis row. `seed_with_analysis/2` is the production path.
"""
@spec seed(DateTime.t(), keyword()) :: {:ok, non_neg_integer()} | {:error, term()}
def seed(%DateTime{} = run_time, opts \\ []) do
source = Keyword.get(opts, :source, "hrrr")
run_time = DateTime.truncate(run_time, :second)
now = DateTime.truncate(DateTime.utc_now(), :microsecond)
rows =
for fh <- 1..@max_forecast_hour do
%{
id: Ecto.UUID.bingenerate(),
run_time: run_time,
forecast_hour: fh,
valid_time: DateTime.add(run_time, fh * 3600, :second),
status: "queued",
attempt: 0,
kind: "forecast",
source: source,
claimed_at: nil,
completed_at: nil,
error: nil,
inserted_at: now,
updated_at: now
}
end
do_insert(rows, run_time, source)
end
@doc """
Reclaim orphan `running` rows — Rust workers that died mid-claim
(SIGKILL, OOM, node drain) leave `status='running'` sticky, which
makes `claim_next` skip the row forever and the status panel shows
a permanent-looking spinner. Flip rows whose `claimed_at` is older
than `cutoff_seconds` back to `queued`. Rows that have already
burned through `@max_reclaim_attempts` become `failed` so the next
hourly seed replaces them instead of looping.
Called from `PropagationGridWorker.seed_chain/0` at :05 each hour.
"""
@spec reclaim_stale_running(pos_integer()) :: %{requeued: non_neg_integer(), failed: non_neg_integer()}
def reclaim_stale_running(cutoff_seconds \\ @stale_running_cutoff_seconds) do
cutoff = DateTime.utc_now() |> DateTime.add(-cutoff_seconds, :second) |> DateTime.truncate(:second)
now = DateTime.truncate(DateTime.utc_now(), :second)
{failed_n, _} =
"grid_tasks"
|> where(
[t],
t.status == "running" and t.claimed_at < ^cutoff and t.attempt >= ^@max_reclaim_attempts
)
|> Repo.update_all(
set: [
status: "failed",
error: "reclaim-orphan: exceeded max_reclaim_attempts",
claimed_at: nil,
updated_at: now
]
)
{requeued_n, _} =
"grid_tasks"
|> where([t], t.status == "running" and t.claimed_at < ^cutoff)
|> Repo.update_all(set: [status: "queued", claimed_at: nil, updated_at: now])
if failed_n > 0 or requeued_n > 0 do
Logger.info("GridTaskEnqueuer: reclaimed stale-running grid_tasks (requeued=#{requeued_n}, failed=#{failed_n})")
end
%{requeued: requeued_n, failed: failed_n}
end
defp do_insert(rows, run_time, source) do
{count, _} =
Repo.insert_all("grid_tasks", rows,
on_conflict: :nothing,
conflict_target: [:run_time, :forecast_hour, :kind, :source]
)
if count > 0 do
Logger.info("GridTaskEnqueuer: seeded #{count} grid_tasks (source=#{source}) for run_time=#{run_time}")
else
Logger.info("GridTaskEnqueuer: run_time=#{run_time} source=#{source} already seeded")
end
{:ok, count}
rescue
e ->
Logger.error("GridTaskEnqueuer: seed failed: #{inspect(e)}")
{:error, e}
end
end