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};
10
11#[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#[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#[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#[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#[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 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}