Skip to content
Merged
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
2 changes: 1 addition & 1 deletion crates/cellule-runtime/docs/runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

<a id="drain"></a>
## Drain in ownership order
Expand Down
8 changes: 5 additions & 3 deletions crates/cellule-runtime/src/cell/actor/task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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;
};
Expand Down
68 changes: 59 additions & 9 deletions crates/cellule-runtime/src/fleet/pressure.rs
Original file line number Diff line number Diff line change
Expand Up @@ -160,20 +160,35 @@ 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,
used: u32,
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<Self> {
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<Self> {
if limit == 0 || requested_limit == 0 || interval_ms <= 0 {
return Err(Error::Capacity("movement budget"));
}
Ok(Self {
Expand All @@ -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<MovementPermit> {
if now_ms < self.window_started_ms {
return Err(Error::Control("movement time regressed"));
Expand All @@ -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<MovementPermit> {
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.
Expand All @@ -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
Expand All @@ -224,4 +273,5 @@ impl MovementBudget {
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct MovementPermit {
completed: bool,
requested: bool,
}
23 changes: 23 additions & 0 deletions crates/cellule-runtime/tests/fleet/pressure.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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]);
Expand Down
Loading