use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::Mutex;
use crate::manager::TransactionDefinition;
use crate::{Propagation, TransactionError, TransactionResult, TransactionStatus};
pub trait LiveTransaction: Send {
fn commit_boxed(
self: Box<Self>,
) -> std::pin::Pin<Box<dyn Future<Output = TransactionResult<()>> + Send>>;
fn rollback_boxed(
self: Box<Self>,
) -> std::pin::Pin<Box<dyn Future<Output = TransactionResult<()>> + Send>>;
}
tokio::task_local! {
static ACTIVE_TX: Arc<Mutex<HashMap<String, (TransactionStatus, Box<dyn LiveTransaction>)>>>;
}
fn active_map() -> Arc<Mutex<HashMap<String, (TransactionStatus, Box<dyn LiveTransaction>)>>> {
ACTIVE_TX.try_with(Clone::clone).expect(
"No active transaction scope: wrap your transactional code with `with_transaction_scope`",
)
}
pub async fn with_transaction_scope<F, R>(f: F) -> R
where
F: Future<Output = R>,
{
ACTIVE_TX
.scope(Arc::new(Mutex::new(HashMap::new())), f)
.await
}
pub async fn bind_transaction(status: TransactionStatus, tx: Box<dyn LiveTransaction>) {
let map = active_map();
let mut guard = map.lock().await;
guard.insert(status.name().to_string(), (status, tx));
}
pub async fn take_transaction(status: &TransactionStatus) -> Option<Box<dyn LiveTransaction>> {
let map = active_map();
let mut guard = map.lock().await;
guard.remove(status.name()).map(|(_, tx)| tx)
}
pub async fn current_status() -> Option<TransactionStatus> {
let map = active_map();
let guard = map.lock().await;
guard.values().next().map(|(s, _)| s.clone())
}
pub async fn resolve_propagation(
definition: &TransactionDefinition,
) -> Result<Option<TransactionStatus>, TransactionError> {
match definition.propagation {
Propagation::Required | Propagation::Supports => Ok(current_status().await),
Propagation::Mandatory => match current_status().await {
Some(s) => Ok(Some(s)),
None => {
Err(TransactionError::InvalidState("No existing transaction for MANDATORY".into()))
},
},
Propagation::Never => {
if current_status().await.is_some() {
Err(TransactionError::InvalidState("Existing transaction present for NEVER".into()))
} else {
Ok(None)
}
},
Propagation::NotSupported | Propagation::RequiresNew | Propagation::Nested => Ok(None),
}
}