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