diff --git a/lib/microwaveprop/weather/hrdps_profile.ex b/lib/microwaveprop/weather/hrdps_profile.ex new file mode 100644 index 00000000..63757fe8 --- /dev/null +++ b/lib/microwaveprop/weather/hrdps_profile.ex @@ -0,0 +1,42 @@ +defmodule Microwaveprop.Weather.HrdpsProfile do + @moduledoc false + use Ecto.Schema + + import Ecto.Changeset + + @primary_key {:id, :binary_id, autogenerate: true} + @foreign_key_type :binary_id + + schema "hrdps_profiles" do + field :valid_time, :utc_datetime + field :lat, :float + field :lon, :float + field :run_time, :utc_datetime + field :profile, {:array, :map} + field :hpbl_m, :float + field :pwat_mm, :float + field :surface_temp_c, :float + field :surface_dewpoint_c, :float + field :surface_pressure_mb, :float + field :surface_refractivity, :float + field :min_refractivity_gradient, :float + field :ducting_detected, :boolean, default: false + field :duct_characteristics, {:array, :map} + field :is_grid_point, :boolean, default: false + + timestamps(type: :utc_datetime) + end + + @type t :: %__MODULE__{} + + @required_fields ~w(valid_time lat lon)a + @optional_fields ~w(run_time profile hpbl_m pwat_mm surface_temp_c surface_dewpoint_c surface_pressure_mb surface_refractivity min_refractivity_gradient ducting_detected duct_characteristics is_grid_point)a + + @spec changeset(t() | Ecto.Changeset.t(), map()) :: Ecto.Changeset.t() + def changeset(hrdps_profile, attrs) do + hrdps_profile + |> cast(attrs, @required_fields ++ @optional_fields) + |> validate_required(@required_fields) + |> unique_constraint([:lat, :lon, :valid_time]) + end +end diff --git a/priv/repo/migrations/20260429220349_create_hrdps_profiles.exs b/priv/repo/migrations/20260429220349_create_hrdps_profiles.exs new file mode 100644 index 00000000..582d3916 --- /dev/null +++ b/priv/repo/migrations/20260429220349_create_hrdps_profiles.exs @@ -0,0 +1,66 @@ +defmodule Microwaveprop.Repo.Migrations.CreateHrdpsProfiles do + use Ecto.Migration + + @disable_ddl_transaction true + @disable_migration_lock true + + # Mirror of hrrr_profiles for HRDPS-derived rows. Separate table (not a + # unified nwp_profiles) so HRDPS can be dropped/reprocessed without + # disturbing the 4500+ HRRR profiles already in prod. Same column set + # and indexes; PARTITION BY RANGE (valid_time) with quarterly partitions + # starting from the deploy quarter. + # + # Unlike hrrr_profiles, there's no historical archive to backfill into + # — HRDPS is forward-only — so the partition list starts at the deploy + # quarter and walks forward. + def up do + execute """ + CREATE TABLE hrdps_profiles ( + id uuid NOT NULL DEFAULT gen_random_uuid(), + valid_time timestamp(0) WITHOUT TIME ZONE NOT NULL, + lat double precision NOT NULL, + lon double precision NOT NULL, + run_time timestamp(0) WITHOUT TIME ZONE, + profile jsonb[] DEFAULT ARRAY[]::jsonb[], + hpbl_m double precision, + pwat_mm double precision, + surface_temp_c double precision, + surface_dewpoint_c double precision, + surface_pressure_mb double precision, + surface_refractivity double precision, + min_refractivity_gradient double precision, + ducting_detected boolean DEFAULT false, + duct_characteristics jsonb[], + is_grid_point boolean DEFAULT false, + inserted_at timestamp(0) WITHOUT TIME ZONE NOT NULL, + updated_at timestamp(0) WITHOUT TIME ZONE NOT NULL + ) PARTITION BY RANGE (valid_time) + """ + + partitions = [ + {"2026-04-01", "2026-07-01"}, + {"2026-07-01", "2026-10-01"}, + {"2026-10-01", "2027-01-01"}, + {"2027-01-01", "2027-04-01"}, + {"2027-04-01", "2027-07-01"}, + {"2027-07-01", "2027-10-01"}, + {"2027-10-01", "2028-01-01"} + ] + + for {from_date, to_date} <- partitions do + suffix = from_date |> String.replace("-", "_") |> String.slice(0..6) + + execute """ + CREATE TABLE hrdps_profiles_#{suffix} PARTITION OF hrdps_profiles + FOR VALUES FROM ('#{from_date}') TO ('#{to_date}') + """ + end + + execute "CREATE UNIQUE INDEX hrdps_profiles_lat_lon_valid_time_index ON hrdps_profiles (lat, lon, valid_time)" + execute "CREATE INDEX hrdps_profiles_valid_time_index ON hrdps_profiles (valid_time)" + end + + def down do + execute "DROP TABLE IF EXISTS hrdps_profiles CASCADE" + end +end diff --git a/test/microwaveprop/weather/hrdps_profile_test.exs b/test/microwaveprop/weather/hrdps_profile_test.exs new file mode 100644 index 00000000..175eb4b0 --- /dev/null +++ b/test/microwaveprop/weather/hrdps_profile_test.exs @@ -0,0 +1,83 @@ +defmodule Microwaveprop.Weather.HrdpsProfileTest do + use Microwaveprop.DataCase, async: false + + alias Microwaveprop.Repo + alias Microwaveprop.Weather.HrdpsProfile + + describe "changeset/2" do + test "valid with required fields" do + attrs = %{valid_time: ~U[2026-04-29 12:00:00Z], lat: 53.32, lon: -60.42} + cs = HrdpsProfile.changeset(%HrdpsProfile{}, attrs) + assert cs.valid? + end + + test "rejects missing required fields" do + cs = HrdpsProfile.changeset(%HrdpsProfile{}, %{}) + refute cs.valid? + assert %{valid_time: ["can't be blank"]} = errors_on(cs) + assert %{lat: ["can't be blank"]} = errors_on(cs) + assert %{lon: ["can't be blank"]} = errors_on(cs) + end + + test "casts the full set of derived fields" do + attrs = %{ + valid_time: ~U[2026-04-29 12:00:00Z], + lat: 53.32, + lon: -60.42, + run_time: ~U[2026-04-29 12:00:00Z], + profile: [%{"pres" => 1000.0, "tmpc" => 12.0, "dwpc" => 7.0, "hght" => 100.0}], + hpbl_m: 800.0, + pwat_mm: nil, + surface_temp_c: 12.0, + surface_dewpoint_c: 7.0, + surface_pressure_mb: 1000.0, + surface_refractivity: 312.0, + min_refractivity_gradient: -50.0, + ducting_detected: false, + duct_characteristics: [], + is_grid_point: true + } + + cs = HrdpsProfile.changeset(%HrdpsProfile{}, attrs) + assert cs.valid? + assert Ecto.Changeset.get_field(cs, :surface_refractivity) == 312.0 + assert Ecto.Changeset.get_field(cs, :is_grid_point) == true + end + end + + describe "persistence" do + test "inserts and reads back via the partitioned hrdps_profiles table" do + attrs = %{ + valid_time: ~U[2026-04-29 12:00:00Z], + lat: 53.32, + lon: -60.42, + run_time: ~U[2026-04-29 12:00:00Z], + surface_temp_c: 7.8 + } + + {:ok, %HrdpsProfile{id: id}} = + %HrdpsProfile{} |> HrdpsProfile.changeset(attrs) |> Repo.insert() + + reloaded = Repo.get!(HrdpsProfile, id) + assert reloaded.lat == 53.32 + assert reloaded.surface_temp_c == 7.8 + end + + test "enforces unique constraint on (lat, lon, valid_time) at the DB level" do + # The partitioned table propagates the unique index per partition, so + # Postgres rejects duplicate (lat, lon, valid_time) inserts. Schema- + # level `unique_constraint/2` matching against changeset errors + # doesn't apply because the partition's index name varies per + # partition; callers that need upserts go through + # `on_conflict: {:replace_all_except, ...}` in the worker. + attrs = %{valid_time: ~U[2026-04-29 12:00:00Z], lat: 53.32, lon: -60.42} + + assert {:ok, _} = + %HrdpsProfile{} |> HrdpsProfile.changeset(attrs) |> Repo.insert() + + assert_raise Ecto.ConstraintError, fn -> + %HrdpsProfile{} |> HrdpsProfile.changeset(attrs) |> Repo.insert() + end + end + end +end