prop/lib/microwaveprop/workers/iemre_fetch_worker.ex
Graham McIntire 47ba8d3b6d
fix(enrichment): stub IEMRE on permanent failure + bump hot CPU
IemreFetchWorker now writes an empty stub observation when IEM returns
a permanent-failure status (404, 422), mirroring the existing {:ok, []}
path. The backfill cron's next tick sees the stub via
has_iemre_observation?/3, generates no job, and mark_status!/3 flips
the contact's iemre_status from :queued to :complete — draining the
51 contacts that have been stuck behind cancelled out-of-grid IEMRE
jobs.

Raise hot pod CPU limit 2 → 3 so BEAM gets 3 schedulers online.
Observed run queues of [2, 10, 0, 0, 0, 0] on the 2-scheduler
configuration during the :05 propagation chain, starving /live and
/health probe handlers and tripping intermittent liveness timeouts.
2026-04-23 12:26:09 -05:00

84 lines
2.7 KiB
Elixir

defmodule Microwaveprop.Workers.IemreFetchWorker do
@moduledoc false
use Oban.Worker,
queue: :iemre,
max_attempts: 20,
unique: [period: 300, fields: [:args], states: [:scheduled, :available]]
alias Microwaveprop.Weather
alias Microwaveprop.Weather.IemClient
require Logger
@impl Oban.Worker
def backoff(%Oban.Job{attempt: attempt}) do
min(120 * Integer.pow(2, attempt - 1), _six_hours = 21_600)
end
@impl Oban.Worker
def perform(%Oban.Job{args: args}) do
%{"lat" => lat, "lon" => lon, "date" => date_str} = args
date = Date.from_iso8601!(date_str)
if Weather.has_iemre_observation?(lat, lon, date) do
Logger.info("IEMRE observation already exists for #{lat},#{lon} @ #{date_str}")
:ok
else
Logger.info("Fetching IEMRE data for #{lat},#{lon} @ #{date_str}")
fetch_and_store_iemre(lat, lon, date, date_str)
end
end
defp fetch_and_store_iemre(lat, lon, date, date_str) do
case IemClient.fetch_iemre(lat, lon, date) do
{:ok, []} ->
# Store stub so this lat/lon/date isn't retried on future backfills
_ = Weather.upsert_iemre_observation(%{lat: lat, lon: lon, date: date, hourly: []})
Logger.info("IEMRE: no data available for #{lat},#{lon} @ #{date_str}, stored stub")
:ok
{:ok, data} ->
_ =
Weather.upsert_iemre_observation(%{
lat: lat,
lon: lon,
date: date,
hourly: data
})
Logger.info("IEMRE observation saved for #{lat},#{lon} @ #{date_str} (#{length(data)} hours)")
:ok
{:error, reason} ->
handle_error(reason, lat, lon, date, date_str)
end
end
defp handle_error(reason, lat, lon, date, date_str) do
if transient_failure?(reason) do
Logger.error("IEMRE transient error for #{lat},#{lon} @ #{date_str}: #{inspect(reason)}")
{:error, reason}
else
# Permanent upstream failure (e.g. 404, 422 out-of-grid). Stub the
# bucket so future backfill enqueues skip the fetch and let
# ContactWeatherEnqueueWorker reconcile the contacts' :queued →
# :complete transition via the empty-jobs branch.
_ = Weather.upsert_iemre_observation(%{lat: lat, lon: lon, date: date, hourly: []})
Logger.warning("IEMRE permanent failure for #{lat},#{lon} @ #{date_str}: #{inspect(reason)}, stored stub")
{:cancel, reason}
end
end
defp transient_failure?(%{__exception__: true}), do: true
defp transient_failure?("IEM IEMRE HTTP " <> status), do: server_error?(status)
defp transient_failure?(_), do: false
defp server_error?(status) do
case Integer.parse(status) do
{code, _} when code in [429, 500, 502, 503, 504] -> true
_ -> false
end
end
end