Skip to main content

ferrum_interfaces/vnext/resource/
provisioning.rs

1use super::{
2    device_capacity_account, invalid_resource, issue_generation,
3    validate_runtime_descriptor_for_admission, Arc, BTreeMap, BTreeSet, CapacityDomainId,
4    CapacityDomainSpec, CapacityUnits, DeviceCapacityClaim, DeviceId, DeviceRuntime,
5    DynamicBackingPoolSpec, DynamicPoolDomainSpec, DynamicPoolMaintenanceController,
6    DynamicPoolSet, DynamicResourceDescriptor, ExecutionPlan, LogicalAdmissionCoordinator,
7    PlanHash, PlanId, PlanNode, RequestIdentity, ResourcePoolId, ResourcePoolIdentity,
8    ResourceReservationBatch, StaticProvisioningBinding, VNextError,
9};
10use crate::vnext::CheckpointCapacityPolicy;
11
12/// One-shot plan/admission authority. It cannot be constructed, cloned, or
13/// deserialized by product or backend code. `ResourceTransaction::begin`
14/// consumes it, closing the old caller-built reservation bypass.
15#[must_use = "an admission permit must be consumed by ResourceTransaction::begin"]
16pub struct StaticProvisioningPermit<R>
17where
18    R: DeviceRuntime,
19{
20    pub(super) maintenance_controller: DynamicPoolMaintenanceController<R>,
21    pub(super) dynamic_pools: Arc<DynamicPoolSet<R>>,
22    pub(super) reservations: ResourceReservationBatch,
23    pub(super) capacity_claim: DeviceCapacityClaim,
24    pub(super) binding: StaticProvisioningBinding,
25    pub(super) runtime: Arc<R>,
26    pub(super) seal: AdmissionSeal,
27}
28
29pub(super) struct AdmissionSeal;
30
31pub(super) fn plan_dynamic_pool_admission(
32    maximum_active_sequences: u32,
33    pools: &[DynamicBackingPoolSpec],
34    descriptors: &[DynamicResourceDescriptor],
35    checkpoint_capacity: Option<CheckpointCapacityPolicy>,
36) -> Result<(LogicalAdmissionCoordinator, Vec<DynamicPoolDomainSpec>), VNextError> {
37    let mut descriptors_by_id = descriptors
38        .iter()
39        .map(|descriptor| (descriptor.base_resource_id().clone(), descriptor.clone()))
40        .collect::<BTreeMap<_, _>>();
41    if descriptors_by_id.len() != descriptors.len() || pools.is_empty() != descriptors.is_empty() {
42        return Err(invalid_resource(
43            "dynamic pool catalog and descriptor membership are inconsistent",
44        ));
45    }
46    let mut seen_pools = BTreeSet::new();
47    let mut domains = Vec::with_capacity(pools.len());
48    for (index, pool) in pools.iter().enumerate() {
49        if !seen_pools.insert(pool.pool_id().clone()) {
50            return Err(invalid_resource(
51                "dynamic pool catalog contains a duplicate pool",
52            ));
53        }
54        let mut members = Vec::with_capacity(pool.resource_ids().len());
55        for resource_id in pool.resource_ids() {
56            let descriptor = descriptors_by_id.remove(resource_id).ok_or_else(|| {
57                invalid_resource("dynamic pool references an unknown or duplicate descriptor")
58            })?;
59            if descriptor.pool_id() != pool.pool_id() {
60                return Err(invalid_resource(
61                    "dynamic descriptor belongs to another core-derived pool",
62                ));
63            }
64            members.push(descriptor);
65        }
66        members.sort_by(|left, right| left.base_resource_id().cmp(right.base_resource_id()));
67        let domain_id = CapacityDomainId::new(
68            u32::try_from(index + 1)
69                .map_err(|_| invalid_resource("dynamic pool domain id exceeds u32"))?,
70        )?;
71        domains.push(DynamicPoolDomainSpec {
72            domain_id,
73            pool: pool.clone(),
74            descriptors: members,
75        });
76    }
77    if !descriptors_by_id.is_empty() {
78        return Err(invalid_resource(
79            "dynamic descriptors are missing from the core-derived pool catalog",
80        ));
81    }
82    let coordinator_domains = domains
83        .iter()
84        .map(|domain| {
85            Ok((
86                domain.domain_id,
87                CapacityDomainSpec::new(
88                    CapacityUnits::ZERO,
89                    CapacityUnits::new(domain.pool.provisioning().maximum_resident_bytes()),
90                )?,
91            ))
92        })
93        .collect::<Result<Vec<_>, VNextError>>()?;
94    Ok((
95        match checkpoint_capacity {
96            Some(policy) => LogicalAdmissionCoordinator::with_checkpoint_capacity(
97                coordinator_domains,
98                maximum_active_sequences,
99                Some(policy),
100            )?,
101            None => {
102                LogicalAdmissionCoordinator::new(coordinator_domains, maximum_active_sequences)?
103            }
104        },
105        domains,
106    ))
107}
108
109impl<R> StaticProvisioningPermit<R>
110where
111    R: DeviceRuntime,
112{
113    pub fn maintenance_controller(&self) -> &DynamicPoolMaintenanceController<R> {
114        &self.maintenance_controller
115    }
116
117    pub fn binding(&self) -> &StaticProvisioningBinding {
118        &self.binding
119    }
120
121    pub fn reservations(&self) -> &ResourceReservationBatch {
122        &self.reservations
123    }
124}
125
126/// Explicit no-op result for plans that have no plan-lifetime buffers. It
127/// binds the validated plan to one exact runtime without manufacturing an
128/// empty reservation ledger or a zero-byte device-capacity claim.
129#[must_use = "no-static provisioning must be retained while the plan runtime is live"]
130pub struct NoStatic<R>
131where
132    R: DeviceRuntime,
133{
134    pub(super) maintenance_controller: DynamicPoolMaintenanceController<R>,
135    pub(super) dynamic_pools: Arc<DynamicPoolSet<R>>,
136    pub(super) binding: StaticProvisioningBinding,
137    pub(super) runtime: Arc<R>,
138}
139
140impl<R> NoStatic<R>
141where
142    R: DeviceRuntime,
143{
144    pub fn maintenance_controller(&self) -> &DynamicPoolMaintenanceController<R> {
145        &self.maintenance_controller
146    }
147
148    pub fn plan_id(&self) -> &PlanId {
149        self.binding.plan_id()
150    }
151
152    pub fn plan_hash(&self) -> &PlanHash {
153        self.binding.plan_hash()
154    }
155
156    pub fn device_id(&self) -> &DeviceId {
157        self.binding.device_id()
158    }
159
160    pub fn device_runtime_implementation_fingerprint(&self) -> &str {
161        self.binding.device_runtime_implementation_fingerprint()
162    }
163
164    pub const fn device_capacity_bytes(&self) -> u64 {
165        self.binding.device_capacity_bytes()
166    }
167
168    pub const fn usable_capacity_bytes(&self) -> u64 {
169        self.binding.usable_capacity_bytes()
170    }
171
172    pub const fn maximum_active_sequences(&self) -> u32 {
173        self.binding.maximum_active_sequences()
174    }
175}
176
177/// Static provisioning has two physically distinct outcomes. Only `Required`
178/// carries transaction authority; `NoStatic` cannot be passed to
179/// `ResourceTransaction::begin`.
180#[must_use = "static provisioning must be retained or committed"]
181pub enum StaticProvisioning<R>
182where
183    R: DeviceRuntime,
184{
185    NoStatic(NoStatic<R>),
186    Required(StaticProvisioningPermit<R>),
187}
188
189/// The indivisible result of plan provisioning. Product code must consume
190/// this owner through [`Self::into_parts`], which hands out the plan runtime
191/// outcome and its unique maintenance controller together. There is no
192/// controller-less extraction path.
193#[must_use = "provisioned plan resources must be split into their runtime and maintenance owners"]
194pub struct ProvisionedPlanResources<R>
195where
196    R: DeviceRuntime,
197{
198    provisioning: StaticProvisioning<R>,
199}
200
201/// Named result of consuming [`ProvisionedPlanResources`]. Keeping both
202/// fields in product ownership prevents maintenance authority from being
203/// silently discarded while request admission remains live.
204#[must_use = "both plan provisioning and maintenance ownership must be retained"]
205pub struct ProvisionedPlanParts<R>
206where
207    R: DeviceRuntime,
208{
209    pub provisioning: StaticProvisioning<R>,
210}
211
212impl<R> StaticProvisioning<R>
213where
214    R: DeviceRuntime,
215{
216    pub const fn has_static_resources(&self) -> bool {
217        matches!(self, Self::Required(_))
218    }
219}
220
221impl<R> ProvisionedPlanResources<R>
222where
223    R: DeviceRuntime,
224{
225    pub(super) fn new(provisioning: StaticProvisioning<R>) -> Self {
226        Self { provisioning }
227    }
228
229    pub fn provisioning(&self) -> &StaticProvisioning<R> {
230        &self.provisioning
231    }
232
233    pub fn into_parts(self) -> ProvisionedPlanParts<R> {
234        ProvisionedPlanParts {
235            provisioning: self.provisioning,
236        }
237    }
238
239    pub fn into_provisioning(self) -> StaticProvisioning<R> {
240        self.provisioning
241    }
242}
243
244impl ExecutionPlan {
245    /// Provisions only plan-lifetime buffers. Dynamic sequence admission is a
246    /// separate logical authority and is never implied by this result.
247    pub fn provision_static<R>(
248        &self,
249        runtime: Arc<R>,
250        request_id: RequestIdentity,
251    ) -> Result<ProvisionedPlanResources<R>, VNextError>
252    where
253        R: DeviceRuntime,
254    {
255        let payload = self.payload();
256        let memory = payload.memory();
257        if runtime.descriptor().id != *payload.device_id()
258            || runtime.descriptor().runtime_implementation_fingerprint
259                != payload.device_runtime_implementation_fingerprint()
260            || runtime.descriptor().total_memory_bytes != memory.device_capacity_bytes()
261        {
262            return Err(invalid_resource(
263                "static provisioning runtime differs from the execution plan",
264            ));
265        }
266        let (logical_admission, domains) = plan_dynamic_pool_admission(
267            memory.maximum_active_sequences(),
268            memory.dynamic_pools(),
269            memory.dynamic_descriptors(),
270            memory.checkpoint_capacity().copied(),
271        )?;
272        let generation = issue_generation()?;
273        let reservations = ResourceReservationBatch::from_allocations(
274            &request_id,
275            memory.static_allocations(),
276            generation,
277        )?;
278        if memory.device_capacity_bytes() == 0
279            || memory.usable_capacity_bytes() == 0
280            || memory.usable_capacity_bytes() > memory.device_capacity_bytes()
281            || memory.static_bytes() > memory.usable_capacity_bytes()
282            || memory.maximum_active_sequences() == 0
283            || (memory.static_bytes() == 0) != memory.static_allocations().is_empty()
284            || reservations.plan_static_size_bytes() != memory.static_bytes()
285            || reservations.total_size_bytes() != memory.static_bytes()
286        {
287            return Err(invalid_resource(
288                "plan-static reservation or capacity evidence is invalid",
289            ));
290        }
291        let pool_identity = ResourcePoolIdentity {
292            pool_id: ResourcePoolId::issue(generation)?,
293            plan_id: payload.plan_id().clone(),
294            plan_hash: self.plan_hash().clone(),
295            device_id: payload.device_id().clone(),
296            device_runtime_implementation_fingerprint: payload
297                .device_runtime_implementation_fingerprint()
298                .to_owned(),
299            admission_generation: generation,
300        };
301        let binding = StaticProvisioningBinding {
302            pool_identity,
303            plan_id: payload.plan_id().clone(),
304            plan_hash: self.plan_hash().clone(),
305            request_id,
306            device_id: payload.device_id().clone(),
307            device_runtime_implementation_fingerprint: payload
308                .device_runtime_implementation_fingerprint()
309                .to_owned(),
310            device_capacity_bytes: memory.device_capacity_bytes(),
311            usable_capacity_bytes: memory.usable_capacity_bytes(),
312            plan_static_bytes: memory.static_bytes(),
313            admitted_bytes: memory.static_bytes(),
314            maximum_active_sequences: memory.maximum_active_sequences(),
315            admission_generation: generation,
316        };
317        validate_runtime_descriptor_for_admission(
318            runtime.descriptor(),
319            &binding,
320            "resource admission preflight",
321        )?;
322        let account = device_capacity_account(
323            binding.device_id(),
324            binding.device_runtime_implementation_fingerprint(),
325            binding.device_capacity_bytes(),
326        )?;
327        let budget = account.register_budget(binding.usable_capacity_bytes())?;
328        let nodes: Arc<[PlanNode]> = Arc::from(payload.nodes().to_vec());
329        let dynamic_pools = Arc::new(DynamicPoolSet::new(
330            Arc::clone(&runtime),
331            binding.clone(),
332            Arc::clone(&budget),
333            logical_admission,
334            domains,
335            nodes,
336            payload.memory().reusable_execution().cloned(),
337        )?);
338        validate_runtime_descriptor_for_admission(
339            runtime.descriptor(),
340            &binding,
341            "resource admission completion",
342        )?;
343        let maintenance_controller =
344            DynamicPoolMaintenanceController::new(Arc::clone(&dynamic_pools));
345        if memory.static_allocations().is_empty() {
346            return Ok(ProvisionedPlanResources::new(StaticProvisioning::NoStatic(
347                NoStatic {
348                    maintenance_controller,
349                    dynamic_pools,
350                    binding,
351                    runtime,
352                },
353            )));
354        }
355        let capacity_claim = account.claim(&budget, reservations.total_size_bytes())?;
356        Ok(ProvisionedPlanResources::new(StaticProvisioning::Required(
357            StaticProvisioningPermit {
358                maintenance_controller,
359                dynamic_pools,
360                reservations,
361                capacity_claim,
362                binding,
363                runtime,
364                seal: AdmissionSeal,
365            },
366        )))
367    }
368}