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}