From 8b0776af3a34469e20b99408edf241c74dfa428e Mon Sep 17 00:00:00 2001 From: Cody Kickertz Date: Tue, 8 Sep 2026 18:30:56 -0500 Subject: [PATCH 1/4] feat(sched): add selective admission retirement --- crates/sched/src/lib.rs | 327 +++++++++++++++++++++++++++++++++++++--- 1 file changed, 309 insertions(+), 18 deletions(-) diff --git a/crates/sched/src/lib.rs b/crates/sched/src/lib.rs index 1ceb0cb..5f40805 100644 --- a/crates/sched/src/lib.rs +++ b/crates/sched/src/lib.rs @@ -545,27 +545,32 @@ impl Scheduler { } self.revoked = true; for admission in self.admissions.values_mut() { - admission.state = match &admission.state { - AdmissionState::Reserved => AdmissionState::Draining { - resident: None, - loading: false, - }, - AdmissionState::Loading => AdmissionState::Draining { - resident: None, - loading: true, - }, - AdmissionState::Resident(resident) | AdmissionState::InUse(resident) => { - AdmissionState::Draining { - resident: Some(resident.clone()), - loading: false, - } - } - AdmissionState::Draining { .. } | AdmissionState::Evicting(_) => continue, - }; + Self::transition_to_draining(admission); } Ok(()) } + /// Request that one admission stop accepting new uses and release safely. + /// + /// The request denies new uses immediately. A reserved lease releases through + /// [`Self::poll_command`]; a loaded resident releases only after its trusted + /// eviction acknowledgement. Repeated requests for an active or already + /// released local admission are idempotent. This does not revoke the grant or + /// affect unrelated admissions. + /// + /// # Errors + /// + /// Returns an error without mutation for a foreign or stale ticket. + pub fn request_retirement(&mut self, ticket: &AdmissionTicket) -> Result<(), SchedulerError> { + self.ensure_local(&ticket.brand, "admission ticket")?; + self.ensure_current(ticket.generation)?; + let Some(admission) = self.admissions.get_mut(&ticket.admission_id) else { + return Ok(()); + }; + Self::transition_to_draining(admission); + Ok(()) + } + /// Install a new resource grant after the revoked prior grant fully drains. /// /// # Errors @@ -716,7 +721,15 @@ impl Scheduler { location: error_location(), }, )?; - admission.state = if self.revoked { + let drain_after_load = self.revoked + || matches!( + &admission.state, + AdmissionState::Draining { + resident: None, + loading: true, + } + ); + admission.state = if drain_after_load { AdmissionState::Draining { resident: Some(resident), loading: false, @@ -920,6 +933,27 @@ impl Scheduler { } } + fn transition_to_draining(admission: &mut Admission) { + let state = match &admission.state { + AdmissionState::Reserved => AdmissionState::Draining { + resident: None, + loading: false, + }, + AdmissionState::Loading => AdmissionState::Draining { + resident: None, + loading: true, + }, + AdmissionState::Resident(resident) | AdmissionState::InUse(resident) => { + AdmissionState::Draining { + resident: Some(resident.clone()), + loading: false, + } + } + AdmissionState::Draining { .. } | AdmissionState::Evicting(_) => return, + }; + admission.state = state; + } + fn next_command(&self) -> Option<(u64, OperationKind)> { self.admissions .iter() @@ -1577,6 +1611,263 @@ mod tests { Ok(()) } + #[test] + fn retirement_of_reserved_admissions_is_idempotent_and_recovers_the_bound() + -> Result<(), SchedulerError> { + let initial = request(&format!("[{}]", workload("initial", 4)), 20)?; + let mut scheduler = Scheduler::new(&initial, SchedulerLimits::try_new(1, 1, 1)?)?; + for profile in ["first", "second", "third"] { + let current = request(&format!("[{}]", workload(profile, 4)), 20)?; + let ticket = one_ticket(&mut scheduler, ¤t)?; + scheduler.request_retirement(&ticket)?; + scheduler.request_retirement(&ticket)?; + assert!( + matches!( + scheduler.begin_use(&ticket), + Err(SchedulerError::AdmissionNotResident { .. }) + ), + "retirement closes a reserved admission before a load can issue" + ); + assert!( + matches!(scheduler.poll_command()?, PollOutcome::Progressed), + "a reserved admission releases without a physical command" + ); + scheduler.request_retirement(&ticket)?; + assert!( + scheduler.admissions.is_empty(), + "an already released local ticket remains an idempotent retirement request" + ); + } + Ok(()) + } + + #[test] + fn retirement_while_loading_reclaims_a_late_success_only_after_eviction() + -> Result<(), SchedulerError> { + let request = request(&format!("[{}]", workload("main", 4)), 20)?; + let mut scheduler = Scheduler::new(&request, SchedulerLimits::try_new(2, 1, 2)?)?; + let ticket = one_ticket(&mut scheduler, &request)?; + let load = poll_command(&mut scheduler)?; + scheduler.request_retirement(&ticket)?; + scheduler.request_retirement(&ticket)?; + scheduler.complete(RuntimeCompletion::Loaded { + operation: load.operation(), + resident: ResidentHandle::try_new("resident-main")?, + })?; + assert!( + matches!( + scheduler.begin_use(&ticket), + Err(SchedulerError::AdmissionNotResident { .. }) + ), + "a late successful load must remain retired under a live grant" + ); + let evict = poll_command(&mut scheduler)?; + scheduler.request_retirement(&ticket)?; + scheduler.complete(RuntimeCompletion::Evicted { + operation: evict.operation(), + })?; + scheduler.request_retirement(&ticket)?; + assert!( + scheduler.admissions.is_empty(), + "only the eviction acknowledgement releases a late loaded resident" + ); + Ok(()) + } + + #[test] + fn retirement_while_loading_releases_only_after_a_failed_load_acknowledgement() + -> Result<(), SchedulerError> { + let request = request(&format!("[{}]", workload("main", 4)), 20)?; + let mut scheduler = Scheduler::new(&request, SchedulerLimits::default())?; + let ticket = one_ticket(&mut scheduler, &request)?; + let load = poll_command(&mut scheduler)?; + scheduler.request_retirement(&ticket)?; + scheduler.complete(RuntimeCompletion::LoadFailed { + operation: load.operation(), + })?; + assert!( + scheduler.admissions.is_empty(), + "the trusted failed-load acknowledgement proves no allocation remains" + ); + assert!( + matches!( + scheduler.complete(RuntimeCompletion::LoadFailed { + operation: load.operation() + }), + Err(SchedulerError::UnknownOperation { .. }) + ), + "a duplicate late load acknowledgement cannot release accounting twice" + ); + scheduler.request_retirement(&ticket)?; + Ok(()) + } + + #[test] + fn failed_retirement_eviction_retains_a_and_allows_b_load_and_use_progress() + -> Result<(), SchedulerError> { + let initial = request(&format!("[{}]", workload("alpha", 4)), 20)?; + let mut scheduler = Scheduler::new(&initial, SchedulerLimits::try_new(3, 1, 2)?)?; + let alpha = loaded_ticket(&mut scheduler)?; + let beta_request = request(&format!("[{}]", workload("beta", 4)), 20)?; + let beta = one_ticket(&mut scheduler, &beta_request)?; + scheduler.request_retirement(&alpha)?; + let failed_evict = poll_command(&mut scheduler)?; + scheduler.complete(RuntimeCompletion::EvictFailed { + operation: failed_evict.operation(), + })?; + assert_eq!( + scheduler.admissions.len(), + 2, + "failed eviction keeps alpha's accounting lease" + ); + let beta_load = poll_command(&mut scheduler)?; + assert!( + matches!( + beta_load.kind(), + RuntimeCommandKind::Load { profile_id, .. } if profile_id == "beta" + ), + "load-first dispatch lets unrelated reserved work progress after failed eviction" + ); + scheduler.complete(RuntimeCompletion::Loaded { + operation: beta_load.operation(), + resident: ResidentHandle::try_new("resident-beta")?, + })?; + let beta_use = scheduler.begin_use(&beta)?; + scheduler.finish_use(&beta_use)?; + let retry = poll_command(&mut scheduler)?; + scheduler.complete(RuntimeCompletion::Evicted { + operation: retry.operation(), + })?; + let later_beta_use = scheduler.begin_use(&beta)?; + scheduler.finish_use(&later_beta_use)?; + Ok(()) + } + + #[test] + fn retirement_waits_for_live_uses_and_a_dropped_permit_is_not_an_acknowledgement() + -> Result<(), SchedulerError> { + let request = request(&format!("[{}]", workload("main", 4)), 20)?; + let mut scheduler = Scheduler::new(&request, SchedulerLimits::default())?; + let ticket = loaded_ticket(&mut scheduler)?; + let first = scheduler.begin_use(&ticket)?; + let second = scheduler.begin_use(&ticket)?; + scheduler.request_retirement(&ticket)?; + assert!( + matches!( + scheduler.begin_use(&ticket), + Err(SchedulerError::AdmissionNotResident { .. }) + ), + "retirement refuses every new use immediately" + ); + scheduler.finish_use(&first)?; + assert!( + matches!(scheduler.poll_command()?, PollOutcome::Idle), + "one remaining live use retains the resident allocation" + ); + scheduler.finish_use(&second)?; + let evict = poll_command(&mut scheduler)?; + scheduler.complete(RuntimeCompletion::Evicted { + operation: evict.operation(), + })?; + + let ticket = loaded_ticket(&mut scheduler)?; + let abandoned = scheduler.begin_use(&ticket)?; + scheduler.request_retirement(&ticket)?; + drop(abandoned); + assert!( + matches!(scheduler.poll_command()?, PollOutcome::Idle), + "dropping a caller-held permit cannot falsely acknowledge physical use completion" + ); + assert_eq!( + scheduler.admissions.len(), + 1, + "the abandoned permit keeps its retired allocation accounted" + ); + Ok(()) + } + + #[test] + fn retirement_refuses_foreign_and_stale_tickets_without_mutation() -> Result<(), SchedulerError> + { + let grant = request(&format!("[{}]", workload("main", 4)), 20)?; + let mut first = Scheduler::new(&grant, SchedulerLimits::default())?; + let mut second = Scheduler::new(&grant, SchedulerLimits::default())?; + let foreign = one_ticket(&mut second, &grant)?; + assert!( + matches!( + first.request_retirement(&foreign), + Err(SchedulerError::ForeignCapability { .. }) + ), + "a ticket from an identical but distinct controller cannot retire work" + ); + assert!( + first.admissions.is_empty(), + "foreign retirement has no side effect" + ); + + let ticket = one_ticket(&mut first, &grant)?; + first.revoke(&first.generation())?; + assert!(matches!(first.poll_command()?, PollOutcome::Progressed)); + let replacement = first.replace_grant(&grant)?; + assert!( + matches!( + first.request_retirement(&ticket), + Err(SchedulerError::StaleGeneration { .. }) + ), + "a prior generation ticket cannot retire a new grant's admission" + ); + assert!( + replacement.value > 1, + "replacement advanced the grant generation" + ); + Ok(()) + } + + #[test] + fn global_revoke_interleaves_with_selective_retirement_and_drains_all_admissions() + -> Result<(), SchedulerError> { + let grant = request( + &format!("[{},{}]", workload("alpha", 4), workload("beta", 4)), + 20, + )?; + let mut scheduler = Scheduler::new(&grant, SchedulerLimits::try_new(3, 1, 2)?)?; + let prepared = scheduler.prepare(&grant)?; + let mut tickets = scheduler.commit(prepared)?; + let beta = tickets.pop().ok_or(SchedulerError::UnknownAdmission { + location: error_location(), + })?; + let alpha = tickets.pop().ok_or(SchedulerError::UnknownAdmission { + location: error_location(), + })?; + for resident in ["resident-alpha", "resident-beta"] { + let load = poll_command(&mut scheduler)?; + scheduler.complete(RuntimeCompletion::Loaded { + operation: load.operation(), + resident: ResidentHandle::try_new(resident)?, + })?; + } + scheduler.request_retirement(&alpha)?; + scheduler.revoke(&scheduler.generation())?; + assert!( + matches!( + scheduler.begin_use(&beta), + Err(SchedulerError::GrantRevoked { .. }) + ), + "global revocation remains the stronger all-admission authority" + ); + for _ in 0..2 { + let evict = poll_command(&mut scheduler)?; + scheduler.complete(RuntimeCompletion::Evicted { + operation: evict.operation(), + })?; + } + assert!( + scheduler.admissions.is_empty(), + "selective and global drains release every lease exactly once" + ); + Ok(()) + } + #[test] fn bounded_sequence_preserves_drain_and_replace_invariants() -> Result<(), SchedulerError> { let initial = request(&format!("[{}]", workload("alpha", 4)), 20)?; From 3808fc96600116f597cde292db9beed5e2af99e3 Mon Sep 17 00:00:00 2001 From: Cody Kickertz Date: Tue, 8 Sep 2026 18:35:23 -0500 Subject: [PATCH 2/4] test(sched): cover selective resident retirement cycles --- crates/sched/src/lib.rs | 32 ++++++++++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/crates/sched/src/lib.rs b/crates/sched/src/lib.rs index 5f40805..17b0a20 100644 --- a/crates/sched/src/lib.rs +++ b/crates/sched/src/lib.rs @@ -1641,6 +1641,38 @@ mod tests { Ok(()) } + #[test] + fn retirement_cycles_reclaim_resident_admission_bound_without_grant_replacement() + -> Result<(), SchedulerError> { + let initial = request(&format!("[{}]", workload("initial", 4)), 20)?; + let mut scheduler = Scheduler::new(&initial, SchedulerLimits::try_new(1, 1, 1)?)?; + for profile in ["first", "second", "third"] { + let current = request(&format!("[{}]", workload(profile, 4)), 20)?; + let ticket = one_ticket(&mut scheduler, ¤t)?; + let load = poll_command(&mut scheduler)?; + scheduler.complete(RuntimeCompletion::Loaded { + operation: load.operation(), + resident: ResidentHandle::try_new(format!("resident-{profile}"))?, + })?; + let permit = scheduler.begin_use(&ticket)?; + scheduler.finish_use(&permit)?; + scheduler.request_retirement(&ticket)?; + let evict = poll_command(&mut scheduler)?; + scheduler.complete(RuntimeCompletion::Evicted { + operation: evict.operation(), + })?; + assert!( + scheduler.admissions.is_empty(), + "every resident cycle returns its lease before the next admission" + ); + } + assert!( + !scheduler.revoked, + "selective retirement does not replace the grant" + ); + Ok(()) + } + #[test] fn retirement_while_loading_reclaims_a_late_success_only_after_eviction() -> Result<(), SchedulerError> { From 9b570d2259eac661256589eab3ebff153c40ad68 Mon Sep 17 00:00:00 2001 From: Cody Kickertz Date: Tue, 8 Sep 2026 18:39:07 -0500 Subject: [PATCH 3/4] test(sched): order failed eviction before new admission --- crates/sched/src/lib.rs | 14 +++++++++++--- 1 file changed, 11 insertions(+), 3 deletions(-) diff --git a/crates/sched/src/lib.rs b/crates/sched/src/lib.rs index 17b0a20..1a994dc 100644 --- a/crates/sched/src/lib.rs +++ b/crates/sched/src/lib.rs @@ -1740,18 +1740,22 @@ mod tests { let initial = request(&format!("[{}]", workload("alpha", 4)), 20)?; let mut scheduler = Scheduler::new(&initial, SchedulerLimits::try_new(3, 1, 2)?)?; let alpha = loaded_ticket(&mut scheduler)?; - let beta_request = request(&format!("[{}]", workload("beta", 4)), 20)?; - let beta = one_ticket(&mut scheduler, &beta_request)?; scheduler.request_retirement(&alpha)?; let failed_evict = poll_command(&mut scheduler)?; + assert!( + matches!(failed_evict.kind(), RuntimeCommandKind::Evict { .. }), + "the retired resident issues its eviction before unrelated work is admitted" + ); scheduler.complete(RuntimeCompletion::EvictFailed { operation: failed_evict.operation(), })?; assert_eq!( scheduler.admissions.len(), - 2, + 1, "failed eviction keeps alpha's accounting lease" ); + let beta_request = request(&format!("[{}]", workload("beta", 4)), 20)?; + let beta = one_ticket(&mut scheduler, &beta_request)?; let beta_load = poll_command(&mut scheduler)?; assert!( matches!( @@ -1767,6 +1771,10 @@ mod tests { let beta_use = scheduler.begin_use(&beta)?; scheduler.finish_use(&beta_use)?; let retry = poll_command(&mut scheduler)?; + assert!( + matches!(retry.kind(), RuntimeCommandKind::Evict { .. }), + "the retained failed eviction becomes retryable after beta's load completes" + ); scheduler.complete(RuntimeCompletion::Evicted { operation: retry.operation(), })?; From 9a0b233c5491cab94ed5760d2fb7849b0f3dfb12 Mon Sep 17 00:00:00 2001 From: Cody Kickertz Date: Tue, 8 Sep 2026 18:41:00 -0500 Subject: [PATCH 4/4] test(sched): retain peers before grant revoke --- crates/sched/src/lib.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/crates/sched/src/lib.rs b/crates/sched/src/lib.rs index 1a994dc..d227479 100644 --- a/crates/sched/src/lib.rs +++ b/crates/sched/src/lib.rs @@ -1887,6 +1887,8 @@ mod tests { })?; } scheduler.request_retirement(&alpha)?; + let beta_use = scheduler.begin_use(&beta)?; + scheduler.finish_use(&beta_use)?; scheduler.revoke(&scheduler.generation())?; assert!( matches!(