defmodule Microwaveprop.Workers.BackfillEnqueueWorker do @moduledoc """ Runs backfill enrichment enqueue as an Oban job so the backfill dashboard returns immediately instead of blocking on the enqueue loop. """ use Oban.Worker, queue: :backfill_enqueue, max_attempts: 1 import Ecto.Query alias Microwaveprop.Radio.Contact alias Microwaveprop.Repo alias Microwaveprop.Weather.NarrClient alias Microwaveprop.Workers.ContactWeatherEnqueueWorker require Logger @enrichable [:pending, :queued, :failed] @valid_types ~w(hrrr weather terrain iemre narr) # Virtual types have no _status column on contacts; they're synthesized # from other status fields (e.g. :narr candidates are `hrrr_status = :unavailable`). # Skip these in any reconcile / ordering pass that reads a _status field. @virtual_types [:narr] # A :queued contact whose updated_at is older than this is one the cron has # already tried ~144 times over 3 days without the underlying worker landing # data. At that point the data is effectively unreachable (dead ASOS station, # pre-2014 HRRR, out-of-grid IEMRE) and the cron should stop re-enqueueing. # set_enrichment_status!/3 is a no-op when the status is unchanged, so a # stuck :queued row retains a stale updated_at across cron cycles — which is # exactly what we need as the "given up" signal. @stale_queued_cutoff_seconds 3 * 86_400 @impl Oban.Worker def perform(%Oban.Job{args: args}) do types = parse_types(args) limit = Map.get(args, "limit") reconciled = reconcile_stale_queued(types) if reconciled > 0 do Logger.info("BackfillEnqueue: reconciled #{reconciled} stale queued contacts") end stale_unavail = reconcile_stale_queued_to_unavailable(types) if stale_unavail > 0 do Logger.info("BackfillEnqueue: marked #{stale_unavail} stuck-queued contacts as unavailable") end contacts = Contact |> where([c], not is_nil(c.pos1) or not is_nil(c.grid1)) |> where(^type_filter(types)) |> order_by(^status_priority_order(types)) |> maybe_limit(limit) |> Repo.all() count = length(contacts) Logger.info("BackfillEnqueue: processing #{count} contacts (types: #{inspect(types)})") Enum.each(contacts, &ContactWeatherEnqueueWorker.enqueue_for_contact(&1, types)) Phoenix.PubSub.broadcast( Microwaveprop.PubSub, "backfill:enqueue_complete", {:enqueue_complete, count} ) :ok end defp maybe_limit(query, nil), do: query defp maybe_limit(query, n) when is_integer(n), do: limit(query, ^n) defp parse_types(%{"types" => types}) when is_list(types) do types |> Enum.filter(&(&1 in @valid_types)) |> Enum.map(&String.to_existing_atom/1) end defp parse_types(_), do: [:hrrr, :weather, :terrain, :iemre, :narr] defp status_priority_order(types) do # Prioritize pending/failed over queued so real work happens first. # Virtual types (e.g. :narr) have no *_status column — fall back to # time-only ordering when every requested type is virtual. case Enum.find(types, &(&1 not in @virtual_types)) do nil -> [desc: dynamic([c], c.qso_timestamp)] type -> field = :"#{type}_status" [ asc: dynamic( [c], fragment("CASE ? WHEN 'pending' THEN 0 WHEN 'failed' THEN 1 ELSE 2 END", field(c, ^field)) ), desc: dynamic([c], c.qso_timestamp) ] end end # Flip contacts that have been stuck in :queued for longer than # @stale_queued_cutoff_seconds to :unavailable. One Repo.update_all per type, # scoped to only the types passed in args so a partial cron (e.g. hrrr-only) # doesn't touch other statuses. Virtual types (e.g. :narr) have no *_status # column and are skipped. Fast: hits the (, updated_at) combo with a # simple range scan. defp reconcile_stale_queued_to_unavailable(types) do cutoff = DateTime.utc_now() |> DateTime.add(-@stale_queued_cutoff_seconds, :second) |> DateTime.truncate(:second) types |> Enum.reject(&(&1 in @virtual_types)) |> Enum.reduce(0, fn type, acc -> field = :"#{type}_status" {n, _} = Contact |> where([c], field(c, ^field) == :queued and c.updated_at < ^cutoff) |> Repo.update_all(set: [{field, :unavailable}]) acc + n end) end # Fast SQL reconciliation for contacts stuck in "queued" where data already exists. # Terrain profiles have a direct contact_id FK so we can reconcile in bulk. defp reconcile_stale_queued(types) do terrain_count = if :terrain in types do {count, _} = Repo.update_all( from(c in Contact, where: c.terrain_status == :queued, where: c.id in subquery(from(tp in "terrain_profiles", select: tp.contact_id)) ), set: [terrain_status: :complete] ) count else 0 end terrain_count end defp type_filter(types) do Enum.reduce(types, dynamic(false), fn :narr, acc -> # NARR targets pre-2014 contacts where HRRR is unavailable. The NCEI # archive ends 2014-10-02 — post-cutoff fetches would 404, so scope # the candidate set here to keep the cron from scanning contacts we # can't serve. coverage_end = NarrClient.coverage_end() dynamic([c], ^acc or (c.hrrr_status == :unavailable and c.qso_timestamp < ^coverage_end)) type, acc -> field = :"#{type}_status" dynamic([c], ^acc or field(c, ^field) in ^@enrichable) end) end end