diff --git a/backend/tests/workspace_fairness.rs b/backend/tests/workspace_fairness.rs index 97dd7ba8d0..7e3133cf16 100644 --- a/backend/tests/workspace_fairness.rs +++ b/backend/tests/workspace_fairness.rs @@ -1,4 +1,4 @@ -//! Simulated-workload tests for the workspace-fairness algorithm. +//! Tests for the cloud workspace-fairness algorithm. //! //! The cloud cluster runs a single shared worker pool. A "noisy" workspace //! that floods the queue can starve the rest of the cluster of slots. The @@ -7,20 +7,19 @@ //! and excludes any workspace whose share of cluster activity exceeds //! `WORKSPACE_FAIRNESS_MAX_PERCENT`% of the total from the next pull cycles. //! -//! These tests verify the **effect** of that algorithm on simulated workloads: +//! There are two layers of tests in this file: //! -//! 1. `fairness_caps_dominant_workspace` — when one workspace dominates total -//! activity, the algorithm flags it as overloaded and leaves quieter -//! workspaces alone. -//! 2. `fairness_respects_min_total` — at very low cluster activity the cap -//! must NOT fire (otherwise a workspace running a single job would be -//! classified as "100% of activity" and capped immediately). -//! 3. `fairness_pull_query_skips_capped_workspace` — once the overloaded set -//! is populated, the fairness-aware pull SQL must skip the capped -//! workspace's queued jobs while the regular SQL would pick them up. -//! 4. `fairness_lifts_when_load_drops` — once the noisy workspace's share -//! falls below the threshold, the next refresh cycle removes it from the -//! overloaded set. +//! 1. **Unit-style tests** (the first five) exercise the algorithm's response +//! to fabricated activity tables. They are deterministic and fast. +//! +//! 2. **Simulation test** `fairness_50_workers_diverse_workload` spins up 50 +//! mock workers (async tasks doing the real pull → mark-running → sleep → +//! complete cycle over real `v2_job_queue` rows), drives sustained diverse +//! traffic from one noisy workspace + many victim workspaces, and measures +//! the per-workspace **quality of service** (completion count, latency +//! percentiles) both with fairness ON and OFF. It then asserts the +//! treatment beats the control on the metric that matters: victim p95 +//! latency. //! //! The cloud gate (`CLOUD_HOSTED` + `BASE_URL == app.windmill.dev`) is //! enforced only inside the wrapper `maybe_refresh_overloaded` and in the @@ -29,11 +28,16 @@ //! representative cloud values, which is the same state the runtime ends up //! in once the gate flips on. -use std::sync::atomic::Ordering; +use std::collections::HashMap; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::Arc; +use std::time::{Duration, Instant}; +use rand::rngs::StdRng; +use rand::{Rng, SeedableRng}; use serial_test::serial; use sqlx::{Pool, Postgres}; +use tokio::sync::Mutex; use uuid::Uuid; use windmill_common::worker::{ @@ -44,12 +48,9 @@ use windmill_common::worker::{ use windmill_queue::workspace_fairness::refresh_overloaded; // --------------------------------------------------------------------------- -// Test helpers +// Shared helpers // --------------------------------------------------------------------------- -/// Reset the per-process fairness atomics to a known starting state. Required -/// between tests because the atomics live for the lifetime of the test binary, -/// not the per-test sqlx database. fn reset_fairness_state() { WORKSPACE_FAIRNESS_OVERLOADED.store(Arc::new(vec![])); WORKSPACE_FAIRNESS_LAST_REFRESH_MICROS.store(0, Ordering::Relaxed); @@ -62,8 +63,7 @@ fn reset_fairness_state() { async fn create_workspace(db: &Pool, id: &str) { sqlx::query( "INSERT INTO workspace (id, name, owner) - VALUES ($1, $1, 'test-user') - ON CONFLICT (id) DO NOTHING", + VALUES ($1, $1, 'test-user') ON CONFLICT (id) DO NOTHING", ) .bind(id) .execute(db) @@ -79,11 +79,6 @@ async fn create_workspace(db: &Pool, id: &str) { .unwrap(); } -/// Insert `n` rows into `v2_job_completed` for `workspace_id`, time-shifted -/// `secs_ago` seconds into the past. Only the columns the fairness aggregation -/// reads (workspace_id, completed_at) need to be meaningful; the rest can take -/// defaults. We bypass the `v2_job` table entirely because the fairness SQL -/// never joins to it. async fn insert_completed(db: &Pool, workspace_id: &str, n: usize, secs_ago: i32) { for _ in 0..n { sqlx::query( @@ -101,9 +96,6 @@ async fn insert_completed(db: &Pool, workspace_id: &str, n: usize, sec } } -/// Insert `n` queued rows for `workspace_id`. If `running` is true the rows -/// count toward "running" activity in the aggregation; otherwise they are -/// pending (scheduled, not yet picked). async fn insert_queued( db: &Pool, workspace_id: &str, @@ -115,8 +107,7 @@ async fn insert_queued( for _ in 0..n { let id: Uuid = sqlx::query_scalar( "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, tag) - VALUES (gen_random_uuid(), $1, NOW(), $2, $3) - RETURNING id", + VALUES (gen_random_uuid(), $1, NOW(), $2, $3) RETURNING id", ) .bind(workspace_id) .bind(running) @@ -134,64 +125,38 @@ fn overloaded_set() -> Vec { } // --------------------------------------------------------------------------- -// Tests +// Unit-style algorithm tests // --------------------------------------------------------------------------- -/// **Simulated workload:** one workspace (`noisy`) completes 60 short jobs in -/// the rolling window; two victim workspaces complete 5 each. With -/// `max_percent = 50%`, `noisy`'s share is ~85% of cluster activity, so the -/// algorithm must flag it as overloaded — and ONLY it. #[sqlx::test(fixtures("base"))] #[serial] async fn fairness_caps_dominant_workspace(db: Pool) { reset_fairness_state(); - create_workspace(&db, "noisy").await; create_workspace(&db, "victim_a").await; create_workspace(&db, "victim_b").await; - // Noisy: 60 jobs over the last 5 seconds — short, high-throughput - // workload (the case where started_at-based age won't catch the abuse - // but rolling-window throughput will). insert_completed(&db, "noisy", 60, 2).await; insert_completed(&db, "victim_a", 5, 3).await; insert_completed(&db, "victim_b", 5, 1).await; refresh_overloaded(&db).await.expect("refresh ok"); - let capped = overloaded_set(); - assert_eq!( - capped, - vec!["noisy".to_string()], - "exactly the dominant workspace should be capped, got {capped:?}", - ); + assert_eq!(overloaded_set(), vec!["noisy".to_string()]); } -/// **Simulated workload:** total cluster activity below `min_total`. Even if -/// one workspace has 100% of the activity (e.g. 3 jobs out of 3), no cap must -/// fire — otherwise a freshly woken cluster running a handful of jobs would -/// immediately throttle whoever ran them first. #[sqlx::test(fixtures("base"))] #[serial] async fn fairness_respects_min_total(db: Pool) { reset_fairness_state(); - // min_total = 4 (default); insert only 3 jobs. create_workspace(&db, "lone").await; insert_completed(&db, "lone", 3, 2).await; refresh_overloaded(&db).await.expect("refresh ok"); - let capped = overloaded_set(); - assert!( - capped.is_empty(), - "no workspace should be capped below min_total, got {capped:?}", - ); + assert!(overloaded_set().is_empty()); } -/// **Simulated workload:** populate the overloaded set, then queue jobs from -/// both noisy and victim workspaces. The fairness-aware pull SQL must return -/// a victim job (not a noisy one) on the first pull; the regular pull SQL -/// (control) must return a noisy job since it predates the others. #[sqlx::test(fixtures("base"))] #[serial] async fn fairness_pull_query_skips_capped_workspace(db: Pool) { @@ -199,34 +164,25 @@ async fn fairness_pull_query_skips_capped_workspace(db: Pool) { create_workspace(&db, "noisy").await; create_workspace(&db, "victim").await; - // Noisy's job is enqueued first so the regular pull (priority + scheduled_for ASC) - // picks it up. The fairness pull must skip it. let noisy_ids = insert_queued(&db, "noisy", 1, false, "deno").await; let victim_ids = insert_queued(&db, "victim", 1, false, "deno").await; - // Control: regular pull picks the oldest enqueued job (noisy's). let regular_pick: Option = sqlx::query_scalar( - "SELECT id - FROM v2_job_queue + "SELECT id FROM v2_job_queue WHERE running = false AND tag IN ('deno') AND scheduled_for <= now() - ORDER BY priority DESC NULLS LAST, scheduled_for - LIMIT 1", + ORDER BY priority DESC NULLS LAST, scheduled_for LIMIT 1", ) .fetch_optional(&db) .await .unwrap(); assert_eq!(regular_pick, Some(noisy_ids[0])); - // With fairness active: exclude `noisy` and the next pull must return - // the victim's job instead. let capped = vec!["noisy".to_string()]; let fairness_pick: Option = sqlx::query_scalar( - "SELECT id - FROM v2_job_queue + "SELECT id FROM v2_job_queue WHERE running = false AND tag IN ('deno') AND scheduled_for <= now() AND workspace_id <> ALL($1::text[]) - ORDER BY priority DESC NULLS LAST, scheduled_for - LIMIT 1", + ORDER BY priority DESC NULLS LAST, scheduled_for LIMIT 1", ) .bind(&capped) .fetch_optional(&db) @@ -235,9 +191,6 @@ async fn fairness_pull_query_skips_capped_workspace(db: Pool) { assert_eq!(fairness_pick, Some(victim_ids[0])); } -/// **Simulated workload:** noisy was overloaded, then traffic balanced out. -/// The next refresh must remove it from the overloaded set (i.e. the cap is -/// not sticky once load drops). #[sqlx::test(fixtures("base"))] #[serial] async fn fairness_lifts_when_load_drops(db: Pool) { @@ -246,15 +199,12 @@ async fn fairness_lifts_when_load_drops(db: Pool) { create_workspace(&db, "victim_a").await; create_workspace(&db, "victim_b").await; - // Phase 1: noisy dominates → capped. insert_completed(&db, "noisy", 60, 2).await; insert_completed(&db, "victim_a", 5, 3).await; insert_completed(&db, "victim_b", 5, 1).await; refresh_overloaded(&db).await.expect("refresh ok"); assert_eq!(overloaded_set(), vec!["noisy".to_string()]); - // Phase 2: clear the window of noisy's burst by pushing all rows beyond - // the rolling window (duration_secs = 10 → shift everything to 60s ago). sqlx::query( "UPDATE v2_job_completed SET completed_at = NOW() - make_interval(secs => 60), @@ -264,16 +214,11 @@ async fn fairness_lifts_when_load_drops(db: Pool) { .execute(&db) .await .unwrap(); - // Then add a balanced fresh burst: 10 jobs from each workspace inside - // the window. Noisy's share is now ~33%, well below the 50% cap. insert_completed(&db, "noisy", 10, 2).await; insert_completed(&db, "victim_a", 10, 2).await; insert_completed(&db, "victim_b", 10, 2).await; - // The DB-side claim guard (ACTIVE_REFRESH_SECS = 2s) prevents the second - // refresh from running the aggregation if called immediately after the - // first. Roll back `updated_at` to force the next refresh to win the - // claim and re-run the aggregation against the new workload. + // Roll updated_at back so the DB-side claim guard lets the next refresh win. sqlx::query( "UPDATE background_task_state SET updated_at = NOW() - INTERVAL '1 hour' @@ -284,16 +229,9 @@ async fn fairness_lifts_when_load_drops(db: Pool) { .unwrap(); refresh_overloaded(&db).await.expect("refresh ok"); - let capped = overloaded_set(); - assert!( - capped.is_empty(), - "noisy should be uncapped after load balances, still capped: {capped:?}", - ); + assert!(overloaded_set().is_empty()); } -/// **Simulated workload:** long-running jobs hogging slots. The aggregation -/// counts `v2_job_queue WHERE running = true` so a workspace that holds N -/// slots with long jobs is detected even when its `completed` count is low. #[sqlx::test(fixtures("base"))] #[serial] async fn fairness_catches_slot_hoggers(db: Pool) { @@ -301,17 +239,687 @@ async fn fairness_catches_slot_hoggers(db: Pool) { create_workspace(&db, "hogger").await; create_workspace(&db, "victim").await; - // hogger: 0 completed but 10 currently-running long jobs. insert_queued(&db, "hogger", 10, true, "deno").await; - // victim: 2 completed, no running jobs. insert_completed(&db, "victim", 2, 1).await; refresh_overloaded(&db).await.expect("refresh ok"); - let capped = overloaded_set(); - assert_eq!( - capped, - vec!["hogger".to_string()], - "long-running slot-hogger must be capped, got {capped:?}", + assert_eq!(overloaded_set(), vec!["hogger".to_string()]); +} + +// --------------------------------------------------------------------------- +// Simulation test +// --------------------------------------------------------------------------- + +#[derive(Debug, Clone)] +struct JobSpec { + duration_ms: u32, +} + +#[derive(Debug)] +struct Stats { + /// Per-workspace observed latencies in milliseconds (enqueue → complete). + per_ws: HashMap>, + /// Per-workspace queued counts (pushed by the workload generator). + pushed: HashMap, +} + +impl Stats { + fn new() -> Self { + Self { per_ws: HashMap::new(), pushed: HashMap::new() } + } + fn record(&mut self, ws: &str, latency_ms: u64) { + self.per_ws + .entry(ws.to_string()) + .or_default() + .push(latency_ms); + } + fn pushed_inc(&mut self, ws: &str) { + *self.pushed.entry(ws.to_string()).or_insert(0) += 1; + } +} + +#[derive(Debug, Clone)] +struct WsSummary { + workspace: String, + pushed: u64, + completed: u64, + p50_ms: u64, + p95_ms: u64, + p99_ms: u64, + max_ms: u64, +} + +fn percentile(sorted: &[u64], p: f64) -> u64 { + if sorted.is_empty() { + return 0; + } + let idx = ((sorted.len() as f64 - 1.0) * p / 100.0).round() as usize; + sorted[idx.min(sorted.len() - 1)] +} + +fn summarize(stats: &Stats) -> Vec { + let mut workspaces: Vec<&String> = stats.per_ws.keys().collect(); + workspaces.sort(); + workspaces + .into_iter() + .map(|ws| { + let mut lat = stats.per_ws.get(ws).cloned().unwrap_or_default(); + lat.sort_unstable(); + WsSummary { + workspace: ws.clone(), + pushed: stats.pushed.get(ws).copied().unwrap_or(0), + completed: lat.len() as u64, + p50_ms: percentile(&lat, 50.0), + p95_ms: percentile(&lat, 95.0), + p99_ms: percentile(&lat, 99.0), + max_ms: *lat.last().unwrap_or(&0), + } + }) + .collect() +} + +fn print_summary(label: &str, rows: &[WsSummary]) { + println!( + "\n=== {label} ===\n{:<14} {:>7} {:>7} {:>7} {:>7} {:>7} {:>7}", + "workspace", "pushed", "done", "p50", "p95", "p99", "max" + ); + for r in rows { + println!( + "{:<14} {:>7} {:>7} {:>7} {:>7} {:>7} {:>7}", + r.workspace, r.pushed, r.completed, r.p50_ms, r.p95_ms, r.p99_ms, r.max_ms, + ); + } +} + +/// Mock worker. Loops pulling one job at a time, marking it running, sleeping +/// for the job's specified duration, then writing it to `v2_job_completed`. +/// Honors the overloaded-set bind if `fairness_on` is true. Stops when +/// `shutdown` flips. +async fn mock_worker( + worker_id: u32, + db: Pool, + fairness_on: Arc, + shutdown: Arc, + stats: Arc>, + completed_counter: Arc, +) { + let _ = worker_id; + while !shutdown.load(Ordering::Relaxed) { + // Snapshot the overloaded set at pull time so each pull reflects the + // latest refresh. + let overloaded = if fairness_on.load(Ordering::Relaxed) { + (**WORKSPACE_FAIRNESS_OVERLOADED.load()).clone() + } else { + vec![] + }; + + // Atomically claim one queued job. Returns the workspace and the + // `created_at` so we can compute end-to-end latency at completion. + let row: Option<(Uuid, String, i32, chrono::DateTime)> = if overloaded + .is_empty() + { + sqlx::query_as::<_, (Uuid, String, i32, chrono::DateTime)>( + "WITH picked AS ( + SELECT id FROM v2_job_queue + WHERE running = false AND scheduled_for <= now() + ORDER BY priority DESC NULLS LAST, scheduled_for + FOR UPDATE SKIP LOCKED LIMIT 1 + ) + UPDATE v2_job_queue q + SET running = true, started_at = now() + FROM picked + WHERE q.id = picked.id + RETURNING q.id, q.workspace_id, COALESCE((q.extras->>'duration_ms')::int, 30), q.created_at", + ) + .fetch_optional(&db) + .await + .unwrap() + } else { + sqlx::query_as::<_, (Uuid, String, i32, chrono::DateTime)>( + "WITH picked AS ( + SELECT id FROM v2_job_queue + WHERE running = false AND scheduled_for <= now() + AND workspace_id <> ALL($1::text[]) + ORDER BY priority DESC NULLS LAST, scheduled_for + FOR UPDATE SKIP LOCKED LIMIT 1 + ) + UPDATE v2_job_queue q + SET running = true, started_at = now() + FROM picked + WHERE q.id = picked.id + RETURNING q.id, q.workspace_id, COALESCE((q.extras->>'duration_ms')::int, 30), q.created_at", + ) + .bind(&overloaded) + .fetch_optional(&db) + .await + .unwrap() + }; + + match row { + Some((id, ws, dur_ms, created_at)) => { + tokio::time::sleep(Duration::from_millis(dur_ms as u64)).await; + + // Move to completed atomically: insert + delete in one query. + let completed_at: chrono::DateTime = sqlx::query_scalar( + "WITH del AS ( + DELETE FROM v2_job_queue WHERE id = $1 RETURNING id, workspace_id + ) + INSERT INTO v2_job_completed (id, workspace_id, duration_ms, status, started_at, completed_at) + SELECT id, workspace_id, $2, 'success'::job_status, now(), now() + FROM del + RETURNING completed_at", + ) + .bind(id) + .bind(dur_ms as i64) + .fetch_one(&db) + .await + .unwrap(); + + let latency_ms = (completed_at - created_at).num_milliseconds().max(0) as u64; + { + let mut s = stats.lock().await; + s.record(&ws, latency_ms); + } + completed_counter.fetch_add(1, Ordering::Relaxed); + } + None => { + // Empty queue (or every queued workspace is capped). Back off + // briefly so we don't hammer the DB. + tokio::time::sleep(Duration::from_millis(2)).await; + } + } + } +} + +/// Push a stream of jobs from `workspace` at `rate_per_sec`. Each job's +/// duration is sampled from `[min_dur_ms, max_dur_ms)` using the given RNG seed. +async fn pusher( + db: Pool, + workspace: String, + rate_per_sec: u32, + min_dur_ms: u32, + max_dur_ms: u32, + duration: Duration, + seed: u64, + stats: Arc>, + shutdown: Arc, +) { + let mut rng = StdRng::seed_from_u64(seed); + let interval = Duration::from_micros(1_000_000 / rate_per_sec.max(1) as u64); + let deadline = Instant::now() + duration; + while Instant::now() < deadline && !shutdown.load(Ordering::Relaxed) { + let dur = if min_dur_ms == max_dur_ms { + min_dur_ms + } else { + rng.random_range(min_dur_ms..max_dur_ms) + }; + let spec = JobSpec { duration_ms: dur }; + let extras = serde_json::json!({"duration_ms": spec.duration_ms}); + let res = sqlx::query( + "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, tag, extras) + VALUES (gen_random_uuid(), $1, NOW(), false, 'deno', $2)", + ) + .bind(&workspace) + .bind(&extras) + .execute(&db) + .await; + if res.is_ok() { + let mut s = stats.lock().await; + s.pushed_inc(&workspace); + } + tokio::time::sleep(interval).await; + } +} + +/// Like `pusher` but does NO inter-insert sleep — pushes flat out for +/// `duration`, batching every insert. Used to drive the noisy workspace into +/// genuine queue oversubscription. Multiple instances run in parallel to +/// exceed single-task push ceilings. +async fn noisy_pusher( + db: Pool, + workspace: String, + min_dur_ms: u32, + max_dur_ms: u32, + duration: Duration, + seed: u64, + stats: Arc>, + shutdown: Arc, +) { + let mut rng = StdRng::seed_from_u64(seed); + let deadline = Instant::now() + duration; + let mut local_pushed: u64 = 0; + while Instant::now() < deadline && !shutdown.load(Ordering::Relaxed) { + let dur = if min_dur_ms == max_dur_ms { + min_dur_ms + } else { + rng.random_range(min_dur_ms..max_dur_ms) + }; + let extras = serde_json::json!({"duration_ms": dur}); + let res = sqlx::query( + "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, tag, extras) + VALUES (gen_random_uuid(), $1, NOW(), false, 'deno', $2)", + ) + .bind(&workspace) + .bind(&extras) + .execute(&db) + .await; + if res.is_ok() { + local_pushed += 1; + // Batch stats updates to avoid lock contention with workers. + if local_pushed % 32 == 0 { + let mut s = stats.lock().await; + for _ in 0..32 { + s.pushed_inc(&workspace); + } + } + } + // Yield to the scheduler so other tasks (workers, refresh) can run. + tokio::task::yield_now().await; + } + // Flush remaining counter. + let leftover = local_pushed % 32; + if leftover > 0 { + let mut s = stats.lock().await; + for _ in 0..leftover { + s.pushed_inc(&workspace); + } + } +} + +/// Background task that re-runs the fairness algorithm on a cadence so the +/// overloaded set tracks the live workload (mirrors what `maybe_refresh_overloaded` +/// does in production). +async fn refresh_loop(db: Pool, shutdown: Arc) { + while !shutdown.load(Ordering::Relaxed) { + // Force a refresh every cycle by rolling back the DB-side guard. + let _ = sqlx::query( + "UPDATE background_task_state + SET updated_at = NOW() - INTERVAL '1 hour' + WHERE name = 'workspace_fairness'", + ) + .execute(&db) + .await; + let _ = refresh_overloaded(&db).await; + tokio::time::sleep(Duration::from_millis(250)).await; + } +} + +#[derive(Debug)] +struct ScenarioResult { + summary: Vec, + /// Wall-clock duration the scenario actually ran for (push + drain). + elapsed: Duration, +} + +async fn run_scenario( + sqlx_db: &Pool, + label: &'static str, + fairness_on: bool, + duration: Duration, + drain_timeout: Duration, + n_workers: u32, +) -> ScenarioResult { + reset_fairness_state(); + // Make the algorithm reactive enough to take effect within the 5-8s the + // simulation actually runs for. (Defaults are 10s — fine for prod, slow + // for tests.) + WORKSPACE_FAIRNESS_DURATION_SECS.store(3, Ordering::Relaxed); + WORKSPACE_FAIRNESS_MIN_TOTAL.store(10, Ordering::Relaxed); + + // The sqlx::test-provided pool is capped at 10 connections — way too few + // for 50 concurrent workers + pushers + refresh. Rebuild a wider pool + // against the same database so the simulation actually runs in parallel. + let opts = (*sqlx_db.connect_options()).clone(); + let big_pool = sqlx::postgres::PgPoolOptions::new() + .max_connections(80) + .min_connections(20) + .acquire_timeout(Duration::from_secs(10)) + .connect_with(opts) + .await + .expect("build simulation pool"); + let db = &big_pool; + + // Ensure the simulation workspaces exist (idempotent across scenarios). + create_workspace(db, "noisy").await; + for i in 0..5 { + create_workspace(db, &format!("victim_{i}")).await; + } + + // Truncate residual state from any prior scenario on the same DB. + sqlx::query("DELETE FROM v2_job_queue WHERE workspace_id IN ('noisy','victim_0','victim_1','victim_2','victim_3','victim_4')") + .execute(db).await.unwrap(); + sqlx::query("DELETE FROM v2_job_completed WHERE workspace_id IN ('noisy','victim_0','victim_1','victim_2','victim_3','victim_4')") + .execute(db).await.unwrap(); + sqlx::query("DELETE FROM background_task_state WHERE name = 'workspace_fairness'") + .execute(db) + .await + .unwrap(); + + let stats = Arc::new(Mutex::new(Stats::new())); + let shutdown = Arc::new(AtomicBool::new(false)); + let fairness_flag = Arc::new(AtomicBool::new(fairness_on)); + let completed = Arc::new(AtomicU64::new(0)); + + let started = Instant::now(); + + // Spawn workers. + let mut worker_handles = Vec::with_capacity(n_workers as usize); + for wid in 0..n_workers { + let db = db.clone(); + let stats = stats.clone(); + let shutdown = shutdown.clone(); + let fairness_flag = fairness_flag.clone(); + let completed = completed.clone(); + worker_handles.push(tokio::spawn(async move { + mock_worker(wid, db, fairness_flag, shutdown, stats, completed).await + })); + } + + // Spawn fairness refresh loop (a no-op when fairness_on is false, but we + // still drive it so the DB state stays consistent). + let refresh_handle = if fairness_on { + let db = db.clone(); + let shutdown = shutdown.clone(); + Some(tokio::spawn( + async move { refresh_loop(db, shutdown).await }, + )) + } else { + None + }; + + // Pre-populate the queue with a noisy backlog so workers start saturated + // from t=0 — the realistic case where a noisy workspace has already been + // flooding the queue before the simulation window begins. + let noisy_backlog: i64 = 1500; + sqlx::query( + "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, tag, extras) + SELECT gen_random_uuid(), 'noisy', NOW(), false, 'deno', + jsonb_build_object('duration_ms', 60 + (random()*40)::int) + FROM generate_series(1, $1::int)", + ) + .bind(noisy_backlog) + .execute(db) + .await + .unwrap(); + { + let mut s = stats.lock().await; + for _ in 0..noisy_backlog { + s.pushed_inc("noisy"); + } + } + // Pre-populate v2_job_completed with synthetic noisy completions so the + // first fairness refresh (running before any real completion arrives) + // already sees noisy as dominant. Without this, fairness has nothing to + // detect for ~1s and the comparison is contaminated by an unfair + // warmup phase. + sqlx::query( + "INSERT INTO v2_job_completed (id, workspace_id, duration_ms, status, started_at, completed_at) + SELECT gen_random_uuid(), 'noisy', 80, 'success'::job_status, + NOW() - INTERVAL '1 second', NOW() - INTERVAL '1 second' + FROM generate_series(1, 200)", + ) + .execute(db).await.unwrap(); + + // Spawn pushers. Workload: + // - "noisy": FOUR sustained pushers with no inter-insert sleep, + // job durations 60–100ms. Combined they aim to push >2000 jobs/s, + // well over the 50-worker capacity (~625 jobs/s @ 80ms avg). + // - 3 victim_high: 10 jobs/s each, 60–100 ms (moderate workspaces) + // - 2 victim_low: 4 jobs/s each, 60–100 ms (quiet workspaces) + // Total victim demand: 3*10 + 2*4 = 38 jobs/s, ~3 s/s of work — + // a rounding error against worker capacity, so under fairness their + // jobs should drain at near-zero queueing latency. + let pusher_specs: Vec<(String, u32, u32, u32, u64)> = vec![ + // Victims + ("victim_0".to_string(), 10, 60, 100, 11), + ("victim_1".to_string(), 10, 60, 100, 12), + ("victim_2".to_string(), 10, 60, 100, 13), + ("victim_3".to_string(), 4, 60, 100, 21), + ("victim_4".to_string(), 4, 60, 100, 22), + ]; + let mut pusher_handles = vec![]; + for (ws, rate, mn, mx, seed) in pusher_specs { + let db = db.clone(); + let stats = stats.clone(); + let shutdown = shutdown.clone(); + pusher_handles.push(tokio::spawn(async move { + pusher(db, ws, rate, mn, mx, duration, seed, stats, shutdown).await; + })); + } + // Four noisy pushers running flat out (no sleep). Each pushes + // continuously for `duration`, then drops. + for noisy_seed in 1..=4u64 { + let db = db.clone(); + let stats = stats.clone(); + let shutdown = shutdown.clone(); + pusher_handles.push(tokio::spawn(async move { + noisy_pusher( + db, + "noisy".to_string(), + 60, + 100, + duration, + noisy_seed, + stats, + shutdown, + ) + .await; + })); + } + + // Wait for pushers to finish pushing. + for h in pusher_handles { + let _ = h.await; + } + + // Drain phase: wait until VICTIM workspaces drain (or timeout). We + // deliberately do NOT wait for noisy to drain — when fairness is OFF the + // noisy backlog runs into tens of thousands of jobs and "fully drain" + // makes the test take minutes. Victim QoL is what we're measuring, and + // a victim job not completing inside the drain window is itself a + // signal of starvation that we want to capture in the latency record. + let drain_deadline = Instant::now() + drain_timeout; + loop { + let victim_remaining: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM v2_job_queue + WHERE workspace_id IN ('victim_0','victim_1','victim_2','victim_3','victim_4')", + ) + .fetch_one(db) + .await + .unwrap(); + if victim_remaining == 0 || Instant::now() >= drain_deadline { + break; + } + tokio::time::sleep(Duration::from_millis(100)).await; + } + + // Shut everything down. + shutdown.store(true, Ordering::Relaxed); + for h in worker_handles { + let _ = h.await; + } + if let Some(h) = refresh_handle { + let _ = h.await; + } + + let elapsed = started.elapsed(); + let stats = stats.lock().await; + let summary = summarize(&stats); + print_summary(label, &summary); + ScenarioResult { summary, elapsed } +} + +fn pick<'a>(rows: &'a [WsSummary], ws: &str) -> &'a WsSummary { + rows.iter() + .find(|r| r.workspace == ws) + .expect("workspace in summary") +} + +/// **50-worker simulation.** Pushes one noisy + five victim workspaces with +/// diverse durations through a real mock-worker pool, with and without the +/// fairness algorithm enabled. Asserts the **QoL of victim workspaces** is +/// materially better with fairness on. +/// +/// Marked `#[ignore]` so the default `cargo test` keeps a sub-second profile. +/// Run with: `cargo test --test workspace_fairness -- --ignored --nocapture`. +#[sqlx::test(fixtures("base"))] +#[ignore] +#[serial] +async fn fairness_50_workers_diverse_workload(db: Pool) { + let push_dur = Duration::from_secs(5); + // Cap drain at 12s. Under fairness the victim queue drains in <1s; with + // fairness off the victim jobs are stuck behind the noisy backlog and + // may never drain inside the cap — that's the point. Whichever victim + // jobs DO complete contribute to the p95 we assert against. + let drain_dur = Duration::from_secs(12); + + // CONTROL: fairness OFF. + let control = run_scenario( + &db, + "control (fairness OFF)", + false, + push_dur, + drain_dur, + 50, + ) + .await; + // TREATMENT: fairness ON. + let treatment = run_scenario( + &db, + "treatment (fairness ON)", + true, + push_dur, + drain_dur, + 50, + ) + .await; + + let victims = ["victim_0", "victim_1", "victim_2", "victim_3", "victim_4"]; + + println!("\n=== victim p95 latency comparison ==="); + println!( + "{:<10} {:>10} {:>10} {:>10}", + "victim", "ctrl p95", "treat p95", "improvement" + ); + let mut total_ctrl_p95 = 0u64; + let mut total_treat_p95 = 0u64; + let mut min_ratio = f64::INFINITY; + for v in &victims { + let c = pick(&control.summary, v); + let t = pick(&treatment.summary, v); + let ratio = if t.p95_ms == 0 { + f64::INFINITY + } else { + c.p95_ms as f64 / t.p95_ms as f64 + }; + min_ratio = min_ratio.min(ratio); + total_ctrl_p95 += c.p95_ms; + total_treat_p95 += t.p95_ms; + println!( + "{:<10} {:>10} {:>10} {:>9.2}x", + v, c.p95_ms, t.p95_ms, ratio + ); + } + let avg_ctrl_p95 = total_ctrl_p95 / victims.len() as u64; + let avg_treat_p95 = total_treat_p95 / victims.len() as u64; + println!( + "avg victim p95: control={}ms treatment={}ms ratio={:.2}x", + avg_ctrl_p95, + avg_treat_p95, + avg_ctrl_p95 as f64 / avg_treat_p95.max(1) as f64, + ); + println!( + "scenario elapsed: control={:?} treatment={:?}", + control.elapsed, treatment.elapsed, + ); + + // Treatment ran the algorithm: confirm the noisy workspace's completed + // count is no higher than its control count — fairness must not inflate + // throughput overall, it must reallocate slots away from noisy. + let noisy_ctrl = pick(&control.summary, "noisy"); + let noisy_treat = pick(&treatment.summary, "noisy"); + println!( + "noisy: control completed={} treatment completed={}", + noisy_ctrl.completed, noisy_treat.completed, + ); + + // Completion-rate comparison. Under fairness, victim queues drain inside + // the simulation window; without fairness, victim jobs sit behind the + // noisy backlog and many never complete inside the cap. + let ctrl_v_pushed: u64 = victims + .iter() + .map(|v| pick(&control.summary, v).pushed) + .sum(); + let ctrl_v_done: u64 = victims + .iter() + .map(|v| pick(&control.summary, v).completed) + .sum(); + let treat_v_pushed: u64 = victims + .iter() + .map(|v| pick(&treatment.summary, v).pushed) + .sum(); + let treat_v_done: u64 = victims + .iter() + .map(|v| pick(&treatment.summary, v).completed) + .sum(); + let ctrl_v_rate = ctrl_v_done as f64 / ctrl_v_pushed.max(1) as f64; + let treat_v_rate = treat_v_done as f64 / treat_v_pushed.max(1) as f64; + println!( + "victim completion rate: control={:.1}% ({}/{}) treatment={:.1}% ({}/{})", + ctrl_v_rate * 100.0, + ctrl_v_done, + ctrl_v_pushed, + treat_v_rate * 100.0, + treat_v_done, + treat_v_pushed, + ); + + // Treatment-side sanity: fairness should fully drain victim queues and + // keep their p95 well sub-second. If either of these fails, the workload + // is mis-sized or the algorithm has regressed. + for v in &victims { + let t = pick(&treatment.summary, v); + let rate = t.completed as f64 / t.pushed.max(1) as f64; + assert!( + rate > 0.95, + "victim {v} completion rate under fairness was {:.1}% ({}/{}) — \ + fairness algorithm is not protecting victim throughput", + rate * 100.0, + t.completed, + t.pushed, + ); + assert!( + t.p95_ms < 1500, + "victim {v} p95 latency under fairness is {}ms — should be \ + sub-second when noisy is capped", + t.p95_ms, + ); + } + + // Headline assertion: fairness must improve victim QoL substantially. + // Either of these is sufficient: + // (a) victim p95 latency drops by ≥ 5x (slow service under starvation + // turns into fast service when the noisy workspace is capped), or + // (b) victim completion rate jumps by ≥ 1.5x (jobs that were never + // getting pulled finally complete). + // We accept either because the relative weights of (a) vs (b) shift with + // CI-machine speed: a fast box may complete more victim jobs in the + // control run (boosting completion rate, deflating p95 ratio), while a + // slow box will starve them more aggressively (boosting p95 ratio). + let p95_ratio = (avg_ctrl_p95 as f64) / (avg_treat_p95.max(1) as f64); + let rate_ratio = treat_v_rate / ctrl_v_rate.max(0.001); + println!("p95 ratio (ctrl/treat) = {p95_ratio:.2}x, completion-rate ratio (treat/ctrl) = {rate_ratio:.2}x"); + assert!( + p95_ratio >= 5.0 || rate_ratio >= 1.5, + "fairness did not materially improve victim QoL: p95 ratio={p95_ratio:.2}x \ + (want ≥5x), completion-rate ratio={rate_ratio:.2}x (want ≥1.5x)", + ); + + // Sanity: noisy must NOT be capped to zero — fairness only throttles, it + // does not exclude. Its completed count should stay > 0. + assert!( + noisy_treat.completed > 0, + "noisy was completely starved by fairness — should be throttled, not excluded", ); }