From 26ef50c9656f7ea56250952e5a54366c3078fc01 Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Tue, 21 Apr 2026 16:39:12 -0500 Subject: [PATCH] fix(ingest): dedupe ASOS duplicates + shrink SWPC decode errors Two unrelated production warnings from the same log window: WeatherFetchWorker raised Postgrex 21000 cardinality_violation on ASOS upserts when IEM returned two METAR records for the same timestamp (routine + special at the same minute). Dedupe by (station_id, observed_at) in upsert_surface_observations, keeping the last occurrence since IEM's ordering is chronological. SwpcClient.fetch_xrays warned with the full truncated response body (~650 KB of GOES X-ray JSON) dumped into the log aggregator whenever SWPC's CDN cut the response mid-stream. Intercept Jason.DecodeError in fetch_and_parse and format it short: keep the byte position and an 80-char preview, drop the rest. --- .../space_weather/swpc_client.ex | 20 +++++++++++++-- lib/microwaveprop/weather.ex | 13 ++++++++++ .../space_weather/swpc_client_test.exs | 25 +++++++++++++++++++ test/microwaveprop/weather_test.exs | 21 ++++++++++++++++ 4 files changed, 77 insertions(+), 2 deletions(-) diff --git a/lib/microwaveprop/space_weather/swpc_client.ex b/lib/microwaveprop/space_weather/swpc_client.ex index c9d8d396..ad1a1aad 100644 --- a/lib/microwaveprop/space_weather/swpc_client.ex +++ b/lib/microwaveprop/space_weather/swpc_client.ex @@ -158,10 +158,10 @@ defmodule Microwaveprop.SpaceWeather.SwpcClient do parser.(body) {:ok, %{status: status, body: body}} -> - {:error, "SWPC HTTP #{status}: #{inspect(body)}"} + {:error, "SWPC HTTP #{status}: #{inspect(body, printable_limit: 200)}"} {:error, reason} -> - {:error, "SWPC request failed: #{inspect(reason)}"} + {:error, "SWPC request failed: #{format_req_error(reason)}"} end end) end @@ -226,4 +226,20 @@ defmodule Microwaveprop.SpaceWeather.SwpcClient do defp req_options do Application.get_env(:microwaveprop, :swpc_req_options, []) end + + # SWPC's CDN occasionally truncates the response mid-stream. Req's + # auto-JSON-decoder then returns `{:error, %Jason.DecodeError{data: body}}` + # where `:data` is the full raw body (up to ~650 KB for xrays-1-day). + # Inspecting that struct at warn level spams the log aggregator. Keep + # enough context to debug (position + a short preview) and drop the body. + defp format_req_error(%Jason.DecodeError{position: pos, data: data}) do + preview = + data + |> to_string() + |> binary_part(0, min(80, byte_size(to_string(data)))) + + "Jason.DecodeError at byte #{pos} (truncated upstream response; preview=#{inspect(preview)})" + end + + defp format_req_error(reason), do: inspect(reason, printable_limit: 200) end diff --git a/lib/microwaveprop/weather.ex b/lib/microwaveprop/weather.ex index 22480f9d..f4249777 100644 --- a/lib/microwaveprop/weather.ex +++ b/lib/microwaveprop/weather.ex @@ -107,6 +107,7 @@ defmodule Microwaveprop.Weather do rows |> Enum.filter(&row_has_observed_at?/1) |> Enum.map(&surface_observation_entry(&1, station.id, now)) + |> dedupe_last_by_conflict_target() case entries do [] -> @@ -147,6 +148,18 @@ defmodule Microwaveprop.Weather do defp row_has_observed_at?(%{"observed_at" => %DateTime{}}), do: true defp row_has_observed_at?(_), do: false + # IEM ASOS occasionally returns two rows with the same timestamp + # (e.g. routine + special METAR at the same minute). Passing both + # to insert_all triggers Postgres 21000 cardinality_violation. + # Keep the last occurrence — IEM's ordering is chronological, so + # the later row is the final correction. + defp dedupe_last_by_conflict_target(entries) do + entries + |> Enum.reverse() + |> Enum.uniq_by(&{&1.station_id, &1.observed_at}) + |> Enum.reverse() + end + defp surface_observation_entry(row, station_id, now) do %{ id: Ecto.UUID.generate(), diff --git a/test/microwaveprop/space_weather/swpc_client_test.exs b/test/microwaveprop/space_weather/swpc_client_test.exs index acd0faea..3327e326 100644 --- a/test/microwaveprop/space_weather/swpc_client_test.exs +++ b/test/microwaveprop/space_weather/swpc_client_test.exs @@ -22,6 +22,31 @@ defmodule Microwaveprop.SpaceWeather.SwpcClientTest do end end + describe "fetch_xrays/0 upstream truncation" do + # SWPC occasionally returns a truncated JSON body (TCP reset mid-stream). + # Req's auto-decoder raises a Jason.DecodeError carrying the entire + # raw body in :data — historically that dumped hundreds of KB into + # our Logger. The error should be returned, but formatted short. + test "returns {:error, reason} without dumping the full truncated body" do + big_body = "[" <> String.duplicate(~s({"time_tag":"2026-04-20T21:35:00Z","energy":"0.1-0.8nm","flux":1.0},), 20_000) + # Truncate — no closing bracket. + truncated = big_body + + Req.Test.stub(SwpcClient, fn conn -> + conn + |> Plug.Conn.put_resp_content_type("application/json") + |> Plug.Conn.send_resp(200, truncated) + end) + + assert {:error, reason} = SwpcClient.fetch_xrays() + assert is_binary(reason) + # Upper-bound the log-noise blast radius. + assert byte_size(reason) < 500 + # Still descriptive. + assert reason =~ ~r/(decode|JSON|truncated|Jason)/i + end + end + describe "parse_f107/1" do test "decodes the SWPC f107_cm_flux fixture" do body = File.read!("test/fixtures/swpc/f107_sample.json") diff --git a/test/microwaveprop/weather_test.exs b/test/microwaveprop/weather_test.exs index 507077bd..a0e5d952 100644 --- a/test/microwaveprop/weather_test.exs +++ b/test/microwaveprop/weather_test.exs @@ -143,6 +143,27 @@ defmodule Microwaveprop.WeatherTest do assert {1, _} = Weather.upsert_surface_observations(station, rows) assert Repo.aggregate(SurfaceObservation, :count) == 1 end + + test "dedupes rows sharing (station_id, observed_at), keeping the last" do + {:ok, station} = Weather.find_or_create_station(@station_attrs) + ts = ~U[2024-09-22 20:00:00Z] + + # IEM occasionally returns two METAR records with the same + # timestamp. Passing both to a single insert_all triggers + # Postgres 21000 cardinality_violation because ON CONFLICT + # cannot update the same row twice in one command. + rows = [ + %{observed_at: ts, temp_f: 70.0, dewpoint_f: 60.0}, + %{observed_at: ts, temp_f: 72.0, dewpoint_f: 61.0} + ] + + assert {1, _} = Weather.upsert_surface_observations(station, rows) + assert Repo.aggregate(SurfaceObservation, :count) == 1 + + [obs] = Repo.all(SurfaceObservation) + assert obs.temp_f == 72.0 + assert obs.dewpoint_f == 61.0 + end end describe "upsert_sounding/2" do