prop/lib/microwaveprop/radio.ex
Graham McIntire 55a51f69ab
Replace boolean enrichment flags with enum status fields
- New fields: hrrr_status, weather_status, terrain_status, iemre_status
  with values: pending, queued, processing, complete, failed, unavailable
- EnrichmentStatus module defines valid state transitions
- Migration backfills: queued=true → complete, false → pending
- Workers set status on transitions (queued on enqueue, complete on finish)
- Show page reads status from contact struct, not computed dynamically
- Backfill dashboard queries status fields for accurate progress counts
- Remove cron-scheduled ContactWeatherEnqueueWorker (enrichment only on
  submission, page view, or manual backfill)
- Partial indexes on status != 'complete' for fast unprocessed lookups
2026-04-02 12:07:51 -05:00

241 lines
7.5 KiB
Elixir

defmodule Microwaveprop.Radio do
@moduledoc false
import Ecto.Query
alias Microwaveprop.Radio.Contact
alias Microwaveprop.Radio.Maidenhead
alias Microwaveprop.Repo
@per_page 20
@sortable_fields ~w(station1 station2 band mode distance_km qso_timestamp inserted_at)a
def list_contacts(opts \\ []) do
page = max(Keyword.get(opts, :page, 1), 1)
offset = (page - 1) * @per_page
search = Keyword.get(opts, :search)
{sort_field, sort_dir} = sort_opts(opts)
base_query = maybe_search(Contact, search)
total_entries = Repo.aggregate(base_query, :count)
total_pages = max(ceil(total_entries / @per_page), 1)
entries =
base_query
|> order_by([q], [{^sort_dir, field(q, ^sort_field)}])
|> limit(^@per_page)
|> offset(^offset)
|> Repo.all()
%{
entries: entries,
grouped_entries: group_reciprocals(entries),
page: page,
total_pages: total_pages,
total_entries: total_entries
}
end
@doc """
Groups reciprocal contacts (same station pair swapped, same band, same timestamp).
Returns a list of `{primary, reciprocals}` tuples where reciprocals is
a (possibly empty) list of matching contacts.
"""
def group_reciprocals(contacts) do
{groups, _seen_ids} =
Enum.reduce(contacts, {[], MapSet.new()}, fn contact, {groups, seen} ->
if MapSet.member?(seen, contact.id) do
{groups, seen}
else
reciprocals =
Enum.filter(contacts, fn other ->
other.id != contact.id and
not MapSet.member?(seen, other.id) and
other.band == contact.band and
reciprocal_pair?(contact, other)
end)
new_seen = Enum.reduce(reciprocals, MapSet.put(seen, contact.id), &MapSet.put(&2, &1.id))
{[{contact, reciprocals} | groups], new_seen}
end
end)
Enum.reverse(groups)
end
defp reciprocal_pair?(a, b) do
same_stations =
(a.station1 == b.station2 and a.station2 == b.station1) or
(a.station1 == b.station1 and a.station2 == b.station2)
same_stations and same_hour?(a.qso_timestamp, b.qso_timestamp)
end
defp same_hour?(t1, t2) do
t1 |> NaiveDateTime.truncate(:second) |> Map.take([:year, :month, :day, :hour]) ==
t2 |> NaiveDateTime.truncate(:second) |> Map.take([:year, :month, :day, :hour])
end
defp maybe_search(query, nil), do: query
defp maybe_search(query, ""), do: query
defp maybe_search(query, search) do
pattern = "%" <> String.upcase(search) <> "%"
where(query, [q], ilike(q.station1, ^pattern) or ilike(q.station2, ^pattern))
end
defp sort_opts(opts) do
sort_by = Keyword.get(opts, :sort_by, :qso_timestamp)
sort_order = Keyword.get(opts, :sort_order, :desc)
field = if sort_by in @sortable_fields, do: sort_by, else: :qso_timestamp
dir = if sort_order in [:asc, :desc], do: sort_order, else: :desc
{field, dir}
end
# ── Enrichment status queries ────────────────────────────────────
@enrichable_statuses [:pending, :failed]
def contacts_needing_enrichment(field, limit \\ 500, extra_filter \\ nil) do
query =
Contact
|> where([q], field(q, ^field) in ^@enrichable_statuses and not is_nil(q.pos1))
|> order_by([q], asc: q.qso_timestamp)
|> limit(^limit)
query = if extra_filter, do: extra_filter.(query), else: query
Repo.all(query)
end
def unprocessed_contacts(limit \\ 500),
do: contacts_needing_enrichment(:weather_status, limit)
def unprocessed_hrrr_contacts(limit \\ 500),
do: contacts_needing_enrichment(:hrrr_status, limit)
def unprocessed_terrain_contacts(limit \\ 500),
do: contacts_needing_enrichment(:terrain_status, limit, &where(&1, [q], not is_nil(q.pos2)))
def unprocessed_iemre_contacts(limit \\ 500),
do: contacts_needing_enrichment(:iemre_status, limit)
def set_enrichment_status!(ids, field, status) do
Contact
|> where([q], q.id in ^ids)
|> Repo.update_all(set: [{field, status}])
end
def mark_weather_queued!(ids), do: set_enrichment_status!(ids, :weather_status, :queued)
def mark_hrrr_queued!(ids), do: set_enrichment_status!(ids, :hrrr_status, :queued)
def mark_terrain_queued!(ids), do: set_enrichment_status!(ids, :terrain_status, :queued)
def mark_iemre_queued!(ids), do: set_enrichment_status!(ids, :iemre_status, :queued)
@doc """
Returns a list of {lat, lon} points along the contact path: pos1, midpoint, pos2.
Returns [pos1] if pos2 is nil, or [] if pos1 is nil.
"""
def contact_path_points(%{pos1: nil}), do: []
def contact_path_points(%{pos1: pos1, pos2: nil}) do
lat = pos1["lat"]
lon = pos1["lon"] || pos1["lng"]
if lat && lon, do: [{lat, lon}], else: []
end
def contact_path_points(%{pos1: pos1, pos2: pos2}) do
lat1 = pos1["lat"]
lon1 = pos1["lon"] || pos1["lng"]
lat2 = pos2["lat"]
lon2 = pos2["lon"] || pos2["lng"]
if lat1 && lon1 && lat2 && lon2 do
mid_lat = (lat1 + lat2) / 2
mid_lon = (lon1 + lon2) / 2
[{lat1, lon1}, {mid_lat, mid_lon}, {lat2, lon2}]
else
[{lat1, lon1}]
end
end
@earth_radius_km 6371.0
def haversine_km(lat1, lon1, lat2, lon2) do
dlat = deg_to_rad(lat2 - lat1)
dlon = deg_to_rad(lon2 - lon1)
rlat1 = deg_to_rad(lat1)
rlat2 = deg_to_rad(lat2)
a =
:math.sin(dlat / 2) ** 2 +
:math.cos(rlat1) * :math.cos(rlat2) * :math.sin(dlon / 2) ** 2
2 * @earth_radius_km * :math.asin(:math.sqrt(a))
end
def backfill_distances(contacts) do
Enum.each(contacts, fn contact ->
if is_nil(contact.distance_km) && contact.pos1 && contact.pos2 do
lat1 = contact.pos1["lat"]
lon1 = contact.pos1["lon"] || contact.pos1["lng"]
lat2 = contact.pos2["lat"]
lon2 = contact.pos2["lon"] || contact.pos2["lng"]
if lat1 && lon1 && lat2 && lon2 do
dist = lat1 |> haversine_km(lon1, lat2, lon2) |> round() |> Decimal.new()
Contact
|> where([q], q.id == ^contact.id)
|> Repo.update_all(set: [distance_km: dist])
end
end
end)
end
defp deg_to_rad(deg), do: deg * :math.pi() / 180
def get_contact!(id) do
Repo.get!(Contact, id)
end
def toggle_flagged_invalid!(contact) do
contact
|> Ecto.Changeset.change(flagged_invalid: !contact.flagged_invalid)
|> Repo.update!()
end
def change_contact(contact, attrs \\ %{}) do
Contact.submission_changeset(contact, attrs)
end
def create_contact(attrs) do
changeset = Contact.submission_changeset(%Contact{}, attrs)
if changeset.valid? do
grid1 = Ecto.Changeset.get_change(changeset, :grid1)
grid2 = Ecto.Changeset.get_change(changeset, :grid2)
case {Maidenhead.to_latlon(grid1), Maidenhead.to_latlon(grid2)} do
{{:ok, {lat1, lon1}}, {:ok, {lat2, lon2}}} ->
distance = lat1 |> haversine_km(lon1, lat2, lon2) |> round() |> Decimal.new()
changeset
|> Ecto.Changeset.put_change(:pos1, %{"lat" => lat1, "lon" => lon1})
|> Ecto.Changeset.put_change(:pos2, %{"lat" => lat2, "lon" => lon2})
|> Ecto.Changeset.put_change(:distance_km, distance)
|> Ecto.Changeset.put_change(:user_submitted, true)
|> Repo.insert()
_ ->
changeset = Ecto.Changeset.add_error(changeset, :grid1, "could not resolve grid to coordinates")
{:error, %{changeset | action: :insert}}
end
else
{:error, %{changeset | action: :insert}}
end
end
end