Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion io-thread-controller.d/engines/threshold.json
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
256 changes: 254 additions & 2 deletions src/engines/threshold.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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 {
Expand All @@ -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(),
}
}
}
Expand Down Expand Up @@ -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 {
Expand All @@ -183,6 +209,9 @@ struct InstanceState {
baseline_window_start: Option<RateWindow>,
/// 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<PostRevertUpHold>,
}

impl InstanceState {
Expand All @@ -194,6 +223,7 @@ impl InstanceState {
pending_validation: None,
baseline_window_start: None,
current_baseline_iops: 0,
post_revert_up_hold: None,
}
}

Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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,
Expand All @@ -496,20 +598,39 @@ 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;
} else {
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!(
Expand Down Expand Up @@ -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,
Expand All @@ -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(_) => {}
Expand Down Expand Up @@ -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<Instance>) {
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<Instance>) {
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]
Expand Down
Loading