prop/lib/microwaveprop/workers/backfill_enqueue_worker.ex
Graham McIntire bbf3e67d85 Add ERA5 to backfill dashboard and enrichment pipeline
- ERA5 checkbox on /backfill page alongside HRRR/weather/terrain/IEMRE
- BackfillEnqueueWorker handles era5 type (targets hrrr_status=unavailable)
- ContactWeatherEnqueueWorker.build_era5_jobs/1 creates ERA5 fetch jobs
  for path points without existing ERA5 profiles
- Progress bar shows ERA5 candidate count (HRRR-unavailable contacts)
- Status panel shows ERA5 candidates and fetched profile count
2026-04-07 12:24:21 -05:00

100 lines
2.9 KiB
Elixir

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.Workers.ContactWeatherEnqueueWorker
require Logger
@enrichable [:pending, :queued, :failed]
@valid_types ~w(hrrr weather terrain iemre era5)
@impl Oban.Worker
def perform(%Oban.Job{args: %{"limit" => limit} = args}) do
types = parse_types(args)
reconciled = reconcile_stale_queued(types)
if reconciled > 0 do
Logger.info("BackfillEnqueue: reconciled #{reconciled} stale queued contacts")
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))
|> 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 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, :era5]
defp status_priority_order(types) do
# Prioritize pending/failed over queued so real work happens first
field = :"#{hd(types)}_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
# 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
:era5, acc ->
# ERA5 targets contacts where HRRR is unavailable (pre-2014, missing from archive)
dynamic([c], ^acc or c.hrrr_status == :unavailable)
type, acc ->
field = :"#{type}_status"
dynamic([c], ^acc or field(c, ^field) in ^@enrichable)
end)
end
end