From 3cb367c8980bf0b91de0c6a67d5b5355636539d9 Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Fri, 17 Apr 2026 09:54:04 -0500 Subject: [PATCH] fix(import): handle Radio.create_contact :duplicate return in worker MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit apply_valid_rows only matched {:ok, _} and {:error, %Changeset{}}, so when Radio.create_contact returned {:error, :duplicate, existing} — the race window where preview flagged a row as new but another writer inserts the same contact before the worker picks up the chunk — the reduce callback crashed with FunctionClauseError. Oban retried the chunk, rolling back the Postgrex transaction and surfacing the DBConnection errors the operator was seeing. Treat :duplicate as a silent skip: the row is already in the DB, there is nothing to insert, and the counter should not double-count the existing contact as an error either. --- lib/microwaveprop/radio/csv_import.ex | 3 + .../workers/contact_import_worker_test.exs | 125 ++++++++++++++++++ 2 files changed, 128 insertions(+) diff --git a/lib/microwaveprop/radio/csv_import.ex b/lib/microwaveprop/radio/csv_import.ex index 6adb101c..90d616ab 100644 --- a/lib/microwaveprop/radio/csv_import.ex +++ b/lib/microwaveprop/radio/csv_import.ex @@ -326,6 +326,9 @@ defmodule Microwaveprop.Radio.CsvImport do ContactWeatherEnqueueWorker.enqueue_for_contact(contact) %{acc | imported: [contact | acc.imported]} + {:error, :duplicate, _existing} -> + acc + {:error, %Ecto.Changeset{} = changeset} -> %{acc | errors: [{row.row_num, changeset_error_strings(changeset)} | acc.errors]} end diff --git a/test/microwaveprop/workers/contact_import_worker_test.exs b/test/microwaveprop/workers/contact_import_worker_test.exs index 6bbc0814..fae73ac8 100644 --- a/test/microwaveprop/workers/contact_import_worker_test.exs +++ b/test/microwaveprop/workers/contact_import_worker_test.exs @@ -219,4 +219,129 @@ defmodule Microwaveprop.Workers.ContactImportWorkerTest do }) end end + + # Reproduces a bug where the preview screen correctly flags an upload row + # as a duplicate of an existing DB contact (via direction-agnostic matching) + # but the async worker ends up inserting it anyway. The full pipeline — + # preview -> enqueue -> worker -> DB — must honor the preview's exclusion + # so the Imported count on /imports/:id never exceeds what the preview + # claimed as "Valid". + describe "preview-to-worker dedup parity" do + test "worker does not insert rows the preview classified as duplicates" do + {:ok, _existing} = + %Contact{} + |> Contact.submission_changeset(%{ + "station1" => "W5XD", + "station2" => "K5TR", + "grid1" => "EM12", + "grid2" => "EM00", + "band" => "10000", + "mode" => "CW", + "qso_timestamp" => "2026-03-28T18:00:00Z", + "submitter_email" => @submitter + }) + |> Repo.insert() + + # Upload row is a direction-swapped version of the existing contact + # at the same time/band — should be flagged as a duplicate. + csv = build_csv(["K5TR,W5XD,EM00,EM12,10000,CW,2026-03-28T18:30:00Z"]) + + {:ok, preview} = CsvImport.preview(csv, @submitter) + assert preview.valid == [] + assert length(preview.duplicates) == 1 + + {:ok, run_id} = CsvImport.enqueue(preview, @submitter) + + run = Repo.get!(ImportRun, run_id) + assert run.total_rows == 0 + assert run.rows["rows"] == [] + + assert Repo.aggregate(Contact, :count) == 1 + end + + test "worker does not insert in-upload direction-swapped duplicates" do + # Same contact logged from both ends — the CSV has row A->B and + # B->A at the same time/band. Only one should land in the DB. + csv = + build_csv([ + "W5XD,K5TR,EM12,EM00,10000,CW,2026-03-28T18:00:00Z", + "K5TR,W5XD,EM00,EM12,10000,CW,2026-03-28T18:00:00Z" + ]) + + {:ok, preview} = CsvImport.preview(csv, @submitter) + assert length(preview.valid) == 1 + assert length(preview.duplicates) == 1 + + {:ok, run_id} = CsvImport.enqueue(preview, @submitter) + + run = Repo.get!(ImportRun, run_id) + assert run.total_rows == 1 + assert length(run.rows["rows"]) == 1 + + assert Repo.aggregate(Contact, :count) == 1 + end + + # Covers the race window: preview ran when no matching contact existed, + # but by the time the worker picks up the chunk another insert (from a + # concurrent upload, a retry, etc.) has already landed. `Radio.create_contact` + # now returns `{:error, :duplicate, existing}`. The worker must not crash, + # must not double-insert, and must keep the counters in a sensible state. + test "worker survives when a row's duplicate appears in DB after preview" do + {:ok, _existing} = + %Contact{} + |> Contact.submission_changeset(%{ + "station1" => "W5XD", + "station2" => "K5TR", + "grid1" => "EM12", + "grid2" => "EM00", + "band" => "10000", + "mode" => "CW", + "qso_timestamp" => "2026-03-28T18:00:00Z", + "submitter_email" => @submitter + }) + |> Repo.insert() + + # Hand-rolled ImportRun: simulates a preview where the row slipped + # through as "valid" (dup appeared after preview ran). + rows = [ + %{ + "kind" => "insert", + "row_num" => 2, + "attrs" => %{ + "station1" => "W5XD", + "station2" => "K5TR", + "grid1" => "EM12", + "grid2" => "EM00", + "band" => "10000", + "mode" => "CW", + "qso_timestamp" => "2026-03-28T18:00:00Z", + "submitter_email" => @submitter + }, + "timestamp" => "2026-03-28T18:00:00Z" + } + ] + + {:ok, run} = + %ImportRun{} + |> ImportRun.changeset(%{ + submitter_email: @submitter, + rows: %{"rows" => rows, "refinement_ids" => []}, + total_rows: 1 + }) + |> Repo.insert() + + assert :ok = + ContactImportWorker.perform(%Oban.Job{ + args: %{"import_run_id" => run.id, "offset" => 0, "limit" => 1} + }) + + reloaded = Repo.get!(ImportRun, run.id) + assert reloaded.processed_rows == 1 + assert reloaded.imported_count == 0 + assert reloaded.status == "complete" + + # The existing contact must remain; no new row inserted. + assert Repo.aggregate(Contact, :count) == 1 + end + end end