use std::fmt;
use std::net::Ipv6Addr;
use std::sync::Arc;
use std::time::Duration;
use sha2::{Digest, Sha256};
#[cfg(feature = "typedb")]
use super::backend::AnswerCancellation;
#[cfg(test)]
use super::backend::BoxFuture;
use super::backend::{DriverBackend, GivenRowsSpec, QueryResult, TxType};
use super::context::TransactionContext;
use super::transaction::Transaction;
use crate::_registry::DescriptorRegistry;
use crate::error::Result;
use crate::match_request::selected_result_executor::{
ManagerHydratedRoots, ManagerRootSelection, SelectedResultExecutor,
};
use crate::match_request::{MatchExecutionLimits, ValidatedMatchRequest, ValidatedMatchResult};
#[cfg(feature = "typedb")]
use crate::query_execution_limits::QueryExecutionDeadline;
use crate::query_execution_limits::QueryExecutionResourceLimits;
#[derive(Clone)]
pub struct Database {
backend: Arc<dyn DriverBackend>,
connection_authority: DatabaseConnectionAuthority,
database_name: String,
answer_limits: Option<QueryExecutionResourceLimits>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DatabaseCreateOutcome {
Created,
AlreadyExists,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DatabaseDeleteOutcome {
Deleted,
AlreadyAbsent,
}
#[derive(Clone, Eq, PartialEq)]
pub(crate) struct DatabaseExecutionIdentity {
connection_authority: DatabaseConnectionAuthority,
database_name: String,
}
impl DatabaseExecutionIdentity {
#[cfg(test)]
pub(crate) fn isolated(database_name: impl Into<String>) -> Self {
Self {
connection_authority: DatabaseConnectionAuthority::isolated(),
database_name: database_name.into(),
}
}
}
impl fmt::Debug for DatabaseExecutionIdentity {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("DatabaseExecutionIdentity([REDACTED])")
}
}
#[derive(Clone)]
pub struct DatabaseConnectionAuthority(DatabaseConnectionAuthorityKind);
#[derive(Clone)]
enum DatabaseConnectionAuthorityKind {
Provider([u8; 32]),
Custom(Arc<()>),
}
impl DatabaseConnectionAuthority {
#[must_use]
pub fn isolated() -> Self {
Self(DatabaseConnectionAuthorityKind::Custom(Arc::new(())))
}
fn for_typedb_address(address: &str) -> Self {
if !is_identity_safe_provider_address(address) {
return Self::isolated();
}
let mut digest = Sha256::new();
digest.update(b"typebridge.orm.database-connection-authority/v1\0");
digest.update(address.as_bytes());
Self(DatabaseConnectionAuthorityKind::Provider(
digest.finalize().into(),
))
}
}
#[doc(hidden)]
#[must_use]
pub fn is_identity_safe_provider_address(address: &str) -> bool {
!address.is_empty() && address.split(',').all(identity_safe_provider_endpoint)
}
fn identity_safe_provider_endpoint(endpoint: &str) -> bool {
let (host_is_valid, port) = if let Some(bracketed) = endpoint.strip_prefix('[') {
let Some((address, port)) = bracketed.split_once("]:") else {
return false;
};
(
!address.is_empty()
&& !port.contains(['[', ']', ':'])
&& address.parse::<Ipv6Addr>().is_ok(),
port,
)
} else {
let Some((host, port)) = endpoint.rsplit_once(':') else {
return false;
};
let host = host.strip_suffix('.').unwrap_or(host);
(
!host.is_empty()
&& host.len() <= 253
&& !host.contains(['[', ']', ':'])
&& host.split('.').all(|label| {
!label.is_empty()
&& label.len() <= 63
&& label
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || byte == b'-')
&& label
.as_bytes()
.first()
.is_some_and(u8::is_ascii_alphanumeric)
&& label
.as_bytes()
.last()
.is_some_and(u8::is_ascii_alphanumeric)
}),
port,
)
};
host_is_valid
&& !port.is_empty()
&& port.bytes().all(|byte| byte.is_ascii_digit())
&& port.parse::<u16>().is_ok_and(|port| port != 0)
}
impl PartialEq for DatabaseConnectionAuthority {
fn eq(&self, other: &Self) -> bool {
match (&self.0, &other.0) {
(
DatabaseConnectionAuthorityKind::Provider(left),
DatabaseConnectionAuthorityKind::Provider(right),
) => left == right,
(
DatabaseConnectionAuthorityKind::Custom(left),
DatabaseConnectionAuthorityKind::Custom(right),
) => Arc::ptr_eq(left, right),
_ => false,
}
}
}
impl Eq for DatabaseConnectionAuthority {}
impl fmt::Debug for DatabaseConnectionAuthority {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("DatabaseConnectionAuthority([REDACTED])")
}
}
impl Database {
#[cfg(feature = "typedb")]
#[doc(hidden)]
pub async fn connect_direct(
installed: &crate::InstalledRuntimeProjection,
policy: &super::direct_connection::DirectConnectionPolicy,
cancellation: AnswerCancellation,
) -> std::result::Result<Self, type_bridge_contract::sdk_diagnostic::SdkExecutionDiagnostic>
{
let prepared =
super::direct_connection::prepare_direct_connection(installed, policy, &cancellation)?;
let transport = prepared.transport.with_generated_3_12_3_requirement();
let mut database = Self::connect_prepared_secure_with_control(
&prepared.endpoint,
&prepared.database,
&prepared.username,
&prepared.password,
transport,
prepared.connection_limits,
prepared.deadline,
cancellation,
)
.await
.map_err(super::direct_connection::lower_secure_connection)?;
database.answer_limits = Some(prepared.answer_limits);
Ok(database)
}
pub fn with_backend(backend: Box<dyn DriverBackend>, database_name: impl Into<String>) -> Self {
Self::with_backend_authority(
backend,
database_name,
DatabaseConnectionAuthority::isolated(),
)
}
pub fn with_backend_authority(
backend: Box<dyn DriverBackend>,
database_name: impl Into<String>,
connection_authority: DatabaseConnectionAuthority,
) -> Self {
Self {
backend: Arc::from(backend),
connection_authority,
database_name: database_name.into(),
answer_limits: None,
}
}
#[cfg(feature = "typedb")]
pub async fn connect(
address: &str,
database: &str,
username: &str,
password: &str,
) -> Result<Self> {
Self::connect_with_options(
address,
database,
username,
password,
super::real_driver::ConnectOptions::default(),
)
.await
}
#[cfg(feature = "typedb")]
pub async fn connect_with_options(
address: &str,
database: &str,
username: &str,
password: &str,
options: super::real_driver::ConnectOptions,
) -> Result<Self> {
let backend =
super::real_driver::RealBackend::connect(address, username, password, options).await?;
Ok(Self {
backend: Arc::new(backend),
connection_authority: DatabaseConnectionAuthority::for_typedb_address(address),
database_name: database.to_string(),
answer_limits: None,
})
}
#[cfg(feature = "typedb")]
pub async fn connect_secure_with_options(
address: &str,
database: &str,
username: &str,
password: &str,
options: super::real_driver::SecureConnectOptions,
) -> super::real_driver::SecureResult<Self> {
let backend =
super::real_driver::RealBackend::connect_secure(address, username, password, options)
.await?;
Ok(Self {
backend: Arc::new(backend),
connection_authority: DatabaseConnectionAuthority::for_typedb_address(address),
database_name: database.to_string(),
answer_limits: None,
})
}
#[cfg(feature = "typedb")]
#[doc(hidden)]
pub async fn connect_prepared_secure_with_options(
address: &str,
database: &str,
username: &str,
password: &str,
options: super::real_driver::PreparedSecureConnectOptions,
) -> super::real_driver::SecureResult<Self> {
let backend = super::real_driver::RealBackend::connect_prepared_secure(
address, username, password, options,
)
.await?;
Ok(Self {
backend: Arc::new(backend),
connection_authority: DatabaseConnectionAuthority::for_typedb_address(address),
database_name: database.to_string(),
answer_limits: None,
})
}
#[cfg(feature = "typedb")]
#[doc(hidden)]
#[allow(clippy::too_many_arguments)]
pub async fn connect_prepared_secure_with_control(
address: &str,
database: &str,
username: &str,
password: &str,
options: super::real_driver::PreparedSecureConnectOptions,
limits: QueryExecutionResourceLimits,
deadline: QueryExecutionDeadline,
cancellation: AnswerCancellation,
) -> super::real_driver::SecureResult<Self> {
let limits = limits.effective();
let control = type_bridge_typedb_runtime::RuntimeConnectionControl::new(
limits.items,
limits.bytes,
limits.statements,
deadline.instant(),
type_bridge_typedb_runtime::RuntimeAnswerCancellation::from_shared(
cancellation.shared(),
),
);
Self::connect_prepared_secure_with_options(
address,
database,
username,
password,
options.with_connection_control(control),
)
.await
}
pub async fn read_transaction(&self) -> Result<Transaction> {
let tx = self
.backend
.open_transaction(&self.database_name, TxType::Read)
.await?;
Ok(Transaction::new(tx, TxType::Read, self.server_version()))
}
pub async fn write_transaction(&self) -> Result<Transaction> {
let tx = self
.backend
.open_transaction(&self.database_name, TxType::Write)
.await?;
Ok(Transaction::new(tx, TxType::Write, self.server_version()))
}
#[doc(hidden)]
pub async fn schema_fenced_read_transaction(
&self,
timeout: Duration,
) -> Result<(Transaction, String)> {
let fenced = self
.backend
.open_schema_fenced_read_transaction(&self.database_name, timeout)
.await?;
let (transaction, schema_text) = fenced.into_parts();
Ok((
Transaction::new(transaction, TxType::Write, self.server_version()),
schema_text,
))
}
pub async fn schema_transaction(&self) -> Result<Transaction> {
let tx = self
.backend
.open_transaction(&self.database_name, TxType::Schema)
.await?;
Ok(Transaction::new(tx, TxType::Schema, self.server_version()))
}
pub async fn transaction_context(&self, tx_type: TxType) -> Result<TransactionContext> {
let capabilities = self.backend.match_capabilities();
let tx = self
.backend
.open_transaction(&self.database_name, tx_type)
.await?;
Ok(TransactionContext::new(
tx,
tx_type,
capabilities,
self.server_version(),
self.execution_identity(),
self.answer_limits,
))
}
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.backend.match_capabilities(), limits)
.execute_compatible_owned(self, validated)
.await
}
pub(crate) async fn execute_manager_roots_with_limits(
&self,
registry: &DescriptorRegistry,
validated: &ValidatedMatchRequest,
selection: ManagerRootSelection,
limits: MatchExecutionLimits,
) -> Result<ManagerHydratedRoots> {
SelectedResultExecutor::new(registry, self.backend.match_capabilities(), limits)
.execute_manager_roots_owned(self, validated, selection)
.await
}
#[cfg(feature = "integration-tests")]
#[doc(hidden)]
pub async fn execute_match_v1_legacy_for_live_test(
&self,
registry: &DescriptorRegistry,
validated: &ValidatedMatchRequest,
) -> Result<ValidatedMatchResult> {
let registry = registry.owned_registry_snapshot()?;
SelectedResultExecutor::new(
®istry,
self.backend.match_capabilities(),
MatchExecutionLimits::default(),
)
.execute_owned(self, validated)
.await
}
pub fn database_name(&self) -> &str {
&self.database_name
}
#[doc(hidden)]
#[must_use]
pub fn derived_journal_database(&self) -> Self {
Self {
backend: Arc::clone(&self.backend),
connection_authority: self.connection_authority.clone(),
database_name: format!(
"{}{}",
self.database_name,
type_bridge_contract::reserved::TYPEBRIDGE_JOURNAL_DATABASE_SUFFIX
),
answer_limits: self.answer_limits,
}
}
#[doc(hidden)]
#[must_use]
pub const fn answer_limits(&self) -> Option<QueryExecutionResourceLimits> {
self.answer_limits
}
pub(crate) fn execution_identity(&self) -> DatabaseExecutionIdentity {
DatabaseExecutionIdentity {
connection_authority: self.connection_authority.clone(),
database_name: self.database_name.clone(),
}
}
#[must_use]
pub fn shares_connection_authority_with(&self, other: &Self) -> bool {
self.connection_authority == other.connection_authority
}
pub fn is_connected(&self) -> bool {
self.backend.is_open()
}
pub fn close(&self) -> Result<()> {
self.backend.close_connection()
}
pub fn server_version(&self) -> Option<type_bridge_core_lib::version::Version> {
self.backend.server_version()
}
pub fn check_schema_annotation_support(&self, typeql: &str) -> Result<()> {
crate::_schema::annotations::check_schema_annotation_support(typeql, self.server_version())
}
pub fn supports_given_stage(&self) -> bool {
use type_bridge_core_lib::version::{Feature, check_feature_supported};
self.backend.supports_given_rows()
&& self
.server_version()
.is_some_and(|server| check_feature_supported(Feature::GivenStage, &server).is_ok())
}
pub fn check_given_stage_support(&self) -> Result<()> {
use type_bridge_core_lib::version::{Feature, check_feature_supported};
let Some(server) = self.server_version() else {
return Err(crate::error::OrmError::QueryExecution(
"given-stage support cannot be proven because the server version is unknown".into(),
));
};
check_feature_supported(Feature::GivenStage, &server)
.map_err(crate::error::OrmError::UnsupportedVersion)?;
if !self.backend.supports_given_rows() {
return Err(crate::error::OrmError::QueryExecution(
"given-stage input rows require an active band-9 provider; the connected server supports the syntax but the negotiated provider cannot transport rows"
.into(),
));
}
Ok(())
}
pub async fn database_exists(&self) -> Result<bool> {
self.backend.database_exists(&self.database_name).await
}
pub async fn create_database(&self) -> Result<()> {
self.create_database_outcome().await?;
Ok(())
}
pub async fn create_database_outcome(&self) -> Result<DatabaseCreateOutcome> {
if self.database_exists().await? {
return Ok(DatabaseCreateOutcome::AlreadyExists);
}
match self.backend.create_database(&self.database_name).await {
Ok(()) => Ok(DatabaseCreateOutcome::Created),
Err(error) => match self.database_exists().await {
Ok(true) => Ok(DatabaseCreateOutcome::AlreadyExists),
Ok(false) | Err(_) => Err(error),
},
}
}
pub async fn delete_database(&self) -> Result<()> {
self.delete_database_outcome().await?;
Ok(())
}
pub async fn delete_database_outcome(&self) -> Result<DatabaseDeleteOutcome> {
if !self.database_exists().await? {
return Ok(DatabaseDeleteOutcome::AlreadyAbsent);
}
match self.backend.delete_database(&self.database_name).await {
Ok(()) => Ok(DatabaseDeleteOutcome::Deleted),
Err(error) => match self.database_exists().await {
Ok(false) => Ok(DatabaseDeleteOutcome::Deleted),
Ok(true) | Err(_) => Err(error),
},
}
}
pub async fn schema_text(&self) -> Result<String> {
self.backend.schema_text(&self.database_name).await
}
pub fn into_shared(self) -> Arc<Self> {
Arc::new(self)
}
#[tracing::instrument(skip(self, typeql), fields(db = %self.database_name))]
pub async fn execute_raw(&self, typeql: &str, tx_type: TxType) -> Result<QueryResult> {
if tx_type == TxType::Schema {
self.check_schema_annotation_support(typeql)?;
}
let mut tx = self
.backend
.open_transaction(&self.database_name, tx_type)
.await?;
let result = tx.query(typeql).await?;
if matches!(tx_type, TxType::Write | TxType::Schema) {
tx.commit().await?;
}
Ok(result)
}
pub(crate) async fn execute_canonical(
&self,
typeql: &str,
tx_type: TxType,
) -> Result<QueryResult> {
if tx_type == TxType::Schema {
self.check_schema_annotation_support(typeql)?;
}
let mut tx = self
.backend
.open_transaction(&self.database_name, tx_type)
.await?;
let result = tx.query_canonical(typeql).await;
match result {
Ok(value) => {
if matches!(tx_type, TxType::Write | TxType::Schema)
&& let Err(error) = tx.commit().await
{
let _ = tx.close().await;
return Err(error);
}
tx.close().await?;
Ok(value)
}
Err(primary) => {
if matches!(tx_type, TxType::Write | TxType::Schema) {
let _ = tx.rollback().await;
let _ = tx.close().await;
} else {
let _ = tx.close().await;
}
Err(primary)
}
}
}
#[tracing::instrument(skip(self, typeql, rows), fields(db = %self.database_name))]
pub async fn execute_with_rows(
&self,
typeql: &str,
tx_type: TxType,
rows: GivenRowsSpec,
) -> Result<QueryResult> {
if tx_type == TxType::Schema {
self.check_schema_annotation_support(typeql)?;
}
self.check_given_stage_support()?;
let mut tx = self
.backend
.open_transaction(&self.database_name, tx_type)
.await?;
let result = match tx.query_with_rows(typeql, rows).await {
Ok(result) => result,
Err(error) => {
if matches!(tx_type, TxType::Write | TxType::Schema) {
let _ = tx.rollback().await;
}
let _ = tx.close().await;
return Err(error);
}
};
if matches!(tx_type, TxType::Write | TxType::Schema)
&& let Err(error) = tx.commit().await
{
let _ = tx.close().await;
return Err(error);
}
tx.close().await?;
Ok(result)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicBool, Ordering};
use crate::error::OrmError;
struct LifecycleBackend {
exists: AtomicBool,
fail_create_after_effect: bool,
fail_delete_after_effect: bool,
}
impl DriverBackend for LifecycleBackend {
fn open_transaction(
&self,
_database: &str,
_tx_type: TxType,
) -> BoxFuture<'_, Result<Box<dyn super::super::backend::TransactionOps>>> {
Box::pin(async { Err(OrmError::Connection("unused fixture transaction".into())) })
}
fn is_open(&self) -> bool {
true
}
fn database_exists(&self, _database: &str) -> BoxFuture<'_, Result<bool>> {
Box::pin(async { Ok(self.exists.load(Ordering::SeqCst)) })
}
fn create_database(&self, _database: &str) -> BoxFuture<'_, Result<()>> {
Box::pin(async {
self.exists.store(true, Ordering::SeqCst);
if self.fail_create_after_effect {
Err(OrmError::Connection("simulated concurrent create".into()))
} else {
Ok(())
}
})
}
fn delete_database(&self, _database: &str) -> BoxFuture<'_, Result<()>> {
Box::pin(async {
self.exists.store(false, Ordering::SeqCst);
if self.fail_delete_after_effect {
Err(OrmError::Connection("simulated concurrent delete".into()))
} else {
Ok(())
}
})
}
}
fn lifecycle_database(
exists: bool,
fail_create_after_effect: bool,
fail_delete_after_effect: bool,
) -> Database {
Database::with_backend(
Box::new(LifecycleBackend {
exists: AtomicBool::new(exists),
fail_create_after_effect,
fail_delete_after_effect,
}),
"bound-database",
)
}
#[test]
fn derived_journal_binding_retains_exact_provider_authority() {
let database = lifecycle_database(false, false, false);
let journal = database.derived_journal_database();
assert_eq!(journal.database_name(), "bound-database__tbv2_journal");
assert!(database.shares_connection_authority_with(&journal));
}
#[tokio::test]
async fn bound_administration_outcomes_are_idempotent_and_race_normalized() {
let database = lifecycle_database(false, false, false);
assert_eq!(
database.create_database_outcome().await.unwrap(),
DatabaseCreateOutcome::Created
);
assert_eq!(
database.create_database_outcome().await.unwrap(),
DatabaseCreateOutcome::AlreadyExists
);
assert_eq!(
database.delete_database_outcome().await.unwrap(),
DatabaseDeleteOutcome::Deleted
);
assert_eq!(
database.delete_database_outcome().await.unwrap(),
DatabaseDeleteOutcome::AlreadyAbsent
);
let concurrent_create = lifecycle_database(false, true, false);
assert_eq!(
concurrent_create.create_database_outcome().await.unwrap(),
DatabaseCreateOutcome::AlreadyExists
);
let concurrent_delete = lifecycle_database(true, false, true);
assert_eq!(
concurrent_delete.delete_database_outcome().await.unwrap(),
DatabaseDeleteOutcome::Deleted
);
}
#[test]
fn connection_authority_is_opaque_redacted_and_exact() {
const SENTINEL: &str = "TB_AUTHORITY_SECRET_31d7";
let first = DatabaseConnectionAuthority::for_typedb_address("provider.example:1729");
let same = DatabaseConnectionAuthority::for_typedb_address("provider.example:1729");
let different = DatabaseConnectionAuthority::for_typedb_address("provider.example:1730");
assert_eq!(first, same);
assert_ne!(first, different);
let rendered = format!("{first:?}");
assert_eq!(rendered, "DatabaseConnectionAuthority([REDACTED])");
assert!(!rendered.contains(SENTINEL));
assert!(!rendered.contains("provider.example"));
let unsafe_address = format!("admin:{SENTINEL}@provider.example:1729");
let unsafe_first = DatabaseConnectionAuthority::for_typedb_address(&unsafe_address);
let unsafe_second = DatabaseConnectionAuthority::for_typedb_address(&unsafe_address);
assert_ne!(unsafe_first, unsafe_second);
assert!(!format!("{unsafe_first:?}").contains(SENTINEL));
let isolated = DatabaseConnectionAuthority::isolated();
assert_eq!(isolated, isolated.clone());
assert_ne!(isolated, DatabaseConnectionAuthority::isolated());
}
#[test]
fn public_identity_safe_address_validator_reuses_authority_grammar() {
for valid in [
"provider.example:1729",
"provider.example:1729,backup.example:1730",
"[2001:db8::1]:1729",
] {
assert!(is_identity_safe_provider_address(valid), "{valid}");
}
for invalid in [
"",
"typedb://provider.example:1729",
"admin@provider.example:1729",
" provider.example:1729",
"provider.example:0",
] {
assert!(!is_identity_safe_provider_address(invalid), "{invalid}");
}
}
#[cfg(feature = "typedb")]
#[tokio::test]
async fn controlled_prepared_connect_forwards_common_pre_dispatch_cancellation() {
const SENTINEL: &str = "TB_ORM_CONTROLLED_CONNECT_SECRET";
let options = super::super::real_driver::SecureConnectOptions {
http_port: 8123,
tls_mode: super::super::real_driver::TlsMode::Disabled,
server_version: Some(type_bridge_core_lib::version::Version::new(3, 12, 1)),
}
.prepare_transport()
.expect("prepare credential-free transport before controlled connect");
let limits = QueryExecutionResourceLimits::default();
let deadline = QueryExecutionDeadline::for_limits(limits);
let cancellation = AnswerCancellation::default();
cancellation.cancel();
let error = Database::connect_prepared_secure_with_control(
SENTINEL,
"database",
SENTINEL,
SENTINEL,
options,
limits,
deadline,
cancellation,
)
.await
.err()
.expect("pre-cancelled controlled connection must reject without provider I/O");
assert!(matches!(
error,
super::super::real_driver::SecureConnectError::Runtime(
type_bridge_typedb_runtime::RuntimeError::ResourceLimit {
code: "provider_cancelled",
..
}
)
));
assert!(!error.to_string().contains(SENTINEL));
}
#[cfg(feature = "typedb")]
#[tokio::test]
async fn controlled_prepared_connect_forwards_zero_statement_ceiling_before_credentials() {
const SENTINEL: &str = "TB_ORM_CONNECTION_CREDENTIAL_SECRET";
let options = super::super::real_driver::SecureConnectOptions {
http_port: 8123,
tls_mode: super::super::real_driver::TlsMode::Disabled,
server_version: Some(type_bridge_core_lib::version::Version::new(3, 12, 1)),
}
.prepare_transport()
.expect("prepare transport before supplying credentials");
let limits = QueryExecutionResourceLimits {
statements: 0,
..QueryExecutionResourceLimits::default()
};
let deadline = QueryExecutionDeadline::for_limits(limits);
let error = Database::connect_prepared_secure_with_control(
"127.0.0.1:1",
"database",
SENTINEL,
SENTINEL,
options,
limits,
deadline,
AnswerCancellation::default(),
)
.await
.err()
.expect("zero statements must reject before driver credential construction");
assert!(matches!(
error,
super::super::real_driver::SecureConnectError::Runtime(
type_bridge_typedb_runtime::RuntimeError::ResourceLimit {
code: "provider_statement_limit",
..
}
)
));
assert!(!error.to_string().contains(SENTINEL));
}
}