- CacheTest: fix sweep timing race by using negative TTL (-1) instead of positive TTL (1) for already-expired entries - ScoreCache: replace ETS match-spec DateTime comparisons with :ets.foldl + DateTime.compare — DateTime structs are maps in Elixir >= 1.15 and ETS can't compare maps with :< / :> guards - Accounts: drop unsupported returning: true on delete_all, return [] for expired tokens list - Backtest: catch ArgumentError from String.to_existing_atom for unknown feature names, preserving the helpful Mix.Error - ContactLive IndexTest: invalidate monthly_bars cache before assertion so test data is visible - RoverLocationsLive MapTest: invalidate cached points before assertion - StatusLiveTest: add DB cleanup in setup to reduce test interference; parameterize NARR candidate coordinates
264 lines
8.1 KiB
Elixir
264 lines
8.1 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
|
|
|
|
# HRRR forecast hours: f01-f18 hourly, then f21-f48 3-hourly.
|
|
# Rust processes whatever rows are seeded — no code changes needed there.
|
|
@forecast_hours [
|
|
1,
|
|
2,
|
|
3,
|
|
4,
|
|
5,
|
|
6,
|
|
7,
|
|
8,
|
|
9,
|
|
10,
|
|
11,
|
|
12,
|
|
13,
|
|
14,
|
|
15,
|
|
16,
|
|
17,
|
|
18,
|
|
21,
|
|
24,
|
|
27,
|
|
30,
|
|
33,
|
|
36,
|
|
39,
|
|
42,
|
|
45,
|
|
48
|
|
]
|
|
|
|
# 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 kind='forecast' rows
|
|
(f01..f18 hourly, f21..f48 3-hourly) 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 <- @forecast_hours 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. 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 <- @forecast_hours 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
|
|
|
|
@doc """
|
|
Seed exactly one kind='forecast' row whose valid_time is closest to
|
|
`now` for the given `run_time`. Used by `HrdpsGridWorker` because the
|
|
/weather map only consumes the current-hour HRDPS cell — the rest of
|
|
the f01..f48 chain isn't worth the wgrib2 cost (each HRDPS chain step
|
|
is ~20 min wall time at the production point density).
|
|
|
|
Forecast hour is clamped to the available forecast hour list. When
|
|
`now` lands before `run_time + 1h`, we still seed `fh=1`; when it's
|
|
past the max forecast hour, we cap there. The Rust worker rejects
|
|
`forecast_hour=0` (`F00Reserved`) so we never seed the analysis hour
|
|
through this path.
|
|
"""
|
|
@spec seed_current_hour(DateTime.t(), keyword()) :: {:ok, non_neg_integer()} | {:error, term()}
|
|
def seed_current_hour(%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)
|
|
target_valid = now |> DateTime.truncate(:second) |> Map.put(:minute, 0) |> Map.put(:second, 0)
|
|
|
|
fh =
|
|
target_valid
|
|
|> DateTime.diff(run_time, :second)
|
|
|> div(3600)
|
|
|> max(1)
|
|
|> min(Enum.max(@forecast_hours))
|
|
|
|
row = %{
|
|
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
|
|
}
|
|
|
|
do_insert([row], run_time, source)
|
|
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
|