status.rs (17254B)
1 //! Passive synchronization status aggregation and host retry decisions. 2 3 use std::collections::BTreeSet; 4 5 use radroots_protocol::runtime::v1::{ 6 OPERATION_SCHEMA_VERSION, SyncCapabilityState, SyncHealth, SyncOutboxStatus, 7 SyncProjectionStatus, SyncRetryDecision, SyncStatusReceipt, 8 }; 9 use radroots_signing::{SignerStatus, status::SignerAvailability}; 10 use radroots_storage::{ 11 EventStore, Outbox, ProjectionStore, 12 authored_delivery::{AuthoredDeliveryPlan, AuthoredDeliveryState}, 13 outbox::OutboxStatus, 14 projection::{ProjectionHealth, ProjectionId, ProjectionStatus}, 15 status::{ 16 EventStoreHealth, EventStoreStatus, IntegrityHealth, ShutdownState, StorageStatus, 17 StorageStatusProvider, 18 }, 19 }; 20 use radroots_transport::{SinkStatus, SourceStatus, capability::Availability}; 21 22 use crate::{Engine, policy::Error}; 23 24 const STATUS_PROJECTION_LIMIT: usize = 256; 25 26 fn map_storage_error(error: radroots_storage::Error) -> Error { 27 match error { 28 radroots_storage::Error::SpaceInsufficient => Error::StorageSpaceInsufficient, 29 _ => Error::StorageFailed, 30 } 31 } 32 33 /// Typed report for one optional injected host capability. 34 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 35 #[derive(Clone, Debug, Eq, PartialEq)] 36 pub struct CapabilityReport<T> { 37 state: SyncCapabilityState, 38 status: Option<T>, 39 } 40 41 impl<T> CapabilityReport<T> { 42 pub const fn state(&self) -> SyncCapabilityState { 43 self.state 44 } 45 46 pub const fn status(&self) -> Option<&T> { 47 self.status.as_ref() 48 } 49 50 const fn unsupported() -> Self { 51 Self { 52 state: SyncCapabilityState::Unsupported, 53 status: None, 54 } 55 } 56 57 const fn compiled(status: Option<T>) -> Self { 58 Self { 59 state: SyncCapabilityState::Compiled, 60 status, 61 } 62 } 63 64 const fn reported(state: SyncCapabilityState, status: T) -> Self { 65 Self { 66 state, 67 status: Some(status), 68 } 69 } 70 } 71 72 /// Requested projection and its optional durable status record. 73 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 74 #[derive(Clone, Debug, Eq, PartialEq)] 75 pub struct ProjectionReport { 76 projection_id: ProjectionId, 77 status: Option<ProjectionStatus>, 78 } 79 80 impl ProjectionReport { 81 pub const fn projection_id(&self) -> &ProjectionId { 82 &self.projection_id 83 } 84 85 pub const fn status(&self) -> Option<&ProjectionStatus> { 86 self.status.as_ref() 87 } 88 } 89 90 /// One passive, side-effect-free synchronization health snapshot. 91 #[cfg_attr(feature = "serde", derive(serde::Serialize))] 92 #[derive(Clone, Debug, Eq, PartialEq)] 93 pub struct SyncStatus { 94 health: SyncHealth, 95 storage: StorageStatus, 96 events: EventStoreStatus, 97 outbox: OutboxStatus, 98 source: CapabilityReport<SourceStatus>, 99 sink: CapabilityReport<SinkStatus>, 100 signer: CapabilityReport<SignerStatus>, 101 projections: Vec<ProjectionReport>, 102 } 103 104 impl SyncStatus { 105 pub const fn health(&self) -> SyncHealth { 106 self.health 107 } 108 109 pub const fn storage(&self) -> StorageStatus { 110 self.storage 111 } 112 113 pub const fn events(&self) -> &EventStoreStatus { 114 &self.events 115 } 116 117 pub const fn outbox(&self) -> OutboxStatus { 118 self.outbox 119 } 120 121 pub const fn source(&self) -> &CapabilityReport<SourceStatus> { 122 &self.source 123 } 124 125 pub const fn sink(&self) -> &CapabilityReport<SinkStatus> { 126 &self.sink 127 } 128 129 pub const fn signer(&self) -> &CapabilityReport<SignerStatus> { 130 &self.signer 131 } 132 133 pub fn projections(&self) -> &[ProjectionReport] { 134 self.projections.as_slice() 135 } 136 137 /// Converts native reports into the versioned passive protocol receipt. 138 pub fn to_protocol(&self) -> SyncStatusReceipt { 139 let mut projection = SyncProjectionStatus { 140 ready: 0, 141 invalidated: 0, 142 rebuilding: 0, 143 failed: 0, 144 untracked: 0, 145 }; 146 for report in &self.projections { 147 match report.status.as_ref().map(ProjectionStatus::health) { 148 Some(ProjectionHealth::Ready) => projection.ready += 1, 149 Some(ProjectionHealth::Invalidated) => projection.invalidated += 1, 150 Some(ProjectionHealth::Rebuilding) => projection.rebuilding += 1, 151 Some(ProjectionHealth::Failed) => projection.failed += 1, 152 None => projection.untracked += 1, 153 } 154 } 155 SyncStatusReceipt { 156 schema_version: OPERATION_SCHEMA_VERSION, 157 health: self.health, 158 storage: storage_state(self.storage, self.events.health()), 159 source: self.source.state, 160 sink: self.sink.state, 161 signer: self.signer.state, 162 outbox: SyncOutboxStatus { 163 pending: self.outbox.pending, 164 leased: self.outbox.leased, 165 retryable: self.outbox.retryable, 166 satisfied: self.outbox.satisfied, 167 exhausted: self.outbox.exhausted, 168 }, 169 projections: projection, 170 } 171 } 172 } 173 174 impl Engine { 175 /// Aggregates passive status without spawning work or initiating recovery. 176 pub async fn status(&self, projection_ids: &[ProjectionId]) -> Result<SyncStatus, Error> { 177 if [ 178 projection_ids.len() > STATUS_PROJECTION_LIMIT, 179 projection_ids.iter().collect::<BTreeSet<_>>().len() != projection_ids.len(), 180 ] 181 .contains(&true) 182 { 183 return Err(Error::InvalidStatusRequest); 184 } 185 let storage = StorageStatusProvider::storage_status(self.storage.as_ref()) 186 .await 187 .map_err(map_storage_error)?; 188 let events = EventStore::status(self.storage.as_ref()) 189 .await 190 .map_err(map_storage_error)?; 191 let outbox = Outbox::status(self.storage.as_ref()) 192 .await 193 .map_err(map_storage_error)?; 194 if outbox.total().is_none() { 195 return Err(Error::StorageFailed); 196 } 197 let source = source_report(self).await; 198 let sink = sink_report(self).await; 199 let signer = signer_report(self).await; 200 let mut projections = Vec::with_capacity(projection_ids.len()); 201 for projection_id in projection_ids { 202 let status = ProjectionStore::status(self.storage.as_ref(), projection_id.clone()) 203 .await 204 .map_err(map_storage_error)?; 205 projections.push(ProjectionReport { 206 projection_id: projection_id.clone(), 207 status, 208 }); 209 } 210 let health = aggregate_health(storage, &events, &source, &sink, &signer, &projections); 211 Ok(SyncStatus { 212 health, 213 storage, 214 events, 215 outbox, 216 source, 217 sink, 218 signer, 219 projections, 220 }) 221 } 222 223 /// Classifies host action for one durable plan without mutating it. 224 pub fn retry_decision( 225 &self, 226 plan: &AuthoredDeliveryPlan, 227 now_unix_ms: u64, 228 ) -> Result<SyncRetryDecision, Error> { 229 if now_unix_ms == 0 { 230 return Err(Error::ClockUnavailable); 231 } 232 if plan.request().is_some() 233 && plan.delivery_satisfaction().map_err(map_storage_error)? 234 == radroots_transport::policy::SatisfactionState::Satisfied 235 { 236 return Ok(SyncRetryDecision::Satisfied); 237 } 238 match plan.state() { 239 AuthoredDeliveryState::Satisfied => return Ok(SyncRetryDecision::Satisfied), 240 AuthoredDeliveryState::Exhausted 241 | AuthoredDeliveryState::FailedTerminal 242 | AuthoredDeliveryState::Cancelled => return Ok(SyncRetryDecision::Exhausted), 243 AuthoredDeliveryState::Pending | AuthoredDeliveryState::Retryable => {} 244 } 245 if now_unix_ms >= plan.intent().deadline_unix_ms() { 246 return Ok(SyncRetryDecision::Expired); 247 } 248 if let Some(claim) = plan.claim_evidence() 249 && now_unix_ms < claim.expires_at_unix_ms() 250 { 251 return Ok(SyncRetryDecision::InFlightUntil { 252 unix_ms: claim.expires_at_unix_ms(), 253 }); 254 } 255 if let Some(retry) = plan.retry() 256 && now_unix_ms < retry.not_before_unix_ms() 257 { 258 return Ok(SyncRetryDecision::DeferredUntil { 259 unix_ms: retry.not_before_unix_ms(), 260 }); 261 } 262 Ok(SyncRetryDecision::Ready) 263 } 264 } 265 266 async fn source_report(engine: &Engine) -> CapabilityReport<SourceStatus> { 267 let Some(source) = engine.source.as_deref() else { 268 return CapabilityReport::unsupported(); 269 }; 270 match source.status().await { 271 Ok(status) if !status.is_configured() => CapabilityReport::compiled(Some(status)), 272 Ok(status) => { 273 let state = availability_state(status.availability()); 274 CapabilityReport::reported(state, status) 275 } 276 Err(_) => CapabilityReport::compiled(None), 277 } 278 } 279 280 async fn sink_report(engine: &Engine) -> CapabilityReport<SinkStatus> { 281 let Some(sink) = engine.sink.as_deref() else { 282 return CapabilityReport::unsupported(); 283 }; 284 match sink.status().await { 285 Ok(status) if !status.is_configured() => CapabilityReport::compiled(Some(status)), 286 Ok(status) => { 287 let state = availability_state(status.availability()); 288 CapabilityReport::reported(state, status) 289 } 290 Err(_) => CapabilityReport::compiled(None), 291 } 292 } 293 294 async fn signer_report(engine: &Engine) -> CapabilityReport<SignerStatus> { 295 let Some(signer) = engine.signer.as_deref() else { 296 return CapabilityReport::unsupported(); 297 }; 298 match signer.status().await { 299 Ok(status) => { 300 let state = match status.availability() { 301 SignerAvailability::Ready => SyncCapabilityState::Available, 302 SignerAvailability::Busy | SignerAvailability::AwaitingAuthentication => { 303 SyncCapabilityState::Degraded 304 } 305 SignerAvailability::Unavailable => SyncCapabilityState::Configured, 306 _ => SyncCapabilityState::Degraded, 307 }; 308 CapabilityReport::reported(state, status) 309 } 310 Err(_) => CapabilityReport::compiled(None), 311 } 312 } 313 314 const fn availability_state(availability: Availability) -> SyncCapabilityState { 315 match availability { 316 Availability::Available => SyncCapabilityState::Available, 317 Availability::Degraded => SyncCapabilityState::Degraded, 318 Availability::Unavailable => SyncCapabilityState::Configured, 319 } 320 } 321 322 fn aggregate_health( 323 storage: StorageStatus, 324 events: &EventStoreStatus, 325 source: &CapabilityReport<SourceStatus>, 326 sink: &CapabilityReport<SinkStatus>, 327 signer: &CapabilityReport<SignerStatus>, 328 projections: &[ProjectionReport], 329 ) -> SyncHealth { 330 if [ 331 matches!( 332 storage.shutdown(), 333 ShutdownState::Closing | ShutdownState::Closed 334 ), 335 storage.integrity().health() == IntegrityHealth::Corrupt, 336 events.health() == EventStoreHealth::Unavailable, 337 ] 338 .contains(&true) 339 { 340 return SyncHealth::Unavailable; 341 } 342 let capability_degraded = [source.state, sink.state, signer.state] 343 .into_iter() 344 .any(|state| { 345 matches!( 346 state, 347 SyncCapabilityState::Compiled 348 | SyncCapabilityState::Configured 349 | SyncCapabilityState::Degraded 350 ) 351 }); 352 let projection_degraded = projections.iter().any(|projection| { 353 !matches!( 354 projection.status.as_ref().map(ProjectionStatus::health), 355 Some(ProjectionHealth::Ready) 356 ) 357 }); 358 if [ 359 storage.integrity().health() != IntegrityHealth::Healthy, 360 events.health() == EventStoreHealth::Degraded, 361 capability_degraded, 362 projection_degraded, 363 ] 364 .contains(&true) 365 { 366 SyncHealth::Degraded 367 } else { 368 SyncHealth::Healthy 369 } 370 } 371 372 const fn storage_state( 373 storage: StorageStatus, 374 event_health: EventStoreHealth, 375 ) -> SyncCapabilityState { 376 if matches!( 377 storage.shutdown(), 378 ShutdownState::Closing | ShutdownState::Closed 379 ) || matches!(storage.integrity().health(), IntegrityHealth::Corrupt) 380 || matches!(event_health, EventStoreHealth::Unavailable) 381 { 382 SyncCapabilityState::Configured 383 } else if matches!(storage.integrity().health(), IntegrityHealth::Healthy) 384 && matches!(event_health, EventStoreHealth::Available) 385 { 386 SyncCapabilityState::Available 387 } else { 388 SyncCapabilityState::Degraded 389 } 390 } 391 392 #[cfg(test)] 393 #[cfg_attr(coverage_nightly, coverage(off))] 394 mod tests { 395 use super::*; 396 use radroots_storage::{ 397 event::SourceGeneration, 398 status::{EventStoreMode, IntegrityStatus, StorageBackend, StorageOpenMode, WriterPolicy}, 399 }; 400 401 fn storage(health: IntegrityHealth, shutdown: ShutdownState) -> StorageStatus { 402 let failed = u32::from(health == IntegrityHealth::Corrupt); 403 let checked = if health == IntegrityHealth::Unknown { 404 None 405 } else { 406 Some(1) 407 }; 408 StorageStatus::new( 409 StorageBackend::Memory, 410 StorageOpenMode::ReadWriteExisting, 411 WriterPolicy::NoWriter, 412 shutdown, 413 IntegrityStatus::new(health, checked, 1, failed).unwrap(), 414 false, 415 0, 416 ) 417 .unwrap() 418 } 419 420 fn events(health: EventStoreHealth) -> EventStoreStatus { 421 EventStoreStatus::new( 422 SourceGeneration::new([21; 32]).unwrap(), 423 EventStoreMode::ReadWrite, 424 health, 425 0, 426 0, 427 0, 428 ) 429 .unwrap() 430 } 431 432 #[test] 433 fn health_and_protocol_classification_cover_every_state() { 434 assert_eq!( 435 map_storage_error(radroots_storage::Error::SpaceInsufficient), 436 Error::StorageSpaceInsufficient 437 ); 438 assert_eq!( 439 map_storage_error(radroots_storage::Error::BackendUnavailable), 440 Error::StorageFailed 441 ); 442 assert_eq!( 443 availability_state(Availability::Available), 444 SyncCapabilityState::Available 445 ); 446 assert_eq!( 447 availability_state(Availability::Degraded), 448 SyncCapabilityState::Degraded 449 ); 450 assert_eq!( 451 availability_state(Availability::Unavailable), 452 SyncCapabilityState::Configured 453 ); 454 let source = CapabilityReport::<SourceStatus>::unsupported(); 455 let sink = CapabilityReport::<SinkStatus>::unsupported(); 456 let signer = CapabilityReport::<SignerStatus>::unsupported(); 457 for (storage, events, expected) in [ 458 ( 459 storage(IntegrityHealth::Healthy, ShutdownState::Open), 460 events(EventStoreHealth::Available), 461 SyncHealth::Healthy, 462 ), 463 ( 464 storage(IntegrityHealth::Corrupt, ShutdownState::Open), 465 events(EventStoreHealth::Available), 466 SyncHealth::Unavailable, 467 ), 468 ( 469 storage(IntegrityHealth::Healthy, ShutdownState::Closing), 470 events(EventStoreHealth::Available), 471 SyncHealth::Unavailable, 472 ), 473 ( 474 storage(IntegrityHealth::Healthy, ShutdownState::Open), 475 events(EventStoreHealth::Unavailable), 476 SyncHealth::Unavailable, 477 ), 478 ( 479 storage(IntegrityHealth::Degraded, ShutdownState::Open), 480 events(EventStoreHealth::Available), 481 SyncHealth::Degraded, 482 ), 483 ] { 484 assert_eq!( 485 aggregate_health(storage, &events, &source, &sink, &signer, &[]), 486 expected 487 ); 488 } 489 let compiled = CapabilityReport::<SourceStatus>::compiled(None); 490 assert_eq!( 491 aggregate_health( 492 storage(IntegrityHealth::Healthy, ShutdownState::Open), 493 &events(EventStoreHealth::Available), 494 &compiled, 495 &sink, 496 &signer, 497 &[], 498 ), 499 SyncHealth::Degraded 500 ); 501 assert_eq!( 502 storage_state( 503 storage(IntegrityHealth::Healthy, ShutdownState::Open), 504 EventStoreHealth::Available, 505 ), 506 SyncCapabilityState::Available 507 ); 508 assert_eq!( 509 storage_state( 510 storage(IntegrityHealth::Degraded, ShutdownState::Open), 511 EventStoreHealth::Degraded, 512 ), 513 SyncCapabilityState::Degraded 514 ); 515 assert_eq!( 516 storage_state( 517 storage(IntegrityHealth::Healthy, ShutdownState::Closed), 518 EventStoreHealth::Available, 519 ), 520 SyncCapabilityState::Configured 521 ); 522 } 523 }