use super::dynamic_pool::{DynamicDeviceCapacityBlocked, DynamicPoolGrowthIntent};
use super::{
invalid_resource, AdmissionDeferred, CapacityEpochs, CapacityVector, CapacityWaitCondition,
DeviceRuntime, DynamicBackingBlocker, DynamicBackingDeferred, DynamicBackingPackingEnvelope,
DynamicBackingPoolId, DynamicChunkQuarantineReason, DynamicPoolGrowthBatchReceipt,
DynamicPoolGrowthReceipt, DynamicPoolGrowthRequest, DynamicPoolMaintenanceBoundaryReceipt,
DynamicPoolSet, DynamicPoolStatus, VNextError,
};
use crate::vnext::{
CapacityShortfallKind, CapacityWaitSnapshot, DeferredAction, DynamicBackingPressure,
DynamicPoolResidentPressure,
};
use serde::Serialize;
use std::collections::BTreeMap;
use std::sync::Arc;
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct DynamicPoolMaintenanceStatus {
epochs: CapacityEpochs,
maximum_active_sequences: u32,
device_capacity_bytes: u64,
effective_device_usable_ceiling_bytes: u64,
process_claimed_bytes: u64,
budget_device_wide_usable_ceiling_bytes: u64,
budget_claimed_bytes: u64,
pools: Vec<DynamicPoolStatus>,
}
impl DynamicPoolMaintenanceStatus {
pub const fn epochs(&self) -> CapacityEpochs {
self.epochs
}
pub const fn maximum_active_sequences(&self) -> u32 {
self.maximum_active_sequences
}
pub const fn device_capacity_bytes(&self) -> u64 {
self.device_capacity_bytes
}
pub const fn effective_device_usable_ceiling_bytes(&self) -> u64 {
self.effective_device_usable_ceiling_bytes
}
pub const fn process_claimed_bytes(&self) -> u64 {
self.process_claimed_bytes
}
pub const fn budget_device_wide_usable_ceiling_bytes(&self) -> u64 {
self.budget_device_wide_usable_ceiling_bytes
}
pub const fn budget_claimed_bytes(&self) -> u64 {
self.budget_claimed_bytes
}
pub fn pools(&self) -> &[DynamicPoolStatus] {
&self.pools
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub enum DynamicDeferredMaintenanceOutcome {
RetryAdmission {
current_epochs: CapacityEpochs,
},
WaitForRelease {
current_epochs: CapacityEpochs,
wait_condition: CapacityWaitCondition,
pressure: DynamicBackingPressure,
#[serde(skip_serializing_if = "Option::is_none")]
maintenance_boundary: Option<DynamicPoolMaintenanceBoundaryReceipt>,
},
Maintained(DynamicPoolGrowthBatchReceipt),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct DynamicPoolQuarantineRelease {
pool_id: DynamicBackingPoolId,
released_chunks: usize,
released_bytes: u64,
}
impl DynamicPoolQuarantineRelease {
pub fn pool_id(&self) -> &DynamicBackingPoolId {
&self.pool_id
}
pub const fn released_chunks(&self) -> usize {
self.released_chunks
}
pub const fn released_bytes(&self) -> u64 {
self.released_bytes
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct DynamicPoolQuarantineReleaseReceipt {
pools: Vec<DynamicPoolQuarantineRelease>,
released_chunks: usize,
released_bytes: u64,
}
impl DynamicPoolQuarantineReleaseReceipt {
pub fn pools(&self) -> &[DynamicPoolQuarantineRelease] {
&self.pools
}
pub const fn released_chunks(&self) -> usize {
self.released_chunks
}
pub const fn released_bytes(&self) -> u64 {
self.released_bytes
}
}
#[must_use = "dynamic pool maintenance controller must be retained by the plan owner"]
pub struct DynamicPoolMaintenanceController<R>
where
R: DeviceRuntime,
{
pools: Arc<DynamicPoolSet<R>>,
}
impl<R> DynamicPoolMaintenanceController<R>
where
R: DeviceRuntime,
{
pub(in crate::vnext::resource) fn new(pools: Arc<DynamicPoolSet<R>>) -> Self {
Self { pools }
}
pub fn pool_ids(&self) -> impl ExactSizeIterator<Item = &DynamicBackingPoolId> {
self.pools.pools.keys()
}
pub fn status(&self) -> Result<DynamicPoolMaintenanceStatus, VNextError> {
let mut pools = Vec::with_capacity(self.pools.pools.len());
for pool in self.pools.pools.values() {
let mut state = pool
.state
.lock()
.map_err(|_| invalid_resource("dynamic backing pool is poisoned"))?;
let live_segments = state.chunks.values().try_fold(0_u64, |total, chunk| {
total
.checked_add(chunk.live_segments)
.ok_or_else(|| invalid_resource("dynamic live segment count overflows u64"))
})?;
let quarantined_bytes = state.quarantined.iter().try_fold(0_u64, |total, chunk| {
total
.checked_add(chunk.backing._grant.bytes())
.ok_or_else(|| invalid_resource("dynamic quarantine bytes overflow u64"))
})?;
let live_occupancy = state.live_occupancy;
let used_bytes = state
.resident_bytes
.checked_sub(state.allocator.free_bytes)
.ok_or_else(|| invalid_resource("dynamic pool free bytes exceed residency"))?;
if live_occupancy.total().physical_bytes() != used_bytes
|| live_occupancy.total().segment_count() != live_segments
{
state.poisoned = true;
return Err(invalid_resource(
"dynamic pool live-claim ledger differs from allocator occupancy",
));
}
pools.push(DynamicPoolStatus {
pool_id: pool.domain.pool_id().clone(),
domain_id: pool.domain.domain_id,
contract: super::DynamicPoolContractStatus::from_domain(&pool.domain),
storage_profile: pool.domain.pool.compatibility().profile(),
resident_bytes: state.resident_bytes,
pending_growth_bytes: state.pending_growth_bytes,
free_bytes: state.allocator.free_bytes,
largest_contiguous_bytes: state.allocator.largest_contiguous_bytes(),
resident_chunks: state.chunks.len(),
live_segments,
live_occupancy,
quarantined_chunks: state.quarantined.len(),
quarantined_bytes,
descriptor_mismatch_chunks: state
.quarantined
.iter()
.filter(|chunk| {
chunk.reason == DynamicChunkQuarantineReason::DescriptorMismatch
})
.count(),
publication_rejected_chunks: state
.quarantined
.iter()
.filter(|chunk| {
chunk.reason == DynamicChunkQuarantineReason::PublicationRejected
})
.count(),
poisoned: state.poisoned,
});
}
let account = &self.pools.budget.account;
let state = account
.state
.lock()
.map_err(|_| invalid_resource("device capacity account is poisoned"))?;
let effective_device_usable_ceiling_bytes = state
.budgets
.values()
.map(|budget| budget.device_wide_usable_ceiling_bytes)
.min()
.ok_or_else(|| invalid_resource("device capacity account has no live budget"))?;
let budget_claimed_bytes = state
.budgets
.get(&self.pools.budget.budget_id)
.ok_or_else(|| invalid_resource("dynamic pool plan budget is stale"))?
.claimed_bytes;
Ok(DynamicPoolMaintenanceStatus {
epochs: self.pools.logical_admission.epochs()?,
maximum_active_sequences: self.pools.maximum_active_sequences(),
device_capacity_bytes: account.device_capacity_bytes,
effective_device_usable_ceiling_bytes,
process_claimed_bytes: state.claimed_bytes,
budget_device_wide_usable_ceiling_bytes: self
.pools
.budget
.device_wide_usable_ceiling_bytes,
budget_claimed_bytes,
pools,
})
}
pub fn initialize_pool(
&self,
pool_id: &DynamicBackingPoolId,
) -> Result<Option<DynamicPoolGrowthReceipt>, VNextError> {
let mut receipt = self
.pools
.maintain_pools(vec![DynamicPoolGrowthIntent::Minimum(pool_id.clone())])?;
Ok(receipt.growths.pop())
}
pub fn initialize_pools(
&self,
pool_ids: &[DynamicBackingPoolId],
) -> Result<DynamicPoolGrowthBatchReceipt, VNextError> {
self.pools.maintain_pools(
pool_ids
.iter()
.cloned()
.map(DynamicPoolGrowthIntent::Minimum)
.collect(),
)
}
pub fn grow_pool(
&self,
pool_id: &DynamicBackingPoolId,
requested_bytes: u64,
) -> Result<DynamicPoolGrowthReceipt, VNextError> {
let request = DynamicPoolGrowthRequest::new(pool_id.clone(), requested_bytes)?;
let mut receipt = self.grow_pools(vec![request])?;
receipt
.growths
.pop()
.ok_or_else(|| invalid_resource("single-pool growth produced no receipt"))
}
pub fn grow_pools(
&self,
requests: Vec<DynamicPoolGrowthRequest>,
) -> Result<DynamicPoolGrowthBatchReceipt, VNextError> {
self.pools.maintain_pools(
requests
.into_iter()
.map(DynamicPoolGrowthIntent::Additional)
.collect(),
)
}
fn wait_snapshot_for_pool_ids<'a>(
&self,
pool_ids: impl IntoIterator<Item = &'a DynamicBackingPoolId>,
) -> Result<CapacityWaitSnapshot, VNextError> {
self.pools.logical_admission.wait_snapshot_for_domains(
pool_ids
.into_iter()
.map(|pool_id| {
self.pools
.pools
.get(pool_id)
.map(|pool| pool.domain.domain_id)
.ok_or_else(|| {
invalid_resource("dynamic maintenance references an unknown pool")
})
})
.collect::<Result<Vec<_>, _>>()?,
)
}
fn capacity_wait_outcome(
&self,
logical_snapshot: CapacityWaitSnapshot,
blocked: DynamicDeviceCapacityBlocked,
maintenance_boundary: DynamicPoolMaintenanceBoundaryReceipt,
) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
if maintenance_boundary.pressure() != &blocked.pressure
|| maintenance_boundary.plan_device_capacity_epoch()
!= blocked.availability.plan_epoch()
|| maintenance_boundary.process_device_capacity_epoch()
!= blocked.availability.process_epoch()
{
return Err(invalid_resource(
"dynamic maintenance boundary differs from its capacity failure",
));
}
let logical_snapshot = logical_snapshot.narrow_to_domains(blocked.planned_domains)?;
let mut observed = logical_snapshot.wait_condition().observed().to_vec();
observed.push(blocked.availability.epoch_for_pressure(&blocked.pressure));
let wait_condition = CapacityWaitCondition::new(
logical_snapshot.wait_condition().coordinator_id(),
observed,
)?;
Ok(DynamicDeferredMaintenanceOutcome::WaitForRelease {
current_epochs: logical_snapshot.epochs(),
wait_condition,
pressure: blocked.pressure.into(),
maintenance_boundary: Some(maintenance_boundary),
})
}
fn pool_resident_wait_outcome(
&self,
logical_snapshot: CapacityWaitSnapshot,
pressure: DynamicPoolResidentPressure,
) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
let domain = self
.pools
.pools
.get(pressure.pool_id())
.map(|pool| pool.domain.domain_id)
.ok_or_else(|| {
invalid_resource("dynamic pool resident pressure references an unknown pool")
})?;
let logical_snapshot = logical_snapshot.narrow_to_domains(vec![domain])?;
Ok(DynamicDeferredMaintenanceOutcome::WaitForRelease {
current_epochs: logical_snapshot.epochs(),
wait_condition: logical_snapshot.wait_condition().clone(),
pressure: pressure.into(),
maintenance_boundary: None,
})
}
fn maintain_deferred_pools(
&self,
intents: Vec<DynamicPoolGrowthIntent>,
capacity_blocked: &mut Option<DynamicDeviceCapacityBlocked>,
maintenance_boundary: &mut Option<DynamicPoolMaintenanceBoundaryReceipt>,
protected_immediate: &CapacityVector,
protected_packing_envelopes: &[DynamicBackingPackingEnvelope],
) -> Result<DynamicPoolGrowthBatchReceipt, VNextError> {
let retry_intents = intents.clone();
match self
.pools
.maintain_pools_observed(intents, capacity_blocked)
{
Err(VNextError::DeviceCapacityUnavailable(pressure)) => {
let planned_domains = capacity_blocked
.as_ref()
.expect("typed capacity failure retains its exact observation")
.planned_domains
.clone();
let blocked = capacity_blocked
.as_ref()
.expect("typed capacity failure retains its exact observation");
let attempt = self.pools.reclaim_idle_chunks_for_pressure(
&pressure,
blocked.availability,
&planned_domains,
protected_immediate,
protected_packing_envelopes,
)?;
*maintenance_boundary = Some(attempt.boundary);
let Some(rebalance) = attempt.rebalance else {
return Err(VNextError::DeviceCapacityUnavailable(pressure));
};
*capacity_blocked = None;
let mut receipt = self
.pools
.maintain_pools_observed(retry_intents, capacity_blocked)?;
receipt.rebalance = Some(rebalance);
receipt.maintenance_boundary = maintenance_boundary.clone();
Ok(receipt)
}
outcome => outcome,
}
}
pub(super) fn maintain_for_live_deferred(
&self,
deferred: &DynamicBackingDeferred,
) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
let coordinator_id = self.pools.logical_admission.id();
if coordinator_id != deferred.epochs().coordinator_id()
|| coordinator_id != deferred.wait_condition().coordinator_id()
{
return Err(invalid_resource(
"dynamic backing deferral belongs to another admission coordinator",
));
}
if deferred.blockers().is_empty() {
return Err(invalid_resource(
"dynamic backing deferral contains no blocking pool",
));
}
let logical_snapshot = self.wait_snapshot_for_pool_ids(
deferred
.blockers()
.iter()
.map(DynamicBackingBlocker::pool_id),
)?;
let mut capacity_blocked = None;
let mut maintenance_boundary = None;
let growth = self.maintain_deferred_pools(
deferred
.blockers()
.iter()
.cloned()
.map(DynamicPoolGrowthIntent::RevalidatedDeferral)
.collect(),
&mut capacity_blocked,
&mut maintenance_boundary,
deferred.protected_immediate(),
deferred.protected_packing_envelopes(),
);
match growth {
Ok(receipt) if receipt.growths().is_empty() => {
let current_epochs = self.pools.logical_admission.epochs()?;
if current_epochs == deferred.epochs() {
return Err(invalid_resource(
"dynamic backing maintenance made no progress on an unchanged deferral",
));
}
Ok(DynamicDeferredMaintenanceOutcome::RetryAdmission { current_epochs })
}
Ok(receipt) => Ok(DynamicDeferredMaintenanceOutcome::Maintained(receipt)),
Err(VNextError::DeviceCapacityUnavailable(_)) => self.capacity_wait_outcome(
logical_snapshot,
capacity_blocked.expect("typed capacity failure retains its exact observation"),
maintenance_boundary
.expect("typed capacity rebalance failure retains its maintenance boundary"),
),
Err(VNextError::DynamicPoolResidentUnavailable(pressure)) => {
self.pool_resident_wait_outcome(logical_snapshot, pressure)
}
Err(error) => Err(error),
}
}
pub fn maintain_for_admission_deferred(
&self,
deferred: &AdmissionDeferred,
) -> Result<DynamicDeferredMaintenanceOutcome, VNextError> {
let coordinator_id = self.pools.logical_admission.id();
if coordinator_id != deferred.epochs().coordinator_id()
|| coordinator_id != deferred.wait_condition().coordinator_id()
{
return Err(invalid_resource(
"logical admission deferral belongs to another coordinator",
));
}
if deferred.action() != DeferredAction::AwaitBackingGrowth {
return Err(invalid_resource(
"logical admission deferral does not request backing growth",
));
}
let pools_by_domain = self
.pools
.pools
.values()
.map(|pool| (pool.domain.domain_id, pool.domain.pool_id().clone()))
.collect::<BTreeMap<_, _>>();
let current = self.pools.logical_admission.snapshot()?;
let mut requested_by_pool = BTreeMap::<DynamicBackingPoolId, u64>::new();
for blocker in deferred
.blockers()
.iter()
.filter(|blocker| blocker.kind() == CapacityShortfallKind::BackingGrowthRequired)
{
let domain = blocker.domain().ok_or_else(|| {
invalid_resource("backing-growth blocker contains no capacity domain")
})?;
let pool_id = pools_by_domain.get(&domain).ok_or_else(|| {
invalid_resource("backing-growth blocker references a non-pool domain")
})?;
let current_total = current
.domains()
.iter()
.find(|snapshot| snapshot.domain() == domain)
.ok_or_else(|| {
invalid_resource("backing-growth blocker references an unknown domain")
})?
.total()
.get();
let missing = blocker.requested().get().saturating_sub(current_total);
if missing == 0 {
continue;
}
requested_by_pool
.entry(pool_id.clone())
.and_modify(|bytes| *bytes = (*bytes).max(missing))
.or_insert(missing);
}
if requested_by_pool.is_empty() {
return Ok(DynamicDeferredMaintenanceOutcome::RetryAdmission {
current_epochs: self.pools.logical_admission.epochs()?,
});
}
let requests = requested_by_pool
.into_iter()
.map(|(pool_id, bytes)| DynamicPoolGrowthRequest::new(pool_id, bytes))
.collect::<Result<Vec<_>, _>>()?;
let logical_snapshot =
self.wait_snapshot_for_pool_ids(requests.iter().map(|request| request.pool_id()))?;
let mut capacity_blocked = None;
let mut maintenance_boundary = None;
let growth = self.maintain_deferred_pools(
requests
.into_iter()
.map(DynamicPoolGrowthIntent::Additional)
.collect(),
&mut capacity_blocked,
&mut maintenance_boundary,
deferred.immediate_requested(),
&[],
);
match growth {
Ok(receipt) => Ok(DynamicDeferredMaintenanceOutcome::Maintained(receipt)),
Err(VNextError::DeviceCapacityUnavailable(_)) => self.capacity_wait_outcome(
logical_snapshot,
capacity_blocked.expect("typed capacity failure retains its exact observation"),
maintenance_boundary
.expect("typed capacity rebalance failure retains its maintenance boundary"),
),
Err(VNextError::DynamicPoolResidentUnavailable(pressure)) => {
self.pool_resident_wait_outcome(logical_snapshot, pressure)
}
Err(error) => Err(error),
}
}
pub fn release_quarantined_chunks(
&self,
) -> Result<DynamicPoolQuarantineReleaseReceipt, VNextError> {
let pools = self.pools.pools.values().cloned().collect::<Vec<_>>();
let _maintenance = pools
.iter()
.map(|pool| {
pool.maintenance
.lock()
.map_err(|_| invalid_resource("dynamic pool maintenance authority is poisoned"))
})
.collect::<Result<Vec<_>, _>>()?;
let mut states = pools
.iter()
.map(|pool| {
pool.state
.lock()
.map_err(|_| invalid_resource("dynamic backing pool is poisoned"))
})
.collect::<Result<Vec<_>, _>>()?;
let pool_totals = states
.iter()
.map(|state| {
let bytes = state.quarantined.iter().try_fold(0_u64, |total, chunk| {
total
.checked_add(chunk.backing._grant.bytes())
.ok_or_else(|| invalid_resource("released quarantine bytes overflow u64"))
})?;
Ok((state.quarantined.len(), bytes))
})
.collect::<Result<Vec<_>, VNextError>>()?;
let released_chunks = pool_totals.iter().try_fold(0_usize, |total, (count, _)| {
total
.checked_add(*count)
.ok_or_else(|| invalid_resource("released quarantine count overflows usize"))
})?;
let released_bytes = pool_totals.iter().try_fold(0_u64, |total, (_, bytes)| {
total
.checked_add(*bytes)
.ok_or_else(|| invalid_resource("released quarantine bytes overflow u64"))
})?;
let mut released = Vec::with_capacity(pools.len());
let mut receipts = Vec::new();
for ((pool, state), (pool_chunks, pool_bytes)) in
pools.iter().zip(states.iter_mut()).zip(pool_totals)
{
if pool_chunks == 0 {
continue;
}
let chunks = std::mem::take(&mut state.quarantined);
debug_assert_eq!(chunks.len(), pool_chunks);
receipts.push(DynamicPoolQuarantineRelease {
pool_id: pool.domain.pool_id().clone(),
released_chunks: pool_chunks,
released_bytes: pool_bytes,
});
released.push(chunks);
}
drop(states);
drop(released);
Ok(DynamicPoolQuarantineReleaseReceipt {
pools: receipts,
released_chunks,
released_bytes,
})
}
}