lib

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

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 }