From ec230c5ac9fb873d44ebc3fe2aec51d39a2cee79 Mon Sep 17 00:00:00 2001 From: Thanos Makatos Date: Thu, 1 Oct 2026 09:45:35 +0000 Subject: [PATCH] Validate scale actions using IOPS instead of cumulative IOs Post-scale validation compared lifetime I/O totals as if they were rates. Example: baseline 165453639 ops, observed 171620504 (+3.7%) after the validation polls, required 173726321 (+5% of the lifetime total), a steady multi-hundred-kops workload still failed validation and was reverted. Scale-up is still proposed when CPU utilisation is above the scale up threshold; validation is meant to require at least scale_up_min_gain more IOPS afterward. Applying that percent to lifetime totals made even healthy post-scale rates look like failures. Compare IOPS over the pre-scale baseline window (from the last settled baseline reset through the scale decision; not a dedicated duration knob) to IOPS over the post-scale window governed by scale_validation_sample_polls. When validating a scale down that ran because CPU utilisation was low, do not revert just because IOPS fell with demand: if utilisation is still at or below the scale up threshold, keep the smaller pool. Enrich the existing scale up decision log with current_iops, min_target_iops, and validation_after_secs to make it clear what the minimum performance must be achieved and when for the scale up to be considered successful. AI-generated: everything, manual fixup of commit message Signed-off-by: Thanos Makatos --- src/controller.rs | 7 +- src/engines/mod.rs | 20 ++ src/engines/threshold.rs | 473 +++++++++++++++++++++++++++++++++------ 3 files changed, 432 insertions(+), 68 deletions(-) diff --git a/src/controller.rs b/src/controller.rs index 4a31cc7..a96befd 100644 --- a/src/controller.rs +++ b/src/controller.rs @@ -553,6 +553,7 @@ impl Controller { host_cpu_util: self.host_cpu_util, psi: self.psi, tick_index: self.tick_index, + scale_poll_secs: self.cfg.scale_poll_secs, }; let fleet: Vec<_> = self.instances.values().cloned().collect(); let refresh_cgroup = self.cfg.refresh_cgroup_on_each_read; @@ -847,6 +848,7 @@ impl Controller { instance.status.write().await.cooldown_until = Some(Instant::now() + Duration::from_secs_f64(self.cfg.cooldown_secs)); } + let applied_at = Instant::now(); self.engine .on_applied( &instance.id, @@ -854,6 +856,7 @@ impl Controller { action, prev_thread_count: previous_count, prev_io_count_total: previous_io_count, + applied_at, }, ) .await; @@ -862,11 +865,12 @@ impl Controller { match instance.client.set_thread_count(target).await { Ok(()) => { + let applied_at = Instant::now(); let mut status = instance.status.write().await; status.thread_count = target; if matches!(action, ScaleAction::Up(_) | ScaleAction::Down(_)) { status.cooldown_until = - Some(Instant::now() + Duration::from_secs_f64(self.cfg.cooldown_secs)); + Some(applied_at + Duration::from_secs_f64(self.cfg.cooldown_secs)); } drop(status); self.engine @@ -876,6 +880,7 @@ impl Controller { action, prev_thread_count: previous_count, prev_io_count_total: previous_io_count, + applied_at, }, ) .await; diff --git a/src/engines/mod.rs b/src/engines/mod.rs index 0c31ec0..15b8c4c 100644 --- a/src/engines/mod.rs +++ b/src/engines/mod.rs @@ -103,6 +103,22 @@ pub struct EngineTickContext { pub psi: PsiSample, /// Monotonically increasing tick sequence number. pub tick_index: u64, + /// Controller poll interval in seconds (for operator-facing durations). + pub scale_poll_secs: f64, +} + +impl Default for EngineTickContext { + fn default() -> Self { + Self { + now: Instant::now(), + min_thread_count: 2, + max_thread_count: 8, + host_cpu_util: 0.0, + psi: PsiSample::default(), + tick_index: 0, + scale_poll_secs: 10.0, + } + } } /// One VM's entry in a fleet-wide plan. @@ -161,6 +177,8 @@ pub enum AppliedOutcome { prev_thread_count: u32, /// Cumulative backend I/O ops before actuation. prev_io_count_total: u64, + /// Time when the controller recorded the outcome. + applied_at: Instant, }, /// Actuation was intentionally skipped. DryRun { @@ -170,6 +188,8 @@ pub enum AppliedOutcome { prev_thread_count: u32, /// Cumulative backend I/O ops before the proposal. prev_io_count_total: u64, + /// Time when the controller recorded the outcome. + applied_at: Instant, }, /// A controller guard deliberately prevented backend actuation. Blocked { diff --git a/src/engines/threshold.rs b/src/engines/threshold.rs index 486240c..4f9c823 100644 --- a/src/engines/threshold.rs +++ b/src/engines/threshold.rs @@ -5,7 +5,7 @@ //! Threshold engine: size a worker pool to carry the observed CPU load. -use std::{collections::HashMap, sync::Arc}; +use std::{collections::HashMap, sync::Arc, time::Instant}; use async_trait::async_trait; use serde::{Deserialize, Serialize}; @@ -150,12 +150,23 @@ struct PendingValidation { action: ScaleAction, /// Worker count to restore after a regression. baseline_thread_count: u32, - /// IOPS observed immediately before actuation. + /// IOPS observed during the pre-scale baseline window. baseline_iops: u64, + /// Cumulative I/O ops at the start of the validation window. + validation_start_io_count: u64, + /// Instant marking the start of the validation window. + validation_start_at: Instant, /// Complete observations still required before evaluation. polls_remaining: u32, } +/// Start of an IOPS-rate measurement window. +#[derive(Debug, Clone, Copy)] +struct RateWindow { + io_count: u64, + at: Instant, +} + /// Per-VM counters owned by this engine. #[derive(Clone)] struct InstanceState { @@ -168,6 +179,10 @@ struct InstanceState { low_util_polls: u32, /// Successful action waiting for its performance validation sample. pending_validation: Option, + /// Start of the post-settlement baseline window (after last validation). + baseline_window_start: Option, + /// Most recent baseline-window IOPS rate, refreshed each evaluate tick. + current_baseline_iops: u64, } impl InstanceState { @@ -177,6 +192,8 @@ impl InstanceState { instance, low_util_polls: 0, pending_validation: None, + baseline_window_start: None, + current_baseline_iops: 0, } } @@ -185,16 +202,52 @@ impl InstanceState { &mut self, action: ScaleAction, previous_thread_count: u32, - previous_iops: u64, + baseline_iops: u64, + validation_start_io_count: u64, + validation_start_at: Instant, polls_remaining: u32, ) { self.pending_validation = Some(PendingValidation { action, baseline_thread_count: previous_thread_count, - baseline_iops: previous_iops, + baseline_iops, + validation_start_io_count, + validation_start_at, polls_remaining, }); } + + /// Begin a new baseline window after settlement (or a non-validated scale). + fn reset_baseline_window(&mut self, io_count: u64, at: Instant) { + self.baseline_window_start = Some(RateWindow { io_count, at }); + self.current_baseline_iops = 0; + } + + /// Refresh the baseline-window IOPS rate ending at `(io_count, now)`. + fn refresh_baseline_rate(&mut self, io_count: u64, now: Instant) -> u64 { + let rate = match self.baseline_window_start { + None => { + self.baseline_window_start = Some(RateWindow { io_count, at: now }); + 0 + } + Some(start) => iops_rate(start.io_count, start.at, io_count, now), + }; + self.current_baseline_iops = rate; + rate + } +} + +/// Average ops/sec between two cumulative counter samples. +fn iops_rate(start_io: u64, start_at: Instant, end_io: u64, end_at: Instant) -> u64 { + let Some(elapsed) = end_at.checked_duration_since(start_at) else { + return 0; + }; + let secs = elapsed.as_secs_f64(); + if secs <= 0.0 { + return 0; + } + let delta = end_io.saturating_sub(start_io) as f64; + (delta / secs).round() as u64 } /// Scaling engine driven by CPU utilisation thresholds. @@ -266,6 +319,7 @@ impl ThresholdEngine { instance_id: &str, thread_count: u32, observed_iops: u64, + per_thread_util: f64, pending: PendingValidation, ) -> Option { let (required_iops, failed) = match pending.action { @@ -279,6 +333,11 @@ impl ThresholdEngine { ) } ScaleAction::Down(_) => { + // Scale-down was driven by low util. If util is still low, an + // IOPS collapse means demand left — keep the smaller pool. + if per_thread_util <= self.cfg.scale_up_threshold { + return None; + } let required = pending.baseline_iops as f64 * (1.0 - self.cfg.scale_down_revert_drop); ( @@ -304,6 +363,13 @@ impl ThresholdEngine { ); Some(pending.baseline_thread_count) } + + /// Seconds until a post-scale IOPS sample is evaluated. + fn validation_after_secs(&self, scale_poll_secs: f64) -> f64 { + // Countdown consumes `scale_validation_sample_polls` ticks, then the + // next tick performs the comparison. + (f64::from(self.cfg.scale_validation_sample_polls) + 1.0) * scale_poll_secs + } } /// Announce why post-action validation is rolling a worker count back. @@ -336,23 +402,41 @@ fn log_performance_revert( ); } -/// Emit one operator-visible scale decision. -fn log_scale_decision( +/// Emit one operator-visible scale-up decision with validation expectations. +fn log_scale_up_decision( instance: &Instance, per_thread_util: f64, thread_count: u32, - action: ScaleAction, + target: u32, + current_iops: u64, + min_target_iops: u64, + validation_after_secs: f64, +) { + tracing::info!( + target: "engine", + event = "scale_up", + vm = instance.to_string(), + util = per_thread_util, + thr = %format!("{thread_count}->{target}"), + current_iops, + min_target_iops, + validation_after_secs, + "threshold scaling decision" + ); +} + +/// Emit one operator-visible scale-down decision. +fn log_scale_down_decision( + instance: &Instance, + per_thread_util: f64, + thread_count: u32, + target: u32, ) { - let (direction, target) = match action { - ScaleAction::Up(target) => ("up", target), - ScaleAction::Down(target) => ("down", target), - ScaleAction::Revert(_) | ScaleAction::None => return, - }; tracing::info!( target: "engine", + event = "scale_down", vm = instance.to_string(), util = per_thread_util, - action = direction, thr = %format!("{thread_count}->{target}"), "threshold scaling decision" ); @@ -369,7 +453,7 @@ impl ScalingEngine for ThresholdEngine { } async fn evaluate(&self, instance: &Arc, context: &EngineTickContext) -> ScaleAction { - let (per_thread_util, thread_count, iops_total, throttled_recently) = { + let (per_thread_util, thread_count, io_count_total, throttled_recently) = { let status = instance.status.read().await; ( status.per_thread_util, @@ -394,16 +478,31 @@ impl ScalingEngine for ThresholdEngine { instance_state.pending_validation = Some(pending); return ScaleAction::None; } - if let Some(target) = self.validation_revert_target( + let observed_iops = iops_rate( + pending.validation_start_io_count, + pending.validation_start_at, + io_count_total, + context.now, + ); + let revert_target = self.validation_revert_target( &instance_state.instance.id, thread_count, - iops_total, + observed_iops, + per_thread_util, pending, - ) { + ); + // Whether retained or reverted, the next baseline excludes this + // validation window. Defer further scaling until the new baseline + // window has at least one full tick of samples. + instance_state.reset_baseline_window(io_count_total, context.now); + if let Some(target) = revert_target { return ScaleAction::Revert(target); } + return ScaleAction::None; } + instance_state.refresh_baseline_rate(io_count_total, context.now); + if down_target < thread_count { instance_state.low_util_polls += 1; } else { @@ -422,16 +521,33 @@ impl ScalingEngine for ThresholdEngine { } // FIXME The min seems redundant given the thread count will always be less than // or equal to the max_thread_count here. - let action = ScaleAction::Up((thread_count + 1).min(context.max_thread_count)); - log_scale_decision(instance, per_thread_util, thread_count, action); - return action; + let target = (thread_count + 1).min(context.max_thread_count); + let current_iops = instance_state.current_baseline_iops; + let (min_target_iops, validation_after_secs) = if self.cfg.scale_up_min_gain > 0.0 { + ( + ((current_iops as f64) * (1.0 + self.cfg.scale_up_min_gain)).ceil() as u64, + self.validation_after_secs(context.scale_poll_secs), + ) + } else { + (current_iops, 0.0) + }; + log_scale_up_decision( + instance, + per_thread_util, + thread_count, + target, + current_iops, + min_target_iops, + validation_after_secs, + ); + return ScaleAction::Up(target); } if down_target < thread_count && instance_state.low_util_polls >= self.cfg.scale_down_sustain_polls { let action = ScaleAction::Down(down_target); - log_scale_decision(instance, per_thread_util, thread_count, action); + log_scale_down_decision(instance, per_thread_util, thread_count, down_target); action } else { ScaleAction::None @@ -439,17 +555,19 @@ impl ScalingEngine for ThresholdEngine { } async fn on_applied(&self, instance_id: &str, outcome: AppliedOutcome) { - let (action, previous_thread_count, previous_iops) = match outcome { + let (action, previous_thread_count, previous_io_count, applied_at) = match outcome { AppliedOutcome::Success { action, prev_thread_count, prev_io_count_total, + applied_at, } | AppliedOutcome::DryRun { action, prev_thread_count, prev_io_count_total, - } => (action, prev_thread_count, prev_io_count_total), + applied_at, + } => (action, prev_thread_count, prev_io_count_total, applied_at), AppliedOutcome::Blocked { .. } | AppliedOutcome::Failed { .. } => return, }; let mut state = self.state.lock().await; @@ -465,7 +583,9 @@ impl ScalingEngine for ThresholdEngine { instance_state.start_validation( action, previous_thread_count, - previous_iops, + instance_state.current_baseline_iops, + previous_io_count, + applied_at, self.cfg.scale_validation_sample_polls, ); } @@ -476,12 +596,19 @@ impl ScalingEngine for ThresholdEngine { instance_state.start_validation( action, previous_thread_count, - previous_iops, + instance_state.current_baseline_iops, + previous_io_count, + applied_at, self.cfg.scale_validation_sample_polls, ); } - ScaleAction::Up(_) | ScaleAction::Down(_) => { + ScaleAction::Up(_) => { + instance_state.low_util_polls = 0; + instance_state.reset_baseline_window(previous_io_count, applied_at); + } + ScaleAction::Down(_) => { instance_state.low_util_polls = 0; + instance_state.reset_baseline_window(previous_io_count, applied_at); } ScaleAction::None | ScaleAction::Revert(_) => {} } @@ -514,7 +641,7 @@ mod tests { use super::{ThresholdConfig, ThresholdEngine, log_performance_revert}; use crate::{ backends::BackendClientError, - engines::{AppliedOutcome, EngineTickContext, PsiSample, ScaleAction, ScalingEngine}, + engines::{AppliedOutcome, EngineTickContext, ScaleAction, ScalingEngine}, instance::{ Instance, InstanceClient, InstancePerfSample, InstanceStatus, ThreadPoolSnapshot, }, @@ -570,33 +697,48 @@ mod tests { instance } - impl Default for EngineTickContext { - fn default() -> Self { - Self { - now: Instant::now() + Duration::from_secs(1), - min_thread_count: 2, - max_thread_count: 8, - host_cpu_util: 0.0, - psi: PsiSample::default(), - tick_index: 0, - } + fn context_at(now: Instant) -> EngineTickContext { + EngineTickContext { + now, + scale_poll_secs: 10.0, + ..EngineTickContext::default() } } - fn context() -> EngineTickContext { - EngineTickContext::default() - } - async fn set_observation( instance: &Arc, thread_count: u32, per_thread_util: f64, - iops: u64, + io_count: u64, ) { let mut status = instance.status.write().await; status.thread_count = thread_count; status.per_thread_util = per_thread_util; - status.perf.as_mut().unwrap().read_io_count = iops; + status.perf.as_mut().unwrap().read_io_count = io_count; + } + + /// Establish a baseline window of `baseline_iops` ops/sec ending at `t1`. + async fn seed_baseline_rate( + engine: &ThresholdEngine, + instance: &Arc, + thread_count: u32, + baseline_iops: u64, + t0: Instant, + t1: Instant, + ) { + engine.on_instance_added(instance).await; + // Park util exactly at the up threshold so neither up nor down fires. + let steady_util = engine.cfg.scale_up_threshold; + set_observation(instance, thread_count, steady_util, 0).await; + assert_eq!( + engine.evaluate(instance, &context_at(t0)).await, + ScaleAction::None + ); + set_observation(instance, thread_count, steady_util, baseline_iops).await; + assert_eq!( + engine.evaluate(instance, &context_at(t1)).await, + ScaleAction::None + ); } /// Test that performance-revert log line includes VM id and the @@ -644,21 +786,28 @@ mod tests { #[future] instance: Arc, ) { let instance = instance.await; - engine.on_instance_added(&instance).await; + let t0 = Instant::now(); + let t1 = t0 + Duration::from_secs(1); + let t_apply = t1; + let t_validate = t1 + Duration::from_secs(1); + + seed_baseline_rate(&engine, &instance, 5, 100_000, t0, t1).await; engine .on_applied( &instance.id, AppliedOutcome::Success { action: ScaleAction::Up(6), prev_thread_count: 5, - prev_io_count_total: 155_000, + prev_io_count_total: 100_000, + applied_at: t_apply, }, ) .await; - set_observation(&instance, 6, 0.6, 122_000).await; + // Validation window: only 50k ops over 1s → 50k IOPS (< 105k required). + set_observation(&instance, 6, 0.6, 150_000).await; assert_eq!( - engine.evaluate(&instance, &context()).await, + engine.evaluate(&instance, &context_at(t_validate)).await, ScaleAction::Revert(5) ); } @@ -672,7 +821,12 @@ mod tests { #[future] instance: Arc, ) { let instance = instance.await; - engine.on_instance_added(&instance).await; + let t0 = Instant::now(); + let t1 = t0 + Duration::from_secs(1); + let t_apply = t1; + let t_validate = t1 + Duration::from_secs(1); + + seed_baseline_rate(&engine, &instance, 3, 100_000, t0, t1).await; engine .on_applied( &instance.id, @@ -680,17 +834,130 @@ mod tests { action: ScaleAction::Up(4), prev_thread_count: 3, prev_io_count_total: 100_000, + applied_at: t_apply, }, ) .await; + // +1% ops over the validation second → still below 5% gain. set_observation(&instance, 4, 0.6, 101_000).await; assert_eq!( - engine.evaluate(&instance, &context()).await, + engine.evaluate(&instance, &context_at(t_validate)).await, ScaleAction::Revert(3) ); } + /// Test that a sufficient post-scale IOPS-rate gain keeps the scale-up. + #[rstest] + #[tokio::test] + async fn scale_up_rate_gain_keeps_scale( + engine: ThresholdEngine, + #[future] instance: Arc, + ) { + let instance = instance.await; + let t0 = Instant::now(); + let t1 = t0 + Duration::from_secs(1); + let t_apply = t1; + let t_validate = t1 + Duration::from_secs(1); + + seed_baseline_rate(&engine, &instance, 3, 100_000, t0, t1).await; + engine + .on_applied( + &instance.id, + AppliedOutcome::Success { + action: ScaleAction::Up(4), + prev_thread_count: 3, + prev_io_count_total: 100_000, + applied_at: t_apply, + }, + ) + .await; + // +10% ops/sec during validation → above 5% min gain. + set_observation(&instance, 4, 0.6, 210_000).await; + + assert_eq!( + engine.evaluate(&instance, &context_at(t_validate)).await, + ScaleAction::None + ); + } + + /// Test that scale-down validation does not revert when util stays low + /// (demand left; IOPS collapse is expected). + #[rstest] + #[tokio::test] + async fn scale_down_keeps_shrink_when_util_stays_low(#[future] instance: Arc) { + let instance = instance.await; + let engine = ThresholdEngine::new(ThresholdConfig { + scale_up_threshold: 0.8, + scale_down_revert_drop: 0.05, + scale_validation_sample_polls: 0, + scale_down_sustain_polls: 1, + ..Default::default() + }); + let t0 = Instant::now(); + let t1 = t0 + Duration::from_secs(1); + let t_apply = t1; + let t_validate = t1 + Duration::from_secs(1); + + seed_baseline_rate(&engine, &instance, 4, 100_000, t0, t1).await; + engine + .on_applied( + &instance.id, + AppliedOutcome::Success { + action: ScaleAction::Down(1), + prev_thread_count: 4, + prev_io_count_total: 100_000, + applied_at: t_apply, + }, + ) + .await; + // IOPS collapsed, but util remains below the up threshold. + set_observation(&instance, 1, 0.0, 100_001).await; + assert_eq!( + engine.evaluate(&instance, &context_at(t_validate)).await, + ScaleAction::None + ); + } + + /// Test that scale-down validation reverts when util is still high (demand + /// remains) but IOPS fell beyond the tolerated drop. + #[rstest] + #[tokio::test] + async fn scale_down_reverts_when_util_stays_high_and_iops_drop( + #[future] instance: Arc, + ) { + let instance = instance.await; + let engine = ThresholdEngine::new(ThresholdConfig { + scale_up_threshold: 0.8, + scale_down_revert_drop: 0.05, + scale_validation_sample_polls: 0, + ..Default::default() + }); + let t0 = Instant::now(); + let t1 = t0 + Duration::from_secs(1); + let t_apply = t1; + let t_validate = t1 + Duration::from_secs(1); + + seed_baseline_rate(&engine, &instance, 4, 100_000, t0, t1).await; + engine + .on_applied( + &instance.id, + AppliedOutcome::Success { + action: ScaleAction::Down(2), + prev_thread_count: 4, + prev_io_count_total: 100_000, + applied_at: t_apply, + }, + ) + .await; + // Demand still present (util above up threshold) but IOPS collapsed. + set_observation(&instance, 2, 0.95, 110_000).await; + assert_eq!( + engine.evaluate(&instance, &context_at(t_validate)).await, + ScaleAction::Revert(4) + ); + } + /// Test that percent fields serde as human percents on the wire and /// fractions in memory. #[test] @@ -707,8 +974,6 @@ mod tests { ); } - // Test that right after a scale up operation the engine does not ask for - // another scale up during the validation period. /// Test that while a prior scale-up is pending validation, further /// scale-ups are suppressed. #[rstest] @@ -718,33 +983,50 @@ mod tests { #[future] instance: Arc, ) { let instance = instance.await; + let t0 = Instant::now(); + let t1 = t0 + Duration::from_secs(1); + let t_apply = t1; - // inform the engine of the scale up - engine.on_instance_added(&instance).await; + seed_baseline_rate(&engine, &instance, 4, 100_000, t0, t1).await; engine .on_applied( &instance.id, AppliedOutcome::Success { action: ScaleAction::Up(5), prev_thread_count: 4, - prev_io_count_total: 149_000, + prev_io_count_total: 100_000, + applied_at: t_apply, }, ) .await; - set_observation(&instance, 5, 0.88, 160_000).await; + set_observation(&instance, 5, 0.88, 220_000).await; - for _ in 0..engine.cfg.scale_validation_sample_polls { - // there should be no scale up during the validation period + for i in 0..engine.cfg.scale_validation_sample_polls { assert_eq!( - engine.evaluate(&instance, &context()).await, + engine + .evaluate( + &instance, + &context_at(t_apply + Duration::from_secs(u64::from(i) + 1)) + ) + .await, ScaleAction::None ); } - // after the validation period the engine is allowed to scale up + // Validation tick settles the scale-up without proposing another action. + let t_settle = t_apply + Duration::from_secs(1); assert_eq!( - engine.evaluate(&instance, &context()).await, + engine.evaluate(&instance, &context_at(t_settle)).await, + ScaleAction::None + ); + + // Next tick with sustained high util may scale again. + set_observation(&instance, 5, 0.88, 330_000).await; + assert_eq!( + engine + .evaluate(&instance, &context_at(t_settle + Duration::from_secs(1))) + .await, ScaleAction::Up(6) ); } @@ -754,35 +1036,90 @@ mod tests { #[rstest] #[test(tokio::test)] async fn pending_validation_reverts_regressive_scale( - engine: ThresholdEngine, + _engine: ThresholdEngine, #[future] instance: Arc, ) { let instance = instance.await; - engine.on_instance_added(&instance).await; + let engine = ThresholdEngine::new(ThresholdConfig { + scale_up_threshold: 0.8, + scale_up_min_gain: 0.05, + scale_validation_sample_polls: 2, + ..Default::default() + }); + let t0 = Instant::now(); + let t1 = t0 + Duration::from_secs(1); + let t_apply = t1; + + seed_baseline_rate(&engine, &instance, 4, 100_000, t0, t1).await; engine .on_applied( &instance.id, AppliedOutcome::Success { action: ScaleAction::Up(5), prev_thread_count: 4, - prev_io_count_total: 149_000, + prev_io_count_total: 100_000, + applied_at: t_apply, }, ) .await; set_observation(&instance, 5, 0.88, 140_000).await; - for _ in 0..engine.cfg.scale_validation_sample_polls { + for i in 0..engine.cfg.scale_validation_sample_polls { assert_eq!( - engine.evaluate(&instance, &context()).await, + engine + .evaluate( + &instance, + &context_at(t_apply + Duration::from_secs(u64::from(i) + 1)) + ) + .await, ScaleAction::None ); } + // 40k ops over ~3s ≈ 13k IOPS ≪ 105k required. assert_eq!( - engine.evaluate(&instance, &context()).await, + engine + .evaluate(&instance, &context_at(t_apply + Duration::from_secs(3))) + .await, ScaleAction::Revert(4) ); } + /// Test that scale-up decision logs current/min IOPS and validation delay. + #[rstest] + #[tokio::test] + async fn scale_up_log_includes_iops_targets( + engine: ThresholdEngine, + #[future] instance: Arc, + ) { + let instance = instance.await; + let output = Arc::new(Mutex::new(Vec::new())); + let writer_output = Arc::clone(&output); + let subscriber = tracing_subscriber::fmt() + .without_time() + .with_ansi(false) + .with_target(false) + .with_max_level(tracing::Level::INFO) + .with_writer(move || BufferWriter(Arc::clone(&writer_output))) + .finish(); + let _guard = tracing::subscriber::set_default(subscriber); + + let t0 = Instant::now(); + let t1 = t0 + Duration::from_secs(1); + seed_baseline_rate(&engine, &instance, 2, 100_000, t0, t1).await; + set_observation(&instance, 2, 0.9, 100_000).await; + + assert_eq!( + engine.evaluate(&instance, &context_at(t1)).await, + ScaleAction::Up(3) + ); + + let rendered = String::from_utf8(output.lock().unwrap().clone()).unwrap(); + assert!(rendered.contains("event=\"scale_up\"") || rendered.contains("event=scale_up")); + assert!(rendered.contains("current_iops=100000")); + assert!(rendered.contains("min_target_iops=105000")); + assert!(rendered.contains("validation_after_secs=10")); + } + /// Test that a missing `threshold.json` still builds an engine with /// defaults. #[test] @@ -801,7 +1138,9 @@ mod tests { async fn evaluate_holds_when_util_below_thresholds(#[future] instance: Arc) { let instance = instance.await; let engine = ThresholdEngine::new(ThresholdConfig::default()); - let action = engine.evaluate(&instance, &context()).await; + let action = engine + .evaluate(&instance, &context_at(Instant::now())) + .await; assert_eq!(action, ScaleAction::None); } }