#![cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use boatramp_core::compute::{
managed_db_spec, ComputeWorkload, ManagedDbEngine, PlacementConstraints,
};
use boatramp_core::deploy::DeployStore;
use boatramp_core::envelope::KeyEnvelope;
use boatramp_core::kv::KvStore;
use boatramp_core::project::{ProjectRef, DEFAULT_PROJECT};
use boatramp_core::sql::{SqlBackend, SqlError};
use boatramp_storage::sql_compute::{
ComputeEndpointResolver, ComputeResolvedSqlBackend, SESSION_KEY_PROJECT, SESSION_KEY_SITE,
};
use boatramp_storage::sql_sqlx::PerTenantSqlResolver;
use boatramp_storage::tenant_provision::{
grant_app_role_ddl, provision_ddl, recover_soft_deprovision_ddl, sanitize_ident,
soft_deprovision_ddl, tenant_db_name, tenant_role_name,
};
use boatramp_storage::ExternalSqlKind;
use crate::config::{ExternalDatabaseConfig, TenantIsolation, TenantScope};
use crate::managed_sql::{DeployEndpointResolver, ManagedSqlCredentials};
use crate::tenant_tombstone::{self, Tombstone};
const DEFAULT_VOLUME_MIB: u32 = 10 * 1024;
pub const DEFAULT_DEPROVISION_GRACE_SECS: u64 = 7 * 24 * 60 * 60;
pub const TOMBSTONE_REAPER_TICK: std::time::Duration = std::time::Duration::from_secs(3600);
fn now_unix_secs() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
fn engine_of(kind: ExternalSqlKind) -> ManagedDbEngine {
match kind {
ExternalSqlKind::Postgres => ManagedDbEngine::Postgres,
ExternalSqlKind::Mysql => ManagedDbEngine::Mysql,
}
}
fn maintenance_database(kind: ExternalSqlKind) -> &'static str {
match kind {
ExternalSqlKind::Postgres => "postgres",
ExternalSqlKind::Mysql => "mysql",
}
}
pub fn tenant_key(scope: TenantScope, project: &str, site: &str) -> (String, bool) {
match scope {
TenantScope::Project => (project.to_string(), project == DEFAULT_PROJECT),
TenantScope::Site => (format!("{project}/{site}"), false),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TenantNames {
pub database: String,
pub role: String,
pub workload: String,
}
pub(crate) fn tenant_names(
isolation: TenantIsolation,
compute: &str,
database: &str,
tenant_ident_raw: &str,
is_default: bool,
) -> TenantNames {
if is_default {
return TenantNames {
database: database.to_string(),
role: database.to_string(),
workload: compute.to_string(),
};
}
let ident = sanitize_ident(tenant_ident_raw);
let (db, workload) = match isolation {
TenantIsolation::Shared => (tenant_db_name(database, &ident), compute.to_string()),
TenantIsolation::Single => (database.to_string(), format!("{compute}-{ident}")),
};
TenantNames {
database: db,
role: tenant_role_name(compute, &ident),
workload,
}
}
pub(crate) fn credential_workload_key(compute: &str, tenant_ident: &str) -> String {
format!("{compute}/{tenant_ident}")
}
pub(crate) fn single_credential_project(project: &str, is_default: bool) -> String {
if is_default {
DEFAULT_PROJECT.to_string()
} else {
project.to_string()
}
}
fn managed_volume_name(workload: &str, is_default: bool) -> String {
if is_default {
"data".to_string()
} else {
workload.to_string()
}
}
fn managed_bindings(
databases: &std::collections::BTreeMap<String, ExternalDatabaseConfig>,
) -> Vec<(&str, ExternalSqlKind, &ExternalDatabaseConfig)> {
databases
.iter()
.filter_map(|(name, db)| {
db.compute.as_deref().filter(|c| !c.is_empty())?;
let kind = ExternalSqlKind::parse(&db.kind)?;
Some((name.as_str(), kind, db))
})
.collect()
}
pub async fn provision_tenant(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
binding: &ExternalDatabaseConfig,
project: &str,
site: &str,
) -> Result<(), String> {
let Some(compute) = binding.compute.as_deref().filter(|c| !c.is_empty()) else {
return Ok(()); };
let Some(kind) = ExternalSqlKind::parse(&binding.kind) else {
return Ok(()); };
let database = binding.database.as_deref().unwrap_or_default();
let user = binding.user.as_deref().unwrap_or_default();
let (tenant_ident_raw, is_default) = tenant_key(binding.tenant_scope, project, site);
let names = tenant_names(
binding.tenant,
compute,
database,
&tenant_ident_raw,
is_default,
);
let creds = ManagedSqlCredentials::new(kv.clone(), envelope.clone());
match binding.tenant {
TenantIsolation::Single => {
provision_single(deploy, &creds, binding, kind, project, &names, is_default).await
}
TenantIsolation::Shared => {
if is_default {
return Ok(());
}
let ident = sanitize_ident(&tenant_ident_raw);
provision_shared(deploy, &creds, kind, compute, user, project, &names, &ident).await
}
}
}
#[allow(clippy::too_many_arguments)]
async fn provision_single(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
binding: &ExternalDatabaseConfig,
kind: ExternalSqlKind,
project: &str,
names: &TenantNames,
is_default: bool,
) -> Result<(), String> {
let cred_project = single_credential_project(project, is_default);
creds
.password(&cred_project, &names.workload)
.await
.map_err(|e| format!("per-tenant credential ({}): {e}", names.workload))?;
let proj = ProjectRef::new(project);
match deploy.get_compute_workload(proj, &names.workload).await {
Ok(Some(_)) => return Ok(()), Ok(None) => {}
Err(e) => return Err(format!("check workload {}: {e}", names.workload)),
}
let mut spec = managed_db_spec(
engine_of(kind),
binding.image.as_deref(),
binding.volume_size_mib.unwrap_or(DEFAULT_VOLUME_MIB),
);
if let Some(grace) = binding.startup_grace_secs {
spec.startup_grace_secs = grace;
}
if let Some(vol) = spec.volumes.first_mut() {
vol.name = managed_volume_name(&names.workload, is_default);
}
let spec_id = deploy
.put_compute_spec(&spec)
.await
.map_err(|e| format!("store spec for {}: {e}", names.workload))?;
let wl = ComputeWorkload {
version: 1,
name: names.workload.clone(),
active: spec_id,
replicas: 1,
placement: PlacementConstraints::default(),
};
deploy
.set_compute_workload(proj, &wl)
.await
.map_err(|e| format!("register workload {}: {e}", names.workload))?;
Ok(())
}
#[allow(clippy::too_many_arguments)]
async fn provision_shared(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
kind: ExternalSqlKind,
compute: &str,
superuser: &str,
project: &str,
names: &TenantNames,
tenant_ident: &str,
) -> Result<(), String> {
let superuser_pw = creds
.password(DEFAULT_PROJECT, compute)
.await
.map_err(|e| format!("superuser credential ({compute}): {e}"))?;
let cred_workload = credential_workload_key(compute, tenant_ident);
let tenant_pw = creds
.password(project, &cred_workload)
.await
.map_err(|e| format!("per-tenant credential ({}): {e}", names.role))?;
let resolver = Arc::new(DeployEndpointResolver::new(deploy.clone(), DEFAULT_PROJECT));
let admin = ComputeResolvedSqlBackend::new(
resolver,
compute,
kind,
maintenance_database(kind),
superuser,
superuser_pw.clone(),
Some(1),
false,
Some(Duration::from_secs(10)),
);
for stmt in provision_ddl(kind, &names.database, &names.role, &tenant_pw) {
if let Err(e) = admin.run_script(&stmt).await {
let is_create_database = stmt.to_ascii_uppercase().contains("CREATE DATABASE");
if is_create_database && is_database_exists_error(&e) {
continue;
}
return Err(format!("provision {}: {e}", names.database));
}
}
let grants = grant_app_role_ddl(kind, &names.role, superuser);
if !grants.is_empty() {
let tenant_resolver =
Arc::new(DeployEndpointResolver::new(deploy.clone(), DEFAULT_PROJECT));
let tenant_admin = ComputeResolvedSqlBackend::new(
tenant_resolver,
compute,
kind,
names.database.clone(),
superuser,
superuser_pw,
Some(1),
false,
Some(Duration::from_secs(10)),
);
for stmt in grants {
tenant_admin
.run_script(&stmt)
.await
.map_err(|e| format!("grant app role on {}: {e}", names.database))?;
}
}
Ok(())
}
fn is_database_exists_error(err: &SqlError) -> bool {
let msg = err.to_string().to_ascii_lowercase();
msg.contains("already exists") && msg.contains("database")
}
pub async fn provision_binding_for_tenants(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
databases: &std::collections::BTreeMap<String, ExternalDatabaseConfig>,
project: &str,
site: &str,
) -> Result<(), String> {
for (_name, _kind, binding) in managed_bindings(databases) {
provision_tenant(deploy, kv, envelope, binding, project, site).await?;
}
Ok(())
}
async fn shared_admin_backend(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
kind: ExternalSqlKind,
compute: &str,
superuser: &str,
) -> Result<ComputeResolvedSqlBackend, String> {
let superuser_pw = creds
.password(DEFAULT_PROJECT, compute)
.await
.map_err(|e| format!("superuser credential ({compute}): {e}"))?;
let resolver = Arc::new(DeployEndpointResolver::new(deploy.clone(), DEFAULT_PROJECT));
Ok(ComputeResolvedSqlBackend::new(
resolver,
compute,
kind,
maintenance_database(kind),
superuser,
superuser_pw,
Some(1),
false,
Some(Duration::from_secs(10)),
))
}
pub async fn deprovision_tenant(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
binding: &ExternalDatabaseConfig,
project: &str,
site: &str,
grace_secs: u64,
) -> Result<(), String> {
let Some(plan) = plan_deprovision(binding, project, site, grace_secs, now_unix_secs()) else {
return Ok(()); };
let creds = ManagedSqlCredentials::new(kv.clone(), envelope.clone());
match plan {
DeprovisionPlan::SingleDrop { workload } => {
deploy
.delete_compute_workload(ProjectRef::new(project), &workload)
.await
.map_err(|e| format!("delete workload {workload}: {e}"))?;
creds
.delete(project, &workload)
.await
.map_err(|e| format!("delete credential {workload}: {e}"))?;
}
DeprovisionPlan::SharedImmediate {
kind,
compute,
superuser,
ddl,
cred_workload,
database,
} => {
let admin = shared_admin_backend(deploy, &creds, kind, &compute, &superuser).await?;
for stmt in ddl {
admin
.run_script(&stmt)
.await
.map_err(|e| format!("deprovision {database}: {e}"))?;
}
creds
.delete(project, &cred_workload)
.await
.map_err(|e| format!("delete credential {cred_workload}: {e}"))?;
}
DeprovisionPlan::SharedSoftPostgres { ddl, tombstone } => {
let admin = shared_admin_backend(
deploy,
&creds,
ExternalSqlKind::Postgres,
&tombstone.compute,
&tombstone.superuser,
)
.await?;
for stmt in ddl {
admin
.run_script(&stmt)
.await
.map_err(|e| format!("soft-deprovision {}: {e}", tombstone.original_db))?;
}
tenant_tombstone::put(kv, &tombstone).await?;
}
}
Ok(())
}
fn plan_deprovision(
binding: &ExternalDatabaseConfig,
project: &str,
site: &str,
grace_secs: u64,
now: u64,
) -> Option<DeprovisionPlan> {
let compute = binding.compute.as_deref().filter(|c| !c.is_empty())?;
let kind = ExternalSqlKind::parse(&binding.kind)?;
let (tenant_ident_raw, is_default) = tenant_key(binding.tenant_scope, project, site);
if is_default {
return None; }
let database = binding.database.as_deref().unwrap_or_default();
let superuser = binding.user.as_deref().unwrap_or_default();
let ident = sanitize_ident(&tenant_ident_raw);
let names = tenant_names(binding.tenant, compute, database, &tenant_ident_raw, false);
Some(match binding.tenant {
TenantIsolation::Single => DeprovisionPlan::SingleDrop {
workload: names.workload,
},
TenantIsolation::Shared if kind == ExternalSqlKind::Postgres && grace_secs > 0 => {
let renamed_db = format!("{}__deleted_{now}", names.database);
let ddl = soft_deprovision_ddl(&names.database, &renamed_db, &names.role);
let tombstone = Tombstone {
version: 1,
project: project.to_string(),
renamed_db,
original_db: names.database.clone(),
role: names.role.clone(),
engine: "postgres".to_string(),
compute: compute.to_string(),
superuser: superuser.to_string(),
cred_workload: credential_workload_key(compute, &ident),
deleted_at: now,
delete_after: now.saturating_add(grace_secs),
};
DeprovisionPlan::SharedSoftPostgres { ddl, tombstone }
}
TenantIsolation::Shared => DeprovisionPlan::SharedImmediate {
kind,
compute: compute.to_string(),
superuser: superuser.to_string(),
ddl: boatramp_storage::tenant_provision::deprovision_ddl(
kind,
&names.database,
&names.role,
),
cred_workload: credential_workload_key(compute, &ident),
database: names.database,
},
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum DeprovisionPlan {
SingleDrop { workload: String },
SharedImmediate {
kind: ExternalSqlKind,
compute: String,
superuser: String,
ddl: Vec<String>,
cred_workload: String,
database: String,
},
SharedSoftPostgres {
ddl: Vec<String>,
tombstone: Tombstone,
},
}
pub async fn recover_tenant(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
project: &str,
renamed_db: &str,
) -> Result<bool, String> {
let Some(ts) = tenant_tombstone::get(kv, project, renamed_db).await? else {
return Ok(false);
};
let creds = ManagedSqlCredentials::new(kv.clone(), envelope.clone());
let admin = shared_admin_backend(
deploy,
&creds,
ExternalSqlKind::Postgres,
&ts.compute,
&ts.superuser,
)
.await?;
for stmt in recover_soft_deprovision_ddl(&ts.renamed_db, &ts.original_db, &ts.role) {
admin
.run_script(&stmt)
.await
.map_err(|e| format!("recover {}: {e}", ts.original_db))?;
}
tenant_tombstone::delete(kv, &ts).await?;
Ok(true)
}
pub struct NodeTenantDeprovisioner {
deploy: DeployStore,
kv: Arc<dyn KvStore>,
envelope: Arc<dyn KeyEnvelope>,
databases: std::collections::BTreeMap<String, ExternalDatabaseConfig>,
grace_secs: u64,
}
impl NodeTenantDeprovisioner {
pub fn new(
deploy: DeployStore,
kv: Arc<dyn KvStore>,
envelope: Arc<dyn KeyEnvelope>,
databases: std::collections::BTreeMap<String, ExternalDatabaseConfig>,
grace_secs: u64,
) -> Self {
Self {
deploy,
kv,
envelope,
databases,
grace_secs,
}
}
async fn deprovision_scope(&self, scope: TenantScope, project: &str, site: &str) {
if project == DEFAULT_PROJECT {
return;
}
for (name, _kind, binding) in managed_bindings(&self.databases) {
if !binding.is_managed_credential() || binding.tenant_scope != scope {
continue;
}
match deprovision_tenant(
&self.deploy,
&self.kv,
&self.envelope,
binding,
project,
site,
self.grace_secs,
)
.await
{
Ok(()) => {
if matches!(scope, TenantScope::Site) {
tracing::info!(
binding = name,
project,
site,
"deprovisioned managed database for deleted site tenant"
);
} else {
tracing::info!(
binding = name,
project,
"deprovisioned managed database for deleted project tenant"
);
}
}
Err(e) => tracing::warn!(
binding = name,
project,
site,
error = %e,
"managed-database deprovision failed (best-effort; delete not blocked)"
),
}
}
}
}
#[async_trait]
impl boatramp_core::sql::TenantDeprovisioner for NodeTenantDeprovisioner {
async fn deprovision_project(&self, project: &str) {
self.deprovision_scope(TenantScope::Project, project, "")
.await;
}
async fn deprovision_site(&self, project: &str, site: &str) {
self.deprovision_scope(TenantScope::Site, project, site)
.await;
}
}
async fn hard_drop_tombstone(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
ts: &Tombstone,
) -> Result<(), String> {
let creds = ManagedSqlCredentials::new(kv.clone(), envelope.clone());
let admin = shared_admin_backend(
deploy,
&creds,
ExternalSqlKind::Postgres,
&ts.compute,
&ts.superuser,
)
.await?;
for stmt in boatramp_storage::tenant_provision::deprovision_ddl(
ExternalSqlKind::Postgres,
&ts.renamed_db,
&ts.role,
) {
admin
.run_script(&stmt)
.await
.map_err(|e| format!("reap {}: {e}", ts.renamed_db))?;
}
creds
.delete(&ts.project, &ts.cred_workload)
.await
.map_err(|e| format!("reap credential {}: {e}", ts.cred_workload))?;
tenant_tombstone::delete(kv, ts).await
}
fn due_tombstones(all: Vec<Tombstone>, now: u64) -> Vec<Tombstone> {
all.into_iter().filter(|t| t.is_due(now)).collect()
}
async fn reap_due(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
now: u64,
) -> usize {
let tombstones = match tenant_tombstone::list(kv).await {
Ok(t) => t,
Err(e) => {
tracing::warn!(error = %e, "tenant tombstone reaper: could not list tombstones");
return 0;
}
};
let mut reaped = 0;
for ts in due_tombstones(tombstones, now) {
match hard_drop_tombstone(deploy, kv, envelope, &ts).await {
Ok(()) => {
reaped += 1;
tracing::info!(
project = %ts.project,
renamed_db = %ts.renamed_db,
"tenant tombstone reaper: hard-dropped a soft-deleted tenant past its grace window"
);
}
Err(e) => tracing::warn!(
project = %ts.project,
renamed_db = %ts.renamed_db,
error = %e,
"tenant tombstone reaper: hard-drop failed (best-effort; retried next sweep)"
),
}
}
reaped
}
pub fn spawn_tenant_tombstone_reaper(
deploy: DeployStore,
kv: Arc<dyn KvStore>,
envelope: Arc<dyn KeyEnvelope>,
is_leader: boatramp_server::CronLeaderGate,
tick: std::time::Duration,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut interval = tokio::time::interval(tick);
interval.tick().await;
loop {
interval.tick().await;
if !is_leader() {
continue;
}
let n = reap_due(&deploy, &kv, &envelope, now_unix_secs()).await;
if n > 0 {
tracing::info!(reaped = n, "tenant tombstone reaper sweep");
}
}
})
}
pub struct NodeTenantSqlResolver {
deploy: DeployStore,
kv: Arc<dyn KvStore>,
envelope: Arc<dyn KeyEnvelope>,
kind: ExternalSqlKind,
compute: String,
database: String,
user: String,
isolation: TenantIsolation,
scope: TenantScope,
rls_session: bool,
pool_max: Option<u32>,
read_only: bool,
connect_timeout: Option<Duration>,
binding: ExternalDatabaseConfig,
}
impl NodeTenantSqlResolver {
pub fn new(
deploy: DeployStore,
kv: Arc<dyn KvStore>,
envelope: Arc<dyn KeyEnvelope>,
binding: &ExternalDatabaseConfig,
) -> Option<Self> {
let compute = binding.compute.as_deref().filter(|c| !c.is_empty())?;
let kind = ExternalSqlKind::parse(&binding.kind)?;
Some(Self {
deploy,
kv,
envelope,
kind,
compute: compute.to_string(),
database: binding.database.clone().unwrap_or_default(),
user: binding.user.clone().unwrap_or_default(),
isolation: binding.tenant,
scope: binding.tenant_scope,
rls_session: binding.rls_session,
pool_max: binding.pool_max,
read_only: binding.read_only,
connect_timeout: binding.connect_timeout_secs.map(Duration::from_secs),
binding: binding.clone(),
})
}
pub fn site_scoped(&self) -> bool {
matches!(self.scope, TenantScope::Site)
}
}
#[async_trait]
impl PerTenantSqlResolver for NodeTenantSqlResolver {
async fn resolve(&self, project: &str, site: &str) -> Result<Arc<dyn SqlBackend>, SqlError> {
provision_tenant(
&self.deploy,
&self.kv,
&self.envelope,
&self.binding,
project,
site,
)
.await
.map_err(SqlError::other)?;
self.build_backend(project, site).await
}
}
impl NodeTenantSqlResolver {
async fn build_backend(
&self,
project: &str,
site: &str,
) -> Result<Arc<dyn SqlBackend>, SqlError> {
let (tenant_ident_raw, is_default) = tenant_key(self.scope, project, site);
let names = tenant_names(
self.isolation,
&self.compute,
&self.database,
&tenant_ident_raw,
is_default,
);
let (cred_project, cred_workload, user) = match self.isolation {
TenantIsolation::Shared => {
if is_default {
(
DEFAULT_PROJECT.to_string(),
self.compute.clone(),
self.user.clone(),
)
} else {
let ident = sanitize_ident(&tenant_ident_raw);
(
project.to_string(),
credential_workload_key(&self.compute, &ident),
names.role.clone(),
)
}
}
TenantIsolation::Single => (
single_credential_project(project, is_default),
names.workload.clone(),
self.user.clone(),
),
};
let password = ManagedSqlCredentials::new(self.kv.clone(), self.envelope.clone())
.password(&cred_project, &cred_workload)
.await
.map_err(SqlError::other)?;
let endpoint_project = match self.isolation {
TenantIsolation::Single if !is_default => project.to_string(),
_ => DEFAULT_PROJECT.to_string(),
};
let resolver: Arc<dyn ComputeEndpointResolver> = Arc::new(DeployEndpointResolver::new(
self.deploy.clone(),
endpoint_project,
));
let mut backend = ComputeResolvedSqlBackend::new(
resolver,
names.workload.clone(),
self.kind,
names.database.clone(),
user,
password,
self.pool_max,
self.read_only,
self.connect_timeout,
);
if self.rls_session {
backend = backend.with_session_context(vec![
(SESSION_KEY_PROJECT, project.to_string()),
(SESSION_KEY_SITE, site.to_string()),
]);
}
Ok(Arc::new(backend))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn tenant_key_project_scope_marks_default() {
assert_eq!(
tenant_key(TenantScope::Project, "default", "blog"),
("default".to_string(), true)
);
assert_eq!(
tenant_key(TenantScope::Project, "acme", "blog"),
("acme".to_string(), false)
);
}
#[test]
fn tenant_key_site_scope_is_qualified_and_never_default() {
assert_eq!(
tenant_key(TenantScope::Site, "default", "blog"),
("default/blog".to_string(), false)
);
assert_eq!(
tenant_key(TenantScope::Site, "acme", "shop"),
("acme/shop".to_string(), false)
);
}
#[test]
fn default_tenant_uses_plain_names_no_hash() {
let n = tenant_names(TenantIsolation::Shared, "pg", "appdb", "default", true);
assert_eq!(n.database, "appdb");
assert_eq!(n.workload, "pg");
assert!(!n.database.contains('_') || n.database == "appdb");
let s = tenant_names(TenantIsolation::Single, "pg", "appdb", "default", true);
assert_eq!(s.database, "appdb");
assert_eq!(s.workload, "pg");
}
#[test]
fn managed_volume_name_isolates_non_default_tenants() {
assert_eq!(managed_volume_name("pg", true), "data");
let a = tenant_names(TenantIsolation::Single, "pg", "appdb", "acme", false).workload;
let b = tenant_names(TenantIsolation::Single, "pg", "appdb", "globex", false).workload;
assert_eq!(
managed_volume_name(&a, false),
a,
"keyed to its own workload"
);
assert_ne!(managed_volume_name(&a, false), "data");
assert_ne!(
managed_volume_name(&a, false),
managed_volume_name(&b, false),
"distinct tenants must not share a volume"
);
}
#[test]
fn shared_tenant_derives_distinct_db_and_role_from_workload_base() {
let ident = sanitize_ident("acme");
let n = tenant_names(TenantIsolation::Shared, "pg", "appdb", "acme", false);
assert_eq!(n.database, tenant_db_name("appdb", &ident));
assert_eq!(n.role, tenant_role_name("pg", &ident));
assert_eq!(n.workload, "pg");
assert_ne!(n.database, "appdb");
}
#[test]
fn shared_two_bindings_same_server_share_one_role_distinct_dbs() {
let a = tenant_names(TenantIsolation::Shared, "pg", "appdb", "acme", false);
let b = tenant_names(TenantIsolation::Shared, "pg", "analytics", "acme", false);
assert_eq!(a.role, b.role, "same (tenant, server) ⇒ one shared role");
assert_ne!(a.database, b.database, "distinct databases per binding");
}
#[test]
fn shared_distinct_tenants_are_isolated() {
let a = tenant_names(TenantIsolation::Shared, "pg", "appdb", "acme", false);
let b = tenant_names(TenantIsolation::Shared, "pg", "appdb", "globex", false);
assert_ne!(a.database, b.database, "cross-tenant database collision!");
assert_ne!(a.role, b.role, "cross-tenant role collision!");
}
#[test]
fn single_tenant_derives_dedicated_workload_plain_db() {
let ident = sanitize_ident("acme");
let n = tenant_names(TenantIsolation::Single, "pg", "appdb", "acme", false);
assert_eq!(n.workload, format!("pg-{ident}"));
assert_eq!(n.database, "appdb");
}
#[test]
fn single_distinct_tenants_get_distinct_workloads() {
let a = tenant_names(TenantIsolation::Single, "pg", "appdb", "acme", false);
let b = tenant_names(TenantIsolation::Single, "pg", "appdb", "globex", false);
assert_ne!(a.workload, b.workload, "cross-tenant workload collision!");
}
#[test]
fn site_grain_isolates_two_sites_of_one_project() {
let (ta, da) = tenant_key(TenantScope::Site, "acme", "blog");
let (tb, db) = tenant_key(TenantScope::Site, "acme", "shop");
assert!(!da && !db);
let na = tenant_names(TenantIsolation::Shared, "pg", "appdb", &ta, false);
let nb = tenant_names(TenantIsolation::Shared, "pg", "appdb", &tb, false);
assert_ne!(na.database, nb.database);
assert_ne!(na.role, nb.role);
}
#[test]
fn credential_key_folds_in_the_tenant() {
let ident = sanitize_ident("acme");
let k = credential_workload_key("pg", &ident);
assert!(k.starts_with("pg/"));
assert_ne!(
k, "pg",
"must never collide with the superuser workload key"
);
assert_ne!(
credential_workload_key("pg", &sanitize_ident("acme")),
credential_workload_key("pg", &sanitize_ident("globex")),
"distinct tenants ⇒ distinct credential keys"
);
}
#[test]
fn maintenance_database_per_engine() {
assert_eq!(maintenance_database(ExternalSqlKind::Postgres), "postgres");
assert_eq!(maintenance_database(ExternalSqlKind::Mysql), "mysql");
}
#[test]
fn database_exists_error_is_recognized() {
assert!(is_database_exists_error(&SqlError::Other(
"database \"appdb_acme\" already exists".into()
)));
assert!(!is_database_exists_error(&SqlError::Other(
"connection refused".into()
)));
assert!(!is_database_exists_error(&SqlError::Other(
"role \"x\" already exists".into()
)));
}
use async_trait::async_trait as _async_trait;
use boatramp_core::envelope::EnvelopeError;
use boatramp_core::kv::MemoryKv;
use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
struct RevEnvelope;
#[_async_trait]
impl KeyEnvelope for RevEnvelope {
async fn wrap(&self, p: &[u8]) -> std::result::Result<Vec<u8>, EnvelopeError> {
Ok(p.iter().rev().copied().collect())
}
async fn unwrap(&self, w: &[u8]) -> std::result::Result<Vec<u8>, EnvelopeError> {
Ok(w.iter().rev().copied().collect())
}
}
struct NullStorage;
#[_async_trait]
impl Storage for NullStorage {
async fn get(&self, _: &str) -> std::result::Result<GetObject, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn get_range(
&self,
_: &str,
_: u64,
_: Option<u64>,
) -> std::result::Result<GetObject, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn put(
&self,
_: &str,
_: ByteStream,
_: PutMeta,
) -> std::result::Result<ObjectMeta, StorageError> {
Err(StorageError::unsupported("null"))
}
async fn head(&self, _: &str) -> std::result::Result<ObjectMeta, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn delete(&self, _: &str) -> std::result::Result<(), StorageError> {
Ok(())
}
async fn list(&self, _: &str) -> std::result::Result<Vec<ObjectMeta>, StorageError> {
Ok(Vec::new())
}
}
fn shared_binding() -> ExternalDatabaseConfig {
ExternalDatabaseConfig {
kind: "postgres".into(),
compute: Some("pg".into()),
database: Some("appdb".into()),
user: Some("super".into()),
tenant: TenantIsolation::Shared,
tenant_scope: TenantScope::Project,
..Default::default()
}
}
fn build_resolver(
binding: &ExternalDatabaseConfig,
) -> (Arc<dyn KvStore>, NodeTenantSqlResolver) {
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let deploy = DeployStore::new(Arc::new(NullStorage), kv.clone());
let envelope: Arc<dyn KeyEnvelope> = Arc::new(RevEnvelope);
let resolver = NodeTenantSqlResolver::new(deploy, kv.clone(), envelope, binding)
.expect("compute-backed managed binding builds a resolver");
(kv, resolver)
}
#[tokio::test]
async fn shared_resolve_seals_isolated_per_tenant_credentials() {
let binding = shared_binding();
let (kv, resolver) = build_resolver(&binding);
let _ = resolver.build_backend("acme", "blog").await.unwrap();
let _ = resolver.build_backend("globex", "shop").await.unwrap();
let acme_ident = sanitize_ident("acme");
let globex_ident = sanitize_ident("globex");
let acme_key = format!(
"managed-sql-cred/acme/{}",
credential_workload_key("pg", &acme_ident)
);
let globex_key = format!(
"managed-sql-cred/globex/{}",
credential_workload_key("pg", &globex_ident)
);
let acme = kv.get(&acme_key).await.unwrap().expect("acme cred sealed");
let globex = kv
.get(&globex_key)
.await
.unwrap()
.expect("globex cred sealed");
assert_ne!(acme, globex, "cross-tenant credential reuse!");
assert!(
kv.get("managed-sql-cred/default/pg")
.await
.unwrap()
.is_none(),
"a derived tenant must not touch the superuser credential"
);
}
#[tokio::test]
async fn shared_default_tenant_uses_plain_superuser_credential() {
let binding = shared_binding();
let (kv, resolver) = build_resolver(&binding);
let _ = resolver.build_backend("default", "blog").await.unwrap();
assert!(kv
.get("managed-sql-cred/default/pg")
.await
.unwrap()
.is_some());
}
#[tokio::test]
async fn deprovision_project_targets_the_tenant_and_skips_default() {
use boatramp_core::sql::TenantDeprovisioner as _;
let binding = ExternalDatabaseConfig {
tenant: TenantIsolation::Single,
tenant_scope: TenantScope::Project,
..shared_binding()
};
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let deploy = DeployStore::new(Arc::new(NullStorage), kv.clone());
let envelope: Arc<dyn KeyEnvelope> = Arc::new(RevEnvelope);
provision_tenant(&deploy, &kv, &envelope, &binding, "default", "")
.await
.unwrap();
provision_tenant(&deploy, &kv, &envelope, &binding, "acme", "")
.await
.unwrap();
let ident = sanitize_ident("acme");
let acme_wl = format!("pg-{ident}"); let acme_cred = format!("managed-sql-cred/acme/pg-{ident}");
let default_cred = "managed-sql-cred/default/pg";
assert!(deploy
.get_compute_workload(ProjectRef::new("acme"), &acme_wl)
.await
.unwrap()
.is_some());
assert!(kv.get(&acme_cred).await.unwrap().is_some());
assert!(kv.get(default_cred).await.unwrap().is_some());
let deprovisioner = NodeTenantDeprovisioner::new(
deploy.clone(),
kv.clone(),
envelope.clone(),
std::iter::once(("pg".to_string(), binding.clone())).collect(),
DEFAULT_DEPROVISION_GRACE_SECS,
);
deprovisioner.deprovision_project("default").await;
assert!(
kv.get(default_cred).await.unwrap().is_some(),
"default-project guard: the single-tenant install must never be dropped"
);
deprovisioner.deprovision_project("acme").await;
assert!(
deploy
.get_compute_workload(ProjectRef::new("acme"), &acme_wl)
.await
.unwrap()
.is_none(),
"acme's dedicated Single workload must be gone"
);
assert!(
kv.get(&acme_cred).await.unwrap().is_none(),
"acme's sealed credential must be gone"
);
assert!(
kv.get(default_cred).await.unwrap().is_some(),
"acme delete must not touch the default install's credential"
);
}
#[tokio::test]
async fn single_startup_grace_override_flows_into_the_stored_spec() {
let binding = ExternalDatabaseConfig {
tenant: TenantIsolation::Single,
tenant_scope: TenantScope::Project,
startup_grace_secs: Some(77),
..shared_binding()
};
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let deploy = DeployStore::new(Arc::new(NullStorage), kv.clone());
let envelope: Arc<dyn KeyEnvelope> = Arc::new(RevEnvelope);
provision_tenant(&deploy, &kv, &envelope, &binding, "default", "")
.await
.unwrap();
let wl = deploy
.get_compute_workload(ProjectRef::DEFAULT, "pg")
.await
.unwrap()
.expect("default Single workload registered");
let stored = deploy
.get_compute_spec(&wl.active)
.await
.unwrap()
.expect("active spec stored");
assert_eq!(
stored.startup_grace_secs, 77,
"the operator grace overrides the engine default in the stored spec"
);
let mut expected = managed_db_spec(ManagedDbEngine::Postgres, None, DEFAULT_VOLUME_MIB);
expected.startup_grace_secs = binding.startup_grace_secs.unwrap();
assert_eq!(
stored.id(),
expected.id(),
"both managed-registration paths build the identical content-addressed spec"
);
}
#[tokio::test]
async fn single_resolve_keys_default_and_derived_distinctly() {
let binding = ExternalDatabaseConfig {
tenant: TenantIsolation::Single,
..shared_binding()
};
let (kv, resolver) = build_resolver(&binding);
let _ = resolver.build_backend("default", "blog").await.unwrap();
let _ = resolver.build_backend("acme", "blog").await.unwrap();
assert!(kv
.get("managed-sql-cred/default/pg")
.await
.unwrap()
.is_some());
let ident = sanitize_ident("acme");
assert!(kv
.get(&format!("managed-sql-cred/acme/pg-{ident}"))
.await
.unwrap()
.is_some());
}
fn mysql_shared_binding() -> ExternalDatabaseConfig {
ExternalDatabaseConfig {
kind: "mysql".into(),
..shared_binding()
}
}
fn single_binding() -> ExternalDatabaseConfig {
ExternalDatabaseConfig {
tenant: TenantIsolation::Single,
..shared_binding()
}
}
#[test]
fn plan_shared_postgres_soft_deletes_renames_not_drops() {
let binding = shared_binding();
let now = 1_700_000_000;
let grace = DEFAULT_DEPROVISION_GRACE_SECS;
let plan = plan_deprovision(&binding, "acme", "", grace, now).expect("a plan");
let DeprovisionPlan::SharedSoftPostgres { ddl, tombstone } = plan else {
panic!("Shared+Postgres+grace>0 must be a soft delete, got {plan:?}");
};
let joined = ddl.join("\n");
assert!(joined.contains("ALTER DATABASE"), "must RENAME:\n{joined}");
assert!(joined.contains("RENAME TO"), "must RENAME aside:\n{joined}");
assert!(
joined.contains("NOLOGIN"),
"must disable the role:\n{joined}"
);
assert!(
joined.contains("pg_terminate_backend"),
"must evict sessions"
);
assert!(
!joined.to_ascii_uppercase().contains("DROP DATABASE"),
"soft delete must NOT drop the database:\n{joined}"
);
assert!(!joined.to_ascii_uppercase().contains("DROP ROLE"));
assert_eq!(tombstone.delete_after, now + grace);
assert_eq!(tombstone.deleted_at, now);
assert_eq!(tombstone.project, "acme");
assert_eq!(tombstone.engine, "postgres");
assert_eq!(tombstone.compute, "pg");
assert_eq!(tombstone.superuser, "super");
assert!(tombstone.renamed_db.ends_with(&format!("__deleted_{now}")));
assert!(tombstone.renamed_db.starts_with(&tombstone.original_db));
assert!(joined.contains(&tombstone.renamed_db));
}
#[test]
fn soft_delete_frees_original_name_without_aliasing() {
let binding = shared_binding();
let now = 1_700_000_000;
let plan = plan_deprovision(&binding, "acme", "", 60, now).expect("a plan");
let DeprovisionPlan::SharedSoftPostgres { ddl, tombstone } = plan else {
panic!("expected a soft delete");
};
assert_ne!(tombstone.original_db, tombstone.renamed_db);
let joined = ddl.join("\n");
assert!(joined.contains(&format!("RENAME TO \"{}\"", tombstone.renamed_db)));
}
#[test]
fn plan_shared_mysql_hard_drops_immediately() {
let binding = mysql_shared_binding();
let plan = plan_deprovision(&binding, "acme", "", DEFAULT_DEPROVISION_GRACE_SECS, 0)
.expect("a plan");
let DeprovisionPlan::SharedImmediate { ddl, kind, .. } = plan else {
panic!("Shared+MySQL must be an immediate drop, got {plan:?}");
};
assert_eq!(kind, ExternalSqlKind::Mysql);
let joined = ddl.join("\n");
assert!(
joined.contains("DROP DATABASE IF EXISTS"),
"MySQL must hard-drop:\n{joined}"
);
assert!(joined.contains("DROP USER IF EXISTS"));
assert!(!joined.contains("RENAME TO"), "MySQL must NOT soft-rename");
}
#[test]
fn plan_single_drops_the_workload_immediately() {
let binding = single_binding();
let plan = plan_deprovision(&binding, "acme", "", DEFAULT_DEPROVISION_GRACE_SECS, 0)
.expect("a plan");
let ident = sanitize_ident("acme");
assert_eq!(
plan,
DeprovisionPlan::SingleDrop {
workload: format!("pg-{ident}")
}
);
}
#[test]
fn plan_grace_zero_takes_the_immediate_path_for_shared_postgres() {
let binding = shared_binding();
let plan = plan_deprovision(&binding, "acme", "", 0, 1_700_000_000).expect("a plan");
let DeprovisionPlan::SharedImmediate { ddl, kind, .. } = plan else {
panic!("grace=0 must be an immediate drop, got {plan:?}");
};
assert_eq!(kind, ExternalSqlKind::Postgres);
let joined = ddl.join("\n");
assert!(
joined.contains("DROP DATABASE IF EXISTS"),
"grace=0 must hard-drop:\n{joined}"
);
assert!(
!joined.contains("RENAME TO"),
"grace=0 must NOT soft-rename"
);
}
#[test]
fn plan_skips_default_tenant_and_bring_your_own() {
assert!(plan_deprovision(&shared_binding(), "default", "", 60, 0).is_none());
let byo = ExternalDatabaseConfig {
kind: "postgres".into(),
compute: None,
url_env: "PG_URL".into(),
..Default::default()
};
assert!(plan_deprovision(&byo, "acme", "", 60, 0).is_none());
}
#[test]
fn reaper_selects_only_due_tombstones() {
let mk = |renamed: &str, delete_after: u64| Tombstone {
version: 1,
project: "acme".into(),
renamed_db: renamed.into(),
original_db: "appdb_acme".into(),
role: "appdb_acme_role".into(),
engine: "postgres".into(),
compute: "pg".into(),
superuser: "super".into(),
cred_workload: "pg/x".into(),
deleted_at: 0,
delete_after,
};
let now = 1_000;
let past = mk("db__deleted_1", now - 1); let exact = mk("db__deleted_2", now); let future = mk("db__deleted_3", now + 1);
let due = due_tombstones(vec![past.clone(), exact.clone(), future.clone()], now);
assert!(due.contains(&past), "an elapsed tombstone is due");
assert!(due.contains(&exact), "delete_after == now is due");
assert!(
!due.contains(&future),
"a tombstone still inside its grace window must NOT be reaped"
);
assert_eq!(due.len(), 2);
}
#[tokio::test]
async fn recover_absent_tombstone_is_a_noop() {
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let deploy = DeployStore::new(Arc::new(NullStorage), kv.clone());
let envelope: Arc<dyn KeyEnvelope> = Arc::new(RevEnvelope);
let recovered = recover_tenant(&deploy, &kv, &envelope, "acme", "nope__deleted_1")
.await
.unwrap();
assert!(!recovered, "recovering an absent tombstone is a no-op");
let ddl = recover_soft_deprovision_ddl("appdb__deleted_1", "appdb", "appdb_role");
assert!(ddl[0].contains("RENAME TO \"appdb\""), "renames back");
assert!(ddl[1].contains("LOGIN") && !ddl[1].contains("NOLOGIN"));
}
}