lib

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

settling.rs (5307B)


      1 use radroots_storage::backup::BackupCapabilityError as Error;
      2 use sqlx::{Connection, Sqlite, SqlitePool, pool::PoolConnection};
      3 
      4 use crate::SqliteStorage;
      5 
      6 /// PoolConnection's asynchronous return pings its worker before releasing the
      7 /// permit. Holding every configured permit therefore also waits for cancelled
      8 /// executor work on connections other than the next snapshot connection.
      9 pub(super) async fn settle(store: &SqliteStorage) -> Result<(), Error> {
     10     store
     11         .lifecycle
     12         .require_open()
     13         .map_err(|_| Error::Unavailable)?;
     14     let runtime = settle_pool(&store.pool).await?;
     15     let protected = settle_pool(&store.private_pool).await?;
     16     store
     17         .lifecycle
     18         .require_open()
     19         .map_err(|_| Error::Unavailable)?;
     20     // Retain both sets simultaneously. Dropping a cancelled attempt returns
     21     // every acquired permit through the same owner, without inventing success.
     22     drop((runtime, protected));
     23     Ok(())
     24 }
     25 
     26 async fn settle_pool(pool: &SqlitePool) -> Result<Vec<PoolConnection<Sqlite>>, Error> {
     27     let capacity = pool.options().get_max_connections();
     28     let mut connections = Vec::with_capacity(capacity as usize);
     29     for _ in 0..capacity {
     30         let mut connection = pool.acquire().await.map_err(|_| Error::Unavailable)?;
     31         connection.ping().await.map_err(|_| Error::Unavailable)?;
     32         connections.push(connection);
     33     }
     34     Ok(connections)
     35 }
     36 
     37 #[cfg(test)]
     38 mod tests {
     39     use super::*;
     40     use radroots_storage::backup::StorageReliability;
     41     use std::{
     42         future::Future,
     43         task::{Context, Poll, Waker},
     44     };
     45 
     46     async fn fixture() -> (tempfile::TempDir, SqliteStorage) {
     47         let root = tempfile::tempdir().unwrap();
     48         let store = SqliteStorage::open(
     49             crate::OpenOptions::new(
     50                 crate::Paths::from_directory(root.path()).unwrap(),
     51                 crate::OpenMode::Create,
     52             )
     53             .with_source_generation(
     54                 radroots_storage::event::SourceGeneration::new([7; 32]).unwrap(),
     55                 100,
     56             )
     57             .unwrap(),
     58         )
     59         .await
     60         .unwrap();
     61         (root, store)
     62     }
     63 
     64     fn poll<F: Future + ?Sized>(future: std::pin::Pin<&mut F>) -> Poll<F::Output> {
     65         future.poll(&mut Context::from_waker(Waker::noop()))
     66     }
     67 
     68     #[tokio::test]
     69     async fn settling_waits_for_both_full_pools_and_cancelled_attempt_can_retry() {
     70         let (_root, store) = fixture().await;
     71         let runtime = settle_pool(&store.pool).await.unwrap();
     72         let mut protected = settle_pool(&store.private_pool).await.unwrap();
     73         let mut attempt = store.settle_backup_writes();
     74         assert!(poll(attempt.as_mut()).is_pending());
     75         drop(runtime);
     76         for _ in 0..32 {
     77             tokio::task::yield_now().await;
     78             assert!(poll(attempt.as_mut()).is_pending());
     79         }
     80         // An idle connection in one pool is not enough. The last admitted
     81         // protected member must finish before a cross-member inventory begins.
     82         let last_protected = protected.pop().unwrap();
     83         drop(protected);
     84         for _ in 0..32 {
     85             tokio::task::yield_now().await;
     86             assert!(poll(attempt.as_mut()).is_pending());
     87         }
     88         drop(attempt);
     89         drop(last_protected);
     90         store.settle_backup_writes().await.unwrap();
     91         assert_eq!(
     92             settle_pool(&store.pool).await.unwrap().len(),
     93             store.pool.options().get_max_connections() as usize
     94         );
     95         store.close().await.unwrap();
     96         assert_eq!(store.settle_backup_writes().await, Err(Error::Unavailable));
     97     }
     98 
     99     #[tokio::test]
    100     async fn settled_owner_observes_committed_and_abandoned_transactions_without_late_changes() {
    101         let (_root, store) = fixture().await;
    102         sqlx::query("CREATE TABLE backup_settling_fixture(value INTEGER NOT NULL)")
    103             .execute(&store.pool)
    104             .await
    105             .unwrap();
    106         let mut pending = store.pool.begin().await.unwrap();
    107         sqlx::query("INSERT INTO backup_settling_fixture VALUES (1)")
    108             .execute(&mut *pending)
    109             .await
    110             .unwrap();
    111         let mut attempt = store.settle_backup_writes();
    112         assert!(poll(attempt.as_mut()).is_pending());
    113         // Dropping a transaction schedules its rollback. No application future
    114         // remains to represent that work, but the owner must still settle it.
    115         drop(pending);
    116         attempt.await.unwrap();
    117         let count: i64 = sqlx::query_scalar("SELECT count(*) FROM backup_settling_fixture")
    118             .fetch_one(&store.pool)
    119             .await
    120             .unwrap();
    121         assert_eq!(count, 0);
    122         let mut committed = store.pool.begin().await.unwrap();
    123         sqlx::query("INSERT INTO backup_settling_fixture VALUES (2)")
    124             .execute(&mut *committed)
    125             .await
    126             .unwrap();
    127         committed.commit().await.unwrap();
    128         store.settle_backup_writes().await.unwrap();
    129         for _ in 0..2 {
    130             let values: Vec<i64> = sqlx::query_scalar("SELECT value FROM backup_settling_fixture")
    131                 .fetch_all(&store.pool)
    132                 .await
    133                 .unwrap();
    134             assert_eq!(values, [2]);
    135         }
    136         store.close().await.unwrap();
    137     }
    138 }