lib

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

status.rs (13067B)


      1 //! Passive SQLite storage status and explicit close lifecycle.
      2 
      3 use std::{
      4     sync::{
      5         Mutex, RwLock,
      6         atomic::{AtomicU8, Ordering},
      7     },
      8     time::Duration,
      9 };
     10 
     11 use radroots_storage::{
     12     Error,
     13     outbox::BoxFuture,
     14     status::{
     15         IntegrityStatus, ShutdownState, StorageBackend, StorageOpenMode, StorageStatus,
     16         StorageStatusProvider, WriterPolicy,
     17     },
     18 };
     19 
     20 use crate::{OpenMode, SqliteStorage, integrity, lock::WriterLock};
     21 
     22 const OPEN: u8 = 0;
     23 const CLOSING: u8 = 1;
     24 const CLOSED: u8 = 2;
     25 const RESTORING: u8 = 3;
     26 
     27 pub(crate) struct StorageLifecycle {
     28     open_mode: StorageOpenMode,
     29     writer_policy: WriterPolicy,
     30     wal_enabled: bool,
     31     busy_timeout: Duration,
     32     shutdown: AtomicU8,
     33     integrity: RwLock<Option<IntegrityStatus>>,
     34     writer_lock: Mutex<Option<WriterLock>>,
     35 }
     36 
     37 impl StorageLifecycle {
     38     pub(crate) fn new(
     39         mode: OpenMode,
     40         busy_timeout: Duration,
     41         writer_lock: Option<WriterLock>,
     42     ) -> Self {
     43         Self {
     44             open_mode: storage_open_mode(mode),
     45             writer_policy: if mode.is_writable() {
     46                 WriterPolicy::AdvisoryProcessLock
     47             } else {
     48                 WriterPolicy::NoWriter
     49             },
     50             wal_enabled: mode.is_writable(),
     51             busy_timeout,
     52             shutdown: AtomicU8::new(OPEN),
     53             integrity: RwLock::new(None),
     54             writer_lock: Mutex::new(writer_lock),
     55         }
     56     }
     57 
     58     pub(crate) fn scaffold(mode: radroots_storage::status::EventStoreMode) -> Self {
     59         let open_mode = match mode {
     60             radroots_storage::status::EventStoreMode::ReadOnly => OpenMode::ReadOnly,
     61             radroots_storage::status::EventStoreMode::ReadWrite => OpenMode::Create,
     62         };
     63         Self::new(open_mode, Duration::from_secs(5), None)
     64     }
     65 
     66     pub(crate) fn integrity(&self) -> Result<IntegrityStatus, Error> {
     67         self.integrity
     68             .read()
     69             .map_err(|_| Error::BackendUnavailable)?
     70             .map_or_else(integrity::unknown, Ok)
     71     }
     72 
     73     pub(crate) fn require_open(&self) -> Result<(), Error> {
     74         if self.shutdown.load(Ordering::Acquire) == OPEN {
     75             Ok(())
     76         } else {
     77             Err(Error::BackendUnavailable)
     78         }
     79     }
     80 
     81     pub(crate) fn record_integrity(
     82         &self,
     83         status: IntegrityStatus,
     84     ) -> Result<IntegrityStatus, Error> {
     85         let mut recorded = self
     86             .integrity
     87             .write()
     88             .map_err(|_| Error::BackendUnavailable)?;
     89         if let Some(previous) = *recorded {
     90             let previous_time = previous
     91                 .checked_at_unix_ms()
     92                 .ok_or(Error::InvalidIntegrityStatus)?;
     93             let candidate_time = status
     94                 .checked_at_unix_ms()
     95                 .ok_or(Error::InvalidIntegrityStatus)?;
     96             if candidate_time < previous_time
     97                 || (candidate_time == previous_time && status != previous)
     98             {
     99                 return Err(Error::InvalidIntegrityStatus);
    100             }
    101         }
    102         *recorded = Some(status);
    103         Ok(status)
    104     }
    105 
    106     fn shutdown(&self) -> ShutdownState {
    107         match self.shutdown.load(Ordering::Acquire) {
    108             OPEN => ShutdownState::Open,
    109             CLOSING | RESTORING => ShutdownState::Closing,
    110             _ => ShutdownState::Closed,
    111         }
    112     }
    113 
    114     pub(crate) fn begin_close(&self) {
    115         let _ = self
    116             .shutdown
    117             .compare_exchange(OPEN, CLOSING, Ordering::AcqRel, Ordering::Acquire);
    118     }
    119 
    120     pub(crate) fn begin_restore_close(&self) -> Result<RestoreCloseAttempt<'_>, Error> {
    121         self.shutdown
    122             .compare_exchange(OPEN, RESTORING, Ordering::AcqRel, Ordering::Acquire)
    123             .map(|_| RestoreCloseAttempt(self))
    124             .map_err(|_| Error::BackendUnavailable)
    125     }
    126 
    127     pub(crate) fn finish_close(&self) -> Result<(), Error> {
    128         if self.shutdown.load(Ordering::Acquire) == RESTORING {
    129             return Ok(());
    130         }
    131         self.release_writer_and_close()
    132     }
    133 
    134     pub(crate) fn finish_restore_close(&self) -> Result<(), Error> {
    135         if self.shutdown.load(Ordering::Acquire) != RESTORING {
    136             return Err(Error::BackendUnavailable);
    137         }
    138         self.release_writer_and_close()
    139     }
    140 
    141     fn release_writer_and_close(&self) -> Result<(), Error> {
    142         let mut writer_lock = self
    143             .writer_lock
    144             .lock()
    145             .map_err(|_| Error::BackendUnavailable)?;
    146         let release_result = writer_lock
    147             .take()
    148             .map(WriterLock::release)
    149             .transpose()
    150             .map(|_| ())
    151             .map_err(|_| Error::BackendUnavailable);
    152         self.shutdown.store(CLOSED, Ordering::Release);
    153         release_result
    154     }
    155 
    156     fn status(&self) -> Result<StorageStatus, Error> {
    157         StorageStatus::new(
    158             StorageBackend::Sqlite,
    159             self.open_mode,
    160             self.writer_policy,
    161             self.shutdown(),
    162             self.integrity()?,
    163             self.wal_enabled,
    164             u32::try_from(self.busy_timeout.as_millis())
    165                 .map_err(|_| Error::InvalidStorageStatus)?,
    166         )
    167     }
    168 }
    169 
    170 /// Keeps writer authority reserved while finalization owns the close. A lost
    171 /// caller only hands draining back to ordinary close; it never releases a lock
    172 /// while either pool can still have active work.
    173 pub(crate) struct RestoreCloseAttempt<'a>(&'a StorageLifecycle);
    174 
    175 impl RestoreCloseAttempt<'_> {
    176     pub(crate) fn finish(self) -> Result<(), Error> {
    177         self.0.finish_restore_close()
    178     }
    179 }
    180 
    181 impl Drop for RestoreCloseAttempt<'_> {
    182     fn drop(&mut self) {
    183         let _ = self.0.shutdown.compare_exchange(
    184             RESTORING,
    185             CLOSING,
    186             Ordering::AcqRel,
    187             Ordering::Acquire,
    188         );
    189     }
    190 }
    191 
    192 impl SqliteStorage {
    193     /// Returns backend-level status without opening a connection or initiating
    194     /// integrity checks, checkpoints, migrations, or other maintenance.
    195     pub async fn storage_status(&self) -> Result<StorageStatus, Error> {
    196         self.lifecycle.status()
    197     }
    198 
    199     /// Closes both pools, releases writable authority, and returns final
    200     /// passive status. Repeated and concurrent calls are idempotent.
    201     pub async fn close(&self) -> Result<StorageStatus, Error> {
    202         self.lifecycle.begin_close();
    203         self.pool.close().await;
    204         self.private_pool.close().await;
    205         self.lifecycle.finish_close()?;
    206         self.lifecycle.status()
    207     }
    208 }
    209 
    210 impl StorageStatusProvider for SqliteStorage {
    211     fn storage_status(&self) -> BoxFuture<'_, Result<StorageStatus, Error>> {
    212         Box::pin(async move { SqliteStorage::storage_status(self).await })
    213     }
    214 }
    215 
    216 const fn storage_open_mode(mode: OpenMode) -> StorageOpenMode {
    217     match mode {
    218         OpenMode::ReadOnly => StorageOpenMode::ReadOnly,
    219         OpenMode::ReadWriteExisting => StorageOpenMode::ReadWriteExisting,
    220         OpenMode::Create => StorageOpenMode::Create,
    221     }
    222 }
    223 
    224 #[cfg(test)]
    225 #[cfg_attr(coverage_nightly, coverage(off))]
    226 mod tests {
    227     use std::time::Duration;
    228 
    229     use radroots_storage::{
    230         EventStore,
    231         event::SourceGeneration,
    232         status::{IntegrityHealth, ShutdownState, StorageBackend, StorageOpenMode, WriterPolicy},
    233     };
    234 
    235     use crate::{OpenOptions, Paths};
    236 
    237     use super::*;
    238 
    239     fn generation(byte: u8) -> SourceGeneration {
    240         SourceGeneration::new([byte; 32]).expect("source generation")
    241     }
    242 
    243     async fn create(directory: &std::path::Path) -> (Paths, SqliteStorage) {
    244         let paths = Paths::from_directory(directory).expect("owned paths");
    245         let store = SqliteStorage::open(
    246             OpenOptions::new(paths.clone(), OpenMode::Create)
    247                 .with_busy_timeout(Duration::from_millis(250))
    248                 .expect("busy timeout")
    249                 .with_source_generation(generation(73), 7_300)
    250                 .expect("source generation"),
    251         )
    252         .await
    253         .expect("create storage");
    254         (paths, store)
    255     }
    256 
    257     #[tokio::test]
    258     async fn status_and_integrity_are_passive_and_report_governed_configuration() {
    259         let directory = tempfile::tempdir().expect("temporary directory");
    260         let (paths, store) = create(directory.path()).await;
    261 
    262         let integrity = store.integrity().await.expect("integrity status");
    263         assert_eq!(integrity.health(), IntegrityHealth::Unknown);
    264         assert_eq!(integrity.checked_at_unix_ms(), None);
    265         assert_eq!(integrity.verified_members(), 0);
    266         assert_eq!(integrity.failed_members(), 0);
    267 
    268         let status = store.storage_status().await.expect("storage status");
    269         assert_eq!(status.backend(), StorageBackend::Sqlite);
    270         assert_eq!(status.open_mode(), StorageOpenMode::Create);
    271         assert_eq!(status.writer_policy(), WriterPolicy::AdvisoryProcessLock);
    272         assert_eq!(status.shutdown(), ShutdownState::Open);
    273         assert_eq!(status.integrity(), integrity);
    274         assert!(status.wal_enabled());
    275         assert_eq!(status.busy_timeout_ms(), 250);
    276 
    277         let reader = SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadOnly))
    278             .await
    279             .expect("read-only storage");
    280         let reader_status = reader.storage_status().await.expect("reader status");
    281         assert_eq!(reader_status.open_mode(), StorageOpenMode::ReadOnly);
    282         assert_eq!(reader_status.writer_policy(), WriterPolicy::NoWriter);
    283         assert!(!reader_status.wal_enabled());
    284         assert_eq!(reader_status.busy_timeout_ms(), 5_000);
    285 
    286         let healthy = IntegrityStatus::new(IntegrityHealth::Healthy, Some(100), 2, 0)
    287             .expect("healthy integrity");
    288         assert_eq!(reader.lifecycle.record_integrity(healthy), Ok(healthy));
    289         assert_eq!(reader.lifecycle.record_integrity(healthy), Ok(healthy));
    290         let older = IntegrityStatus::new(IntegrityHealth::Healthy, Some(99), 2, 0)
    291             .expect("older integrity");
    292         assert_eq!(
    293             reader.lifecycle.record_integrity(older),
    294             Err(Error::InvalidIntegrityStatus)
    295         );
    296         let conflicting = IntegrityStatus::new(IntegrityHealth::Degraded, Some(100), 1, 1)
    297             .expect("conflicting integrity");
    298         assert_eq!(
    299             reader.lifecycle.record_integrity(conflicting),
    300             Err(Error::InvalidIntegrityStatus)
    301         );
    302 
    303         let restoring = reader.lifecycle.begin_restore_close().unwrap();
    304         assert!(matches!(
    305             reader.lifecycle.begin_restore_close(),
    306             Err(Error::BackendUnavailable)
    307         ));
    308         assert_eq!(reader.lifecycle.finish_close(), Ok(()));
    309         assert_eq!(restoring.finish(), Ok(()));
    310         assert_eq!(
    311             reader.lifecycle.finish_restore_close(),
    312             Err(Error::BackendUnavailable)
    313         );
    314     }
    315 
    316     #[tokio::test]
    317     async fn close_is_observable_shared_idempotent_and_releases_writable_authority() {
    318         let directory = tempfile::tempdir().expect("temporary directory");
    319         let (paths, store) = create(directory.path()).await;
    320         let clone = store.clone();
    321         let held_connection = store.pool.acquire().await.expect("held connection");
    322         let mut close = Box::pin(clone.close());
    323 
    324         tokio::select! {
    325             biased;
    326             result = &mut close => panic!("close completed before the checked-out connection was returned: {result:?}"),
    327             () = tokio::task::yield_now() => {}
    328         }
    329         assert_eq!(
    330             store
    331                 .storage_status()
    332                 .await
    333                 .expect("closing status")
    334                 .shutdown(),
    335             ShutdownState::Closing
    336         );
    337 
    338         drop(held_connection);
    339         assert_eq!(
    340             close.await.expect("first close").shutdown(),
    341             ShutdownState::Closed
    342         );
    343         assert_eq!(
    344             store
    345                 .storage_status()
    346                 .await
    347                 .expect("shared closed status")
    348                 .shutdown(),
    349             ShutdownState::Closed
    350         );
    351         assert_eq!(
    352             store.close().await.expect("idempotent close").shutdown(),
    353             ShutdownState::Closed
    354         );
    355         assert_eq!(
    356             EventStore::status(&store).await,
    357             Err(Error::BackendUnavailable)
    358         );
    359 
    360         let reopened = SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting))
    361             .await
    362             .expect("writer authority released before final clone drop");
    363         assert_eq!(
    364             reopened
    365                 .storage_status()
    366                 .await
    367                 .expect("reopened status")
    368                 .shutdown(),
    369             ShutdownState::Open
    370         );
    371         let reopened_clone = reopened.clone();
    372         let (first, second) = tokio::join!(reopened.close(), reopened_clone.close());
    373         assert_eq!(
    374             first.expect("concurrent close one").shutdown(),
    375             ShutdownState::Closed
    376         );
    377         assert_eq!(
    378             second.expect("concurrent close two").shutdown(),
    379             ShutdownState::Closed
    380         );
    381     }
    382 }