mod async_group;
mod data_acc;
mod data_conn;
mod data_hub;
mod data_src;
use crate::{ErrEntry, SendSyncNonNull, TxnFailureReport};
use std::any;
use std::collections::HashMap;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
pub use data_conn::DataConnError;
pub use data_hub::DataHubError;
pub use data_src::{
create_static_data_src_container, setup_async, setup_with_order_async, uses, uses_async,
DataSrcError,
};
#[doc(inline)]
pub use crate::_logic as logic;
#[doc(inline)]
pub use crate::_uses_for_async as uses;
#[allow(clippy::type_complexity)]
pub struct AsyncGroup {
attrs: Vec<(usize, Arc<str>)>,
tasks: Vec<Pin<Box<dyn Future<Output = errs::Result<()>> + Send + 'static>>>,
_index: usize,
_name: Arc<str>,
}
#[allow(async_fn_in_trait)]
#[allow(unused_variables)] pub trait DataConn {
fn commit_async(
&mut self,
ag: &mut AsyncGroup,
) -> impl Future<Output = errs::Result<()>> + Send;
fn pre_commit_async(
&mut self,
ag: &mut AsyncGroup,
) -> impl Future<Output = errs::Result<()>> + Send {
async { Ok(()) }
}
fn post_commit_async(
&mut self,
ag: &mut AsyncGroup,
) -> impl Future<Output = errs::Result<()>> + Send {
async { Ok(()) }
}
fn is_committed(&self) -> bool;
fn rollback_async(
&mut self,
ag: &mut AsyncGroup,
) -> impl Future<Output = errs::Result<()>> + Send;
#[cfg_attr(coverage_nightly, coverage(off))]
fn on_txn_failure_async(
&mut self,
ag: &mut AsyncGroup,
reports: Arc<[TxnFailureReport]>,
) -> impl Future<Output = ()> + Send {
async {}
}
fn close(&mut self);
}
pub(crate) struct NoopDataConn {}
#[cfg_attr(coverage_nightly, coverage(off))]
impl DataConn for NoopDataConn {
async fn commit_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
Ok(())
}
fn is_committed(&self) -> bool {
false
}
async fn rollback_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
Ok(())
}
fn close(&mut self) {}
}
#[allow(clippy::type_complexity)]
#[repr(C)]
pub(crate) struct DataConnContainer<C = NoopDataConn>
where
C: DataConn + 'static,
{
drop_fn: fn(*const DataConnContainer),
is_fn: fn(any::TypeId) -> bool,
type_fn: fn() -> &'static str,
commit_fn: for<'ag> fn(
*const DataConnContainer,
&'ag mut AsyncGroup,
) -> Pin<Box<dyn Future<Output = errs::Result<()>> + Send + 'ag>>,
pre_commit_fn: for<'ag> fn(
*const DataConnContainer,
&'ag mut AsyncGroup,
) -> Pin<Box<dyn Future<Output = errs::Result<()>> + Send + 'ag>>,
post_commit_fn: for<'ag> fn(
*const DataConnContainer,
&'ag mut AsyncGroup,
)
-> Pin<Box<dyn Future<Output = errs::Result<()>> + Send + 'ag>>,
is_committed_fn: fn(*const DataConnContainer) -> bool,
rollback_fn: for<'ag> fn(
*const DataConnContainer,
&'ag mut AsyncGroup,
) -> Pin<Box<dyn Future<Output = errs::Result<()>> + Send + 'ag>>,
on_txn_failure_fn: for<'ag> fn(
*const DataConnContainer,
&'ag mut AsyncGroup,
Arc<[TxnFailureReport]>,
) -> Pin<Box<dyn Future<Output = ()> + Send + 'ag>>,
close_fn: fn(*const DataConnContainer),
name: Arc<str>,
data_conn: Box<C>,
}
pub(crate) struct DataConnManager {
vec: Vec<Option<SendSyncNonNull<DataConnContainer>>>,
index_map: HashMap<Arc<str>, usize>,
committed: bool,
}
#[trait_variant::make(Send)]
#[allow(unused_variables)] pub trait DataSrc<C>
where
C: DataConn + 'static,
{
async fn setup_async(&mut self, ag: &mut AsyncGroup) -> errs::Result<()>;
fn close(&mut self);
async fn create_data_conn_async(&mut self) -> errs::Result<Box<C>>;
}
pub(crate) struct NoopDataSrc {}
#[cfg_attr(coverage_nightly, coverage(off))]
impl DataSrc<NoopDataConn> for NoopDataSrc {
async fn setup_async(&mut self, _ag: &mut AsyncGroup) -> errs::Result<()> {
Ok(())
}
fn close(&mut self) {}
async fn create_data_conn_async(&mut self) -> errs::Result<Box<NoopDataConn>> {
Ok(Box::new(NoopDataConn {}))
}
}
#[allow(clippy::type_complexity)]
#[repr(C)]
pub(crate) struct DataSrcContainer<S = NoopDataSrc, C = NoopDataConn>
where
S: DataSrc<C>,
C: DataConn + 'static,
{
drop_fn: fn(*const DataSrcContainer),
close_fn: fn(*const DataSrcContainer),
is_data_conn_fn: fn(any::TypeId) -> bool,
setup_fn: for<'ag> fn(
*const DataSrcContainer,
&'ag mut AsyncGroup,
) -> Pin<Box<dyn Future<Output = errs::Result<()>> + Send + 'ag>>,
create_data_conn_fn: fn(
*const DataSrcContainer,
) -> Pin<
Box<dyn Future<Output = errs::Result<Box<DataConnContainer<C>>>> + Send + 'static>,
>,
local: bool,
name: Arc<str>,
data_src: S,
}
pub(crate) struct DataSrcManager {
vec_unready: Vec<SendSyncNonNull<DataSrcContainer>>,
vec_ready: Vec<SendSyncNonNull<DataSrcContainer>>,
local: bool,
}
pub struct AutoShutdown {}
pub struct DataHub {
local_data_src_manager: DataSrcManager,
data_src_map: HashMap<Arc<str>, (bool, usize)>,
data_conn_manager: DataConnManager,
fixed: bool,
}
pub trait DataAcc {
#[allow(async_fn_in_trait)]
fn get_data_conn_async<C: DataConn + 'static>(
&mut self,
name: &str,
) -> impl Future<Output = errs::Result<&mut C>> + Send;
}
#[doc(hidden)]
pub struct StaticDataSrcContainer {
pub(crate) ssnnptr: SendSyncNonNull<DataSrcContainer>,
}
#[doc(hidden)]
pub struct StaticDataSrcRegistration {
pub(crate) factory: fn() -> StaticDataSrcContainer,
}