diff --git a/io-thread-controller.d/engines/threshold.json b/io-thread-controller.d/engines/threshold.json index 032c643..f023663 100644 --- a/io-thread-controller.d/engines/threshold.json +++ b/io-thread-controller.d/engines/threshold.json @@ -5,5 +5,7 @@ "scale_up_min_gain_percent": 5, "scale_down_revert_drop_percent": 3, "scale_validation_sample_polls": 3, - "block_scale_up_when_cgroup_throttled": true + "block_scale_up_when_cgroup_throttled": true, + "post_revert_up_grace_polls": 1, + "post_revert_up_recovery_polls": 2 } diff --git a/src/engines/threshold.rs b/src/engines/threshold.rs index 4f9c823..221f989 100644 --- a/src/engines/threshold.rs +++ b/src/engines/threshold.rs @@ -75,6 +75,13 @@ pub struct ThresholdConfig { /// Suppress scale-up after the VM's cgroup reports new throttled CPU time. #[serde(default = "default_true")] pub block_scale_up_when_cgroup_throttled: bool, + /// Polls after a failed scale-up revert during which metrics are ignored. + #[serde(default = "default_post_revert_up_grace_polls")] + pub post_revert_up_grace_polls: u32, + /// Polls after grace spent measuring a settled baseline while scale-up is + /// still blocked. + #[serde(default = "default_post_revert_up_recovery_polls")] + pub post_revert_up_recovery_polls: u32, } fn default_scale_up_threshold() -> f64 { @@ -101,6 +108,12 @@ fn default_scale_validation_sample_polls() -> u32 { fn default_true() -> bool { true } +fn default_post_revert_up_grace_polls() -> u32 { + 1 +} +fn default_post_revert_up_recovery_polls() -> u32 { + 2 +} impl Default for ThresholdConfig { fn default() -> Self { @@ -112,6 +125,8 @@ impl Default for ThresholdConfig { scale_down_revert_drop: default_scale_down_revert_drop(), scale_validation_sample_polls: default_scale_validation_sample_polls(), block_scale_up_when_cgroup_throttled: true, + post_revert_up_grace_polls: default_post_revert_up_grace_polls(), + post_revert_up_recovery_polls: default_post_revert_up_recovery_polls(), } } } @@ -167,6 +182,17 @@ struct RateWindow { at: Instant, } +/// Hold that blocks scale-up after a failed scale-up was reverted. +#[derive(Debug, Clone, Copy)] +enum PostRevertUpHold { + /// Metrics ignored while the pool settles after the revert. + Grace { polls_remaining: u32 }, + /// Scale-up still blocked while a settled baseline rate is measured. + Recovery { + polls_remaining: u32, + window_start: RateWindow, + }, +} /// Per-VM counters owned by this engine. #[derive(Clone)] struct InstanceState { @@ -183,6 +209,9 @@ struct InstanceState { baseline_window_start: Option, /// Most recent baseline-window IOPS rate, refreshed each evaluate tick. current_baseline_iops: u64, + /// After a failed scale-up revert, suppress further ups through grace then + /// recovery while a settled baseline is established. + post_revert_up_hold: Option, } impl InstanceState { @@ -194,6 +223,7 @@ impl InstanceState { pending_validation: None, baseline_window_start: None, current_baseline_iops: 0, + post_revert_up_hold: None, } } @@ -235,6 +265,66 @@ impl InstanceState { self.current_baseline_iops = rate; rate } + + /// Start grace→recovery hold after a failed scale-up revert. + fn begin_post_revert_up_hold(&mut self, grace_polls: u32) { + self.baseline_window_start = None; + self.current_baseline_iops = 0; + self.post_revert_up_hold = Some(PostRevertUpHold::Grace { + polls_remaining: grace_polls, + }); + } + + /// Advance grace/recovery state. Returns true while scale-up must stay + /// blocked. + fn advance_post_revert_up_hold( + &mut self, + io_count: u64, + now: Instant, + recovery_polls: u32, + ) -> bool { + let Some(hold) = self.post_revert_up_hold.take() else { + return false; + }; + match hold { + PostRevertUpHold::Grace { polls_remaining } => { + if polls_remaining > 0 { + self.post_revert_up_hold = Some(PostRevertUpHold::Grace { + polls_remaining: polls_remaining - 1, + }); + return true; + } + // Grace finished: anchor the recovery baseline window. + if recovery_polls == 0 { + self.reset_baseline_window(io_count, now); + return false; + } + self.post_revert_up_hold = Some(PostRevertUpHold::Recovery { + polls_remaining: recovery_polls, + window_start: RateWindow { io_count, at: now }, + }); + true + } + PostRevertUpHold::Recovery { + polls_remaining, + window_start, + } => { + let rate = iops_rate(window_start.io_count, window_start.at, io_count, now); + self.current_baseline_iops = rate; + if polls_remaining > 1 { + self.post_revert_up_hold = Some(PostRevertUpHold::Recovery { + polls_remaining: polls_remaining - 1, + window_start, + }); + return true; + } + // Final recovery sample: publish the window and allow scale-up. + self.baseline_window_start = Some(window_start); + self.current_baseline_iops = rate; + false + } + } + } } /// Average ops/sec between two cumulative counter samples. @@ -442,6 +532,17 @@ fn log_scale_down_decision( ); } +/// Announce that scale-up is held after a failed-up revert. +fn log_post_revert_up_hold(instance_id: &str, grace_polls: u32, recovery_polls: u32) { + tracing::info!( + target: "engine", + event = "post_revert_up_hold", + vm = %instance_id, + grace_polls, + recovery_polls, + "blocking scale-up after failed scale-up revert" + ); +} #[async_trait] impl ScalingEngine for ThresholdEngine { fn name(&self) -> &'static str { @@ -484,6 +585,7 @@ impl ScalingEngine for ThresholdEngine { io_count_total, context.now, ); + let failed_up = matches!(pending.action, ScaleAction::Up(_)); let revert_target = self.validation_revert_target( &instance_state.instance.id, thread_count, @@ -496,12 +598,29 @@ impl ScalingEngine for ThresholdEngine { // 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 { + if failed_up { + instance_state.begin_post_revert_up_hold(self.cfg.post_revert_up_grace_polls); + log_post_revert_up_hold( + &instance.id, + self.cfg.post_revert_up_grace_polls, + self.cfg.post_revert_up_recovery_polls, + ); + } return ScaleAction::Revert(target); } return ScaleAction::None; } - instance_state.refresh_baseline_rate(io_count_total, context.now); + let block_scale_up = instance_state.advance_post_revert_up_hold( + io_count_total, + context.now, + self.cfg.post_revert_up_recovery_polls, + ); + // Grace ignores metrics; recovery updates the rate inside advance. + // Once the hold ends (or if there was none), refresh the normal baseline. + if !block_scale_up { + instance_state.refresh_baseline_rate(io_count_total, context.now); + } if down_target < thread_count { instance_state.low_util_polls += 1; @@ -509,7 +628,9 @@ impl ScalingEngine for ThresholdEngine { instance_state.low_util_polls = 0; } - if thread_count < context.max_thread_count && per_thread_util > self.cfg.scale_up_threshold + if !block_scale_up + && thread_count < context.max_thread_count + && per_thread_util > self.cfg.scale_up_threshold { if self.cfg.block_scale_up_when_cgroup_throttled && throttled_recently { tracing::info!( @@ -593,6 +714,7 @@ impl ScalingEngine for ThresholdEngine { // Successful actuation consumes the sustained low-utilisation // run; subsequent downscaling must establish a fresh run. instance_state.low_util_polls = 0; + instance_state.post_revert_up_hold = None; instance_state.start_validation( action, previous_thread_count, @@ -608,6 +730,7 @@ impl ScalingEngine for ThresholdEngine { } ScaleAction::Down(_) => { instance_state.low_util_polls = 0; + instance_state.post_revert_up_hold = None; instance_state.reset_baseline_window(previous_io_count, applied_at); } ScaleAction::None | ScaleAction::Revert(_) => {} @@ -881,6 +1004,135 @@ mod tests { ); } + /// Test that a failed scale-up enters grace then recovery: scale-up is + /// blocked while a settled baseline is measured, then the next up uses + /// that recovery-window IOPS (not a depressed post-revert blip). + #[rstest] + #[tokio::test] + async fn failed_scale_up_grace_then_recovery_baseline(#[future] instance: Arc) { + let instance = instance.await; + let engine = ThresholdEngine::new(ThresholdConfig { + scale_up_threshold: 0.8, + scale_up_min_gain: 0.05, + scale_validation_sample_polls: 0, + post_revert_up_grace_polls: 1, + post_revert_up_recovery_polls: 2, + ..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, 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; + set_observation(&instance, 4, 0.95, 101_000).await; + assert_eq!( + engine.evaluate(&instance, &context_at(t_validate)).await, + ScaleAction::Revert(3) + ); + + // Grace: util blip / depressed IOPS must not produce a scale-up. + set_observation(&instance, 3, 0.4, 150_000).await; + assert_eq!( + engine + .evaluate(&instance, &context_at(t_validate + Duration::from_secs(1))) + .await, + ScaleAction::None + ); + + // Recovery tick 1: enter recovery window (anchor), still blocked. + set_observation(&instance, 3, 0.95, 150_000).await; + assert_eq!( + engine + .evaluate(&instance, &context_at(t_validate + Duration::from_secs(2))) + .await, + ScaleAction::None + ); + + // Recovery tick 2: measure ~100k IOPS over 1s, still blocked. + set_observation(&instance, 3, 0.95, 250_000).await; + assert_eq!( + engine + .evaluate(&instance, &context_at(t_validate + Duration::from_secs(3))) + .await, + ScaleAction::None + ); + + // Hold complete: scale-up allowed; baseline is the recovery-window rate + // (~100k), so min_target reflects that — not the grace-period blip. + set_observation(&instance, 3, 0.95, 350_000).await; + assert_eq!( + engine + .evaluate(&instance, &context_at(t_validate + Duration::from_secs(4))) + .await, + ScaleAction::Up(4) + ); + { + let state = engine.state.lock().await; + let instance_state = state.get(&instance.id).unwrap(); + // Recovery anchored at 150k @ +2s through 350k @ +4s → 100k IOPS. + assert_eq!(instance_state.current_baseline_iops, 100_000); + assert!(instance_state.post_revert_up_hold.is_none()); + } + } + + /// Test that scale-down remains allowed during the post-revert up hold. + #[rstest] + #[tokio::test] + async fn scale_down_allowed_during_post_revert_up_hold(#[future] instance: Arc) { + let instance = instance.await; + let engine = ThresholdEngine::new(ThresholdConfig { + scale_up_threshold: 0.8, + scale_up_min_gain: 0.05, + scale_validation_sample_polls: 0, + scale_down_sustain_polls: 1, + post_revert_up_grace_polls: 2, + post_revert_up_recovery_polls: 2, + ..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, 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; + set_observation(&instance, 4, 0.95, 101_000).await; + assert_eq!( + engine.evaluate(&instance, &context_at(t_validate)).await, + ScaleAction::Revert(3) + ); + + // During grace, low util can still scale down. + set_observation(&instance, 3, 0.0, 101_000).await; + assert_eq!( + engine + .evaluate(&instance, &context_at(t_validate + Duration::from_secs(1))) + .await, + ScaleAction::Down(2) + ); + } /// Test that scale-down validation does not revert when util stays low /// (demand left; IOPS collapse is expected). #[rstest]