diff --git a/crates/cellule-runtime/docs/runtime.md b/crates/cellule-runtime/docs/runtime.md index 4c626cb..9d5f6c2 100644 --- a/crates/cellule-runtime/docs/runtime.md +++ b/crates/cellule-runtime/docs/runtime.md @@ -490,7 +490,7 @@ The node scheduler: | Enter shedding | 80 percent on any dimension | Sustained for one second | | Return to normal | Below 60 percent | Sustained for one second | -**Eviction.** While shedding, the actor starts one bounded eviction per sample through the same movement budget and victim selection a transfer uses, so a hot node releases settled Cells instead of admitting work it cannot hold. A brief spike never triggers a move. +**Eviction.** While shedding, the actor starts one bounded eviction per sample using its two-in-flight, two-completions-per-second pressure budget and the same victim selection a transfer uses. A brief spike never triggers a move. Explicit `release_idle_cell` calls have a separate 32-in-flight, 32-completions-per-second budget, so a cold repository scan cannot spend the automatic pressure-shedding allowance. Both paths still require the same generation, settled-work and authoritative-release checks; the larger requested-release budget is an admission bound, not a measured sustainable object-store rate. ## Drain in ownership order diff --git a/crates/cellule-runtime/src/cell/actor/task.rs b/crates/cellule-runtime/src/cell/actor/task.rs index 2174cb4..79265d5 100644 --- a/crates/cellule-runtime/src/cell/actor/task.rs +++ b/crates/cellule-runtime/src/cell/actor/task.rs @@ -24,7 +24,9 @@ pub(super) async fn run( Ok(classifier) => classifier, Err(_) => return, }; - let mut movement = match MovementBudget::new(2, 1_000) { + // Requested releases have their own bounded capacity so a cold scan does + // not consume the two-per-second automatic pressure-shedding allowance. + let mut movement = match MovementBudget::with_requested_limit(2, 32, 1_000) { Ok(budget) => budget, Err(_) => return, }; @@ -246,7 +248,7 @@ fn classify_pressure_sample( let state = pressure.observe(sample)?; telemetry.pressure_state(state); if matches!(state, PressureState::Shedding | PressureState::Critical) - && movement_permits.len() < 2 + && movement.in_flight() < 2 { let _ = start_bounded_evictions( 1, @@ -709,7 +711,7 @@ pub(super) fn handle_message( let _ = reply.send(Err(Error::CellDraining)); return; } - let Ok(mut permit) = movement.try_start(unix_millis()) else { + let Ok(mut permit) = movement.try_start_requested(unix_millis()) else { let _ = reply.send(Err(Error::Capacity("movement budget"))); return; }; diff --git a/crates/cellule-runtime/src/fleet/pressure.rs b/crates/cellule-runtime/src/fleet/pressure.rs index 8c41b9f..682fa9c 100644 --- a/crates/cellule-runtime/src/fleet/pressure.rs +++ b/crates/cellule-runtime/src/fleet/pressure.rs @@ -160,7 +160,7 @@ pub enum MovementKind { Receive, } -/// Synchronous concurrency/rate budget for pressure movement. +/// Synchronous concurrency/rate budgets for pressure and requested movement. #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub struct MovementBudget { limit: u32, @@ -168,12 +168,27 @@ pub struct MovementBudget { interval_ms: i64, window_started_ms: i64, completed_in_window: u32, + requested_limit: u32, + requested_used: u32, + requested_window_started_ms: i64, + requested_completed_in_window: u32, } impl MovementBudget { - /// Creates a budget of `limit` movements per `interval_ms` window. + /// Creates equal budgets for pressure and requested movement. pub fn new(limit: u32, interval_ms: i64) -> Result { - if limit == 0 || interval_ms <= 0 { + Self::with_requested_limit(limit, limit, interval_ms) + } + + /// Sets independent budgets for automatic pressure movement and explicit + /// release requests. Each limit bounds both in-flight work and completions + /// per interval; neither class can consume the other's reserved capacity. + pub fn with_requested_limit( + limit: u32, + requested_limit: u32, + interval_ms: i64, + ) -> Result { + if limit == 0 || requested_limit == 0 || interval_ms <= 0 { return Err(Error::Capacity("movement budget")); } Ok(Self { @@ -182,11 +197,14 @@ impl MovementBudget { interval_ms, window_started_ms: 0, completed_in_window: 0, + requested_limit, + requested_used: 0, + requested_window_started_ms: 0, + requested_completed_in_window: 0, }) } - /// Reserves one movement, failing when the window is spent or time ran - /// backwards. + /// Reserves one automatic pressure movement. pub fn try_start(&mut self, now_ms: i64) -> Result { if now_ms < self.window_started_ms { return Err(Error::Control("movement time regressed")); @@ -199,7 +217,32 @@ impl MovementBudget { return Err(Error::Capacity("movement budget")); } self.used += 1; - Ok(MovementPermit { completed: false }) + Ok(MovementPermit { + completed: false, + requested: false, + }) + } + + /// Reserves one explicit release without consuming pressure-shedding + /// capacity. The caller still performs its generation and preflight checks. + pub fn try_start_requested(&mut self, now_ms: i64) -> Result { + if now_ms < self.requested_window_started_ms { + return Err(Error::Control("movement time regressed")); + } + if now_ms.saturating_sub(self.requested_window_started_ms) >= self.interval_ms { + self.requested_window_started_ms = now_ms; + self.requested_completed_in_window = 0; + } + if self.requested_used >= self.requested_limit + || self.requested_completed_in_window >= self.requested_limit + { + return Err(Error::Capacity("movement budget")); + } + self.requested_used += 1; + Ok(MovementPermit { + completed: false, + requested: true, + }) } /// Finishes a reservation: frees its slot and counts it in the window. @@ -208,11 +251,17 @@ impl MovementBudget { return; } permit.completed = true; - self.used = self.used.saturating_sub(1); - self.completed_in_window = self.completed_in_window.saturating_add(1); + if permit.requested { + self.requested_used = self.requested_used.saturating_sub(1); + self.requested_completed_in_window = + self.requested_completed_in_window.saturating_add(1); + } else { + self.used = self.used.saturating_sub(1); + self.completed_in_window = self.completed_in_window.saturating_add(1); + } } - /// Returns the reservations that have not completed yet. + /// Returns automatic pressure movements that have not completed yet. #[must_use] pub const fn in_flight(self) -> u32 { self.used @@ -224,4 +273,5 @@ impl MovementBudget { #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub struct MovementPermit { completed: bool, + requested: bool, } diff --git a/crates/cellule-runtime/tests/fleet/pressure.rs b/crates/cellule-runtime/tests/fleet/pressure.rs index cef96cb..4d3ffe3 100644 --- a/crates/cellule-runtime/tests/fleet/pressure.rs +++ b/crates/cellule-runtime/tests/fleet/pressure.rs @@ -97,6 +97,29 @@ fn movement_budget_bounds_concurrency_and_rate() { assert_eq!(budget.in_flight(), 1); } +#[test] +fn requested_releases_do_not_spend_pressure_shedding_budget() { + assert!(MovementBudget::with_requested_limit(2, 0, 100).is_err()); + let mut budget = MovementBudget::with_requested_limit(1, 2, 100).unwrap(); + let mut requested_a = budget.try_start_requested(0).unwrap(); + let mut requested_b = budget.try_start_requested(0).unwrap(); + assert!(budget.try_start_requested(0).is_err()); + + let mut pressure = budget.try_start(0).unwrap(); + assert_eq!(budget.in_flight(), 1); + assert!(budget.try_start(0).is_err()); + budget.complete(&mut requested_a); + budget.complete(&mut requested_a); + budget.complete(&mut requested_b); + assert!(budget.try_start_requested(0).is_err()); + + budget.complete(&mut pressure); + assert_eq!(budget.in_flight(), 0); + assert!(budget.try_start(0).is_err()); + let _requested = budget.try_start_requested(100).unwrap(); + let _pressure = budget.try_start(100).unwrap(); +} + #[tokio::test] async fn node_reports_every_pressure_tier_it_classifies() { let session = SessionId::from_bytes([123; 16]);