1use std::collections::{BTreeMap, BTreeSet};
4use std::sync::{Arc, Mutex, MutexGuard};
5
6use thiserror::Error;
7use vyre_driver::backend::{ArtifactInstance, ArtifactMaterializer, BackendError, Resource};
8
9const ZERO_UPLOAD_CHUNK_BYTES: usize = 1024 * 1024;
10
11#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
13pub struct ResourceSetKey {
14 pub source_digest: [u8; 32],
16 pub artifact_digest: [u8; 32],
18}
19
20#[derive(Debug, Clone, Copy, PartialEq, Eq)]
22pub struct ImmutableResourceUpload<'a> {
23 pub name: &'a str,
25 pub bytes: &'a [u8],
27 pub blake3: [u8; 32],
29}
30
31#[derive(Clone)]
33pub struct ArtifactInstanceBinding {
34 name: String,
35 instance: Arc<dyn ArtifactInstance>,
36 byte_len: u64,
37}
38
39impl std::fmt::Debug for ArtifactInstanceBinding {
40 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
41 formatter
42 .debug_struct("ArtifactInstanceBinding")
43 .field("name", &self.name)
44 .field("artifact", &self.instance.artifact())
45 .field("payload", &self.instance.payload())
46 .field("device", self.instance.device())
47 .field("byte_len", &self.byte_len)
48 .finish()
49 }
50}
51
52impl ArtifactInstanceBinding {
53 #[must_use]
55 pub fn new(
56 name: impl Into<String>,
57 instance: Arc<dyn ArtifactInstance>,
58 byte_len: u64,
59 ) -> Self {
60 Self {
61 name: name.into(),
62 instance,
63 byte_len,
64 }
65 }
66
67 #[must_use]
69 pub fn name(&self) -> &str {
70 &self.name
71 }
72
73 #[must_use]
75 pub fn instance(&self) -> &Arc<dyn ArtifactInstance> {
76 &self.instance
77 }
78
79 #[must_use]
81 pub const fn byte_len(&self) -> u64 {
82 self.byte_len
83 }
84}
85
86#[derive(Debug)]
88pub struct ResourceSetAdmission<'a> {
89 pub key: ResourceSetKey,
91 pub immutable_resources: Vec<ImmutableResourceUpload<'a>>,
93 pub artifacts: Vec<ArtifactInstanceBinding>,
95}
96
97#[derive(Debug, Clone, Copy, PartialEq, Eq)]
99pub enum ResourceAdmissionStatus {
100 Cold,
102 Warm,
104}
105
106#[derive(Debug, Clone, Copy, PartialEq, Eq)]
108pub struct ResourceSetLease {
109 pub key: ResourceSetKey,
111 pub status: ResourceAdmissionStatus,
113}
114
115#[derive(Debug, Clone, Copy, PartialEq, Eq)]
117pub struct MutableStateSpec<'a> {
118 pub name: &'a str,
120 pub byte_len: usize,
122}
123
124#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
126pub struct StateId(pub u64);
127
128#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
130pub struct StateLease {
131 pub id: StateId,
133 pub generation: u64,
135}
136
137pub trait ResidentResourceDevice: Send + Sync {
142 fn allocate(&self, byte_len: usize) -> Result<Resource, BackendError>;
144 fn upload_many(&self, uploads: &[(&Resource, &[u8])]) -> Result<(), BackendError>;
146 fn upload_at(
148 &self,
149 resource: &Resource,
150 offset: usize,
151 bytes: &[u8],
152 ) -> Result<(), BackendError>;
153 fn free(&self, resource: Resource) -> Result<(), BackendError>;
155}
156
157#[derive(Clone)]
159pub struct MaterializerResourceDevice {
160 materializer: Arc<dyn ArtifactMaterializer>,
161}
162
163impl MaterializerResourceDevice {
164 #[must_use]
166 pub fn new(materializer: Arc<dyn ArtifactMaterializer>) -> Self {
167 Self { materializer }
168 }
169}
170
171impl ResidentResourceDevice for MaterializerResourceDevice {
172 fn allocate(&self, byte_len: usize) -> Result<Resource, BackendError> {
173 self.materializer.allocate_resident(byte_len)
174 }
175
176 fn upload_many(&self, uploads: &[(&Resource, &[u8])]) -> Result<(), BackendError> {
177 for (resource, bytes) in uploads {
178 self.materializer.upload_resident(resource, bytes)?;
179 }
180 Ok(())
181 }
182
183 fn upload_at(
184 &self,
185 resource: &Resource,
186 offset: usize,
187 bytes: &[u8],
188 ) -> Result<(), BackendError> {
189 self.materializer
190 .upload_resident_at(resource, offset, bytes)
191 }
192
193 fn free(&self, resource: Resource) -> Result<(), BackendError> {
194 self.materializer.free_resident(resource)
195 }
196}
197
198#[derive(Clone)]
199struct ResidentImmutableResource {
200 resource: Resource,
201 byte_len: u64,
202 digest: [u8; 32],
203}
204
205struct ResidentArtifact {
206 instance: Arc<dyn ArtifactInstance>,
207 artifact: [u8; 32],
208 byte_len: u64,
209}
210
211struct ResidentResourceSet {
212 immutable_resources: BTreeMap<String, ResidentImmutableResource>,
213 artifacts: BTreeMap<String, ResidentArtifact>,
214 accounted_bytes: u64,
215 active_states: u64,
216}
217
218struct ResidentStateSet {
219 resource_set: ResourceSetKey,
220 generation: u64,
221 states: BTreeMap<String, Resource>,
222 state_sizes: BTreeMap<String, usize>,
223 accounted_bytes: u64,
224}
225
226struct ResidencyState {
227 resource_sets: BTreeMap<ResourceSetKey, ResidentResourceSet>,
228 states: BTreeMap<StateId, ResidentStateSet>,
229 next_state: u64,
230 used_bytes: u64,
231}
232
233impl Default for ResidencyState {
234 fn default() -> Self {
235 Self {
236 resource_sets: BTreeMap::new(),
237 states: BTreeMap::new(),
238 next_state: 1,
239 used_bytes: 0,
240 }
241 }
242}
243
244pub struct ResourceResidency {
247 device: Arc<dyn ResidentResourceDevice>,
248 budget_bytes: u64,
249 state: Mutex<ResidencyState>,
250}
251
252impl std::fmt::Debug for ResourceResidency {
253 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
254 formatter
255 .debug_struct("ResourceResidency")
256 .field("budget_bytes", &self.budget_bytes)
257 .finish_non_exhaustive()
258 }
259}
260
261impl ResourceResidency {
262 #[must_use]
264 pub fn new(materializer: Arc<dyn ArtifactMaterializer>, budget_bytes: u64) -> Self {
265 Self::with_device(
266 Arc::new(MaterializerResourceDevice::new(materializer)),
267 budget_bytes,
268 )
269 }
270
271 #[must_use]
273 pub fn with_device(device: Arc<dyn ResidentResourceDevice>, budget_bytes: u64) -> Self {
274 Self {
275 device,
276 budget_bytes,
277 state: Mutex::new(ResidencyState::default()),
278 }
279 }
280
281 pub fn admit_resource_set(
283 &self,
284 request: ResourceSetAdmission<'_>,
285 ) -> Result<ResourceSetLease, ResourceResidencyError> {
286 validate_key(request.key)?;
287 let prepared_immutable_resources =
288 validate_immutable_resources(&request.immutable_resources)?;
289 let prepared_artifacts = validate_artifacts(&request.artifacts)?;
290 let requested_bytes =
291 accounted_resource_set_bytes(&prepared_immutable_resources, &prepared_artifacts)?;
292 let mut state = self.lock_state()?;
293
294 if let Some(resident) = state.resource_sets.get(&request.key) {
295 validate_warm_resource_set(
296 resident,
297 &prepared_immutable_resources,
298 &prepared_artifacts,
299 )?;
300 return Ok(ResourceSetLease {
301 key: request.key,
302 status: ResourceAdmissionStatus::Warm,
303 });
304 }
305 ensure_budget(
306 state.used_bytes,
307 requested_bytes,
308 self.budget_bytes,
309 "resource-set admission",
310 )?;
311
312 let mut allocations = Vec::with_capacity(prepared_immutable_resources.len());
313 for immutable_resource in &prepared_immutable_resources {
314 let byte_len = usize::try_from(immutable_resource.byte_len).map_err(|_| {
315 ResourceResidencyError::ByteLengthOverflow {
316 context: format!("immutable_resource `{}`", immutable_resource.name),
317 }
318 })?;
319 match self.device.allocate(byte_len) {
320 Ok(resource) => allocations.push(resource),
321 Err(error) => {
322 return Err(self.rollback_error(
323 allocations,
324 "immutable_resource allocation",
325 error,
326 ));
327 }
328 }
329 }
330 let uploads = allocations
331 .iter()
332 .zip(request.immutable_resources.iter())
333 .map(|(resource, immutable_resource)| (resource, immutable_resource.bytes))
334 .collect::<Vec<_>>();
335 if let Err(error) = self.device.upload_many(&uploads) {
336 return Err(self.rollback_error(allocations, "immutable_resource batch upload", error));
337 }
338
339 let immutable_resources = prepared_immutable_resources
340 .into_iter()
341 .zip(allocations)
342 .map(|(immutable_resource, resource)| {
343 (
344 immutable_resource.name,
345 ResidentImmutableResource {
346 resource,
347 byte_len: immutable_resource.byte_len,
348 digest: immutable_resource.digest,
349 },
350 )
351 })
352 .collect();
353 let artifacts = request
354 .artifacts
355 .into_iter()
356 .map(|artifact| {
357 (
358 artifact.name,
359 ResidentArtifact {
360 artifact: artifact.instance.artifact().0,
361 instance: artifact.instance,
362 byte_len: artifact.byte_len,
363 },
364 )
365 })
366 .collect();
367 state.used_bytes = state.used_bytes.checked_add(requested_bytes).ok_or(
368 ResourceResidencyError::ByteLengthOverflow {
369 context: "committed residency bytes".into(),
370 },
371 )?;
372 state.resource_sets.insert(
373 request.key,
374 ResidentResourceSet {
375 immutable_resources,
376 artifacts,
377 accounted_bytes: requested_bytes,
378 active_states: 0,
379 },
380 );
381 Ok(ResourceSetLease {
382 key: request.key,
383 status: ResourceAdmissionStatus::Cold,
384 })
385 }
386
387 pub fn start_state(
389 &self,
390 resource_set: ResourceSetKey,
391 specs: &[MutableStateSpec<'_>],
392 ) -> Result<StateLease, ResourceResidencyError> {
393 let prepared = validate_state_specs(specs)?;
394 let requested_bytes = prepared.iter().try_fold(0_u64, |total, (_, bytes)| {
395 total
396 .checked_add(*bytes as u64)
397 .ok_or(ResourceResidencyError::ByteLengthOverflow {
398 context: "state-state byte total".into(),
399 })
400 })?;
401 let mut state = self.lock_state()?;
402 if !state.resource_sets.contains_key(&resource_set) {
403 return Err(ResourceResidencyError::ResourceSetNotResident { key: resource_set });
404 }
405 ensure_budget(
406 state.used_bytes,
407 requested_bytes,
408 self.budget_bytes,
409 "state admission",
410 )?;
411 let next_active_states = state
412 .resource_sets
413 .get(&resource_set)
414 .ok_or(ResourceResidencyError::ResourceSetNotResident { key: resource_set })?
415 .active_states
416 .checked_add(1)
417 .ok_or(ResourceResidencyError::StateIdentityOverflow)?;
418 let id = StateId(state.next_state);
419 state.next_state = state
420 .next_state
421 .checked_add(1)
422 .ok_or(ResourceResidencyError::StateIdentityOverflow)?;
423
424 let mut resources = Vec::with_capacity(prepared.len());
425 for (_, byte_len) in &prepared {
426 match self.device.allocate(*byte_len) {
427 Ok(resource) => {
428 if let Err(error) = self.zero_resource(&resource, *byte_len) {
429 resources.push(resource);
430 return Err(self.rollback_error(resources, "state zeroing", error));
431 }
432 resources.push(resource);
433 }
434 Err(error) => {
435 return Err(self.rollback_error(resources, "state allocation", error));
436 }
437 }
438 }
439 let states = prepared
440 .iter()
441 .map(|(name, _)| name.clone())
442 .zip(resources)
443 .collect();
444 let state_sizes = prepared.into_iter().collect();
445 state.used_bytes = state.used_bytes.checked_add(requested_bytes).ok_or(
446 ResourceResidencyError::ByteLengthOverflow {
447 context: "committed state bytes".into(),
448 },
449 )?;
450 state
451 .resource_sets
452 .get_mut(&resource_set)
453 .ok_or(ResourceResidencyError::ResourceSetNotResident { key: resource_set })?
454 .active_states = next_active_states;
455 state.states.insert(
456 id,
457 ResidentStateSet {
458 resource_set,
459 generation: 0,
460 states,
461 state_sizes,
462 accounted_bytes: requested_bytes,
463 },
464 );
465 Ok(StateLease { id, generation: 0 })
466 }
467
468 pub fn mutable_state(
470 &self,
471 lease: StateLease,
472 name: &str,
473 ) -> Result<Resource, ResourceResidencyError> {
474 let residency = self.lock_state()?;
475 let resident_state = validate_state(&residency, lease)?;
476 resident_state.states.get(name).cloned().ok_or_else(|| {
477 ResourceResidencyError::StateNotFound {
478 state: lease.id,
479 name: name.to_string(),
480 }
481 })
482 }
483
484 pub fn reset_state(&self, lease: StateLease) -> Result<StateLease, ResourceResidencyError> {
486 let mut residency = self.lock_state()?;
487 let resident_state = validate_state(&residency, lease)?;
488 let reset_inputs = resident_state
489 .states
490 .iter()
491 .map(|(name, resource)| {
492 let byte_len = resident_state.state_sizes[name];
493 (resource.clone(), byte_len)
494 })
495 .collect::<Vec<_>>();
496 let next_generation = resident_state
497 .generation
498 .checked_add(1)
499 .ok_or(ResourceResidencyError::StateGenerationOverflow { state: lease.id })?;
500 for (resource, byte_len) in &reset_inputs {
501 if let Err(error) = self.zero_resource(resource, *byte_len) {
502 let removed = residency
503 .states
504 .remove(&lease.id)
505 .ok_or(ResourceResidencyError::StateLeaseNotFound { state: lease.id })?;
506 release_state_accounting(&mut residency, &removed)?;
507 let resources = removed.states.into_values().collect();
508 return Err(self.rollback_error(resources, "state reset", error));
509 }
510 }
511 residency
512 .states
513 .get_mut(&lease.id)
514 .ok_or(ResourceResidencyError::StateLeaseNotFound { state: lease.id })?
515 .generation = next_generation;
516 Ok(StateLease {
517 id: lease.id,
518 generation: next_generation,
519 })
520 }
521
522 pub fn cancel_state(&self, lease: StateLease) -> Result<(), ResourceResidencyError> {
524 self.release_state(lease, "state cancellation")
525 }
526
527 pub fn finish_state(&self, lease: StateLease) -> Result<(), ResourceResidencyError> {
529 self.release_state(lease, "state completion")
530 }
531
532 pub fn evict_resource_set(&self, key: ResourceSetKey) -> Result<(), ResourceResidencyError> {
534 let mut state = self.lock_state()?;
535 let resource_set = state
536 .resource_sets
537 .get(&key)
538 .ok_or(ResourceResidencyError::ResourceSetNotResident { key })?;
539 if resource_set.active_states != 0 {
540 return Err(ResourceResidencyError::ResourceSetInUse {
541 key,
542 active_states: resource_set.active_states,
543 });
544 }
545 let removed = state
546 .resource_sets
547 .remove(&key)
548 .ok_or(ResourceResidencyError::ResourceSetNotResident { key })?;
549 state.used_bytes = state
550 .used_bytes
551 .checked_sub(removed.accounted_bytes)
552 .ok_or(ResourceResidencyError::AccountingUnderflow)?;
553 let resources = removed
554 .immutable_resources
555 .into_values()
556 .map(|immutable_resource| immutable_resource.resource)
557 .collect::<Vec<_>>();
558 self.release_resources(resources, "resource_set eviction")
559 }
560
561 pub fn immutable_resource(
563 &self,
564 key: ResourceSetKey,
565 name: &str,
566 ) -> Result<Resource, ResourceResidencyError> {
567 let state = self.lock_state()?;
568 state
569 .resource_sets
570 .get(&key)
571 .ok_or(ResourceResidencyError::ResourceSetNotResident { key })?
572 .immutable_resources
573 .get(name)
574 .map(|immutable_resource| immutable_resource.resource.clone())
575 .ok_or_else(|| ResourceResidencyError::ImmutableResourceNotFound {
576 key,
577 name: name.to_string(),
578 })
579 }
580
581 pub fn artifact(
583 &self,
584 key: ResourceSetKey,
585 name: &str,
586 ) -> Result<Arc<dyn ArtifactInstance>, ResourceResidencyError> {
587 let state = self.lock_state()?;
588 state
589 .resource_sets
590 .get(&key)
591 .ok_or(ResourceResidencyError::ResourceSetNotResident { key })?
592 .artifacts
593 .get(name)
594 .map(|artifact| Arc::clone(&artifact.instance))
595 .ok_or_else(|| ResourceResidencyError::ArtifactNotFound {
596 key,
597 name: name.to_string(),
598 })
599 }
600
601 pub fn replace_artifact_instance(
608 &self,
609 key: ResourceSetKey,
610 name: &str,
611 instance: Arc<dyn ArtifactInstance>,
612 ) -> Result<(), ResourceResidencyError> {
613 let mut state = self.lock_state()?;
614 let artifact = state
615 .resource_sets
616 .get_mut(&key)
617 .ok_or(ResourceResidencyError::ResourceSetNotResident { key })?
618 .artifacts
619 .get_mut(name)
620 .ok_or_else(|| ResourceResidencyError::ArtifactNotFound {
621 key,
622 name: name.to_string(),
623 })?;
624 if artifact.artifact != instance.artifact().0 {
625 return Err(ResourceResidencyError::WarmResourceSetMismatch);
626 }
627 artifact.instance = instance;
628 Ok(())
629 }
630
631 pub fn used_bytes(&self) -> Result<u64, ResourceResidencyError> {
633 Ok(self.lock_state()?.used_bytes)
634 }
635
636 pub fn active_states(&self, key: ResourceSetKey) -> Result<u64, ResourceResidencyError> {
638 Ok(self
639 .lock_state()?
640 .resource_sets
641 .get(&key)
642 .ok_or(ResourceResidencyError::ResourceSetNotResident { key })?
643 .active_states)
644 }
645
646 fn release_state(
647 &self,
648 lease: StateLease,
649 context: &'static str,
650 ) -> Result<(), ResourceResidencyError> {
651 let mut state = self.lock_state()?;
652 validate_state(&state, lease)?;
653 let removed = state
654 .states
655 .remove(&lease.id)
656 .ok_or(ResourceResidencyError::StateLeaseNotFound { state: lease.id })?;
657 release_state_accounting(&mut state, &removed)?;
658 self.release_resources(removed.states.into_values().collect(), context)
659 }
660
661 fn zero_resource(&self, resource: &Resource, byte_len: usize) -> Result<(), BackendError> {
662 let zeroes = vec![0_u8; ZERO_UPLOAD_CHUNK_BYTES.min(byte_len)];
663 let mut offset = 0_usize;
664 while offset < byte_len {
665 let chunk = (byte_len - offset).min(zeroes.len());
666 self.device.upload_at(resource, offset, &zeroes[..chunk])?;
667 offset += chunk;
668 }
669 Ok(())
670 }
671
672 fn rollback_error(
673 &self,
674 resources: Vec<Resource>,
675 operation: &'static str,
676 source: BackendError,
677 ) -> ResourceResidencyError {
678 match self.release_resources(resources, "admission rollback") {
679 Ok(()) => ResourceResidencyError::Backend {
680 operation,
681 detail: source.to_string(),
682 },
683 Err(cleanup) => ResourceResidencyError::Rollback {
684 operation,
685 detail: source.to_string(),
686 cleanup: cleanup.to_string(),
687 },
688 }
689 }
690
691 fn release_resources(
692 &self,
693 resources: Vec<Resource>,
694 context: &'static str,
695 ) -> Result<(), ResourceResidencyError> {
696 let mut failures = Vec::new();
697 for resource in resources {
698 if let Err(error) = self.device.free(resource) {
699 failures.push(error.to_string());
700 }
701 }
702 if failures.is_empty() {
703 Ok(())
704 } else {
705 Err(ResourceResidencyError::Release {
706 context,
707 details: failures.join("; "),
708 })
709 }
710 }
711
712 fn lock_state(&self) -> Result<MutexGuard<'_, ResidencyState>, ResourceResidencyError> {
713 self.state
714 .lock()
715 .map_err(|_| ResourceResidencyError::LockPoisoned)
716 }
717}
718
719impl Drop for ResourceResidency {
720 fn drop(&mut self) {
721 let state = match self.state.get_mut() {
722 Ok(state) => state,
723 Err(poisoned) => poisoned.into_inner(),
724 };
725 let mut resources = std::mem::take(&mut state.states)
726 .into_values()
727 .flat_map(|state| state.states.into_values())
728 .collect::<Vec<_>>();
729 resources.extend(
730 std::mem::take(&mut state.resource_sets)
731 .into_values()
732 .flat_map(|resource_set| resource_set.immutable_resources.into_values())
733 .map(|immutable_resource| immutable_resource.resource),
734 );
735 for resource in resources {
736 if let Err(error) = self.device.free(resource) {
737 tracing::error!(
738 error = %error,
739 "resource_set residency drop could not release a backend resource"
740 );
741 }
742 }
743 }
744}
745
746struct PreparedImmutableResource {
747 name: String,
748 byte_len: u64,
749 digest: [u8; 32],
750}
751
752struct ValidatedArtifact {
753 name: String,
754 artifact: [u8; 32],
755 byte_len: u64,
756}
757
758fn validate_key(key: ResourceSetKey) -> Result<(), ResourceResidencyError> {
759 if key.source_digest == [0; 32] || key.artifact_digest == [0; 32] {
760 return Err(ResourceResidencyError::ZeroIdentity);
761 }
762 Ok(())
763}
764
765fn validate_immutable_resources(
766 immutable_resources: &[ImmutableResourceUpload<'_>],
767) -> Result<Vec<PreparedImmutableResource>, ResourceResidencyError> {
768 let mut names = BTreeSet::new();
769 let mut prepared = Vec::with_capacity(immutable_resources.len());
770 for immutable_resource in immutable_resources {
771 if immutable_resource.name.is_empty() || !names.insert(immutable_resource.name) {
772 return Err(ResourceResidencyError::DuplicateOrEmptyName {
773 kind: "immutable_resource",
774 name: immutable_resource.name.to_string(),
775 });
776 }
777 let actual = *blake3::hash(immutable_resource.bytes).as_bytes();
778 if actual != immutable_resource.blake3 {
779 return Err(ResourceResidencyError::ImmutableResourceDigestMismatch {
780 name: immutable_resource.name.to_string(),
781 actual,
782 expected: immutable_resource.blake3,
783 });
784 }
785 prepared.push(PreparedImmutableResource {
786 name: immutable_resource.name.to_string(),
787 byte_len: immutable_resource.bytes.len() as u64,
788 digest: immutable_resource.blake3,
789 });
790 }
791 Ok(prepared)
792}
793
794fn validate_artifacts(
795 artifacts: &[ArtifactInstanceBinding],
796) -> Result<Vec<ValidatedArtifact>, ResourceResidencyError> {
797 let mut names = BTreeSet::new();
798 let mut validated = Vec::with_capacity(artifacts.len());
799 for artifact in artifacts {
800 if artifact.name.is_empty() || !names.insert(artifact.name.as_str()) {
801 return Err(ResourceResidencyError::DuplicateOrEmptyName {
802 kind: "artifact",
803 name: artifact.name.clone(),
804 });
805 }
806 validated.push(ValidatedArtifact {
807 name: artifact.name.clone(),
808 artifact: artifact.instance.artifact().0,
809 byte_len: artifact.byte_len,
810 });
811 }
812 Ok(validated)
813}
814
815fn validate_state_specs(
816 specs: &[MutableStateSpec<'_>],
817) -> Result<Vec<(String, usize)>, ResourceResidencyError> {
818 let mut names = BTreeSet::new();
819 let mut prepared = Vec::with_capacity(specs.len());
820 for spec in specs {
821 if spec.name.is_empty() || !names.insert(spec.name) {
822 return Err(ResourceResidencyError::DuplicateOrEmptyName {
823 kind: "mutable state",
824 name: spec.name.to_string(),
825 });
826 }
827 if spec.byte_len == 0 {
828 return Err(ResourceResidencyError::ZeroStateBytes {
829 name: spec.name.to_string(),
830 });
831 }
832 prepared.push((spec.name.to_string(), spec.byte_len));
833 }
834 Ok(prepared)
835}
836
837fn accounted_resource_set_bytes(
838 immutable_resources: &[PreparedImmutableResource],
839 artifacts: &[ValidatedArtifact],
840) -> Result<u64, ResourceResidencyError> {
841 immutable_resources
842 .iter()
843 .map(|immutable_resource| immutable_resource.byte_len)
844 .chain(artifacts.iter().map(|artifact| artifact.byte_len))
845 .try_fold(0_u64, |total, bytes| {
846 total
847 .checked_add(bytes)
848 .ok_or(ResourceResidencyError::ByteLengthOverflow {
849 context: "resource-set and artifact byte total".into(),
850 })
851 })
852}
853
854fn validate_warm_resource_set(
855 resident: &ResidentResourceSet,
856 immutable_resources: &[PreparedImmutableResource],
857 artifacts: &[ValidatedArtifact],
858) -> Result<(), ResourceResidencyError> {
859 if resident.immutable_resources.len() != immutable_resources.len()
860 || resident.artifacts.len() != artifacts.len()
861 {
862 return Err(ResourceResidencyError::WarmResourceSetMismatch);
863 }
864 for immutable_resource in immutable_resources {
865 let Some(existing) = resident.immutable_resources.get(&immutable_resource.name) else {
866 return Err(ResourceResidencyError::WarmResourceSetMismatch);
867 };
868 if existing.byte_len != immutable_resource.byte_len
869 || existing.digest != immutable_resource.digest
870 {
871 return Err(ResourceResidencyError::WarmResourceSetMismatch);
872 }
873 }
874 for artifact in artifacts {
875 let Some(existing) = resident.artifacts.get(&artifact.name) else {
876 return Err(ResourceResidencyError::WarmResourceSetMismatch);
877 };
878 if existing.artifact != artifact.artifact || existing.byte_len != artifact.byte_len {
879 return Err(ResourceResidencyError::WarmResourceSetMismatch);
880 }
881 }
882 Ok(())
883}
884
885fn ensure_budget(
886 used: u64,
887 requested: u64,
888 budget: u64,
889 context: &'static str,
890) -> Result<(), ResourceResidencyError> {
891 let required =
892 used.checked_add(requested)
893 .ok_or(ResourceResidencyError::ByteLengthOverflow {
894 context: context.into(),
895 })?;
896 if required > budget {
897 return Err(ResourceResidencyError::OutOfMemory {
898 context,
899 used,
900 requested,
901 budget,
902 });
903 }
904 Ok(())
905}
906
907fn validate_state(
908 state: &ResidencyState,
909 lease: StateLease,
910) -> Result<&ResidentStateSet, ResourceResidencyError> {
911 let state = state
912 .states
913 .get(&lease.id)
914 .ok_or(ResourceResidencyError::StateLeaseNotFound { state: lease.id })?;
915 if state.generation != lease.generation {
916 return Err(ResourceResidencyError::StaleStateLease {
917 state: lease.id,
918 expected_generation: state.generation,
919 actual_generation: lease.generation,
920 });
921 }
922 Ok(state)
923}
924
925fn release_state_accounting(
926 residency: &mut ResidencyState,
927 resident_state: &ResidentStateSet,
928) -> Result<(), ResourceResidencyError> {
929 residency.used_bytes = residency
930 .used_bytes
931 .checked_sub(resident_state.accounted_bytes)
932 .ok_or(ResourceResidencyError::AccountingUnderflow)?;
933 let resource_set = residency
934 .resource_sets
935 .get_mut(&resident_state.resource_set)
936 .ok_or(ResourceResidencyError::ResourceSetNotResident {
937 key: resident_state.resource_set,
938 })?;
939 resource_set.active_states = resource_set
940 .active_states
941 .checked_sub(1)
942 .ok_or(ResourceResidencyError::AccountingUnderflow)?;
943 Ok(())
944}
945
946#[derive(Debug, Clone, PartialEq, Eq, Error)]
948pub enum ResourceResidencyError {
949 #[error("resident resource-set identity is zero. Fix: use verified source and compiler artifact digests")]
951 ZeroIdentity,
952 #[error(
954 "{kind} name `{name}` is empty or duplicated. Fix: provide one stable name per binding"
955 )]
956 DuplicateOrEmptyName {
957 kind: &'static str,
959 name: String,
961 },
962 #[error("immutable_resource `{name}` does not match its trusted BLAKE3 digest")]
964 ImmutableResourceDigestMismatch {
965 name: String,
967 actual: [u8; 32],
969 expected: [u8; 32],
971 },
972 #[error("residency byte arithmetic overflowed for {context}. Fix: shard the admission")]
974 ByteLengthOverflow {
975 context: String,
977 },
978 #[error("{context} needs {requested} additional bytes with {used} already used, over budget {budget}. Fix: evict idle resource sets or reduce state capacity")]
980 OutOfMemory {
981 context: &'static str,
983 used: u64,
985 requested: u64,
987 budget: u64,
989 },
990 #[error("{operation} failed: {detail}")]
992 Backend {
993 operation: &'static str,
995 detail: String,
997 },
998 #[error("{operation} failed: {detail}; rollback also failed: {cleanup}")]
1000 Rollback {
1001 operation: &'static str,
1003 detail: String,
1005 cleanup: String,
1007 },
1008 #[error("{context} could not release all resident resources: {details}")]
1010 Release {
1011 context: &'static str,
1013 details: String,
1015 },
1016 #[error("warm resource-set request disagrees with resident immutable resource or artifact bindings. Fix: use a new artifact digest for a changed plan")]
1018 WarmResourceSetMismatch,
1019 #[error(
1021 "resource set {key:?} is not resident. Fix: admit the resource_set before starting or binding a state"
1022 )]
1023 ResourceSetNotResident {
1024 key: ResourceSetKey,
1026 },
1027 #[error("resource set {key:?} has {active_states} active states. Fix: finish or cancel them before eviction")]
1029 ResourceSetInUse {
1030 key: ResourceSetKey,
1032 active_states: u64,
1034 },
1035 #[error("state identity space is exhausted. Fix: restart the residency manager rather than reusing stale identities")]
1037 StateIdentityOverflow,
1038 #[error("state {state:?} generation space is exhausted. Fix: finish it and start a new state")]
1040 StateGenerationOverflow {
1041 state: StateId,
1043 },
1044 #[error("state {state:?} is not active. Fix: discard stale leases and start a new state")]
1046 StateLeaseNotFound {
1047 state: StateId,
1049 },
1050 #[error("state {state:?} lease generation {actual_generation} is stale; current generation is {expected_generation}")]
1052 StaleStateLease {
1053 state: StateId,
1055 expected_generation: u64,
1057 actual_generation: u64,
1059 },
1060 #[error("state {state:?} has no state `{name}`")]
1062 StateNotFound {
1063 state: StateId,
1065 name: String,
1067 },
1068 #[error("mutable state `{name}` is zero bytes. Fix: omit unused state or provide its exact positive size")]
1070 ZeroStateBytes {
1071 name: String,
1073 },
1074 #[error("resident resource set {key:?} has no immutable resource `{name}`")]
1076 ImmutableResourceNotFound {
1077 key: ResourceSetKey,
1079 name: String,
1081 },
1082 #[error("resident resource set {key:?} has no artifact `{name}`")]
1084 ArtifactNotFound {
1085 key: ResourceSetKey,
1087 name: String,
1089 },
1090 #[error(
1092 "residency accounting underflowed. Fix: stop using the manager and rebuild residency state"
1093 )]
1094 AccountingUnderflow,
1095 #[error(
1097 "residency state lock is poisoned. Fix: rebuild the manager before admitting more work"
1098 )]
1099 LockPoisoned,
1100}