use crate::migrate::{Migration, MigrationScope};
use std::path::Path;
#[cfg(feature = "postgres")]
use std::sync::Arc;
use tracing::{info, warn};
use crate::core::Column as _;
use crate::migrate;
#[cfg(feature = "postgres")]
use crate::sql::sqlx::postgres::PgPoolOptions;
#[cfg(feature = "postgres")]
use crate::sql::sqlx::PgPool;
use sqlx::Database;
use super::error::TenancyError;
use super::org::{Org, StorageMode};
use super::pools::TenantPools;
#[derive(Debug, Default)]
pub struct TenantMigrationReport {
pub tenants: Vec<TenantMigrationOutcome>,
}
impl TenantMigrationReport {
#[must_use]
pub fn all_ok(&self) -> bool {
self.tenants.iter().all(|t| t.error.is_none())
}
#[must_use]
pub fn failure_count(&self) -> usize {
self.tenants.iter().filter(|t| t.error.is_some()).count()
}
}
#[derive(Debug)]
pub struct TenantMigrationOutcome {
pub slug: String,
pub applied: Vec<Migration>,
pub error: Option<TenancyError>,
}
const SYSTEM_LEDGER: &str = "__rustango_system_migrations__";
async fn apply_system_migrations(
pool: &crate::sql::Pool,
dir: &Path,
scope: crate::core::ModelScope,
) -> Result<Vec<Migration>, TenancyError> {
let project_root = if dir.file_name().and_then(|n| n.to_str()) == Some("migrations") {
dir.parent().unwrap_or(dir)
} else {
dir
};
let _ = crate::migrate::make_migrations_system(project_root, scope, None);
let system_dir = project_root.join("system").join("migrations");
if !system_dir.is_dir() {
return Ok(Vec::new());
}
let migration_scope = match scope {
crate::core::ModelScope::Registry => MigrationScope::Registry,
crate::core::ModelScope::Tenant => MigrationScope::Tenant,
};
let applied = match scoped_subset(&system_dir, migration_scope).await? {
ScopedDir::Owned(temp) => {
let r =
crate::migrate::migrate_pool_with_ledger(pool, temp.path(), SYSTEM_LEDGER).await?;
drop(temp);
r
}
ScopedDir::Original => {
crate::migrate::migrate_pool_with_ledger(pool, &system_dir, SYSTEM_LEDGER).await?
}
};
Ok(applied)
}
pub async fn migrate_registry<DB: Database>(
pools: &TenantPools<DB>,
dir: &Path,
) -> Result<Vec<Migration>, TenancyError>
where
crate::sql::Pool: From<sqlx::Pool<DB>>,
{
migrate_registry_pool(&pools.registry_pool(), dir).await
}
pub async fn migrate_registry_pool(
registry: &crate::sql::Pool,
dir: &Path,
) -> Result<Vec<Migration>, TenancyError> {
info!(target: "crate::tenancy", "applying registry-scoped migrations");
let scoped_dir = scoped_subset(dir, MigrationScope::Registry).await?;
let mut applied = match scoped_dir {
ScopedDir::Owned(temp) => {
let result = crate::migrate::migrate_pool(registry, temp.path()).await?;
drop(temp);
result
}
ScopedDir::Original => crate::migrate::migrate_pool(registry, dir).await?,
};
applied
.extend(apply_system_migrations(registry, dir, crate::core::ModelScope::Registry).await?);
if let Err(e) = crate::contenttypes::ensure_seeded(registry).await {
tracing::warn!(
target: "crate::tenancy",
error = %e,
"contenttypes::ensure_seeded failed for registry pool",
);
}
info!(
target: "crate::tenancy",
applied = applied.len(),
"registry migrations done"
);
Ok(applied)
}
#[cfg(feature = "postgres")]
pub async fn migrate_tenants(
pools: &TenantPools,
dir: &Path,
registry_url: &str,
) -> Result<TenantMigrationReport, TenancyError> {
let scoped = scoped_subset(dir, MigrationScope::Tenant).await?;
let scoped_path = match &scoped {
ScopedDir::Owned(temp) => temp.path().to_path_buf(),
ScopedDir::Original => dir.to_path_buf(),
};
let migrations_in_scope = rustango::migrate::file::list_dir(&scoped_path)
.map(|m| m.len())
.unwrap_or(0);
let orgs: Vec<Org> = Org::objects()
.where_(Org::active.eq(true))
.fetch_on(pools.registry())
.await?;
info!(
target: "crate::tenancy",
tenants = orgs.len(),
migrations = migrations_in_scope,
dir = %dir.display(),
"applying tenant-scoped migrations"
);
if migrations_in_scope == 0 && !orgs.is_empty() {
warn!(
target: "crate::tenancy",
dir = %dir.display(),
"no tenant-scoped migrations found in dir; tenants will record applied=0 — \
pass the flat migrations directory or a project root containing one"
);
}
let mut report = TenantMigrationReport::default();
for org in &orgs {
let outcome = run_for_one_tenant(pools, org, &scoped_path, registry_url).await;
match &outcome {
Ok(applied) => info!(
target: "crate::tenancy",
slug = %org.slug,
applied = applied.len(),
"tenant migrations done"
),
Err(e) => warn!(
target: "crate::tenancy",
slug = %org.slug,
error = %e,
"tenant migration failed; continuing with remaining tenants"
),
}
report.tenants.push(match outcome {
Ok(applied) => TenantMigrationOutcome {
slug: org.slug.clone(),
applied,
error: None,
},
Err(error) => TenantMigrationOutcome {
slug: org.slug.clone(),
applied: Vec::new(),
error: Some(error),
},
});
}
Ok(report)
}
pub async fn migrate_tenants_db<DB: Database>(
pools: &TenantPools<DB>,
dir: &Path,
_registry_url: &str,
) -> Result<TenantMigrationReport, TenancyError>
where
crate::sql::Pool: From<sqlx::Pool<DB>>,
{
use crate::sql::FetcherPool as _;
let scoped = scoped_subset(dir, MigrationScope::Tenant).await?;
let scoped_path = match &scoped {
ScopedDir::Owned(temp) => temp.path().to_path_buf(),
ScopedDir::Original => dir.to_path_buf(),
};
let registry_pool = pools.registry_pool();
let orgs: Vec<Org> = Org::objects()
.where_(Org::active.eq(true))
.fetch(®istry_pool)
.await?;
info!(
target: "crate::tenancy",
tenants = orgs.len(),
dir = %dir.display(),
"applying tenant-scoped migrations (db-mode only)"
);
let mut report = TenantMigrationReport::default();
for org in &orgs {
let outcome = run_for_one_tenant_db(pools, org, &scoped_path).await;
match &outcome {
Ok(applied) => info!(
target: "crate::tenancy",
slug = %org.slug,
applied = applied.len(),
"tenant migrations done"
),
Err(e) => warn!(
target: "crate::tenancy",
slug = %org.slug,
error = %e,
"tenant migration failed; continuing with remaining tenants"
),
}
report.tenants.push(match outcome {
Ok(applied) => TenantMigrationOutcome {
slug: org.slug.clone(),
applied,
error: None,
},
Err(error) => TenantMigrationOutcome {
slug: org.slug.clone(),
applied: Vec::new(),
error: Some(error),
},
});
}
Ok(report)
}
async fn run_for_one_tenant_db<DB: Database>(
pools: &TenantPools<DB>,
org: &Org,
dir: &Path,
) -> Result<Vec<Migration>, TenancyError>
where
crate::sql::Pool: From<sqlx::Pool<DB>>,
{
let mode = StorageMode::parse(&org.storage_mode).map_err(|got| {
TenancyError::Validation(format!(
"org `{}` has unknown storage_mode `{got}`",
org.slug
))
})?;
if !matches!(mode, StorageMode::Database) {
return Err(TenancyError::Validation(format!(
"org `{}` is schema-mode but migrate_tenants_db only handles \
database-mode tenants (schema-mode is PG-only by language)",
org.slug,
)));
}
let tenant_pool = pools.database_pool_for_org(org).await?;
let inner_pool = match &tenant_pool {
super::pools::TenantPool::Database { pool } => crate::sql::Pool::from((**pool).clone()),
#[cfg(feature = "postgres")]
super::pools::TenantPool::Schema { .. } => {
unreachable!("database_pool_for_org rejects schema-mode")
}
};
let mut applied = migrate::migrate_pool(&inner_pool, dir).await?;
applied
.extend(apply_system_migrations(&inner_pool, dir, crate::core::ModelScope::Tenant).await?);
if let Err(e) = super::permissions::auto_create_permissions_pool(&inner_pool).await {
tracing::warn!(target: "crate::tenancy", slug = %org.slug, error = %e, "auto_create_permissions_pool failed for database-mode tenant");
}
if let Err(e) = crate::contenttypes::ensure_seeded(&inner_pool).await {
tracing::warn!(target: "crate::tenancy", slug = %org.slug, error = %e, "contenttypes::ensure_seeded failed for database-mode tenant");
}
Ok(applied)
}
pub async fn migrate_tenants_dyn<DB: Database>(
pools: &TenantPools<DB>,
dir: &Path,
registry_url: &str,
) -> Result<TenantMigrationReport, TenancyError>
where
crate::sql::Pool: From<sqlx::Pool<DB>>,
{
#[cfg(feature = "postgres")]
if let Some(pg) = (pools as &dyn std::any::Any).downcast_ref::<TenantPools<sqlx::Postgres>>() {
return migrate_tenants(pg, dir, registry_url).await;
}
migrate_tenants_db(pools, dir, registry_url).await
}
#[cfg(feature = "postgres")]
async fn run_for_one_tenant(
pools: &TenantPools,
org: &Org,
dir: &Path,
registry_url: &str,
) -> Result<Vec<Migration>, TenancyError> {
let mode = StorageMode::parse(&org.storage_mode).map_err(|got| {
TenancyError::Validation(format!(
"org `{}` has unknown storage_mode `{got}`",
org.slug
))
})?;
match mode {
StorageMode::Schema => {
let schema = org.schema_name.clone().unwrap_or_else(|| org.slug.clone());
let pool = build_schema_scoped_pool(registry_url, &schema).await?;
let mut applied = migrate::migrate(&pool, dir).await?;
let dbpool: crate::sql::Pool = pool.clone().into();
applied.extend(
apply_system_migrations(&dbpool, dir, crate::core::ModelScope::Tenant).await?,
);
if let Err(e) = super::permissions::auto_create_permissions(&pool).await {
tracing::warn!(target: "crate::tenancy", slug = %org.slug, error = %e, "auto_create_permissions failed for schema-mode tenant");
}
if let Err(e) = crate::contenttypes::ensure_seeded(&dbpool).await {
tracing::warn!(target: "crate::tenancy", slug = %org.slug, error = %e, "contenttypes::ensure_seeded failed for schema-mode tenant");
}
pool.close().await;
Ok(applied)
}
StorageMode::Database => {
let tenant_pool = pools.pool_for_org(org).await?;
let mut applied = migrate::migrate(tenant_pool.pool(), dir).await?;
let dbpool: crate::sql::Pool = tenant_pool.pool().clone().into();
applied.extend(
apply_system_migrations(&dbpool, dir, crate::core::ModelScope::Tenant).await?,
);
if let Err(e) = super::permissions::auto_create_permissions(tenant_pool.pool()).await {
tracing::warn!(target: "crate::tenancy", slug = %org.slug, error = %e, "auto_create_permissions failed for database-mode tenant");
}
if let Err(e) = crate::contenttypes::ensure_seeded(&dbpool).await {
tracing::warn!(target: "crate::tenancy", slug = %org.slug, error = %e, "contenttypes::ensure_seeded failed for database-mode tenant");
}
Ok(applied)
}
}
}
#[cfg(feature = "postgres")]
async fn build_schema_scoped_pool(
registry_url: &str,
schema: &str,
) -> Result<PgPool, TenancyError> {
let bootstrap = PgPool::connect(registry_url).await?;
let create_sql = format!(
"CREATE SCHEMA IF NOT EXISTS {}",
quote_ident_for_schema(schema)
);
rustango::sql::sqlx::query(&create_sql)
.execute(&bootstrap)
.await?;
bootstrap.close().await;
let schema_owned: Arc<str> = Arc::from(schema);
let pool = PgPoolOptions::new()
.max_connections(2)
.after_connect(move |conn, _meta| {
let schema = Arc::clone(&schema_owned);
Box::pin(async move {
let stmt = format!(
"SET search_path TO {}, public",
quote_ident_for_schema(&schema)
);
rustango::sql::sqlx::query(&stmt).execute(conn).await?;
Ok(())
})
})
.connect(registry_url)
.await?;
Ok(pool)
}
async fn scoped_subset(dir: &Path, scope: MigrationScope) -> Result<ScopedDir, TenancyError> {
let all = rustango::migrate::file::list_dir(dir)?;
if all.iter().all(|m| m.scope == scope) {
return Ok(ScopedDir::Original);
}
let temp = tempdir_under_target()?;
let temp_path = temp.path().to_path_buf();
for mig in &all {
if mig.scope == scope {
let target = temp_path.join(format!("{}.json", mig.name));
rustango::migrate::file::write(&target, mig)?;
}
}
Ok(ScopedDir::Owned(temp))
}
enum ScopedDir {
Original,
Owned(TempDir),
}
struct TempDir(std::path::PathBuf);
impl TempDir {
fn path(&self) -> &Path {
&self.0
}
}
impl Drop for TempDir {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
fn tempdir_under_target() -> Result<TempDir, TenancyError> {
use std::sync::atomic::{AtomicU64, Ordering};
static COUNTER: AtomicU64 = AtomicU64::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let pid = std::process::id();
let mut p = std::env::temp_dir();
p.push(format!("rustango_tenancy_scoped_{pid}_{n}"));
std::fs::create_dir_all(&p).map_err(|e| {
TenancyError::Validation(format!(
"could not create scoped-migration tempdir at {}: {e}",
p.display()
))
})?;
Ok(TempDir(p))
}
#[cfg(feature = "postgres")]
fn quote_ident_for_schema(name: &str) -> String {
let escaped = name.replace('"', "\"\"");
format!("\"{escaped}\"")
}