use super::*;
use crate::vnext::{
CheckpointCapacityPolicy, DeviceCapacityPressure, DynamicPoolGrowthBatchReceipt,
DynamicPoolResidentPressure, ExecutionPlan, SequenceCheckpointBytePlan,
};
#[must_use = "checkpoint maintenance must be attempted or dropped"]
pub struct CheckpointCapacityMaintenance<R: DeviceRuntime> {
binding: TrustedPlanRuntimeBinding<R>,
requests: CheckpointBackingRequests,
policy: Option<CheckpointCapacityPolicy>,
}
#[derive(Debug)]
pub enum CheckpointCapacityMaintenanceSkipReason {
Retention(CheckpointRetentionSkipReason),
DeviceCapacity(DeviceCapacityPressure),
PoolResident(DynamicPoolResidentPressure),
}
#[derive(Debug)]
pub enum CheckpointCapacityMaintenanceOutcome {
Ready(DynamicPoolGrowthBatchReceipt),
Skipped(CheckpointCapacityMaintenanceSkipReason),
}
#[derive(Clone, Copy)]
enum CheckpointMaintenanceMode {
AvailableCapacityOnly,
ReclaimIdleNonTargets,
}
impl<R: DeviceRuntime> PlanRuntimeResources<R> {
pub fn try_maintain_checkpoint_with_idle_reclaim(
self: &Arc<Self>,
maintenance: CheckpointCapacityMaintenance<R>,
) -> Result<CheckpointCapacityMaintenanceOutcome, VNextError> {
if !Arc::ptr_eq(self, &maintenance.binding.resources) {
return Err(invalid_resource(
"checkpoint maintenance belongs to another resource root",
));
}
maintenance.maintain_with_mode(CheckpointMaintenanceMode::ReclaimIdleNonTargets)
}
}
impl<R: DeviceRuntime> TrustedPlanRuntimeBinding<R> {
pub fn prepare_checkpoint_capacity_maintenance(
&self,
plan: &ExecutionPlan,
byte_plan: &SequenceCheckpointBytePlan,
) -> Result<CheckpointCapacityMaintenance<R>, VNextError> {
let _lifecycle = self
.resources
.read_lifecycle("prepare checkpoint maintenance")?;
if plan.plan_hash() != self.plan_hash()
|| byte_plan.plan_hash() != self.plan_hash()
|| plan.checkpoint_byte_plan(byte_plan.boundary())? != *byte_plan
{
return Err(invalid_resource(
"checkpoint maintenance requires this plan's certified byte layout",
));
}
let requests = byte_plan.backing_requests()?;
self.evaluate_checkpoint_backing(&requests)?;
Ok(CheckpointCapacityMaintenance {
binding: TrustedPlanRuntimeBinding {
resources: Arc::clone(&self.resources),
},
requests,
policy: plan.payload().memory().checkpoint_capacity().copied(),
})
}
}
impl<R: DeviceRuntime> CheckpointCapacityMaintenance<R> {
pub fn try_maintain(self) -> Result<CheckpointCapacityMaintenanceOutcome, VNextError> {
self.maintain_with_mode(CheckpointMaintenanceMode::AvailableCapacityOnly)
}
fn maintain_with_mode(
self,
mode: CheckpointMaintenanceMode,
) -> Result<CheckpointCapacityMaintenanceOutcome, VNextError> {
let _lifecycle = self
.binding
.resources
.read_lifecycle("maintain checkpoint capacity")?;
let Some(policy) = self.policy else {
return Ok(CheckpointCapacityMaintenanceOutcome::Skipped(
CheckpointCapacityMaintenanceSkipReason::Retention(
CheckpointRetentionSkipReason::Disabled,
),
));
};
let evaluated = self.binding.evaluate_checkpoint_backing(&self.requests)?;
let retained = self
.binding
.logical_admission()
.checkpoint_retained_bytes()?;
let maximum = policy.maximum_retained_bytes();
let remaining = maximum.checked_sub(retained).ok_or_else(|| {
invalid_resource("checkpoint retained bytes exceed the bound plan policy")
})?;
if evaluated.extent_bytes > remaining {
return Ok(CheckpointCapacityMaintenanceOutcome::Skipped(
CheckpointCapacityMaintenanceSkipReason::Retention(
CheckpointRetentionSkipReason::Capacity {
requested_bytes: evaluated.extent_bytes,
retained_bytes: retained,
maximum_bytes: maximum,
},
),
));
}
let pools = self.binding.dynamic_pools();
let result = match mode {
CheckpointMaintenanceMode::AvailableCapacityOnly => {
pools.maintain_checkpoint_capacity(&evaluated.slices)
}
CheckpointMaintenanceMode::ReclaimIdleNonTargets => {
pools.maintain_checkpoint_capacity_with_idle_reclaim(&evaluated.slices)
}
};
match result {
Ok(receipt) => Ok(CheckpointCapacityMaintenanceOutcome::Ready(receipt)),
Err(VNextError::DeviceCapacityUnavailable(pressure)) => {
Ok(CheckpointCapacityMaintenanceOutcome::Skipped(
CheckpointCapacityMaintenanceSkipReason::DeviceCapacity(pressure),
))
}
Err(VNextError::DynamicPoolResidentUnavailable(pressure)) => {
Ok(CheckpointCapacityMaintenanceOutcome::Skipped(
CheckpointCapacityMaintenanceSkipReason::PoolResident(pressure),
))
}
Err(error) => Err(error),
}
}
}