diff --git a/lib/microwaveprop/workers/backfill_enqueue_worker.ex b/lib/microwaveprop/workers/backfill_enqueue_worker.ex index 7dfdaba5..d3d3bd46 100644 --- a/lib/microwaveprop/workers/backfill_enqueue_worker.ex +++ b/lib/microwaveprop/workers/backfill_enqueue_worker.ex @@ -14,7 +14,7 @@ defmodule Microwaveprop.Workers.BackfillEnqueueWorker do require Logger @enrichable [:pending, :queued, :failed] - @valid_types ~w(hrrr weather terrain iemre) + @valid_types ~w(hrrr weather terrain iemre era5) @impl Oban.Worker def perform(%Oban.Job{args: %{"limit" => limit} = args}) do @@ -52,7 +52,7 @@ defmodule Microwaveprop.Workers.BackfillEnqueueWorker do types |> Enum.filter(&(&1 in @valid_types)) |> Enum.map(&String.to_existing_atom/1) end - defp parse_types(_), do: [:hrrr, :weather, :terrain, :iemre] + defp parse_types(_), do: [:hrrr, :weather, :terrain, :iemre, :era5] defp status_priority_order(types) do # Prioritize pending/failed over queued so real work happens first @@ -87,9 +87,14 @@ defmodule Microwaveprop.Workers.BackfillEnqueueWorker do end defp type_filter(types) do - Enum.reduce(types, dynamic(false), fn type, acc -> - field = :"#{type}_status" - dynamic([c], ^acc or field(c, ^field) in ^@enrichable) + Enum.reduce(types, dynamic(false), fn + :era5, acc -> + # ERA5 targets contacts where HRRR is unavailable (pre-2014, missing from archive) + dynamic([c], ^acc or c.hrrr_status == :unavailable) + + type, acc -> + field = :"#{type}_status" + dynamic([c], ^acc or field(c, ^field) in ^@enrichable) end) end end diff --git a/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex b/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex index 727e40f8..ad168155 100644 --- a/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex +++ b/lib/microwaveprop/workers/contact_weather_enqueue_worker.ex @@ -5,6 +5,7 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorker do alias Microwaveprop.Radio alias Microwaveprop.Weather alias Microwaveprop.Weather.HrrrClient + alias Microwaveprop.Workers.Era5FetchWorker alias Microwaveprop.Workers.HrrrFetchWorker alias Microwaveprop.Workers.IemreFetchWorker alias Microwaveprop.Workers.TerrainProfileWorker @@ -26,8 +27,9 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorker do hrrr_jobs = if :hrrr in types, do: build_hrrr_jobs([contact]), else: [] terrain_jobs = if :terrain in types, do: build_terrain_jobs([contact]), else: [] iemre_jobs = if :iemre in types, do: build_iemre_jobs([contact]), else: [] + era5_jobs = if :era5 in types, do: build_era5_jobs([contact]), else: [] - all_jobs = weather_jobs ++ hrrr_jobs ++ terrain_jobs ++ iemre_jobs + all_jobs = weather_jobs ++ hrrr_jobs ++ terrain_jobs ++ iemre_jobs ++ era5_jobs if all_jobs != [] do insert_all_chunked(all_jobs) @@ -155,6 +157,37 @@ defmodule Microwaveprop.Workers.ContactWeatherEnqueueWorker do |> Enum.uniq_by(fn changeset -> changeset.changes.args end) end + def build_era5_jobs(contacts) do + contacts + |> Enum.flat_map(&era5_jobs_for_contact/1) + |> Enum.uniq_by(fn changeset -> changeset.changes.args end) + end + + defp era5_jobs_for_contact(%{pos1: nil}), do: [] + + defp era5_jobs_for_contact(contact) do + valid_time = %{DateTime.truncate(contact.qso_timestamp, :second) | minute: 0, second: 0} + + contact + |> Radio.contact_path_points() + |> Enum.flat_map(fn {lat, lon} -> + rlat = Float.round(lat * 4) / 4 + rlon = Float.round(lon * 4) / 4 + + if Weather.find_nearest_era5(rlat, rlon, valid_time) do + [] + else + [ + Era5FetchWorker.new(%{ + "lat" => rlat, + "lon" => rlon, + "valid_time" => DateTime.to_iso8601(valid_time) + }) + ] + end + end) + end + defp jobs_for_contact(contact) do contact |> Radio.contact_path_points() diff --git a/lib/microwaveprop_web/live/backfill_live.ex b/lib/microwaveprop_web/live/backfill_live.ex index 5d776266..c20e3f48 100644 --- a/lib/microwaveprop_web/live/backfill_live.ex +++ b/lib/microwaveprop_web/live/backfill_live.ex @@ -31,7 +31,7 @@ defmodule MicrowavepropWeb.BackfillLive do assign(socket, page_title: "Backfill", limit: limit, - types: MapSet.new(~w(hrrr weather terrain iemre)), + types: MapSet.new(~w(hrrr weather terrain iemre era5)), stats: stats, unprocessed: unprocessed, db_stats: db_stats, @@ -46,10 +46,10 @@ defmodule MicrowavepropWeb.BackfillLive do limit = String.to_integer(limit_str) types = - ~w(hrrr weather terrain iemre) + ~w(hrrr weather terrain iemre era5) |> Enum.filter(&(params[&1] == "true")) |> then(fn - [] -> ~w(hrrr weather terrain iemre) + [] -> ~w(hrrr weather terrain iemre era5) types -> types end) @@ -114,6 +114,10 @@ defmodule MicrowavepropWeb.BackfillLive do iemre = Repo.one(from(c in Contact, where: c.iemre_status in ^incomplete and not is_nil(c.pos1), select: count())) + # ERA5 targets contacts where HRRR is unavailable + era5 = + Repo.one(from(c in Contact, where: c.hrrr_status == :unavailable and not is_nil(c.pos1), select: count())) + all_done = Repo.one( from(c in Contact, @@ -132,6 +136,7 @@ defmodule MicrowavepropWeb.BackfillLive do hrrr: hrrr, weather: weather, iemre: iemre, + era5: era5, remaining: max(terrain, max(hrrr, max(weather, iemre))), complete: all_done, total: total @@ -233,6 +238,18 @@ defmodule MicrowavepropWeb.BackfillLive do # Per-table row counts from pg_stat estimates (already fetched above, avoids slow count(*)) table_counts = Map.new(tables, fn %{table: t, rows: r} -> {t, r} end) + # ERA5 doesn't have its own status column — synthesize from HRRR unavailable + era5_profiles count + era5_count = Map.get(table_counts, "era5_profiles", 0) + + hrrr_unavailable = + statuses |> Map.get("hrrr", []) |> Enum.find_value(0, fn {s, c} -> if s == "unavailable", do: c end) + + statuses = + Map.put(statuses, "era5", [ + {"candidates (HRRR unavail)", hrrr_unavailable}, + {"profiles fetched", era5_count} + ]) + hrrr_count = Map.get(table_counts, "hrrr_profiles", 0) terrain_count = Map.get(table_counts, "terrain_profiles", 0) obs_count = Map.get(table_counts, "surface_observations", 0) @@ -246,6 +263,7 @@ defmodule MicrowavepropWeb.BackfillLive do statuses: statuses, totals: %{ hrrr_profiles: hrrr_count, + era5_profiles: Map.get(table_counts, "era5_profiles", 0), terrain_profiles: terrain_count, surface_observations: obs_count, soundings: sounding_count, @@ -306,7 +324,7 @@ defmodule MicrowavepropWeb.BackfillLive do {format_number(@unprocessed.complete)} / {format_number(@unprocessed.total)} fully enriched
- <%= for {label, complete} <- [{"HRRR", @unprocessed.total - @unprocessed.hrrr}, {"Weather", @unprocessed.total - @unprocessed.weather}, {"Terrain", @unprocessed.total - @unprocessed.terrain}, {"IEMRE", @unprocessed.total - @unprocessed.iemre}] do %> + <%= for {label, complete} <- [{"HRRR", @unprocessed.total - @unprocessed.hrrr}, {"Weather", @unprocessed.total - @unprocessed.weather}, {"Terrain", @unprocessed.total - @unprocessed.terrain}, {"IEMRE", @unprocessed.total - @unprocessed.iemre}, {"ERA5", @unprocessed.total - @unprocessed.era5}] do %> <% pct = if @unprocessed.total > 0, do: round(complete * 100 / @unprocessed.total), else: 0 %>
@@ -353,7 +371,7 @@ defmodule MicrowavepropWeb.BackfillLive do />
- <%= for type <- ~w(hrrr weather terrain iemre) do %> + <%= for type <- ~w(hrrr weather terrain iemre era5) do %>