From ee7809e17b53aa593c20c839d60d022497423f32 Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Tue, 21 Apr 2026 18:02:44 -0500 Subject: [PATCH] hardening: Rust worker circuit-breaker + delete wall-clock sleep from beacon test MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit **Rust workers exit after 30 consecutive claim_next failures.** Both `prop-grid-rs` and `hrrr-point-worker` ran a `while` loop that slept 10s and retried forever on `claim_next` errors. If Postgres went unreachable, pods would stay alive but useless — no circuit breaker, no liveness probe, just an infinite stream of error logs. Now we increment a consecutive-failure counter and exit with an error once it hits 30 (~5 minutes of DB unavailability). k8s restart policy then cycles the pod, which both triggers an operator alert and gives us a fresh sqlx pool on the way back up. A successful claim (Some or None) resets the counter so one flaky query doesn't poison it. **BeaconMonitorsTest: drop the Process.sleep(1100).** The test relied on wall-clock gaps to disambiguate `inserted_at` values on the `:utc_datetime` (second-precision) timestamps — one full second of sleep per run, and on shared CI it could still flake when both inserts landed in the same second. Replaced with an explicit `Repo.update_all` that back-dates the first row by one second. 1.1s faster, deterministic. --- .../prop_grid_rs/src/bin/hrrr_point_worker.rs | 21 ++++++++++++++- rust/prop_grid_rs/src/bin/worker.rs | 27 ++++++++++++++++++- test/microwaveprop/beacon_monitors_test.exs | 16 ++++++++++- 3 files changed, 61 insertions(+), 3 deletions(-) diff --git a/rust/prop_grid_rs/src/bin/hrrr_point_worker.rs b/rust/prop_grid_rs/src/bin/hrrr_point_worker.rs index d58ce795..fff65ad2 100644 --- a/rust/prop_grid_rs/src/bin/hrrr_point_worker.rs +++ b/rust/prop_grid_rs/src/bin/hrrr_point_worker.rs @@ -12,6 +12,9 @@ use prop_grid_rs::{db, fetcher::HrrrClient, hrrr_points, telemetry}; use tracing::{error, info, warn}; const IDLE_SLEEP: Duration = Duration::from_secs(10); +/// ~5 minutes of consecutive Postgres-unavailable failures before we +/// exit and let the kubelet restart us. See worker.rs for rationale. +const MAX_CONSECUTIVE_CLAIM_FAILURES: u32 = 30; #[tokio::main] async fn main() -> Result<(), Box> { @@ -35,9 +38,11 @@ async fn main() -> Result<(), Box> { info!(pg_conns = max_conn, "hrrr-point-worker started"); + let mut claim_failures: u32 = 0; while !shutdown.load(Ordering::Relaxed) { match db::claim_next_hrrr_task(&pool).await { Ok(Some(task)) => { + claim_failures = 0; let task_id = task.id; let valid_time = task.valid_time; let points = task.points.clone(); @@ -68,10 +73,24 @@ async fn main() -> Result<(), Box> { } } Ok(None) => { + claim_failures = 0; tokio::time::sleep(IDLE_SLEEP).await; } Err(e) => { - error!(error = %e, "claim_next_hrrr_task failed; backing off"); + claim_failures = claim_failures.saturating_add(1); + error!( + error = %e, + consecutive_failures = claim_failures, + "claim_next_hrrr_task failed; backing off" + ); + if claim_failures >= MAX_CONSECUTIVE_CLAIM_FAILURES { + error!( + consecutive_failures = claim_failures, + "too many consecutive claim failures — exiting so kubelet restarts the pod" + ); + shutdown.store(true, Ordering::Relaxed); + break; + } tokio::time::sleep(IDLE_SLEEP).await; } } diff --git a/rust/prop_grid_rs/src/bin/worker.rs b/rust/prop_grid_rs/src/bin/worker.rs index 6fdec80a..41ba6bea 100644 --- a/rust/prop_grid_rs/src/bin/worker.rs +++ b/rust/prop_grid_rs/src/bin/worker.rs @@ -98,6 +98,13 @@ async fn main() -> Result<(), Box> { Ok(()) } +/// Give up and exit the process after this many consecutive +/// `claim_next` failures. At `IDLE_SLEEP = 10s` that's ~5 minutes of +/// a broken Postgres before the pod restart-loops and lets k8s + +/// operators notice. Below this threshold the worker stays alive — +/// a short DB blip shouldn't cycle the pod. +const MAX_CONSECUTIVE_CLAIM_FAILURES: u32 = 30; + async fn worker_loop( worker_id: usize, pool: sqlx::PgPool, @@ -105,6 +112,7 @@ async fn worker_loop( scores_dir: PathBuf, shutdown: Arc, ) { + let mut claim_failures: u32 = 0; while !shutdown.load(Ordering::Relaxed) { // Analysis tasks (kind='analysis', f00) get priority — a /map // user values the current hour more than +1..+18h. Forecast is @@ -117,6 +125,7 @@ async fn worker_loop( match claimed { Ok(Some(task)) => { + claim_failures = 0; let task_id = task.id; let _guard = metrics::InFlightGuard::new(); let started = Instant::now(); @@ -188,10 +197,26 @@ async fn worker_loop( } } Ok(None) => { + claim_failures = 0; tokio::time::sleep(IDLE_SLEEP).await; } Err(e) => { - error!(worker_id, error = %e, "claim_next failed; backing off"); + claim_failures = claim_failures.saturating_add(1); + error!( + worker_id, + error = %e, + consecutive_failures = claim_failures, + "claim_next failed; backing off" + ); + if claim_failures >= MAX_CONSECUTIVE_CLAIM_FAILURES { + error!( + worker_id, + consecutive_failures = claim_failures, + "too many consecutive claim failures — exiting so kubelet restarts the pod" + ); + shutdown.store(true, Ordering::Relaxed); + return; + } tokio::time::sleep(IDLE_SLEEP).await; } } diff --git a/test/microwaveprop/beacon_monitors_test.exs b/test/microwaveprop/beacon_monitors_test.exs index 7244b968..bede86aa 100644 --- a/test/microwaveprop/beacon_monitors_test.exs +++ b/test/microwaveprop/beacon_monitors_test.exs @@ -41,8 +41,22 @@ defmodule Microwaveprop.BeaconMonitorsTest do user = user_fixture() other = user_fixture() + # beacon_monitors.inserted_at is :utc_datetime (second precision). + # Back-date the first row one second via a direct update so the + # ordering assertion is deterministic without burning ~1s on a + # Process.sleep — this is a slow-test-suite footgun that also + # makes the test flaky on a shared CI runner where the two + # inserts can land in the same wall-clock second. {:ok, first} = BeaconMonitors.create_monitor(user, %{"name" => "First"}) - Process.sleep(1100) + + backdated = DateTime.add(first.inserted_at, -1, :second) + + {1, _} = + Microwaveprop.Repo.update_all( + from(m in BeaconMonitor, where: m.id == ^first.id), + set: [inserted_at: backdated] + ) + {:ok, second} = BeaconMonitors.create_monitor(user, %{"name" => "Second"}) {:ok, _other} = BeaconMonitors.create_monitor(other, %{"name" => "Nope"})