diff --git a/config/config.exs b/config/config.exs index cfe32727..522cea89 100644 --- a/config/config.exs +++ b/config/config.exs @@ -45,7 +45,7 @@ config :microwaveprop, MicrowavepropWeb.Endpoint, config :microwaveprop, Oban, engine: Oban.Pro.Engines.Smart, repo: Microwaveprop.Repo, - queues: [solar: 1, weather: 20, enqueue: 1, hrrr: 5, terrain: 4, commercial: 2, iemre: 10, propagation: 1], + queues: [solar: 1, weather: 20, enqueue: 1, hrrr: 5, terrain: 4, commercial: 2, iemre: 10, propagation: 1, admin: 1], plugins: [ {Oban.Plugins.Pruner, max_age: 3600 * 24}, {Oban.Plugins.Lifeline, rescue_after: to_timeout(minute: 30)}, diff --git a/config/runtime.exs b/config/runtime.exs index 8ba1734a..1fd043da 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -150,7 +150,8 @@ if config_env() == :prod do rate_limit: [allowed: 10, period: {1, :hour}] ], rtma: 2, - backfill_enqueue: 1 + backfill_enqueue: 1, + admin: 1 ], plugins: [ {Oban.Plugins.Pruner, max_age: 3600 * 24}, diff --git a/lib/microwaveprop/release.ex b/lib/microwaveprop/release.ex index 2f41dd5c..b5f71889 100644 --- a/lib/microwaveprop/release.ex +++ b/lib/microwaveprop/release.ex @@ -2,6 +2,9 @@ defmodule Microwaveprop.Release do @moduledoc """ Release tasks that run in production without Mix installed. + Long-running tasks are enqueued as Oban jobs so `eval` returns + immediately. Check Oban Web at /admin/oban for progress. + Usage via `bin/microwaveprop eval`: bin/microwaveprop eval "Microwaveprop.Release.migrate" @@ -13,6 +16,8 @@ defmodule Microwaveprop.Release do """ @app :microwaveprop + alias Microwaveprop.Workers.AdminTaskWorker + def migrate do load_app() @@ -26,126 +31,29 @@ defmodule Microwaveprop.Release do {:ok, _, _} = Ecto.Migrator.with_repo(repo, &Ecto.Migrator.run(&1, :down, to: version)) end - @doc """ - Run a single feature backtest. Prints markdown report to stdout. - """ def backtest(feature_name) when is_binary(feature_name) do start_app() - - alias Microwaveprop.Backtest - alias Microwaveprop.Backtest.Features - - fun_atom = String.to_atom(feature_name) - - unless function_exported?(Features, fun_atom, 3) do - available = - Features.__info__(:functions) - |> Enum.filter(fn {_, a} -> a == 3 end) - |> Enum.map_join(", ", fn {f, _} -> to_string(f) end) - - IO.puts("Unknown feature: #{feature_name}\nAvailable: #{available}") - System.halt(1) - end - - feature_fun = &apply(Features, fun_atom, [&1, &2, &3]) - name = "Microwaveprop.Backtest.Features.#{feature_name}" - - report = Backtest.evaluate(feature_fun, sample_size: 5000, feature_name: name) - distance_bins = Backtest.lift_by_distance(feature_fun, sample_size: 5000) - band_stats = Backtest.lift_by_band(feature_fun, sample_size: 5000) - - IO.puts(Backtest.to_markdown(report, distance_bins: distance_bins, band_stats: band_stats)) + {:ok, job} = Oban.insert(AdminTaskWorker.new(%{task: "backtest", feature: feature_name})) + IO.puts("Enqueued backtest for #{feature_name} (job #{job.id})") end - @doc """ - Run consolidated backtest across all features. Prints markdown to stdout. - """ def backtest_all do start_app() - - alias Microwaveprop.Backtest - alias Microwaveprop.Backtest.Features - - features = Features.all_features() - IO.puts("Running consolidated backtest for #{map_size(features)} features...") - - results = Backtest.consolidated_report(features, sample_size: 5000) - IO.puts(Backtest.to_consolidated_markdown(results)) + {:ok, job} = Oban.insert(AdminTaskWorker.new(%{task: "backtest_all"})) + IO.puts("Enqueued consolidated backtest (job #{job.id})") end - @doc """ - Build surface temperature climatology from hrrr_profiles. - """ def climatology(min_samples \\ 3) do start_app() - - alias Microwaveprop.Repo - - %{rows: combos} = - Repo.query!( - """ - SELECT EXTRACT(MONTH FROM valid_time)::int AS month, - EXTRACT(HOUR FROM valid_time)::int AS hour - FROM hrrr_profiles - WHERE surface_temp_c IS NOT NULL AND is_grid_point = true - GROUP BY 1, 2 - ORDER BY 1, 2 - """, - [], - timeout: 120_000 - ) - - IO.puts("Building climatology (min_samples=#{min_samples}, #{length(combos)} batches)...") - - total = - combos - |> Enum.with_index(1) - |> Enum.reduce(0, fn {[month, hour], idx}, acc -> - %{num_rows: count} = - Repo.query!( - """ - INSERT INTO hrrr_climatology (id, lat, lon, month, hour, - mean_surface_temp_c, stddev_surface_temp_c, sample_count) - SELECT gen_random_uuid(), lat, lon, - $2 AS month, - $3 AS hour, - AVG(surface_temp_c), - STDDEV_SAMP(surface_temp_c), - COUNT(*) - FROM hrrr_profiles - WHERE surface_temp_c IS NOT NULL - AND is_grid_point = true - AND EXTRACT(MONTH FROM valid_time)::int = $2 - AND EXTRACT(HOUR FROM valid_time)::int = $3 - GROUP BY lat, lon - HAVING COUNT(*) >= $1 - ON CONFLICT (lat, lon, month, hour) - DO UPDATE SET - mean_surface_temp_c = EXCLUDED.mean_surface_temp_c, - stddev_surface_temp_c = EXCLUDED.stddev_surface_temp_c, - sample_count = EXCLUDED.sample_count - """, - [min_samples, month, hour], - timeout: 300_000 - ) - - IO.puts(" [#{idx}/#{length(combos)}] month=#{month} hour=#{hour}: #{count} rows") - acc + count - end) - - IO.puts("Upserted #{total} climatology records total.") + {:ok, job} = Oban.insert(AdminTaskWorker.new(%{task: "climatology", min_samples: min_samples})) + IO.puts("Enqueued climatology build (job #{job.id})") end - @doc """ - Enqueue native HRRR backfill jobs for the top-N hours by contact count. - """ def native_backfill(limit \\ 500) do start_app() - alias Microwaveprop.Repo - %{rows: rows} = - Repo.query!( + Microwaveprop.Repo.query!( """ SELECT EXTRACT(YEAR FROM qso_timestamp)::int AS year, EXTRACT(MONTH FROM qso_timestamp)::int AS month, @@ -172,84 +80,17 @@ defmodule Microwaveprop.Release do Microwaveprop.Workers.HrrrNativeGridWorker.new(args, unique: [period: :infinity]) ) do {:ok, _} -> IO.puts(" #{year}-#{month}-#{day} #{hour}Z (#{cnt} contacts)") - {:error, _} -> IO.puts(" #{year}-#{month}-#{day} #{hour}Z (skipped, already enqueued)") + {:error, _} -> IO.puts(" #{year}-#{month}-#{day} #{hour}Z (skipped)") end end IO.puts("Done.") end - @doc """ - Compute derived fields (inversion, Ri, ducts) on native profiles that lack them. - """ def native_derive(limit \\ 10_000) do start_app() - - import Ecto.Query - - alias Microwaveprop.Propagation.Duct - alias Microwaveprop.Propagation.Inversion - alias Microwaveprop.Repo - alias Microwaveprop.Weather.HrrrNativeProfile - alias Microwaveprop.Weather.ThetaE - - profiles = - HrrrNativeProfile - |> where([p], is_nil(p.bulk_richardson) and p.level_count > 2) - |> limit(^limit) - |> Repo.all() - - IO.puts("Deriving fields for #{length(profiles)} profiles...") - - count = - Enum.count(profiles, fn profile -> - p = %{ - heights_m: profile.heights_m, - temp_k: profile.temp_k, - spfh: profile.spfh, - pressure_pa: profile.pressure_pa, - u_wind_ms: profile.u_wind_ms, - v_wind_ms: profile.v_wind_ms, - tke_m2s2: profile.tke_m2s2, - level_count: profile.level_count - } - - inversion_fields = - case Inversion.find_inversion_top(p) do - {:ok, %{height_m: top_h, level_idx: top_idx, base_idx: base_idx}} -> - ri = Inversion.bulk_richardson(p, base_idx, top_idx) - shear = Inversion.shear_magnitude(p, base_idx, top_idx) - theta_e_jump = ThetaE.theta_e_jump(p, base_idx, top_idx) - - [ - inversion_top_m: top_h, - bulk_richardson: ri, - shear_at_top_ms: shear, - theta_e_jump_k: theta_e_jump - ] - - :none -> - [ - inversion_top_m: nil, - bulk_richardson: nil, - shear_at_top_ms: nil, - theta_e_jump_k: nil - ] - end - - duct_result = Duct.analyze(p) - duct_fields = [ducts: duct_result.ducts, best_duct_band_ghz: duct_result.best_duct_band_ghz] - fields = inversion_fields ++ duct_fields - - {1, _} = - HrrrNativeProfile - |> where([pr], pr.id == ^profile.id) - |> Repo.update_all(set: fields) - - true - end) - - IO.puts("Updated #{count} profiles with derived fields.") + {:ok, job} = Oban.insert(AdminTaskWorker.new(%{task: "native_derive", limit: limit})) + IO.puts("Enqueued native derive (job #{job.id})") end defp repos do @@ -263,6 +104,5 @@ defmodule Microwaveprop.Release do defp start_app do Application.ensure_all_started(@app) - Oban.pause_all_queues(Oban) end end diff --git a/lib/microwaveprop/workers/admin_task_worker.ex b/lib/microwaveprop/workers/admin_task_worker.ex new file mode 100644 index 00000000..852e0399 --- /dev/null +++ b/lib/microwaveprop/workers/admin_task_worker.ex @@ -0,0 +1,190 @@ +defmodule Microwaveprop.Workers.AdminTaskWorker do + @moduledoc """ + Oban worker for long-running admin tasks (backtests, climatology, etc). + + Enqueued via `Microwaveprop.Release` so that `eval` returns immediately + and the work runs in the normal Oban pipeline without blocking. + + Each task type is dispatched by the `"task"` arg: + + Microwaveprop.Workers.AdminTaskWorker.new(%{task: "backtest_all"}) + Microwaveprop.Workers.AdminTaskWorker.new(%{task: "backtest", feature: "naive_gradient"}) + Microwaveprop.Workers.AdminTaskWorker.new(%{task: "climatology", min_samples: 3}) + Microwaveprop.Workers.AdminTaskWorker.new(%{task: "native_derive", limit: 10000}) + """ + use Oban.Worker, queue: :admin, max_attempts: 1, unique: [period: 60] + + require Logger + + alias Microwaveprop.Backtest + alias Microwaveprop.Backtest.Features + alias Microwaveprop.Propagation.Duct + alias Microwaveprop.Propagation.Inversion + alias Microwaveprop.Repo + alias Microwaveprop.Weather.HrrrNativeProfile + alias Microwaveprop.Weather.ThetaE + + import Ecto.Query + + @impl Oban.Worker + def perform(%Oban.Job{args: %{"task" => "backtest_all"} = args}) do + sample_size = Map.get(args, "sample_size", 5000) + features = Features.all_features() + + Logger.info("AdminTask: running consolidated backtest for #{map_size(features)} features") + + results = Backtest.consolidated_report(features, sample_size: sample_size) + markdown = Backtest.to_consolidated_markdown(results) + + path = "priv/backtest_reports/consolidated.md" + File.mkdir_p!(Path.dirname(path)) + File.write!(path, markdown) + + Logger.info("AdminTask: consolidated backtest complete, wrote #{path}") + :ok + end + + def perform(%Oban.Job{args: %{"task" => "backtest", "feature" => feature_name} = args}) do + sample_size = Map.get(args, "sample_size", 5000) + fun_atom = String.to_atom(feature_name) + + unless function_exported?(Features, fun_atom, 3) do + Logger.error("AdminTask: unknown feature #{feature_name}") + {:error, "unknown feature: #{feature_name}"} + end + + feature_fun = &apply(Features, fun_atom, [&1, &2, &3]) + name = "Microwaveprop.Backtest.Features.#{feature_name}" + + Logger.info("AdminTask: running backtest for #{feature_name}") + + report = Backtest.evaluate(feature_fun, sample_size: sample_size, feature_name: name) + distance_bins = Backtest.lift_by_distance(feature_fun, sample_size: sample_size) + band_stats = Backtest.lift_by_band(feature_fun, sample_size: sample_size) + markdown = Backtest.to_markdown(report, distance_bins: distance_bins, band_stats: band_stats) + + path = "priv/backtest_reports/#{feature_name}.md" + File.mkdir_p!(Path.dirname(path)) + File.write!(path, markdown) + + Logger.info("AdminTask: backtest for #{feature_name} complete, wrote #{path}") + :ok + end + + def perform(%Oban.Job{args: %{"task" => "climatology"} = args}) do + min_samples = Map.get(args, "min_samples", 3) + + %{rows: combos} = + Repo.query!( + """ + SELECT EXTRACT(MONTH FROM valid_time)::int AS month, + EXTRACT(HOUR FROM valid_time)::int AS hour + FROM hrrr_profiles + WHERE surface_temp_c IS NOT NULL AND is_grid_point = true + GROUP BY 1, 2 + ORDER BY 1, 2 + """, + [], + timeout: 120_000 + ) + + Logger.info("AdminTask: building climatology (min_samples=#{min_samples}, #{length(combos)} batches)") + + total = + combos + |> Enum.with_index(1) + |> Enum.reduce(0, fn {[month, hour], idx}, acc -> + %{num_rows: count} = + Repo.query!( + """ + INSERT INTO hrrr_climatology (id, lat, lon, month, hour, + mean_surface_temp_c, stddev_surface_temp_c, sample_count) + SELECT gen_random_uuid(), lat, lon, + $2 AS month, + $3 AS hour, + AVG(surface_temp_c), + STDDEV_SAMP(surface_temp_c), + COUNT(*) + FROM hrrr_profiles + WHERE surface_temp_c IS NOT NULL + AND is_grid_point = true + AND EXTRACT(MONTH FROM valid_time)::int = $2 + AND EXTRACT(HOUR FROM valid_time)::int = $3 + GROUP BY lat, lon + HAVING COUNT(*) >= $1 + ON CONFLICT (lat, lon, month, hour) + DO UPDATE SET + mean_surface_temp_c = EXCLUDED.mean_surface_temp_c, + stddev_surface_temp_c = EXCLUDED.stddev_surface_temp_c, + sample_count = EXCLUDED.sample_count + """, + [min_samples, month, hour], + timeout: 300_000 + ) + + if rem(idx, 10) == 0, do: Logger.info("AdminTask: climatology [#{idx}/#{length(combos)}]") + acc + count + end) + + Logger.info("AdminTask: climatology complete, upserted #{total} records") + :ok + end + + def perform(%Oban.Job{args: %{"task" => "native_derive"} = args}) do + limit = Map.get(args, "limit", 10_000) + + profiles = + HrrrNativeProfile + |> where([p], is_nil(p.bulk_richardson) and p.level_count > 2) + |> limit(^limit) + |> Repo.all() + + Logger.info("AdminTask: deriving fields for #{length(profiles)} native profiles") + + count = + Enum.count(profiles, fn profile -> + p = %{ + heights_m: profile.heights_m, + temp_k: profile.temp_k, + spfh: profile.spfh, + pressure_pa: profile.pressure_pa, + u_wind_ms: profile.u_wind_ms, + v_wind_ms: profile.v_wind_ms, + tke_m2s2: profile.tke_m2s2, + level_count: profile.level_count + } + + inversion_fields = + case Inversion.find_inversion_top(p) do + {:ok, %{height_m: top_h, level_idx: top_idx, base_idx: base_idx}} -> + [ + inversion_top_m: top_h, + bulk_richardson: Inversion.bulk_richardson(p, base_idx, top_idx), + shear_at_top_ms: Inversion.shear_magnitude(p, base_idx, top_idx), + theta_e_jump_k: ThetaE.theta_e_jump(p, base_idx, top_idx) + ] + + :none -> + [inversion_top_m: nil, bulk_richardson: nil, shear_at_top_ms: nil, theta_e_jump_k: nil] + end + + duct_result = Duct.analyze(p) + fields = inversion_fields ++ [ducts: duct_result.ducts, best_duct_band_ghz: duct_result.best_duct_band_ghz] + + {1, _} = + HrrrNativeProfile + |> where([pr], pr.id == ^profile.id) + |> Repo.update_all(set: fields) + + true + end) + + Logger.info("AdminTask: derived fields for #{count} profiles") + :ok + end + + def perform(%Oban.Job{args: %{"task" => task}}) do + Logger.error("AdminTask: unknown task #{task}") + {:error, "unknown task: #{task}"} + end +end