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 }