use std::sync::Arc;
use tokio::sync::Mutex;
use super::backend::{
AnswerConsumer, BoundedAnswerLimits, BoundedAnswerStats, GivenRowsSpec, QueryResult,
TransactionOps, TxType,
};
use crate::error::{OrmError, Result};
use crate::match_request::selected_result_executor::SelectedResultExecutor;
use crate::match_request::{
CapabilitySet, MatchExecutionLimits, ValidatedMatchRequest, ValidatedMatchResult,
};
use crate::registry::DescriptorRegistry;
use type_bridge_core_lib::ast::{
TypedFetchRows, TypedHydrateThings, TypedPageRematch, TypedRootScan,
};
pub struct TransactionContext {
inner: Arc<Mutex<Option<Box<dyn TransactionOps>>>>,
tx_type: TxType,
match_capabilities: CapabilitySet,
}
impl TransactionContext {
pub(crate) fn new(
inner: Box<dyn TransactionOps>,
tx_type: TxType,
match_capabilities: CapabilitySet,
) -> Self {
Self {
inner: Arc::new(Mutex::new(Some(inner))),
tx_type,
match_capabilities,
}
}
pub async fn query(&self, typeql: &str) -> Result<QueryResult> {
let mut guard = self.inner.lock().await;
let tx = guard
.as_mut()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.query(typeql).await
}
pub async fn query_with_rows(&self, typeql: &str, rows: GivenRowsSpec) -> Result<QueryResult> {
let mut guard = self.inner.lock().await;
let tx = guard
.as_mut()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.query_with_rows(typeql, rows).await
}
pub(crate) async fn query_typed_bounded(
&self,
query: &TypedFetchRows,
limits: BoundedAnswerLimits,
consumer: &mut dyn AnswerConsumer,
) -> Result<BoundedAnswerStats> {
let mut guard = self.inner.lock().await;
let tx = guard
.as_mut()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.query_typed_bounded(query, limits, consumer).await
}
pub(crate) async fn hydrate_typed_bounded(
&self,
query: &TypedHydrateThings,
limits: BoundedAnswerLimits,
consumer: &mut dyn AnswerConsumer,
) -> Result<BoundedAnswerStats> {
let mut guard = self.inner.lock().await;
let tx = guard
.as_mut()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.hydrate_typed_bounded(query, limits, consumer).await
}
pub(crate) async fn query_root_typed_bounded(
&self,
query: &TypedRootScan,
limits: BoundedAnswerLimits,
consumer: &mut dyn AnswerConsumer,
) -> Result<BoundedAnswerStats> {
let mut guard = self.inner.lock().await;
let tx = guard
.as_mut()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.query_root_typed_bounded(query, limits, consumer).await
}
pub(crate) async fn rematch_page_typed_bounded(
&self,
query: &TypedPageRematch,
limits: BoundedAnswerLimits,
consumer: &mut dyn AnswerConsumer,
) -> Result<BoundedAnswerStats> {
let mut guard = self.inner.lock().await;
let tx = guard
.as_mut()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.rematch_page_typed_bounded(query, limits, consumer).await
}
pub async fn commit(&self) -> Result<()> {
let mut guard = self.inner.lock().await;
let mut tx = guard
.take()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.commit().await
}
pub async fn rollback(&self) -> Result<()> {
let mut guard = self.inner.lock().await;
let mut tx = guard
.take()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.rollback().await
}
pub async fn close(&self) -> Result<()> {
let mut guard = self.inner.lock().await;
let Some(mut tx) = guard.take() else {
return Ok(());
};
tx.close().await
}
pub fn tx_type(&self) -> TxType {
self.tx_type
}
pub async fn execute_match(
&self,
registry: &DescriptorRegistry,
validated: &ValidatedMatchRequest,
) -> Result<ValidatedMatchResult> {
self.execute_match_with_limits(registry, validated, MatchExecutionLimits::default())
.await
}
pub async fn execute_match_with_limits(
&self,
registry: &DescriptorRegistry,
validated: &ValidatedMatchRequest,
limits: MatchExecutionLimits,
) -> Result<ValidatedMatchResult> {
SelectedResultExecutor::new(registry, self.match_capabilities.clone(), limits)
.execute_borrowed(self, validated)
.await
}
}
impl Clone for TransactionContext {
fn clone(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
tx_type: self.tx_type,
match_capabilities: self.match_capabilities.clone(),
}
}
}