ferrum_interfaces/vnext/resource/
provisioning.rs1use 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#[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#[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#[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#[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#[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 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}