saluki_core/runtime/state/resources/mod.rs
1//! External resource management.
2
3use std::{
4 any::{Any, TypeId},
5 fmt, mem,
6 ops::{Deref, DerefMut},
7 sync::{Arc, Mutex},
8};
9
10use async_trait::async_trait;
11use saluki_common::collections::FastHashMap;
12use saluki_error::GenericError;
13use serde::Serialize;
14use snafu::Snafu;
15use stringtheory::MetaString;
16use tracing::{debug, warn};
17
18use crate::{runtime::process::Id as ProcessId, support::SubsystemIdentifier};
19
20mod api;
21pub use self::api::{ResourceRegistryAPIHandler, ResourceRegistryState};
22
23mod sublease;
24use self::sublease::SubleaseLedger;
25pub use self::sublease::{Sublease, Subleases};
26
27mod worker;
28pub use self::worker::ResourceRegistryWorker;
29
30#[cfg(test)]
31mod tests;
32
33/// The kind of an external resource.
34///
35/// Deliberately a closed set: the registry coordinates a known, bounded collection of scarce things, so introducing a
36/// kind is a design decision that belongs here rather than in a downstream crate. This names the kind only; it implies
37/// no dependency on the types that implement it.
38#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
39#[serde(rename_all = "snake_case")]
40pub enum ResourceKind {
41 /// A bound network socket.
42 Socket,
43
44 /// A synthetic kind for exercising cross-kind behavior in tests.
45 #[cfg(test)]
46 Test,
47}
48
49impl ResourceKind {
50 /// Returns the string representation of this kind.
51 pub const fn as_str(&self) -> &'static str {
52 match self {
53 Self::Socket => "socket",
54 #[cfg(test)]
55 Self::Test => "test",
56 }
57 }
58}
59
60impl fmt::Display for ResourceKind {
61 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
62 f.write_str(self.as_str())
63 }
64}
65
66/// A specification naming one resource, and the blueprint for creating it.
67///
68/// A specification describes *what* to create without creating it, so it can be built and compared long before the
69/// underlying resource exists. It also names exactly one resource type, which is what keeps two callers from acquiring
70/// the same resource as two different types: there is no way to express the mismatch.
71#[async_trait]
72pub trait ResourceSpecification: Clone + fmt::Debug + Send + Sync + 'static {
73 /// Type of the resource this specification creates.
74 type Resource: Send + 'static;
75
76 /// Kind of the resource this specification names.
77 const KIND: ResourceKind;
78
79 /// Returns the key for the resource this specification names.
80 ///
81 /// Keys must be unique within their kind: two specifications of the same kind that yield the same key are
82 /// understood to name the same underlying scarce thing and will conflict with each other. Keys of different kinds
83 /// never collide, so a key only has to distinguish a resource from its siblings.
84 ///
85 /// Only the parts of a specification that identify the underlying thing belong in the key. Settings that do not
86 /// change *which* resource is named should be left out, so that two specifications differing only in their
87 /// settings are correctly recognized as naming the same resource.
88 fn key(&self) -> MetaString;
89
90 /// Creates the resource.
91 ///
92 /// A resource that is only meaningful as a set of underlying handles -- several sockets bound to one address with
93 /// `SO_REUSEPORT`, say -- holds all of them itself. Creating them together in a single call is what makes them
94 /// atomic: if any one fails, this returns an error and nothing is registered.
95 ///
96 /// `subleases` issues subleases on the resource, for a resource that hands out subresources able to outlive the
97 /// lease they came from -- a connectionless listener lending the bound socket underneath it to every stream it
98 /// yields, say. Keep it on the resource and issue one per subresource; outstanding subleases keep the resource
99 /// from being handed to another acquirer. A resource that is never subdivided has no use for it.
100 ///
101 /// # Errors
102 ///
103 /// If the resource can't be created, an error is returned and nothing is registered.
104 async fn create(&self, subleases: Subleases) -> Result<Self::Resource, GenericError>;
105
106 /// Prepares a returning resource for its next holder.
107 ///
108 /// Called once a returned resource has no subleases outstanding, immediately before it is handed to its next
109 /// holder -- so it runs knowing the previous holder is genuinely finished, including with anything it lent out.
110 /// Implement this only for a resource
111 /// that accumulates state over the course of a single lease and must start clean for the next one -- a listener
112 /// tracking how many of its pre-bound sockets it has handed out, for example. Everything the resource is *for*,
113 /// such as the sockets themselves, must survive: the point of the registry is that it outlives its holders.
114 ///
115 /// Defaults to doing nothing, which is right for a resource that carries no per-lease state.
116 ///
117 /// This is deliberately infallible. A resource that can't be made fit for reuse should be
118 /// [`discard`][ResourceLease::discard]ed by its holder instead, so the next acquisition builds a fresh one.
119 ///
120 /// # Panics
121 ///
122 /// Don't. A panic here aborts the acquisition that triggered it, and the resource goes back to the registry
123 /// intact but only partly reset, so the next acquisition gets the same resource and runs the same reset again.
124 fn reset(_resource: &mut Self::Resource) {}
125}
126
127/// An error that occurred while acquiring a resource.
128#[derive(Debug, Snafu)]
129#[snafu(context(suffix(false)))]
130pub enum AcquireError {
131 /// The resource is already leased by something else in this process.
132 #[snafu(display(
133 "{} resource '{}' is already leased by '{}' (acquired by process {})",
134 kind,
135 key,
136 owner,
137 acquisition_process_id.as_usize()
138 ))]
139 AlreadyLeased {
140 /// Kind of the resource.
141 kind: ResourceKind,
142
143 /// Key of the resource.
144 key: MetaString,
145
146 /// Rendered identity of the subsystem holding the resource.
147 ///
148 /// This is the authoritative answer to who holds the resource. Rendered rather than kept as a
149 /// [`SubsystemIdentifier`], which stores enough segments inline to make this error large enough to slow down
150 /// every `Result` that carries it.
151 owner: MetaString,
152
153 /// Identifier of the process that acquired the resource.
154 acquisition_process_id: ProcessId,
155 },
156
157 /// The key is registered, but holds a different type of resource.
158 ///
159 /// Two specifications of the same kind produced the same key while naming different resource types, which means
160 /// their keys are not as unique as they need to be.
161 #[snafu(display(
162 "{} resource '{}' is registered as `{}`, but was requested as `{}`",
163 kind,
164 key,
165 existing_type,
166 requested_type
167 ))]
168 TypeMismatch {
169 /// Kind of the resource.
170 kind: ResourceKind,
171
172 /// Key of the resource.
173 key: MetaString,
174
175 /// Type the resource was registered as.
176 existing_type: &'static str,
177
178 /// Type the resource was requested as.
179 requested_type: &'static str,
180 },
181
182 /// The resource could not be created.
183 #[snafu(display("failed to create {} resource '{}': {}", kind, key, source))]
184 CreationFailed {
185 /// Kind of the resource.
186 kind: ResourceKind,
187
188 /// Key of the resource.
189 key: MetaString,
190
191 /// Source of the error.
192 source: GenericError,
193 },
194}
195
196/// Identifies one registry entry.
197///
198/// Keying on the kind as well as the key means a kind namespaces its own keys, so two unrelated resource families can
199/// never collide on a coincidentally equal key. It stays deliberately coarser than the resource type: two
200/// representations of the same scarce thing must collide, and keying by type would let each of them claim it.
201#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
202struct EntryKey {
203 kind: ResourceKind,
204 key: MetaString,
205}
206
207impl EntryKey {
208 fn new<S: ResourceSpecification>(spec: &S) -> Self {
209 Self {
210 kind: S::KIND,
211 key: spec.key(),
212 }
213 }
214}
215
216impl fmt::Display for EntryKey {
217 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
218 write!(f, "{}:{}", self.kind, self.key)
219 }
220}
221
222/// Identity of whatever currently holds a resource.
223#[derive(Clone, Debug)]
224struct LeaseInfo {
225 owner: SubsystemIdentifier,
226 acquisition_process_id: ProcessId,
227}
228
229impl LeaseInfo {
230 fn from_owner(owner: &SubsystemIdentifier) -> Self {
231 Self {
232 owner: owner.clone(),
233 acquisition_process_id: ProcessId::current(),
234 }
235 }
236}
237
238/// RAII guard to disarm in-progress leases when resource creation fails to run to completion.
239///
240/// [`ResourceRegistry::acquire`] marks an entry as [`EntryState::Creating`] before awaiting
241/// [`ResourceSpecification::create`], and
242/// that await is a cancellation point: a supervisor aborts a worker that is still initializing when shutdown arrives,
243/// so a builder's acquisition can be dropped partway through. Without this guard the claim would outlive the
244/// acquisition and block the key for the life of the process.
245struct ClaimCreationGuard<'a> {
246 registry: &'a ResourceRegistry,
247 key: &'a EntryKey,
248 armed: bool,
249}
250
251impl<'a> ClaimCreationGuard<'a> {
252 /// Creates a new guard for the given key in the armed state.
253 fn from_key(registry: &'a ResourceRegistry, key: &'a EntryKey) -> Self {
254 Self {
255 registry,
256 key,
257 armed: true,
258 }
259 }
260
261 /// Disarm and consume the guard.
262 fn disarm(mut self) {
263 self.armed = false;
264 }
265}
266
267impl Drop for ClaimCreationGuard<'_> {
268 fn drop(&mut self) {
269 if !self.armed {
270 return;
271 }
272
273 let mut state = self.registry.inner.lock().unwrap();
274
275 // Only take back a claim that is still ours to take back.
276 if matches!(
277 state.entries.get(self.key).map(|entry| &entry.state),
278 Some(EntryState::Creating(_))
279 ) {
280 debug!(key = %self.key, "Resource creation was cancelled. Releasing the claim on its key.");
281 state.entries.remove(self.key);
282 }
283 }
284}
285
286/// RAII guard to release a claim on an existing entry when an acquisition doesn't run to completion.
287///
288/// [`ResourceRegistry::acquire`] claims an idle or discarded entry before awaiting its outstanding subleases, and that
289/// await is a cancellation point. Without this guard the claim would outlive the acquisition and block the key for the
290/// life of the process.
291///
292/// Unlike [`ClaimCreationGuard`], this only ever clears a flag: the entry keeps whatever it was holding for the whole
293/// wait, so there is no path here that can drop a resource and release the underlying resource. A tombstone left
294/// behind by a cancelled acquisition costs nothing beyond its entry, and the next acquirer clears it.
295struct ClaimedEntryGuard<'a> {
296 registry: &'a ResourceRegistry,
297 key: &'a EntryKey,
298 armed: bool,
299}
300
301impl<'a> ClaimedEntryGuard<'a> {
302 /// Creates a new guard for the given key in the armed state.
303 fn from_key(registry: &'a ResourceRegistry, key: &'a EntryKey) -> Self {
304 Self {
305 registry,
306 key,
307 armed: true,
308 }
309 }
310
311 /// Disarm and consume the guard.
312 fn disarm(mut self) {
313 self.armed = false;
314 }
315}
316
317impl Drop for ClaimedEntryGuard<'_> {
318 fn drop(&mut self) {
319 if !self.armed {
320 return;
321 }
322
323 let mut state = self.registry.inner.lock().unwrap();
324 if let Some(entry) = state.entries.get_mut(self.key) {
325 debug!(key = %self.key, "Resource acquisition was cancelled. Releasing the claim on its entry.");
326 entry.release_claim();
327 }
328 }
329}
330
331/// Lifecycle state of a registry entry.
332enum EntryState {
333 /// The resource is held by the registry.
334 ///
335 /// `claimed_by` is set while an acquisition waits for outstanding subleases to be returned: the resource is still
336 /// right here, but it is already spoken for. Keeping the resource in the entry rather than moving it into the
337 /// waiting acquisition is deliberate -- it means no code path, including a cancelled acquisition, can drop it and
338 /// release the underlying resource.
339 Idle {
340 value: Box<dyn Any + Send>,
341 claimed_by: Option<LeaseInfo>,
342 },
343
344 /// Creation is in flight. The entry holds nothing yet, but is already spoken for.
345 Creating(LeaseInfo),
346
347 /// The resource is lent out.
348 Leased(LeaseInfo),
349
350 /// The resource was discarded while it still had subleases outstanding.
351 ///
352 /// A tombstone: the resource itself is gone, but the things it lent out are not, and for a subdivided resource
353 /// those are what hold the underlying resource open -- a connectionless listener's socket stays bound and
354 /// receiving for as long as a stream is reading it. The entry stays behind to keep the key reserved until they
355 /// come back, because building a replacement in the meantime would put two of the resource in the world at once.
356 ///
357 /// `claimed_by` works as it does for [`Idle`][Self::Idle], serializing acquirers waiting on the same key. The
358 /// acquisition that wins the wait replaces this entry with one of its own.
359 Discarded { claimed_by: Option<LeaseInfo> },
360}
361
362impl EntryState {
363 fn holder(&self) -> Option<&LeaseInfo> {
364 match self {
365 Self::Idle { claimed_by, .. } | Self::Discarded { claimed_by } => claimed_by.as_ref(),
366 Self::Creating(info) | Self::Leased(info) => Some(info),
367 }
368 }
369}
370
371/// One resource registered under a single key.
372struct Entry {
373 type_id: TypeId,
374 type_name: &'static str,
375 spec_desc: String,
376 state: EntryState,
377 acquisitions: u64,
378
379 /// Subleases taken out on this resource.
380 ///
381 /// Lives on the entry rather than on any one lease: a sublease can outlive the head lease that issued it, and the
382 /// resource keeps its [`Subleases`] across every lease it is handed out under.
383 subleases: Arc<SubleaseLedger>,
384}
385
386impl Entry {
387 fn new<S: ResourceSpecification>(spec: &S, lease_info: LeaseInfo) -> Self {
388 Self {
389 type_id: TypeId::of::<S::Resource>(),
390 type_name: std::any::type_name::<S::Resource>(),
391 spec_desc: format!("{:?}", spec),
392 state: EntryState::Creating(lease_info),
393 acquisitions: 0,
394 subleases: SubleaseLedger::new(),
395 }
396 }
397
398 /// Reported lifecycle state of this entry.
399 fn state_name(&self) -> &'static str {
400 match &self.state {
401 // The head lease is back, but the resource isn't available until its subleases are too.
402 EntryState::Idle { .. } if self.subleases.outstanding() > 0 => "subleased",
403 EntryState::Idle { .. } => "idle",
404 EntryState::Creating(_) => "creating",
405 EntryState::Leased(_) => "leased",
406 EntryState::Discarded { .. } => "discarded",
407 }
408 }
409
410 /// Claims this entry for an acquisition, or explains why it can't be claimed.
411 ///
412 /// Nothing is handed over here. Either way the caller has to wait for any outstanding subleases first, and that
413 /// wait happens with the registry lock released, so this returns the entry's sublease ledger to wait on alongside
414 /// what to do once it settles.
415 fn claim<S: ResourceSpecification>(
416 &mut self, key: &EntryKey, new_spec: &S, lease_info: LeaseInfo,
417 ) -> Result<Claim, AcquireError> {
418 let already_leased = |holder: &LeaseInfo| AcquireError::AlreadyLeased {
419 kind: S::KIND,
420 key: key.key.clone(),
421 owner: MetaString::from(holder.owner.to_string()),
422 acquisition_process_id: holder.acquisition_process_id,
423 };
424
425 // A discarded entry holds no resource, only the key, so there is nothing here to type-check against or to
426 // hand over: whoever wins the wait builds fresh from their own specification. Checking the type of a resource
427 // that has already been dropped would refuse an acquisition over a resource that no longer exists.
428 if let EntryState::Discarded { claimed_by } = &mut self.state {
429 if let Some(holder) = claimed_by.as_ref() {
430 return Err(already_leased(holder));
431 }
432
433 *claimed_by = Some(lease_info);
434
435 return Ok(Claim::Rebuild(Arc::clone(&self.subleases)));
436 }
437
438 if self.type_id != TypeId::of::<S::Resource>() {
439 return Err(AcquireError::TypeMismatch {
440 kind: S::KIND,
441 key: key.key.clone(),
442 existing_type: self.type_name,
443 requested_type: std::any::type_name::<S::Resource>(),
444 });
445 }
446
447 if let Some(holder) = self.state.holder() {
448 return Err(already_leased(holder));
449 }
450
451 // The key identifies the resource, so a specification differing only in its settings still names this same
452 // resource. Hand back what exists rather than rebuilding it, but say so, since the new settings have no effect.
453 let new_spec_desc = format!("{:?}", new_spec);
454 if self.spec_desc != new_spec_desc {
455 warn!(
456 %key,
457 existing = %self.spec_desc,
458 requested = %new_spec_desc,
459 "Resource acquired with a different specification than it was created with. Using the existing resource; \
460 the requested specification has no effect."
461 );
462 }
463
464 // `holder` returned `None` just above, and the discarded case returned earlier, so the entry is idle and
465 // unclaimed.
466 match &mut self.state {
467 EntryState::Idle { claimed_by, .. } => *claimed_by = Some(lease_info),
468 _ => unreachable!("entry without a holder is idle or discarded"),
469 }
470
471 Ok(Claim::HandOver(Arc::clone(&self.subleases)))
472 }
473
474 /// Hands a claimed resource to its acquirer.
475 ///
476 /// Runs once every sublease has been returned, so the previous holder is genuinely finished with the resource --
477 /// including with anything it lent out.
478 ///
479 /// Resetting the resource is the caller's job, not this method's: [`ResourceSpecification::reset`] is
480 /// implementor-supplied code, and running it here would run it while the registry lock is held and while the
481 /// resource is owned by nothing but a local. See [`ResourceRegistry::acquire`].
482 fn hand_over<S: ResourceSpecification>(
483 &mut self, registry: &ResourceRegistry, key: &EntryKey,
484 ) -> ResourceLease<S::Resource> {
485 // Take the claim first, so the state can be replaced wholesale without needing a placeholder to stand in for
486 // the resource while it moves.
487 let lease_info = match &mut self.state {
488 EntryState::Idle { claimed_by, .. } => claimed_by.take().expect("a claimed entry holds its claim"),
489 _ => unreachable!("a claimed entry is idle"),
490 };
491
492 let value = match mem::replace(&mut self.state, EntryState::Leased(lease_info)) {
493 EntryState::Idle { value, .. } => value,
494 _ => unreachable!("a claimed entry is idle and holds its value"),
495 };
496
497 self.acquisitions += 1;
498
499 ResourceLease {
500 value: Some(
501 *value
502 .downcast::<S::Resource>()
503 .expect("entry type checked when claimed"),
504 ),
505 registry: registry.clone(),
506 key: key.clone(),
507 }
508 }
509
510 /// Releases a claim without handing the resource over.
511 fn release_claim(&mut self) {
512 match &mut self.state {
513 EntryState::Idle { claimed_by, .. } | EntryState::Discarded { claimed_by } => *claimed_by = None,
514 EntryState::Creating(_) | EntryState::Leased(_) => {}
515 }
516 }
517}
518
519/// What an acquisition that claimed an existing entry has to do once the entry's subleases settle.
520enum Claim {
521 /// The resource is in the entry, waiting to be handed over.
522 HandOver(Arc<SubleaseLedger>),
523
524 /// The entry is a tombstone for a discarded resource, and a replacement has to be built in its place.
525 Rebuild(Arc<SubleaseLedger>),
526}
527
528impl Claim {
529 /// The ledger whose subleases have to settle before this claim can be acted on.
530 fn subleases(&self) -> &SubleaseLedger {
531 match self {
532 Self::HandOver(subleases) | Self::Rebuild(subleases) => subleases,
533 }
534 }
535}
536
537#[derive(Default)]
538struct RegistryState {
539 entries: FastHashMap<EntryKey, Entry>,
540}
541
542impl RegistryState {
543 fn snapshot(&self) -> Vec<ResourceStatus> {
544 let mut statuses = self
545 .entries
546 .iter()
547 .map(|(key, entry)| ResourceStatus {
548 kind: key.kind,
549 key: key.key.to_string(),
550 spec: entry.spec_desc.clone(),
551 state: entry.state_name(),
552 owner: entry.state.holder().map(|info| info.owner.to_string()),
553 acquisition_process_id: entry.state.holder().map(|info| info.acquisition_process_id.as_usize()),
554 acquisitions: entry.acquisitions,
555 outstanding_subleases: entry.subleases.outstanding(),
556 })
557 .collect::<Vec<_>>();
558 statuses.sort_by(|a, b| (a.kind, &a.key).cmp(&(b.kind, &b.key)));
559
560 statuses
561 }
562}
563
564/// A registry for scarce, externally backed resources.
565///
566/// In many cases, data planes will have to interact with the outside world by way of exposing network endpoints, or
567/// exposing files, and so on... referred to here as "resources." These resources are unique, or are conceptually meant
568/// to be unique: there should be no other OS processes trying to take ownership of them, and only one part of the code
569/// in the data plane should own them.
570///
571/// A [`ResourceRegistry`] owns those resources on behalf of the entire data plane and lends them out. A child process
572/// never owns a resource, but instead holds a [`ResourceLease`]. When the lease drops -- including when the child
573/// process holding it dies -- the resource returns to the registry intact and still live, ready for the next acquirer.
574/// Since the registry outlives the components that use its resources, a component can be torn down and rebuilt without
575/// the underlying resource being automatically released back to the operating system due to typical Rust drop
576/// semantics.
577///
578/// # Groups
579///
580/// Resources are named by a [`ResourceSpecification`], which provides the blueprint for how to create a particular
581/// resource, such as a network socket, when a caller attempts to acquire it. Resource specifications are generally tied
582/// one-to-one with a particular type.
583///
584/// The specification provides both the information necessary to properly determine one unique resource from
585/// another, as well as a mechanism for consistent creation of potentially complex resources, including asynchronous
586/// initialization.
587///
588/// # Keys and conflicts
589///
590/// Entries are keyed by kind and key together, deliberately *not* by resource type. Two different Rust types can
591/// easily describe the same scarce thing -- a connection-oriented listener and a general one over the same address --
592/// and keying by type would let each of them claim it. Kind is coarse enough that such representations still collide,
593/// while keeping unrelated resource families from colliding on a coincidentally equal key. A key therefore only has to
594/// be unique within its own kind.
595#[derive(Clone, Default)]
596pub struct ResourceRegistry {
597 inner: Arc<Mutex<RegistryState>>,
598}
599
600impl ResourceRegistry {
601 /// Creates an empty registry.
602 pub fn new() -> Self {
603 Self::default()
604 }
605
606 /// Acquires the resource named by `spec`, creating it if it isn't registered yet.
607 ///
608 /// `owner` identifies the subsystem taking the lease and is recorded, alongside the current process identifier, for
609 /// accounting.
610 ///
611 /// If the resource is already registered, the existing one is handed back rather than a new one being created. This
612 /// is the mechanism by which a resource outlives the components that use it.
613 ///
614 /// # Errors
615 ///
616 /// If the resource is already leased, if the key is registered to a different type of resource, or if creation
617 /// fails, an error is returned.
618 pub async fn acquire<S: ResourceSpecification>(
619 &self, owner: &SubsystemIdentifier, spec: S,
620 ) -> Result<ResourceLease<S::Resource>, AcquireError> {
621 let key = EntryKey::new(&spec);
622 let lease_info = LeaseInfo::from_owner(owner);
623
624 // Claim the resource if it's already registered.
625 //
626 // Otherwise, start the registration process by inserting an uninitialized entry that gives us lease ownership
627 // prior to actually creating the resource and finalizing it.
628 let claim = {
629 let mut state = self.inner.lock().unwrap();
630 match state.entries.get_mut(&key) {
631 Some(entry) => Some(entry.claim::<S>(&key, &spec, lease_info.clone())?),
632 None => {
633 let new_entry = Entry::new(&spec, lease_info.clone());
634 state.entries.insert(key.clone(), new_entry);
635 None
636 }
637 }
638 };
639
640 if let Some(claim) = claim {
641 // Whichever way the claim goes, an entry that already existed has to settle its subleases first:
642 // whatever the previous holder lent out is still in use, and for a subdivided resource that means the
643 // underlying resource is still in use. A connectionless listener's socket is still bound and receiving
644 // until the stream holding a sublease on it is dropped, so neither handing that listener over nor binding
645 // a replacement for it is safe while one is outstanding.
646 //
647 // This resolves immediately for a resource that was never subdivided, which is most of them.
648 let claim_guard = ClaimedEntryGuard::from_key(self, &key);
649 claim.subleases().settled().await;
650 claim_guard.disarm();
651
652 let mut state = self.inner.lock().unwrap();
653 match claim {
654 Claim::HandOver(_) => {
655 let entry = state
656 .entries
657 .get_mut(&key)
658 .expect("entry was claimed above and a claimed entry is only removed by its holder");
659
660 let mut lease = entry.hand_over::<S>(self, &key);
661
662 // Clear whatever the previous holder accumulated, now that it and its subleases are all gone.
663 //
664 // Deliberately done here rather than in `hand_over`: the reset is implementor-supplied code, and
665 // two things have to be true before it runs. The lock has to be released, or a panic in it would
666 // poison the registry for every other key. And the resource has to already be owned by its lease,
667 // so that the same panic unwinds through `ResourceLease::drop` and returns the resource to the
668 // registry, still live, instead of dropping it and releasing the underlying resource.
669 drop(state);
670 S::reset(&mut lease);
671
672 debug!(%key, %owner, "Acquired resource.");
673
674 return Ok(lease);
675 }
676 Claim::Rebuild(_) => {
677 // The discarded resource is finally gone in full, so the tombstone has done its job. Replace it
678 // with an uninitialized entry of our own and fall through to creation, exactly as if the key had
679 // been free all along.
680 debug!(%key, "Discarded resource has fully returned. Recreating it for the next holder.");
681 state.entries.insert(key.clone(), Entry::new(&spec, lease_info.clone()));
682 }
683 }
684 }
685
686 // Create the resource.
687 //
688 // We establish a "creation guard" which is a drop guard that ensures we remove our pending entry if we fail to
689 // create the resource, including if this asynchronous call is cancelled, so that we don't permanently tie up
690 // the resource in an uninitialized state.
691 let subleases = {
692 let state = self.inner.lock().unwrap();
693 let entry = state.entries.get(&key).expect("entry was just inserted");
694 Subleases::from_ledger(&entry.subleases)
695 };
696
697 let claim_guard = ClaimCreationGuard::from_key(self, &key);
698 let created = spec.create(subleases).await;
699 claim_guard.disarm();
700
701 let mut state = self.inner.lock().unwrap();
702 match created {
703 Ok(value) => {
704 let entry = state
705 .entries
706 .get_mut(&key)
707 .expect("entry was inserted before creation and is only removed by this function");
708
709 entry.acquisitions += 1;
710 entry.state = EntryState::Leased(lease_info);
711
712 debug!(%key, %owner, "Created resource.");
713
714 Ok(ResourceLease {
715 value: Some(value),
716 registry: self.clone(),
717 key,
718 })
719 }
720 Err(source) => {
721 // Drop the claim so that a later acquire can retry.
722 state.entries.remove(&key);
723 Err(AcquireError::CreationFailed {
724 kind: S::KIND,
725 key: key.key,
726 source,
727 })
728 }
729 }
730 }
731
732 /// Returns a snapshot of every registered resource, ordered by key.
733 pub fn snapshot(&self) -> Vec<ResourceStatus> {
734 let state = self.inner.lock().unwrap();
735 state.snapshot()
736 }
737
738 /// Creates an API handler for reporting the state of all registered resources.
739 pub fn api_handler(&self) -> ResourceRegistryAPIHandler {
740 ResourceRegistryAPIHandler::from_registry(self.clone())
741 }
742
743 /// Creates a [`ResourceRegistryWorker`] that publishes the registry over the control plane.
744 pub fn worker(&self) -> ResourceRegistryWorker {
745 ResourceRegistryWorker::new(self.clone())
746 }
747
748 /// Returns a resource to the registry, marking its entry idle.
749 fn return_value(&self, key: &EntryKey, value: Box<dyn Any + Send>) {
750 let mut state = self.inner.lock().unwrap();
751 if let Some(entry) = state.entries.get_mut(key) {
752 debug!(%key, outstanding_subleases = entry.subleases.outstanding(), "Resource returned to registry.");
753
754 // Deliberately not reset here: subleases issued by the departing holder can still be outstanding, so the
755 // resource isn't finished being used yet. `Entry::hand_over` resets it once they have all come back.
756 entry.state = EntryState::Idle {
757 value,
758 claimed_by: None,
759 };
760 }
761 }
762
763 /// Drops a resource instead of returning it, so that the next acquisition creates a fresh one.
764 ///
765 /// The resource itself is already gone by the time this runs -- [`ResourceLease::discard`] drops it -- but a
766 /// subdivided resource isn't released by that alone: its subresources hold the underlying resource open, so a
767 /// connectionless listener's socket stays bound until the last stream reading it is dropped. Dropping the entry
768 /// now would let the next acquisition bind a replacement alongside that socket, and with `SO_REUSEPORT` the two
769 /// would quietly split incoming datagrams between them. So the entry stays as a tombstone until its subleases
770 /// come back, keeping the key reserved without keeping anything alive.
771 fn discard_value(&self, key: &EntryKey) {
772 let mut state = self.inner.lock().unwrap();
773 let Some(entry) = state.entries.get_mut(key) else {
774 return;
775 };
776
777 let outstanding = entry.subleases.outstanding();
778 if outstanding == 0 {
779 state.entries.remove(key);
780 debug!(%key, "Resource discarded; it will be recreated on the next acquisition.");
781 } else {
782 entry.state = EntryState::Discarded { claimed_by: None };
783 debug!(
784 %key,
785 outstanding_subleases = outstanding,
786 "Resource discarded with subleases outstanding. Its key stays reserved until they are returned, \
787 after which it will be recreated on the next acquisition."
788 );
789 }
790 }
791}
792
793/// Reported state of a single resource.
794#[derive(Clone, Debug, Serialize)]
795pub struct ResourceStatus {
796 /// Kind of the resource.
797 pub kind: ResourceKind,
798
799 /// Key the resource is registered under, unique within its kind.
800 pub key: String,
801
802 /// Rendered specification the resource was created from.
803 pub spec: String,
804
805 /// Lifecycle state of the resource: `creating`, `leased`, `idle`, `subleased`, or `discarded`.
806 ///
807 /// `subleased` is idle-but-unavailable: the head lease is back, but something the holder lent out is still in
808 /// use. `discarded` means the resource is gone and only its key is still reserved, for the same reason.
809 pub state: &'static str,
810
811 /// Subsystem holding the resource, if any.
812 ///
813 /// This is the authoritative answer to who holds the resource.
814 pub owner: Option<String>,
815
816 /// Process that acquired the resource, if any.
817 ///
818 /// This is the process that ran the acquisition, which is not necessarily the one using the resource now: a
819 /// component acquires while it is being built, and only afterwards does it get a process of its own.
820 ///
821 /// Use [`owner`][Self::owner] to identify the holder.
822 pub acquisition_process_id: Option<usize>,
823
824 /// Number of times the resource has been acquired.
825 pub acquisitions: u64,
826
827 /// Number of subleases currently outstanding on the resource.
828 ///
829 /// Non-zero once a holder has released the resource while something it lent out is still in use. The resource
830 /// isn't handed to its next acquirer until this reaches zero.
831 pub outstanding_subleases: usize,
832}
833
834/// An exclusive lease on a resource.
835///
836/// Dereferences to the resource itself. Dropping the lease returns the resource to the registry still live, so a lease
837/// is a loan, never ownership. See [`ResourceRegistry`] for the full model.
838///
839/// This is the *head* lease. Dropping it relinquishes the holder's claim, but doesn't end the resource's lease while
840/// [`Sublease`]s issued against it are still outstanding: the registry won't hand the resource to another acquirer
841/// until those come back too.
842pub struct ResourceLease<R: Send + 'static> {
843 value: Option<R>,
844 registry: ResourceRegistry,
845 key: EntryKey,
846}
847
848impl<R: Send + 'static> ResourceLease<R> {
849 /// Returns the kind of this resource.
850 pub fn kind(&self) -> ResourceKind {
851 self.key.kind
852 }
853
854 /// Returns the key this resource is registered under, unique within its kind.
855 pub fn key(&self) -> &MetaString {
856 &self.key.key
857 }
858
859 /// Returns the resource to the registry.
860 ///
861 /// Equivalent to dropping the lease; useful where the return should be obvious at the call site.
862 pub fn release(self) {}
863
864 /// Drops the resource instead of returning it, so the next acquisition creates a fresh one.
865 ///
866 /// Use this when the resource has hit an error it can't recover from and handing it to the next acquirer would pass
867 /// the problem along.
868 ///
869 /// Discarding doesn't free the key any sooner than releasing would. Outstanding [`Sublease`]s are what hold a
870 /// subdivided resource open, so the key stays reserved until they come back, and only then does the next
871 /// acquisition build a replacement -- otherwise the replacement would exist alongside the thing being discarded.
872 pub fn discard(mut self) {
873 // Dropping the value here is the point: it is what releases the underlying resource.
874 let _ = self.value.take();
875 self.registry.discard_value(&self.key);
876 }
877}
878
879impl<R: Send + 'static> Deref for ResourceLease<R> {
880 type Target = R;
881
882 fn deref(&self) -> &Self::Target {
883 self.value.as_ref().expect("lease holds its value until dropped")
884 }
885}
886
887impl<R: Send + 'static> DerefMut for ResourceLease<R> {
888 fn deref_mut(&mut self) -> &mut Self::Target {
889 self.value.as_mut().expect("lease holds its value until dropped")
890 }
891}
892
893impl<R: Send + 'static> Drop for ResourceLease<R> {
894 fn drop(&mut self) {
895 if let Some(value) = self.value.take() {
896 self.registry.return_value(&self.key, Box::new(value));
897 }
898 }
899}
900
901impl<R: Send + fmt::Debug + 'static> fmt::Debug for ResourceLease<R> {
902 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
903 f.debug_struct("ResourceLease")
904 .field("key", &self.key)
905 .field("value", &self.value)
906 .finish()
907 }
908}