use std::sync::Arc;
use tokio::sync::Mutex;
use super::backend::{
AnswerConsumer, BoundedAnswerLimits, BoundedAnswerStats, GivenRowsSpec, QueryResult,
TransactionOps, TxType,
};
use crate::_registry::DescriptorRegistry;
use crate::error::{OrmError, Result};
use crate::match_request::selected_result_executor::SelectedResultExecutor;
use crate::match_request::{
CapabilitySet, MatchExecutionLimits, ValidatedMatchRequest, ValidatedMatchResult,
};
use type_bridge_core_lib::ast::{
TypedFetchRows, TypedHydrateThings, TypedPageRematch, TypedRootScan,
};
use type_bridge_core_lib::version::Version;
pub struct TransactionContext {
inner: Arc<Mutex<Option<Box<dyn TransactionOps>>>>,
tx_type: TxType,
match_capabilities: CapabilitySet,
server_version: Option<Version>,
}
impl TransactionContext {
pub(crate) fn new(
inner: Box<dyn TransactionOps>,
tx_type: TxType,
match_capabilities: CapabilitySet,
server_version: Option<Version>,
) -> Self {
Self {
inner: Arc::new(Mutex::new(Some(inner))),
tx_type,
match_capabilities,
server_version,
}
}
pub async fn query(&self, typeql: &str) -> Result<QueryResult> {
self.check_schema_annotation_support(typeql)?;
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
}
#[doc(hidden)]
pub(crate) async fn query_canonical(&self, typeql: &str) -> Result<QueryResult> {
self.check_schema_annotation_support(typeql)?;
let mut guard = self.inner.lock().await;
let tx = guard
.as_mut()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.query_canonical(typeql).await
}
pub(crate) async fn schema_snapshot(&self) -> Result<Option<String>> {
let mut guard = self.inner.lock().await;
let tx = guard
.as_mut()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
tx.schema_snapshot().await
}
pub async fn query_with_rows(&self, typeql: &str, rows: GivenRowsSpec) -> Result<QueryResult> {
self.check_schema_annotation_support(typeql)?;
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 supports_exactly_one_tuple_proof(&self) -> Result<bool> {
let guard = self.inner.lock().await;
let tx = guard
.as_ref()
.ok_or_else(|| OrmError::Transaction("Transaction already consumed".into()))?;
Ok(tx.supports_exactly_one_tuple_proof())
}
pub(crate) async fn query_tuple_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_tuple_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
}
fn check_schema_annotation_support(&self, typeql: &str) -> Result<()> {
if self.tx_type == TxType::Schema {
crate::_schema::annotations::check_schema_annotation_support(
typeql,
self.server_version,
)?;
}
Ok(())
}
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_compatible_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(),
server_version: self.server_version,
}
}
}