use std::fmt;
use std::sync::Arc;
#[allow(unused_imports)]
use crate::error::{Error, Result};
use crate::schema::{Schema, SchemaPackage, Unbound};
use type_bridge_orm::_registry::DescriptorRegistry;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DatabaseCreateOutcome {
Created,
AlreadyExists,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DatabaseDeleteOutcome {
Deleted,
AlreadyAbsent,
}
#[cfg(feature = "typedb")]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ManagedDatabasePairState {
Absent,
StandaloneManaged,
OwnedPair,
OwnedJournalOrphan,
}
#[cfg(feature = "typedb")]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ManagedDatabaseDeleteOutcome {
AlreadyAbsent,
DeletedStandaloneManaged,
DeletedOwnedPair,
DeletedOwnedJournalOrphan,
}
#[cfg(feature = "typedb")]
pub struct ManagedDatabaseDeletionPlan {
inner: type_bridge_schema_migration_typedb::ManagedDatabasePairDeletionPlan,
}
#[cfg(feature = "typedb")]
impl ManagedDatabaseDeletionPlan {
#[must_use]
pub fn inspected_state(&self) -> ManagedDatabasePairState {
map_pair_state(self.inner.inspected_state())
}
pub async fn execute(self) -> Result<ManagedDatabaseDeleteOutcome> {
self.inner
.execute()
.await
.map(map_delete_outcome)
.map_err(administration_error)
}
pub async fn execute_controlled(
self,
control: &crate::MigrationExecutionControl,
) -> Result<ManagedDatabaseDeleteOutcome> {
self.inner
.execute_controlled(control)
.await
.map(map_delete_outcome)
.map_err(administration_error)
}
}
#[derive(Clone, PartialEq, Eq)]
pub struct ConnectionOptions {
address: String,
database: String,
username: Option<String>,
password: Option<String>,
http_port: u16,
tls: bool,
}
impl fmt::Debug for ConnectionOptions {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("ConnectionOptions")
.field("address", &self.address)
.field("database", &self.database)
.field("username", &self.username)
.field("password", &self.password.as_ref().map(|_| "[REDACTED]"))
.field("http_port", &self.http_port)
.field("tls", &self.tls)
.finish()
}
}
impl ConnectionOptions {
#[must_use]
pub fn new(address: impl Into<String>, database: impl Into<String>) -> Self {
Self {
address: address.into(),
database: database.into(),
username: None,
password: None,
http_port: 8000,
tls: false,
}
}
#[must_use]
pub fn credentials(mut self, username: impl Into<String>, password: impl Into<String>) -> Self {
self.username = Some(username.into());
self.password = Some(password.into());
self
}
#[must_use]
pub fn http_port(mut self, port: u16) -> Self {
self.http_port = port;
self
}
#[must_use]
pub fn tls(mut self, enabled: bool) -> Self {
self.tls = enabled;
self
}
#[must_use]
pub fn address(&self) -> &str {
&self.address
}
#[must_use]
pub fn database(&self) -> &str {
&self.database
}
#[must_use]
pub fn get_http_port(&self) -> u16 {
self.http_port
}
#[must_use]
pub fn is_tls(&self) -> bool {
self.tls
}
}
impl From<(&str, &str)> for ConnectionOptions {
fn from((address, database): (&str, &str)) -> Self {
Self::new(address, database)
}
}
pub struct Database<S: Schema = Unbound> {
inner: type_bridge_orm::Database,
installed_schema: Option<Arc<type_bridge_orm::InstalledRuntimeProjection>>,
match_registry: Option<Arc<DescriptorRegistry>>,
#[cfg_attr(not(feature = "typedb"), allow(dead_code))]
managed_scope_id: Option<type_bridge_contract::managed_scope::ManagedScopeId>,
marker: std::marker::PhantomData<fn() -> S>,
}
pub(crate) fn build_match_registry(
installed: &type_bridge_orm::InstalledRuntimeProjection,
) -> Result<Arc<DescriptorRegistry>> {
installed
.match_registry()
.map(Arc::new)
.map_err(Error::from_orm)
}
impl Database<Unbound> {
#[cfg(feature = "typedb")]
pub async fn connect(options: impl Into<ConnectionOptions>) -> Result<Database<Unbound>> {
let opts = options.into();
let username = opts.username.as_deref().unwrap_or("admin");
let password = opts.password.as_deref().unwrap_or("password");
let orm_opts = type_bridge_orm::ConnectOptions {
http_port: opts.http_port,
tls: opts.tls,
..type_bridge_orm::ConnectOptions::default()
};
let inner = type_bridge_orm::Database::connect_with_options(
&opts.address,
&opts.database,
username,
password,
orm_opts,
)
.await
.map_err(Error::from_orm)?;
Ok(Database {
inner,
installed_schema: None,
match_registry: None,
managed_scope_id: None,
marker: std::marker::PhantomData,
})
}
#[allow(dead_code)]
pub(crate) fn from_orm_database(inner: type_bridge_orm::Database) -> Self {
Self {
inner,
installed_schema: None,
match_registry: None,
managed_scope_id: None,
marker: std::marker::PhantomData,
}
}
pub fn with_schema<S: Schema>(self, schema: SchemaPackage<S>) -> Result<Database<S>> {
let (installed, authority) = schema.verify_and_install_with_authority()?;
let match_registry = build_match_registry(&installed)?;
Ok(Database::from_bound_parts(
self.inner,
installed,
match_registry,
authority.map(|authority| authority.managed_scope().id().clone()),
))
}
}
impl<S: Schema> Database<S> {
pub(crate) fn from_bound_parts(
inner: type_bridge_orm::Database,
installed: Arc<type_bridge_orm::InstalledRuntimeProjection>,
match_registry: Arc<DescriptorRegistry>,
managed_scope_id: Option<type_bridge_contract::managed_scope::ManagedScopeId>,
) -> Self {
Self {
inner,
installed_schema: Some(installed),
match_registry: Some(match_registry),
managed_scope_id,
marker: std::marker::PhantomData,
}
}
#[cfg(test)]
pub(crate) fn from_test_parts(
inner: type_bridge_orm::Database,
installed: type_bridge_orm::InstalledRuntimeProjection,
) -> Self {
let installed = Arc::new(installed);
let match_registry =
build_match_registry(&installed).expect("test projection descriptors register");
Self {
inner,
installed_schema: Some(installed),
match_registry: Some(match_registry),
managed_scope_id: None,
marker: std::marker::PhantomData,
}
}
#[cfg(test)]
pub(crate) fn from_test_unbound_parts(inner: type_bridge_orm::Database) -> Self {
Self {
inner,
installed_schema: None,
match_registry: None,
managed_scope_id: None,
marker: std::marker::PhantomData,
}
}
pub fn entities<M>(&self) -> crate::entity_manager::EntityManager<'_, S, M>
where
M: crate::__codegen::EntityModel<Schema = S>,
{
crate::entity_manager::EntityManager::new(self)
}
pub fn relations<M>(&self) -> crate::relation_manager::RelationManager<'_, S, M>
where
M: crate::__codegen::RelationModel<Schema = S>,
{
crate::relation_manager::RelationManager::new(self)
}
pub async fn write(&self) -> Result<crate::transaction::WriteTransaction<'_, S>> {
crate::transaction::WriteTransaction::open(self).await
}
pub async fn read(&self) -> Result<crate::transaction::ReadTransaction<'_, S>> {
crate::transaction::ReadTransaction::open(self).await
}
#[must_use]
pub fn database_name(&self) -> &str {
self.inner.database_name()
}
pub async fn database_exists(&self) -> Result<bool> {
self.inner.database_exists().await.map_err(Error::from_orm)
}
#[cfg(feature = "typedb")]
pub async fn database_exists_controlled(
&self,
control: &crate::MigrationExecutionControl,
) -> Result<bool> {
if let Some(scope) = &self.managed_scope_id {
return self
.pair_administrator(scope.clone())?
.database_exists_controlled(control)
.await
.map_err(administration_error);
}
control.check().map_err(administration_error)?;
self.database_exists().await
}
pub async fn create_database(&self) -> Result<DatabaseCreateOutcome> {
#[cfg(feature = "typedb")]
if let Some(scope) = &self.managed_scope_id {
return self
.pair_administrator(scope.clone())?
.create_database_outcome()
.await
.map(|outcome| match outcome {
type_bridge_schema_migration_typedb::ManagedDatabasePairCreateOutcome::Created => DatabaseCreateOutcome::Created,
type_bridge_schema_migration_typedb::ManagedDatabasePairCreateOutcome::AlreadyExists => DatabaseCreateOutcome::AlreadyExists,
})
.map_err(administration_error);
}
self.inner
.create_database_outcome()
.await
.map(|outcome| match outcome {
type_bridge_orm::session::DatabaseCreateOutcome::Created => {
DatabaseCreateOutcome::Created
}
type_bridge_orm::session::DatabaseCreateOutcome::AlreadyExists => {
DatabaseCreateOutcome::AlreadyExists
}
})
.map_err(Error::from_orm)
}
#[cfg(feature = "typedb")]
pub async fn create_database_controlled(
&self,
control: &crate::MigrationExecutionControl,
) -> Result<DatabaseCreateOutcome> {
if let Some(scope) = &self.managed_scope_id {
return self
.pair_administrator(scope.clone())?
.create_database_outcome_controlled(control)
.await
.map(|outcome| match outcome {
type_bridge_schema_migration_typedb::ManagedDatabasePairCreateOutcome::Created => DatabaseCreateOutcome::Created,
type_bridge_schema_migration_typedb::ManagedDatabasePairCreateOutcome::AlreadyExists => DatabaseCreateOutcome::AlreadyExists,
})
.map_err(administration_error);
}
control.check().map_err(administration_error)?;
self.create_database().await
}
pub async fn delete_database(&self) -> Result<DatabaseDeleteOutcome> {
#[cfg(feature = "typedb")]
if self.managed_scope_id.is_some() {
return Err(Error::Database {
message: "managed database deletion requires plan_database_delete() and explicit plan execution"
.to_owned(),
source: None,
});
}
self.inner
.delete_database_outcome()
.await
.map(|outcome| match outcome {
type_bridge_orm::session::DatabaseDeleteOutcome::Deleted => {
DatabaseDeleteOutcome::Deleted
}
type_bridge_orm::session::DatabaseDeleteOutcome::AlreadyAbsent => {
DatabaseDeleteOutcome::AlreadyAbsent
}
})
.map_err(Error::from_orm)
}
#[cfg(feature = "typedb")]
pub async fn inspect_database_pair(&self) -> Result<ManagedDatabasePairState> {
let scope = self
.managed_scope_id
.clone()
.ok_or_else(|| Error::Database {
message:
"managed database administration requires verified generated schema authority"
.to_owned(),
source: None,
})?;
self.pair_administrator(scope)?
.inspect()
.await
.map(map_pair_state)
.map_err(administration_error)
}
#[cfg(feature = "typedb")]
pub async fn inspect_database_pair_controlled(
&self,
control: &crate::MigrationExecutionControl,
) -> Result<ManagedDatabasePairState> {
let scope = self
.managed_scope_id
.clone()
.ok_or_else(|| Error::Database {
message:
"managed database administration requires verified generated schema authority"
.to_owned(),
source: None,
})?;
self.pair_administrator(scope)?
.inspect_controlled(control)
.await
.map(map_pair_state)
.map_err(administration_error)
}
#[cfg(feature = "typedb")]
pub async fn plan_database_delete(&self) -> Result<ManagedDatabaseDeletionPlan> {
let scope = self
.managed_scope_id
.clone()
.ok_or_else(|| Error::Database {
message: "managed database deletion requires verified generated schema authority"
.to_owned(),
source: None,
})?;
self.pair_administrator(scope)?
.plan_delete()
.await
.map(|inner| ManagedDatabaseDeletionPlan { inner })
.map_err(administration_error)
}
#[cfg(feature = "typedb")]
pub async fn plan_database_delete_controlled(
&self,
control: &crate::MigrationExecutionControl,
) -> Result<ManagedDatabaseDeletionPlan> {
let scope = self
.managed_scope_id
.clone()
.ok_or_else(|| Error::Database {
message: "managed database deletion requires verified generated schema authority"
.to_owned(),
source: None,
})?;
self.pair_administrator(scope)?
.plan_delete_controlled(control)
.await
.map(|inner| ManagedDatabaseDeletionPlan { inner })
.map_err(administration_error)
}
#[cfg(feature = "typedb")]
fn pair_administrator(
&self,
scope: type_bridge_contract::managed_scope::ManagedScopeId,
) -> Result<type_bridge_schema_migration_typedb::ManagedDatabasePairAdministrator> {
type_bridge_schema_migration_typedb::ManagedDatabasePairAdministrator::from_managed_database(
Arc::new(self.inner.clone()),
scope,
)
.map_err(administration_error)
}
pub fn close(&self) -> Result<()> {
self.inner.close().map_err(Error::from_orm)
}
#[must_use]
pub fn is_schema_bound(&self) -> bool {
self.installed_schema.is_some()
}
#[allow(dead_code)]
pub(crate) fn inner_orm(&self) -> &type_bridge_orm::Database {
&self.inner
}
pub(crate) fn operation_limits(
&self,
requested: type_bridge_orm::QueryExecutionResourceLimits,
) -> type_bridge_orm::QueryExecutionResourceLimits {
self.inner.answer_limits().map_or_else(
|| requested.effective(),
|ceiling| requested.constrained_by(ceiling),
)
}
#[allow(dead_code)]
pub(crate) fn installed_schema(
&self,
) -> Option<&Arc<type_bridge_orm::InstalledRuntimeProjection>> {
self.installed_schema.as_ref()
}
#[allow(dead_code)]
pub(crate) fn match_registry(&self) -> Option<&Arc<DescriptorRegistry>> {
self.match_registry.as_ref()
}
}
#[cfg(feature = "typedb")]
fn map_pair_state(
state: type_bridge_schema_migration_typedb::ManagedDatabasePairState,
) -> ManagedDatabasePairState {
match state {
type_bridge_schema_migration_typedb::ManagedDatabasePairState::Absent => {
ManagedDatabasePairState::Absent
}
type_bridge_schema_migration_typedb::ManagedDatabasePairState::StandaloneManaged => {
ManagedDatabasePairState::StandaloneManaged
}
type_bridge_schema_migration_typedb::ManagedDatabasePairState::OwnedPair => {
ManagedDatabasePairState::OwnedPair
}
type_bridge_schema_migration_typedb::ManagedDatabasePairState::OwnedJournalOrphan => {
ManagedDatabasePairState::OwnedJournalOrphan
}
}
}
#[cfg(feature = "typedb")]
fn map_delete_outcome(
outcome: type_bridge_schema_migration_typedb::ManagedDatabasePairDeleteOutcome,
) -> ManagedDatabaseDeleteOutcome {
match outcome {
type_bridge_schema_migration_typedb::ManagedDatabasePairDeleteOutcome::AlreadyAbsent => {
ManagedDatabaseDeleteOutcome::AlreadyAbsent
}
type_bridge_schema_migration_typedb::ManagedDatabasePairDeleteOutcome::DeletedStandaloneManaged => ManagedDatabaseDeleteOutcome::DeletedStandaloneManaged,
type_bridge_schema_migration_typedb::ManagedDatabasePairDeleteOutcome::DeletedOwnedPair => {
ManagedDatabaseDeleteOutcome::DeletedOwnedPair
}
type_bridge_schema_migration_typedb::ManagedDatabasePairDeleteOutcome::DeletedOwnedJournalOrphan => ManagedDatabaseDeleteOutcome::DeletedOwnedJournalOrphan,
}
}
#[cfg(feature = "typedb")]
fn administration_error(error: type_bridge_contract::diagnostic::Diagnostic) -> Error {
Error::from_contract_diagnostic(error)
}
#[cfg(test)]
mod tests {
use std::collections::BTreeSet;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use type_bridge_orm::error::OrmError;
use type_bridge_orm::session::backend::{BoxFuture, DriverBackend, TransactionOps, TxType};
use super::Database;
use crate::schema::Unbound;
struct CloseBackend {
closed: Arc<AtomicBool>,
close_calls: Arc<AtomicUsize>,
}
struct AdministrationBackend {
databases: Arc<Mutex<BTreeSet<String>>>,
}
impl DriverBackend for AdministrationBackend {
fn open_transaction(
&self,
_database: &str,
_tx_type: TxType,
) -> BoxFuture<'_, std::result::Result<Box<dyn TransactionOps>, OrmError>> {
Box::pin(async { Err(OrmError::Connection("unexpected transaction".into())) })
}
fn is_open(&self) -> bool {
true
}
fn database_exists(
&self,
database: &str,
) -> BoxFuture<'_, std::result::Result<bool, OrmError>> {
let exists = self.databases.lock().unwrap().contains(database);
Box::pin(async move { Ok(exists) })
}
fn create_database(
&self,
database: &str,
) -> BoxFuture<'_, std::result::Result<(), OrmError>> {
self.databases.lock().unwrap().insert(database.to_owned());
Box::pin(async { Ok(()) })
}
fn delete_database(
&self,
database: &str,
) -> BoxFuture<'_, std::result::Result<(), OrmError>> {
self.databases.lock().unwrap().remove(database);
Box::pin(async { Ok(()) })
}
fn schema_text(
&self,
_database: &str,
) -> BoxFuture<'_, std::result::Result<String, OrmError>> {
Box::pin(async { Ok(String::new()) })
}
}
impl DriverBackend for CloseBackend {
fn open_transaction(
&self,
_database: &str,
_tx_type: TxType,
) -> BoxFuture<'_, std::result::Result<Box<dyn TransactionOps>, OrmError>> {
Box::pin(async {
Err(OrmError::Connection(
"closed test backend cannot open transactions".into(),
))
})
}
fn is_open(&self) -> bool {
!self.closed.load(Ordering::SeqCst)
}
fn close_connection(&self) -> std::result::Result<(), OrmError> {
self.close_calls.fetch_add(1, Ordering::SeqCst);
self.closed.store(true, Ordering::SeqCst);
Ok(())
}
}
#[test]
fn explicit_database_close_is_idempotent() {
let closed = Arc::new(AtomicBool::new(false));
let close_calls = Arc::new(AtomicUsize::new(0));
let inner = type_bridge_orm::Database::with_backend(
Box::new(CloseBackend {
closed: Arc::clone(&closed),
close_calls: Arc::clone(&close_calls),
}),
"app",
);
let database: Database<Unbound> = Database::from_test_unbound_parts(inner);
database.close().unwrap();
database.close().unwrap();
assert!(closed.load(Ordering::SeqCst));
assert_eq!(close_calls.load(Ordering::SeqCst), 2);
}
#[cfg(feature = "typedb")]
#[tokio::test]
async fn generated_database_administration_is_pair_aware_and_plan_owned() {
let databases = Arc::new(Mutex::new(BTreeSet::new()));
let inner = type_bridge_orm::Database::with_backend(
Box::new(AdministrationBackend {
databases: Arc::clone(&databases),
}),
"app",
);
let database: Database<Unbound> = Database {
inner,
installed_schema: None,
match_registry: None,
managed_scope_id: Some(
type_bridge_contract::managed_scope::ManagedScopeId::new("generated-scope")
.unwrap(),
),
marker: std::marker::PhantomData,
};
let cancellation = type_bridge_schema_migration::MigrationCancellation::default();
cancellation.cancel();
let control = type_bridge_schema_migration::MigrationExecutionControl::new(
cancellation,
None,
type_bridge_schema_migration::MigrationExecutionResourceLimits::default(),
);
let cancelled = database
.database_exists_controlled(&control)
.await
.expect_err("pre-cancelled administration rejects before an effect");
assert_eq!(cancelled.code(), Some("migration_execution_cancelled"));
assert_eq!(cancelled.category(), crate::ErrorCategory::Cancelled);
assert_eq!(
database.create_database().await.unwrap(),
super::DatabaseCreateOutcome::Created
);
assert_eq!(
database.inspect_database_pair().await.unwrap(),
super::ManagedDatabasePairState::StandaloneManaged
);
assert!(database.delete_database().await.is_err());
let plan = database.plan_database_delete().await.unwrap();
drop(database);
assert_eq!(
plan.execute().await.unwrap(),
super::ManagedDatabaseDeleteOutcome::DeletedStandaloneManaged
);
assert!(databases.lock().unwrap().is_empty());
}
}