lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

cache.rs (11003B)


      1 //! Single-writer, latest-value publication for passive operations reads.
      2 
      3 use core::fmt;
      4 use std::sync::Arc;
      5 
      6 use tokio::sync::watch;
      7 
      8 use super::{ServiceOperationalState, StatusContractError};
      9 
     10 /// One immutable cached observation and its service-owned metrics metadata.
     11 ///
     12 /// The holder retains only the latest `Arc`, so publication has a fixed
     13 /// one-snapshot capacity. The metadata remains typed and is completed by the
     14 /// bounded metrics contract rather than interpreted by the host cache.
     15 #[derive(Debug)]
     16 pub struct CachedServiceState<M> {
     17     operational: ServiceOperationalState,
     18     metrics: M,
     19 }
     20 
     21 impl<M> CachedServiceState<M> {
     22     #[must_use]
     23     pub const fn new(operational: ServiceOperationalState, metrics: M) -> Self {
     24         Self {
     25             operational,
     26             metrics,
     27         }
     28     }
     29 
     30     #[must_use]
     31     pub const fn operational(&self) -> &ServiceOperationalState {
     32         &self.operational
     33     }
     34 
     35     #[must_use]
     36     pub const fn metrics(&self) -> &M {
     37         &self.metrics
     38     }
     39 }
     40 
     41 /// The sole update authority for one cached service-state channel.
     42 ///
     43 /// This handle deliberately does not implement `Clone`. Services may create
     44 /// any number of readers, but ownership of lifecycle publication stays
     45 /// explicit and singular.
     46 pub struct CachedServiceStatePublisher<M> {
     47     sender: watch::Sender<Arc<CachedServiceState<M>>>,
     48 }
     49 
     50 impl<M> CachedServiceStatePublisher<M> {
     51     /// Publishes one legal lifecycle update and atomically replaces the cache.
     52     pub fn publish(&mut self, next: CachedServiceState<M>) -> Result<(), StatusContractError> {
     53         let current_phase = self.sender.borrow().operational().phase();
     54         let next_phase = next.operational().phase();
     55         if !current_phase.can_transition_to(next_phase) {
     56             return Err(StatusContractError::IllegalTransition {
     57                 from: current_phase,
     58                 to: next_phase,
     59             });
     60         }
     61         self.sender.send_replace(Arc::new(next));
     62         Ok(())
     63     }
     64 
     65     /// Adds a passive reader without sharing publication authority.
     66     #[must_use]
     67     pub fn subscribe(&self) -> CachedServiceStateReader<M> {
     68         CachedServiceStateReader {
     69             receiver: self.sender.subscribe(),
     70         }
     71     }
     72 }
     73 
     74 /// A cloneable passive view of the latest cached service state.
     75 pub struct CachedServiceStateReader<M> {
     76     receiver: watch::Receiver<Arc<CachedServiceState<M>>>,
     77 }
     78 
     79 impl<M> Clone for CachedServiceStateReader<M> {
     80     fn clone(&self) -> Self {
     81         Self {
     82             receiver: self.receiver.clone(),
     83         }
     84     }
     85 }
     86 
     87 impl<M> CachedServiceStateReader<M> {
     88     /// Returns the latest snapshot without executing a probe or awaiting I/O.
     89     #[must_use]
     90     pub fn snapshot(&self) -> Arc<CachedServiceState<M>> {
     91         Arc::clone(&self.receiver.borrow())
     92     }
     93 
     94     /// Waits for a later publication and then returns that latest snapshot.
     95     pub async fn changed(&mut self) -> Result<Arc<CachedServiceState<M>>, StatusPublisherDropped> {
     96         self.receiver
     97             .changed()
     98             .await
     99             .map_err(|_| StatusPublisherDropped)?;
    100         Ok(self.snapshot_and_mark_seen())
    101     }
    102 
    103     fn snapshot_and_mark_seen(&mut self) -> Arc<CachedServiceState<M>> {
    104         Arc::clone(&self.receiver.borrow_and_update())
    105     }
    106 }
    107 
    108 /// Creates a bounded latest-value cache with one publisher and one reader.
    109 #[must_use]
    110 pub fn cached_service_state<M>(
    111     initial: CachedServiceState<M>,
    112 ) -> (CachedServiceStatePublisher<M>, CachedServiceStateReader<M>) {
    113     let (sender, receiver) = watch::channel(Arc::new(initial));
    114     (
    115         CachedServiceStatePublisher { sender },
    116         CachedServiceStateReader { receiver },
    117     )
    118 }
    119 
    120 /// Indicates that no owner remains able to publish another snapshot.
    121 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    122 pub struct StatusPublisherDropped;
    123 
    124 impl fmt::Display for StatusPublisherDropped {
    125     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    126         formatter.write_str("cached service-state publisher was dropped")
    127     }
    128 }
    129 
    130 impl std::error::Error for StatusPublisherDropped {}
    131 
    132 #[cfg(test)]
    133 mod tests {
    134     use std::sync::{
    135         Arc,
    136         atomic::{AtomicUsize, Ordering},
    137     };
    138 
    139     use super::*;
    140     use crate::{CommonReasonCode, Readiness, ReasonCodes, ServicePhase};
    141 
    142     #[derive(Debug)]
    143     struct MetricsMetadata {
    144         revision: u64,
    145         probe_calls: Arc<AtomicUsize>,
    146     }
    147 
    148     fn state(
    149         phase: ServicePhase,
    150         readiness: Readiness,
    151         reasons: ReasonCodes,
    152     ) -> ServiceOperationalState {
    153         ServiceOperationalState::new(phase, readiness, reasons).unwrap()
    154     }
    155 
    156     fn metadata(revision: u64, probe_calls: &Arc<AtomicUsize>) -> MetricsMetadata {
    157         MetricsMetadata {
    158             revision,
    159             probe_calls: Arc::clone(probe_calls),
    160         }
    161     }
    162 
    163     #[tokio::test]
    164     async fn concurrent_readers_observe_the_latest_atomic_publication() {
    165         let probe_calls = Arc::new(AtomicUsize::new(0));
    166         let (mut publisher, reader) = cached_service_state(CachedServiceState::new(
    167             state(
    168                 ServicePhase::Starting,
    169                 Readiness::NOT_READY,
    170                 ReasonCodes::empty(),
    171             ),
    172             metadata(0, &probe_calls),
    173         ));
    174         let readers: Vec<_> = (0..8)
    175             .map(|_| {
    176                 let mut reader = reader.clone();
    177                 tokio::spawn(async move {
    178                     loop {
    179                         let snapshot = reader.changed().await.unwrap();
    180                         if snapshot.metrics().revision == 100 {
    181                             return (
    182                                 snapshot.operational().phase(),
    183                                 snapshot.operational().readiness(),
    184                             );
    185                         }
    186                     }
    187                 })
    188             })
    189             .collect();
    190 
    191         publisher
    192             .publish(CachedServiceState::new(
    193                 state(ServicePhase::Ready, Readiness::READY, ReasonCodes::empty()),
    194                 metadata(1, &probe_calls),
    195             ))
    196             .unwrap();
    197         for revision in 2..=100 {
    198             publisher
    199                 .publish(CachedServiceState::new(
    200                     state(ServicePhase::Ready, Readiness::READY, ReasonCodes::empty()),
    201                     metadata(revision, &probe_calls),
    202                 ))
    203                 .unwrap();
    204         }
    205 
    206         for reader in readers {
    207             assert_eq!(
    208                 reader.await.unwrap(),
    209                 (ServicePhase::Ready, Readiness::READY)
    210             );
    211         }
    212         assert_eq!(probe_calls.load(Ordering::SeqCst), 0);
    213     }
    214 
    215     #[test]
    216     fn shutdown_state_and_illegal_updates_are_explicit() {
    217         let probe_calls = Arc::new(AtomicUsize::new(0));
    218         let (mut publisher, reader) = cached_service_state(CachedServiceState::new(
    219             state(ServicePhase::Ready, Readiness::READY, ReasonCodes::empty()),
    220             metadata(4, &probe_calls),
    221         ));
    222         let shutdown_reasons =
    223             ReasonCodes::new([CommonReasonCode::ShutdownInProgress.into()]).unwrap();
    224         publisher
    225             .publish(CachedServiceState::new(
    226                 state(
    227                     ServicePhase::Stopping,
    228                     Readiness::NOT_READY,
    229                     shutdown_reasons,
    230                 ),
    231                 metadata(5, &probe_calls),
    232             ))
    233             .unwrap();
    234 
    235         let stopping = reader.snapshot();
    236         assert_eq!(stopping.operational().phase(), ServicePhase::Stopping);
    237         assert!(!stopping.operational().readiness().is_ready());
    238         assert_eq!(
    239             stopping.operational().reasons().as_slice()[0].as_str(),
    240             CommonReasonCode::ShutdownInProgress.as_str()
    241         );
    242 
    243         assert_eq!(
    244             publisher.publish(CachedServiceState::new(
    245                 state(ServicePhase::Ready, Readiness::READY, ReasonCodes::empty(),),
    246                 metadata(6, &probe_calls),
    247             )),
    248             Err(StatusContractError::IllegalTransition {
    249                 from: ServicePhase::Stopping,
    250                 to: ServicePhase::Ready,
    251             })
    252         );
    253         assert_eq!(reader.snapshot().metrics().revision, 5);
    254     }
    255 
    256     #[test]
    257     fn repeated_snapshot_reads_are_passive_and_preserve_one_arc() {
    258         let probe_calls = Arc::new(AtomicUsize::new(0));
    259         let (_publisher, reader) = cached_service_state(CachedServiceState::new(
    260             state(ServicePhase::Ready, Readiness::READY, ReasonCodes::empty()),
    261             metadata(1, &probe_calls),
    262         ));
    263         let first = reader.snapshot();
    264         for _ in 0..1_000 {
    265             let next = reader.snapshot();
    266             assert!(Arc::ptr_eq(&first, &next));
    267             assert_eq!(next.metrics().probe_calls.load(Ordering::SeqCst), 0);
    268         }
    269     }
    270 
    271     #[tokio::test]
    272     async fn reader_clone_and_publisher_drop_retain_the_last_snapshot() {
    273         let probe_calls = Arc::new(AtomicUsize::new(0));
    274         let (mut publisher, reader) = cached_service_state(CachedServiceState::new(
    275             state(ServicePhase::Ready, Readiness::READY, ReasonCodes::empty()),
    276             metadata(1, &probe_calls),
    277         ));
    278         let mut clone = reader.clone();
    279         publisher
    280             .publish(CachedServiceState::new(
    281                 state(
    282                     ServicePhase::Degraded,
    283                     Readiness::READY,
    284                     ReasonCodes::new([CommonReasonCode::DatabaseLowDisk.into()]).unwrap(),
    285                 ),
    286                 metadata(2, &probe_calls),
    287             ))
    288             .unwrap();
    289         drop(publisher);
    290 
    291         assert_eq!(reader.snapshot().metrics().revision, 2);
    292         assert_eq!(clone.changed().await.unwrap().metrics().revision, 2);
    293         assert!(matches!(clone.changed().await, Err(StatusPublisherDropped)));
    294         assert_eq!(
    295             clone.snapshot().operational().phase(),
    296             ServicePhase::Degraded
    297         );
    298     }
    299 
    300     #[tokio::test]
    301     async fn race_between_wakeup_and_borrow_marks_the_returned_latest_value_seen() {
    302         let probe_calls = Arc::new(AtomicUsize::new(0));
    303         let (mut publisher, mut reader) = cached_service_state(CachedServiceState::new(
    304             state(ServicePhase::Ready, Readiness::READY, ReasonCodes::empty()),
    305             metadata(0, &probe_calls),
    306         ));
    307         publisher
    308             .publish(CachedServiceState::new(
    309                 state(ServicePhase::Ready, Readiness::READY, ReasonCodes::empty()),
    310                 metadata(1, &probe_calls),
    311             ))
    312             .unwrap();
    313 
    314         reader.receiver.changed().await.unwrap();
    315         publisher
    316             .publish(CachedServiceState::new(
    317                 state(ServicePhase::Ready, Readiness::READY, ReasonCodes::empty()),
    318                 metadata(2, &probe_calls),
    319             ))
    320             .unwrap();
    321 
    322         assert_eq!(reader.snapshot_and_mark_seen().metrics().revision, 2);
    323         assert!(!reader.receiver.has_changed().unwrap());
    324     }
    325 }