saddle-db 0.3.8

Saddle managed asynchronous database access and transactions
Documentation
use std::{
    future::{Future, poll_fn},
    pin::pin,
    sync::atomic::Ordering,
    task::Poll,
};

use saddle_core::{DbPhysicalDisposition, DbPhysicalExecutionHalf};
use saddle_runtime::profusegw::{
    ProfuseGwDatabaseFinalizationCompletion, finish_profusegw_database_disposition,
};
use sqlx::{MySql, pool::PoolConnection};

use crate::{
    observation::DatabaseOperationObservation,
    startup_pool::StartupManagedDatabaseProcessCapability,
};

/// Terminal business value issued only after the physical connection and the
/// admitted request have reached one sealed disposition.
#[doc(hidden)]
pub struct DatabasePhysicalTerminal<T> {
    disposition: DbPhysicalDisposition,
    value: T,
}

impl<T> DatabasePhysicalTerminal<T> {
    pub fn into_outcome(self) -> (DbPhysicalDisposition, T) {
        (self.disposition, self.value)
    }
}

#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum RequestedPhysicalDisposition {
    Return,
    Discard,
}

impl StartupManagedDatabaseProcessCapability {
    pub(crate) fn ensure_live(&self) -> bool {
        self.process_live.load(Ordering::Acquire)
    }

    pub(crate) async fn finish_physical<T>(
        &self,
        mut connection: Option<PoolConnection<MySql>>,
        mut completion: ProfuseGwDatabaseFinalizationCompletion,
        execution: DbPhysicalExecutionHalf,
        requested: RequestedPhysicalDisposition,
        value: T,
        observation: DatabaseOperationObservation,
    ) -> DatabasePhysicalTerminal<T> {
        let disposition = match (requested, connection.as_mut()) {
            (RequestedPhysicalDisposition::Return, Some(connection_ref)) => {
                let returned = {
                    let return_to_pool = connection_ref.return_to_pool();
                    let mut return_to_pool = pin!(return_to_pool);
                    poll_fn(|context| {
                        if return_to_pool.as_mut().poll(context).is_ready() {
                            return Poll::Ready(true);
                        }
                        if completion.poll_physical_deadline(context).is_ready() {
                            return Poll::Ready(false);
                        }
                        Poll::Pending
                    })
                    .await
                };
                if returned {
                    drop(connection.take());
                    DbPhysicalDisposition::Returned
                } else {
                    drop(
                        connection
                            .take()
                            .expect("physical connection exists")
                            .detach(),
                    );
                    DbPhysicalDisposition::Discarded
                }
            }
            (_, Some(_)) => {
                drop(
                    connection
                        .take()
                        .expect("physical connection exists")
                        .detach(),
                );
                DbPhysicalDisposition::Discarded
            }
            _ => DbPhysicalDisposition::Discarded,
        };
        let physical = match disposition {
            DbPhysicalDisposition::Returned => self.physical.connection_returned(execution, value),
            DbPhysicalDisposition::Discarded => {
                self.physical.connection_discarded(execution, value)
            }
        };
        let physical = match physical {
            Ok(physical) => physical,
            Err(_) => std::process::abort(),
        };
        let value = match finish_profusegw_database_disposition(completion, physical) {
            Ok(value) => value,
            Err(_) => std::process::abort(),
        };
        observation.finish(disposition);
        DatabasePhysicalTerminal { disposition, value }
    }
}