Hourly ContactPositionBackfillWorker fills null pos1/pos2
Safety net for rows that land via direct DB writes (manual fixes, bulk imports) and bypass Radio.create_contact's grid-resolution requirement. Runs every hour on the :admin queue with a 500-row per-invocation cap. Fills whichever side resolves — if one grid is invalid, the other still gets populated; distance_km is only written when both positions end up set.
This commit is contained in:
parent
1a8e55dfb0
commit
01d554f2e7
3 changed files with 249 additions and 0 deletions
|
|
@ -223,6 +223,10 @@ if config_env() == :prod do
|
|||
# NarrFetchWorker against NCEI. See narr_jobs_for_contact/1.
|
||||
{"*/30 * * * *", Microwaveprop.Workers.BackfillEnqueueWorker,
|
||||
args: %{"limit" => 500, "types" => ["hrrr", "weather", "terrain", "iemre", "narr"]}},
|
||||
# Hourly safety net for pos1/pos2/distance_km. Normally every contact
|
||||
# gets positions at insert time via Radio.resolve_grids_and_insert/1,
|
||||
# but direct DB writes (manual fixes, bulk imports) can bypass that.
|
||||
{"0 * * * *", Microwaveprop.Workers.ContactPositionBackfillWorker, args: %{"limit" => 500}},
|
||||
# UWYO publishes the 00Z/12Z Canadian radiosondes ~90 minutes after launch
|
||||
{"30 1 * * *", CanadianSoundingFetchWorker},
|
||||
{"30 13 * * *", CanadianSoundingFetchWorker},
|
||||
|
|
|
|||
105
lib/microwaveprop/workers/contact_position_backfill_worker.ex
Normal file
105
lib/microwaveprop/workers/contact_position_backfill_worker.ex
Normal file
|
|
@ -0,0 +1,105 @@
|
|||
defmodule Microwaveprop.Workers.ContactPositionBackfillWorker do
|
||||
@moduledoc """
|
||||
Hourly safety net that fills in `pos1` / `pos2` / `distance_km` for contacts
|
||||
whose grids resolve but whose positions are null.
|
||||
|
||||
New contacts always have positions because `Radio.create_contact/1` refuses
|
||||
to insert without resolvable grids. This worker exists for rows that arrive
|
||||
via direct DB writes (bulk imports, manual fixes) which bypass that path.
|
||||
"""
|
||||
use Oban.Worker, queue: :admin, max_attempts: 1, unique: [period: 60]
|
||||
|
||||
import Ecto.Query
|
||||
|
||||
alias Microwaveprop.Radio
|
||||
alias Microwaveprop.Radio.Contact
|
||||
alias Microwaveprop.Radio.Maidenhead
|
||||
alias Microwaveprop.Repo
|
||||
|
||||
require Logger
|
||||
|
||||
@default_limit 500
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) do
|
||||
limit = Map.get(args, "limit", @default_limit)
|
||||
|
||||
candidates =
|
||||
Contact
|
||||
|> where([c], (is_nil(c.pos1) and not is_nil(c.grid1)) or (is_nil(c.pos2) and not is_nil(c.grid2)))
|
||||
|> limit(^limit)
|
||||
|> Repo.all()
|
||||
|
||||
count =
|
||||
Enum.count(candidates, fn contact ->
|
||||
updates = compute_updates(contact)
|
||||
|
||||
if updates == [] do
|
||||
false
|
||||
else
|
||||
{n, _} =
|
||||
Contact
|
||||
|> where(id: ^contact.id)
|
||||
|> Repo.update_all(set: updates)
|
||||
|
||||
n == 1
|
||||
end
|
||||
end)
|
||||
|
||||
if count > 0 do
|
||||
Logger.info("ContactPositionBackfill: populated positions for #{count} contacts")
|
||||
end
|
||||
|
||||
{:ok, %{updated: count}}
|
||||
end
|
||||
|
||||
defp compute_updates(contact) do
|
||||
pos1 = maybe_pos(contact.pos1, contact.grid1)
|
||||
pos2 = maybe_pos(contact.pos2, contact.grid2)
|
||||
|
||||
[]
|
||||
|> add_pos_change(:pos1, contact.pos1, pos1)
|
||||
|> add_pos_change(:pos2, contact.pos2, pos2)
|
||||
|> add_distance_change(contact, pos1, pos2)
|
||||
|> maybe_touch_updated_at()
|
||||
end
|
||||
|
||||
defp maybe_pos(existing, _grid) when not is_nil(existing), do: existing
|
||||
defp maybe_pos(_nil, nil), do: nil
|
||||
|
||||
defp maybe_pos(_nil, grid) do
|
||||
case Maidenhead.to_latlon(grid) do
|
||||
{:ok, {lat, lon}} -> %{"lat" => lat, "lon" => lon}
|
||||
:error -> nil
|
||||
end
|
||||
end
|
||||
|
||||
defp add_pos_change(updates, _field, existing, _new) when not is_nil(existing), do: updates
|
||||
defp add_pos_change(updates, _field, _existing, nil), do: updates
|
||||
defp add_pos_change(updates, field, _nil, new), do: [{field, new} | updates]
|
||||
|
||||
defp add_distance_change(updates, contact, pos1, pos2) do
|
||||
cond do
|
||||
pos1 == nil or pos2 == nil ->
|
||||
updates
|
||||
|
||||
not is_nil(contact.distance_km) and contact.pos1 != nil and contact.pos2 != nil ->
|
||||
updates
|
||||
|
||||
true ->
|
||||
distance =
|
||||
pos1["lat"]
|
||||
|> Radio.haversine_km(pos1["lon"], pos2["lat"], pos2["lon"])
|
||||
|> round()
|
||||
|> Decimal.new()
|
||||
|
||||
[{:distance_km, distance} | updates]
|
||||
end
|
||||
end
|
||||
|
||||
defp maybe_touch_updated_at([]), do: []
|
||||
|
||||
defp maybe_touch_updated_at(updates) do
|
||||
[{:updated_at, DateTime.truncate(DateTime.utc_now(), :second)} | updates]
|
||||
end
|
||||
end
|
||||
|
|
@ -0,0 +1,140 @@
|
|||
defmodule Microwaveprop.Workers.ContactPositionBackfillWorkerTest do
|
||||
use Microwaveprop.DataCase, async: true
|
||||
|
||||
alias Microwaveprop.Radio.Contact
|
||||
alias Microwaveprop.Repo
|
||||
alias Microwaveprop.Workers.ContactPositionBackfillWorker
|
||||
|
||||
defp insert_raw_contact!(attrs) do
|
||||
defaults = %{
|
||||
station1: "W5XD",
|
||||
station2: "K5TR",
|
||||
band: Decimal.new(10_000),
|
||||
qso_timestamp: ~U[2026-03-28 18:00:00Z],
|
||||
submitter_email: "t@example.com",
|
||||
user_submitted: true
|
||||
}
|
||||
|
||||
%Contact{}
|
||||
|> Map.merge(defaults)
|
||||
|> Map.merge(attrs)
|
||||
|> Repo.insert!()
|
||||
end
|
||||
|
||||
describe "perform/1" do
|
||||
test "backfills pos1 from grid1 when pos1 is null" do
|
||||
contact =
|
||||
insert_raw_contact!(%{grid1: "EM12", grid2: "EM00", pos1: nil, pos2: nil, distance_km: nil})
|
||||
|
||||
assert {:ok, %{updated: 1}} = ContactPositionBackfillWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
reloaded = Repo.reload!(contact)
|
||||
assert %{"lat" => _, "lon" => _} = reloaded.pos1
|
||||
assert %{"lat" => _, "lon" => _} = reloaded.pos2
|
||||
assert reloaded.distance_km
|
||||
end
|
||||
|
||||
test "backfills only pos1 when pos2 is already set" do
|
||||
pos2 = %{"lat" => 40.0, "lon" => -100.0}
|
||||
|
||||
contact =
|
||||
insert_raw_contact!(%{grid1: "EM12", grid2: "EM00", pos1: nil, pos2: pos2, distance_km: nil})
|
||||
|
||||
assert {:ok, %{updated: 1}} = ContactPositionBackfillWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
reloaded = Repo.reload!(contact)
|
||||
assert %{"lat" => _, "lon" => _} = reloaded.pos1
|
||||
assert reloaded.pos2 == pos2
|
||||
assert reloaded.distance_km
|
||||
end
|
||||
|
||||
test "is a no-op for contacts that already have both positions" do
|
||||
pos1 = %{"lat" => 32.0, "lon" => -97.0}
|
||||
pos2 = %{"lat" => 40.0, "lon" => -100.0}
|
||||
|
||||
contact =
|
||||
insert_raw_contact!(%{
|
||||
grid1: "EM12",
|
||||
grid2: "EM00",
|
||||
pos1: pos1,
|
||||
pos2: pos2,
|
||||
distance_km: Decimal.new(1234)
|
||||
})
|
||||
|
||||
assert {:ok, %{updated: 0}} = ContactPositionBackfillWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
reloaded = Repo.reload!(contact)
|
||||
assert reloaded.pos1 == pos1
|
||||
assert reloaded.pos2 == pos2
|
||||
assert Decimal.equal?(reloaded.distance_km, Decimal.new(1234))
|
||||
end
|
||||
|
||||
test "skips contact with null pos and null grid" do
|
||||
contact = insert_raw_contact!(%{grid1: nil, grid2: nil, pos1: nil, pos2: nil, distance_km: nil})
|
||||
|
||||
assert {:ok, %{updated: 0}} = ContactPositionBackfillWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
reloaded = Repo.reload!(contact)
|
||||
assert is_nil(reloaded.pos1)
|
||||
end
|
||||
|
||||
test "fills only the resolvable side when one grid is invalid" do
|
||||
contact =
|
||||
insert_raw_contact!(%{grid1: "ZZZZ99", grid2: "EM00", pos1: nil, pos2: nil, distance_km: nil})
|
||||
|
||||
assert {:ok, %{updated: 1}} = ContactPositionBackfillWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
reloaded = Repo.reload!(contact)
|
||||
assert is_nil(reloaded.pos1)
|
||||
assert %{"lat" => _, "lon" => _} = reloaded.pos2
|
||||
assert is_nil(reloaded.distance_km)
|
||||
end
|
||||
|
||||
test "skips contact where both grids are unresolvable" do
|
||||
contact =
|
||||
insert_raw_contact!(%{grid1: "ZZZZ99", grid2: "!!!!!!", pos1: nil, pos2: nil, distance_km: nil})
|
||||
|
||||
assert {:ok, %{updated: 0}} = ContactPositionBackfillWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
reloaded = Repo.reload!(contact)
|
||||
assert is_nil(reloaded.pos1)
|
||||
assert is_nil(reloaded.pos2)
|
||||
end
|
||||
|
||||
test "respects limit arg" do
|
||||
for i <- 1..3 do
|
||||
insert_raw_contact!(%{
|
||||
grid1: "EM12",
|
||||
grid2: "EM00",
|
||||
pos1: nil,
|
||||
pos2: nil,
|
||||
distance_km: nil,
|
||||
qso_timestamp: DateTime.add(~U[2026-03-28 18:00:00Z], i * 3600, :second)
|
||||
})
|
||||
end
|
||||
|
||||
assert {:ok, %{updated: 2}} =
|
||||
ContactPositionBackfillWorker.perform(%Oban.Job{args: %{"limit" => 2}})
|
||||
|
||||
remaining =
|
||||
Contact
|
||||
|> Ecto.Query.from(where: [grid1: "EM12"])
|
||||
|> Ecto.Query.where([c], is_nil(c.pos1))
|
||||
|> Repo.aggregate(:count)
|
||||
|
||||
assert remaining == 1
|
||||
end
|
||||
|
||||
test "computes distance_km from both grids" do
|
||||
contact =
|
||||
insert_raw_contact!(%{grid1: "EM12", grid2: "EM00", pos1: nil, pos2: nil, distance_km: nil})
|
||||
|
||||
assert {:ok, _} = ContactPositionBackfillWorker.perform(%Oban.Job{args: %{}})
|
||||
|
||||
reloaded = Repo.reload!(contact)
|
||||
# EM12 and EM00 span ~333 km in the Central US
|
||||
distance = Decimal.to_integer(reloaded.distance_km)
|
||||
assert distance > 0 and distance < 1000
|
||||
end
|
||||
end
|
||||
end
|
||||
Loading…
Add table
Reference in a new issue