FreshnessMonitor's comment claimed 'Oban unique constraint on
PropagationGridWorker prevents duplicates' — but the worker declared
no unique: option. During a long outage the monitor's 5-min stale-check
tick would pile up identical jobs (a 2h outage = 24 stacked jobs, each
doing the full f00-f18 HRRR chain).
Put the dedup on the worker (vs. the monitor's insert call) so anything
enqueuing it benefits — FreshnessMonitor, the hourly cron, or a manual
mix propagation_grid run.
unique: [period: 3600, states: [:available, :scheduled, :executing,
:retryable]] aligns with the hourly cron cadence. The seed job args
(%{}) collapse with themselves while chain-step jobs keep distinct
run_time/forecast_hour and remain enqueuable.
64 lines
1.9 KiB
Elixir
64 lines
1.9 KiB
Elixir
defmodule Microwaveprop.Propagation.FreshnessMonitor do
|
|
@moduledoc """
|
|
Monitors propagation score freshness and enqueues grid worker jobs
|
|
when data is stale. Checks every 5 minutes. Covers missed cron ticks,
|
|
slow deploys, worker crashes, and any other gap in hourly scoring.
|
|
"""
|
|
|
|
use GenServer
|
|
|
|
alias Microwaveprop.Propagation
|
|
alias Microwaveprop.Workers.PropagationGridWorker
|
|
|
|
require Logger
|
|
|
|
@check_interval to_timeout(minute: 5)
|
|
@stale_threshold_minutes 120
|
|
|
|
@spec start_link(keyword()) :: GenServer.on_start() | :ignore
|
|
def start_link(_opts) do
|
|
if Application.get_env(:microwaveprop, :start_freshness_monitor, true) do
|
|
GenServer.start_link(__MODULE__, :ok, name: __MODULE__)
|
|
else
|
|
:ignore
|
|
end
|
|
end
|
|
|
|
@impl true
|
|
def init(:ok) do
|
|
send(self(), :check)
|
|
{:ok, %{}}
|
|
end
|
|
|
|
@impl true
|
|
def handle_info(:check, state) do
|
|
check_freshness()
|
|
Process.send_after(self(), :check, @check_interval)
|
|
{:noreply, state}
|
|
end
|
|
|
|
defp check_freshness do
|
|
case Propagation.latest_valid_time() do
|
|
nil ->
|
|
Logger.info("FreshnessMonitor: no scores found, enqueuing grid worker")
|
|
enqueue_if_not_queued()
|
|
|
|
latest ->
|
|
age = DateTime.diff(DateTime.utc_now(), latest, :minute)
|
|
|
|
if age > @stale_threshold_minutes do
|
|
Logger.info("FreshnessMonitor: scores are #{age}m old, enqueuing grid worker")
|
|
enqueue_if_not_queued()
|
|
end
|
|
end
|
|
end
|
|
|
|
defp enqueue_if_not_queued do
|
|
# PropagationGridWorker declares `unique:` over a 1-hour window on
|
|
# the seed args (`%{}`), so repeated 5-minute ticks during a long
|
|
# outage collapse into a single seed job — not one stacked chain
|
|
# per tick. `Oban.insert` returns `{:ok, job}` either way; the
|
|
# `conflict?` field on the returned job distinguishes the two.
|
|
Oban.insert(PropagationGridWorker.new(%{}))
|
|
end
|
|
end
|