prop/lib/microwaveprop/workers/backfill_enqueue_worker.ex
Graham McIntire 622edee180
feat(propagation): per-contact mechanism classification
Classifies every contact's likely non-LOS propagation mechanism and
persists the result on contacts.propagation_mechanism. Mechanism is
determined in priority order:

  1. user_declared_prop_mode (ADIF PROP_MODE from the operator log)
  2. EME — moon-ephemeris check, ≥2m band, >1800 km path
  3. aurora — Kp≥5 + high-lat path, 50-432 MHz
  4. sporadic-E — foEs × 5 ≥ band_mhz, 400-2500 km path
  5. meteor_scatter — ±3 days of a shower peak, VHF/UHF
  6. rain_scatter — common-volume radar heavy rain, 5-11 GHz ≤800 km
  7. tropo_duct — HRRR native_best_duct ≥ band or ducting_detected
  8. line_of_sight — ≤50 km path
  9. troposcatter — default

Persisted via MechanismClassifyWorker (queue: :mechanism, unique on
contact_id). Submit-time enqueue path includes :mechanism by default;
BackfillEnqueueWorker cron now handles :mechanism alongside existing
types so prod continuously backfills any contact with
propagation_mechanism_status in (:pending, :queued, :failed). Also
added :radar to the cron's type list so common-volume radar backfill
runs automatically rather than only via `mix radar_backfill`.

New modules:
- Microwaveprop.Propagation.MoonEphemeris — Meeus low-precision moon
  position, accuracy ±1° — enough for the mutual-visibility EME test
- Microwaveprop.Propagation.MechanismClassifier — plug-in priority
  chain over the evidence map
- Microwaveprop.Workers.MechanismClassifyWorker — assembles inputs
  from HRRR / native profiles / common-volume radar / solar_indices /
  ionosonde + calls the classifier

ADIF importer now reads PROP_MODE into user_declared_prop_mode so
operator-tagged mechanisms (EME/ES/MS/RS/AS/AUR) become ground truth.
2026-04-18 10:42:08 -05:00

175 lines
5.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.Weather.NarrClient
alias Microwaveprop.Workers.ContactWeatherEnqueueWorker
require Logger
@enrichable [:pending, :queued, :failed]
@valid_types ~w(hrrr weather terrain iemre narr radar mechanism)
# Virtual types have no <type>_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 <type>_status field.
@virtual_types [:narr]
# `radar` and `mechanism` use differently-named *_status columns than
# `<type>_status`, so we remap them when building queries and update_alls.
@status_column_overrides %{
mechanism: :propagation_mechanism_status
}
# 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, :radar, :mechanism]
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 = status_column(type)
[
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
defp status_column(type), do: Map.get(@status_column_overrides, type, :"#{type}_status")
# 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 (<field>, 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 = status_column(type)
{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 = status_column(type)
dynamic([c], ^acc or field(c, ^field) in ^@enrichable)
end)
end
end