#![cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use std::sync::Arc;
use async_trait::async_trait;
use boatramp_core::deploy::DeployStore;
use boatramp_core::envelope::KeyEnvelope;
use boatramp_core::kv::KvStore;
use boatramp_core::sql::{
RepairCheck, RepairError, RepairMode, RepairReport, RepairStatus, SqlBackend, SqlError,
SqlValue,
};
use boatramp_storage::tenant_provision::{
grant_app_role_ddl, provision_ddl, quote_ident, sanitize_ident,
};
use boatramp_storage::ExternalSqlKind;
use boatramp_core::project::ProjectRef;
use crate::config::{ExternalDatabaseConfig, TenantIsolation, TenantScope};
use crate::managed_sql::{DeployEndpointResolver, ManagedSqlCredentials};
use crate::tenant_sql::{
owner_credential_workload_key, shared_admin_backend_for_db, single_credential_project,
tenant_key, tenant_names,
};
const LEDGER_SCHEMA: &str = "boatramp_migrations";
const LEDGER_TABLE: &str = "schema_migrations";
pub struct NodeTenantRepair {
databases: std::collections::BTreeMap<String, ExternalDatabaseConfig>,
deploy: DeployStore,
kv: Arc<dyn KvStore>,
envelope: Option<Arc<dyn KeyEnvelope>>,
}
impl NodeTenantRepair {
pub fn new(
databases: std::collections::BTreeMap<String, ExternalDatabaseConfig>,
deploy: DeployStore,
kv: Arc<dyn KvStore>,
envelope: Option<Arc<dyn KeyEnvelope>>,
) -> Self {
Self {
databases,
deploy,
kv,
envelope,
}
}
}
#[async_trait]
impl boatramp_core::sql::TenantRepair for NodeTenantRepair {
async fn repair(
&self,
project: &str,
db: &str,
mode: RepairMode,
) -> Result<RepairReport, RepairError> {
boatramp_core::project::validate_resource_name("database", db)
.map_err(|e| RepairError::Other(e.to_string()))?;
let binding = self.databases.get(db).ok_or(RepairError::NotConfigured)?;
let envelope = self.envelope.clone().ok_or_else(|| {
RepairError::Other(format!(
"managed database {db:?} needs a [secrets] envelope to repair its sealed credentials"
))
})?;
let creds = ManagedSqlCredentials::new(self.kv.clone(), envelope);
repair_tenant(&self.deploy, &creds, binding, db, project, mode).await
}
}
struct Derived {
kind: ExternalSqlKind,
compute: String,
superuser: String,
ident: String,
database: String,
runtime_role: String,
owner_role: String,
project: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum RepairModel {
SharedPostgres,
DedicatedPostgres,
Mysql,
Libsql,
External,
}
fn classify_model(binding: &ExternalDatabaseConfig) -> RepairModel {
let compute_backed = binding.compute.as_deref().is_some_and(|c| !c.is_empty());
match ExternalSqlKind::parse(&binding.kind) {
Some(ExternalSqlKind::Postgres) if compute_backed => match binding.tenant {
TenantIsolation::Shared => RepairModel::SharedPostgres,
TenantIsolation::Single => RepairModel::DedicatedPostgres,
},
Some(ExternalSqlKind::Postgres) => RepairModel::External,
Some(ExternalSqlKind::Mysql) => RepairModel::Mysql,
None => {
if is_libsql_kind(&binding.kind) {
RepairModel::Libsql
} else {
RepairModel::External
}
}
}
}
pub async fn repair_tenant(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
binding: &ExternalDatabaseConfig,
db_binding_name: &str,
project: &str,
mode: RepairMode,
) -> Result<RepairReport, RepairError> {
let backend_class = classify_backend(binding);
let model = classify_model(binding);
tracing::info!(
target: "boatramp::repair",
project = %project,
binding = %db_binding_name,
mode = %mode.as_str(),
backend = %backend_class,
model = ?model,
"repair: classified backend model"
);
match model {
RepairModel::SharedPostgres => {
repair_shared_postgres(deploy, creds, binding, db_binding_name, project, mode).await
}
RepairModel::DedicatedPostgres => {
Ok(
repair_dedicated_postgres(deploy, creds, binding, project, mode, &backend_class)
.await,
)
}
RepairModel::Mysql => {
Ok(repair_mysql(deploy, creds, binding, project, mode, &backend_class).await)
}
RepairModel::Libsql => Ok(repair_libsql(binding, project, mode, &backend_class).await),
RepairModel::External => Ok(repair_external(binding, project, mode, &backend_class).await),
}
}
async fn repair_shared_postgres(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
binding: &ExternalDatabaseConfig,
db_binding_name: &str,
project: &str,
mode: RepairMode,
) -> Result<RepairReport, RepairError> {
let backend_class = classify_backend(binding);
let compute = binding
.compute
.as_deref()
.filter(|c| !c.is_empty())
.unwrap_or_default();
if matches!(binding.tenant_scope, TenantScope::Site) {
let mut report = base_report(binding, project, &backend_class, mode);
report.tenant = format!("{}@{compute} (site-scoped)", binding_db(binding));
report.checks.push(skip_check(
"topology",
"site-scoped managed database; operator repair is project-level and cannot target a \
specific site's database",
));
return Ok(report);
}
let (tenant_ident_raw, is_default) = tenant_key(binding.tenant_scope, project, "");
if is_default {
let mut report = base_report(binding, project, &backend_class, mode);
report.tenant = format!("{}@{compute} (default tenant)", binding_db(binding));
report.checks.push(skip_check(
"topology",
"reserved default tenant (single-tenant install); no per-tenant owner model to \
reconcile",
));
return Ok(report);
}
let database = binding_db(binding);
let names = tenant_names(
binding.tenant,
compute,
&database,
&tenant_ident_raw,
is_default,
);
let derived = Derived {
kind: ExternalSqlKind::Postgres,
compute: compute.to_string(),
superuser: binding.user.as_deref().unwrap_or_default().to_string(),
ident: sanitize_ident(&tenant_ident_raw),
database: names.database.clone(),
runtime_role: names.role.clone(),
owner_role: names.owner_role.clone(),
project: project.to_string(),
};
tracing::info!(
target: "boatramp::repair",
project = %derived.project,
binding = %db_binding_name,
mode = %mode.as_str(),
derived_db = %derived.database,
derived_runtime_role = %derived.runtime_role,
derived_owner_role = %derived.owner_role,
maintenance_user = %derived.superuser,
compute = %derived.compute,
"repair: resolved derived tenant identity"
);
let mut report = base_report(binding, project, &backend_class, mode);
report.tenant = derived.database.clone();
run_shared_postgres_checks(deploy, creds, &derived, mode, &mut report).await;
Ok(report)
}
async fn run_shared_postgres_checks(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
d: &Derived,
mode: RepairMode,
report: &mut RepairReport,
) {
let maint = match build_backend(deploy, creds, d, maintenance_db(d.kind)).await {
Ok(b) => b,
Err(e) => {
report.checks.push(error_check(
"connectivity",
format!("could not build the maintenance connection: {e}"),
));
return;
}
};
match probe_live_or_soft_deleted(&maint, d).await {
Ok(LiveState::Live) => {}
Ok(LiveState::SoftDeletedOnly) => {
report.checks.push(skip_check(
"soft-delete",
format!(
"only a soft-deleted sibling of {:?} exists (recover it first); repair changed \
nothing",
d.database
),
));
return;
}
Ok(LiveState::Absent) => {
report.checks.push(skip_check(
"database",
format!(
"no database {:?} (nor a soft-deleted sibling) exists — this tenant was never \
provisioned; use the provisioning path, not repair",
d.database
),
));
return;
}
Err(e) => {
report.checks.push(error_check(
"soft-delete",
format!("could not probe the database's live/soft-deleted state: {e}"),
));
return;
}
}
let maint_is_superuser = match probe_current_user_superuser(&maint).await {
Ok(v) => v,
Err(e) => {
report.checks.push(error_check(
"superuser-precondition",
format!("could not determine whether the maintenance identity is a superuser: {e}"),
));
false
}
};
let owner_ready = check_owner_role_exists(creds, d, &maint, mode, report).await;
if owner_ready {
check_owner_role_attrs(&maint, d, mode, report).await;
} else {
report.checks.push(skip_check(
"owner-role-attrs",
"owner role is not (yet) present (see owner-role-exists)",
));
}
if owner_ready {
check_db_owner(&maint, d, mode, report).await;
} else {
report.checks.push(skip_check(
"db-owner",
"owner role is not (yet) present, so `ALTER DATABASE … OWNER TO <owner>` is not emitted",
));
}
if owner_ready {
check_object_ownership(deploy, creds, d, maint_is_superuser, mode, report).await;
} else {
report.checks.push(skip_check(
"object-ownership",
"owner role is not (yet) present, so object re-ownership is not emitted",
));
}
check_connect_grants(&maint, d, mode, report).await;
check_runtime_dml_grants(deploy, creds, d, mode, report).await;
if owner_ready {
check_ledger(deploy, creds, d, mode, report).await;
} else {
report.checks.push(skip_check(
"ledger",
"owner role is not (yet) present, so ledger scaffolding + re-ownership is deferred",
));
}
check_owner_credential_sealed(creds, d, mode, report).await;
check_connectivity(deploy, creds, d, mode, report).await;
}
async fn repair_dedicated_postgres(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
binding: &ExternalDatabaseConfig,
project: &str,
mode: RepairMode,
backend_class: &str,
) -> RepairReport {
let mut report = base_report(binding, project, backend_class, mode);
let compute = binding
.compute
.as_deref()
.filter(|c| !c.is_empty())
.unwrap_or_default();
let database = binding_db(binding);
let user = binding.user.as_deref().unwrap_or_default().to_string();
let kind = ExternalSqlKind::Postgres;
let (tenant_ident_raw, is_default) = tenant_key(binding.tenant_scope, project, "");
if matches!(binding.tenant_scope, TenantScope::Site) {
report.tenant = format!("{database}@{compute} (site-scoped)");
report.checks.push(skip_check(
"topology",
"site-scoped managed database; operator repair is project-level and cannot target a \
specific site's database",
));
return report;
}
let names = tenant_names(
binding.tenant,
compute,
&database,
&tenant_ident_raw,
is_default,
);
let cred_project = single_credential_project(project, is_default);
let cred_workload = names.workload.clone();
report.tenant = format!("{}@{}", names.database, names.workload);
let boundary = "dedicated Postgres — the container is the isolation boundary (no shared \
owner/runtime role split, no RLS)";
for check in [
"owner-role",
"object-ownership",
"connect-grants",
"runtime-grants",
] {
report.checks.push(skip_check(check, boundary));
}
match deploy
.get_compute_workload(ProjectRef::new(project), &names.workload)
.await
{
Ok(Some(_)) => report.checks.push(ok_check(
"compute-workload",
format!("dedicated workload {:?} is registered", names.workload),
)),
Ok(None) => report.checks.push(RepairCheck {
check: "compute-workload".to_string(),
status: RepairStatus::Drift,
detail: format!(
"dedicated workload {:?} is not registered; provision the tenant (repair does not \
spawn a server — it reconciles an EXISTING one)",
names.workload
),
ddl: None,
}),
Err(e) => report.checks.push(error_check(
"compute-workload",
format!(
"could not check the compute workload {:?}: {e}",
names.workload
),
)),
}
check_credential_sealed(
creds,
&cred_project,
&cred_workload,
mode,
"credential-sealed",
&mut report,
)
.await;
let endpoint_project = if is_default {
boatramp_core::project::DEFAULT_PROJECT.to_string()
} else {
project.to_string()
};
let backend = build_compute_backend(
deploy,
creds,
kind,
&names.workload,
&names.database,
&user,
&cred_project,
&cred_workload,
&endpoint_project,
mode,
)
.await;
match backend {
Ok(backend) => {
check_pg_ledger_exists(&backend, kind, mode, &mut report).await;
match backend.run_query("SELECT 1;").await {
Ok(_) => report.checks.push(ok_check(
"connectivity",
format!("connected to {:?} as {:?}", names.database, user),
)),
Err(e) => report.checks.push(error_check(
"connectivity",
format!(
"could not connect to {:?} as {:?}: {e}",
names.database, user
),
)),
}
}
Err(BackendBuildError::Unsealed) => {
report.checks.push(skip_check(
"ledger",
"skipped — the server credential is not yet sealed; a dry-run does not seal it \
(the ledger + connectivity are verified after apply)",
));
report.checks.push(skip_check(
"connectivity",
"the server credential is not yet sealed — connectivity verified after apply",
));
}
Err(e) => {
report.checks.push(skip_check(
"ledger",
format!("skipped — could not build the tenant connection: {e}"),
));
report.checks.push(error_check(
"connectivity",
format!("could not build the tenant connection: {e}"),
));
}
}
report
}
async fn repair_mysql(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
binding: &ExternalDatabaseConfig,
project: &str,
mode: RepairMode,
backend_class: &str,
) -> RepairReport {
let mut report = base_report(binding, project, backend_class, mode);
let kind = ExternalSqlKind::Mysql;
let database = binding_db(binding);
let compute_backed = binding.compute.as_deref().is_some_and(|c| !c.is_empty());
let runtime = build_runtime_backend(deploy, creds, binding, kind, project, mode).await;
if compute_backed {
if matches!(binding.tenant_scope, TenantScope::Site) {
report.tenant = format!("{database}@mysql (site-scoped)");
report.checks.push(skip_check(
"topology",
"site-scoped managed database; operator repair is project-level and cannot target \
a specific site's database",
));
return report;
}
let (tenant_ident_raw, is_default) = tenant_key(binding.tenant_scope, project, "");
let compute = binding.compute.as_deref().unwrap_or_default();
let names = tenant_names(
binding.tenant,
compute,
&database,
&tenant_ident_raw,
is_default,
);
report.tenant = format!("{}@{}", names.database, names.workload);
} else {
report.tenant = format!("{database}@{} (external mysql)", binding.url_env);
}
check_mysql_ddl_identity(binding, &database, compute_backed, &mut report);
let runtime = match runtime {
Ok(b) => b,
Err(BackendBuildError::Unsealed) => {
for check in ["runtime-user", "database", "ledger"] {
report.checks.push(skip_check(
check,
"skipped — the runtime credential is not yet sealed; a dry-run does not seal \
it (verified after apply)",
));
}
report.checks.push(skip_check(
"connectivity",
"the runtime credential is not yet sealed — connectivity verified after apply",
));
return report;
}
Err(e) => {
for check in ["runtime-user", "database", "ledger"] {
report.checks.push(skip_check(
check,
format!("skipped — could not build the runtime connection: {e}"),
));
}
report.checks.push(error_check(
"connectivity",
format!("could not build the runtime connection: {e}"),
));
return report;
}
};
check_mysql_runtime_grant(&runtime, &database, &mut report).await;
check_mysql_database_exists(&runtime, &database, &mut report).await;
check_mysql_ledger_exists(&runtime, &mut report).await;
match runtime.run_query("SELECT 1;").await {
Ok(_) => report.checks.push(ok_check(
"connectivity",
format!("connected to the MySQL runtime for {database:?}"),
)),
Err(e) => report.checks.push(error_check(
"connectivity",
format!("could not connect to the MySQL runtime for {database:?}: {e}"),
)),
}
report
}
async fn repair_libsql(
binding: &ExternalDatabaseConfig,
project: &str,
mode: RepairMode,
backend_class: &str,
) -> RepairReport {
let mut report = base_report(binding, project, backend_class, mode);
let boundary = "libsql/SQLite — the file is the trust boundary; there are no roles/RLS to \
reconcile";
for check in [
"owner-role",
"object-ownership",
"connect-grants",
"runtime-grants",
] {
report.checks.push(skip_check(check, boundary));
}
libsql_file_reconcile(binding, mode, report).await
}
#[cfg(feature = "migrate")]
async fn libsql_file_reconcile(
binding: &ExternalDatabaseConfig,
mode: RepairMode,
mut report: RepairReport,
) -> RepairReport {
let Some(path) = binding
.path
.as_deref()
.filter(|p| !p.as_os_str().is_empty())
else {
report.tenant = format!("{}@libsql (remote sqld)", binding_db(binding));
report.checks.push(skip_check(
"db-file",
"remote-sqld libsql binding (url_env, no single-node `path`); the namespace lives on \
the sqld server — not a boatramp-owned local file to reconcile",
));
report.checks.push(skip_check(
"ledger",
"skipped — a remote-sqld libsql binding has no local file to probe",
));
report.checks.push(skip_check(
"connectivity",
"skipped — a remote-sqld libsql binding is not a single-node file target",
));
return report;
};
report.tenant = format!("{}@{}", binding_db(binding), path.display());
let exists = path.is_file();
if exists {
report.checks.push(ok_check(
"db-file",
format!("db file {} exists", path.display()),
));
} else {
report.checks.push(RepairCheck {
check: "db-file".to_string(),
status: RepairStatus::Drift,
detail: format!(
"db file {} does not exist; it is created lazily on first use / migrate — repair \
does not create it (a dry-run must be side-effect-free, and creating an empty db \
is provisioning, not repair)",
path.display()
),
ddl: None,
});
}
if !exists && matches!(mode, RepairMode::DryRun) {
report.checks.push(skip_check(
"ledger",
"skipped — the db file does not exist and a dry-run must not create it",
));
report.checks.push(skip_check(
"connectivity",
"skipped — the db file does not exist (dry-run does not create it)",
));
return report;
}
match boatramp_storage::LibsqlSql::open_local(path).await {
Ok(sql) => {
check_libsql_ledger(&sql, mode, &mut report).await;
use boatramp_core::sql::SqlBackend;
match sql.run_query("SELECT 1;").await {
Ok(_) => report.checks.push(ok_check(
"connectivity",
format!("opened {} and ran SELECT 1", path.display()),
)),
Err(e) => report.checks.push(error_check(
"connectivity",
format!("could not query {}: {e}", path.display()),
)),
}
}
Err(e) => {
report.checks.push(error_check(
"ledger",
format!("could not open {} to probe the ledger: {e}", path.display()),
));
report.checks.push(error_check(
"connectivity",
format!("could not open {}: {e}", path.display()),
));
}
}
report
}
#[cfg(not(feature = "migrate"))]
async fn libsql_file_reconcile(
binding: &ExternalDatabaseConfig,
_mode: RepairMode,
mut report: RepairReport,
) -> RepairReport {
report.tenant = format!("{}@libsql", binding_db(binding));
report.checks.push(skip_check(
"topology",
"libsql repair needs the `migrate` feature (the embedded libsql substrate); this build has \
no libsql runner to open the file",
));
report
}
async fn repair_external(
binding: &ExternalDatabaseConfig,
project: &str,
mode: RepairMode,
backend_class: &str,
) -> RepairReport {
let _ = project;
let mut report = base_report(binding, project, backend_class, mode);
report.tenant = format!("{}@{} (external)", binding_db(binding), binding.url_env);
let operator_owned =
"operator-owned binding (bring-your-own url_env); boatramp owns no roles / \
ownership / grants here — nothing to reconcile";
for check in [
"owner-role",
"object-ownership",
"connect-grants",
"runtime-grants",
] {
report.checks.push(skip_check(check, operator_owned));
}
let kind = ExternalSqlKind::parse(&binding.kind);
match kind {
Some(kind) => match build_external_backend(binding, kind) {
Ok(backend) => {
check_pg_or_mysql_ledger_exists(&backend, kind, mode, &mut report).await;
match backend.run_query("SELECT 1;").await {
Ok(_) => report.checks.push(ok_check(
"connectivity",
format!(
"connected to the external {} via {}",
binding.kind, binding.url_env
),
)),
Err(e) => report.checks.push(error_check(
"connectivity",
format!("could not connect to the external database: {e}"),
)),
}
}
Err(e) => {
report.checks.push(skip_check(
"ledger",
format!("skipped — could not build the external connection: {e}"),
));
report.checks.push(error_check(
"connectivity",
format!("could not build the external connection: {e}"),
));
}
},
None => {
report.checks.push(skip_check(
"ledger",
format!(
"engine {:?} is not a sqlx engine this build reconciles the ledger for; only \
connectivity/topology is inspectable",
binding.kind
),
));
report.checks.push(skip_check(
"connectivity",
format!("engine {:?} is not a recognized sqlx engine", binding.kind),
));
}
}
report
}
async fn check_owner_role_exists(
creds: &ManagedSqlCredentials,
d: &Derived,
maint: &Arc<dyn SqlBackend>,
mode: RepairMode,
report: &mut RepairReport,
) -> bool {
let present = match probe_role_exists(maint, &d.owner_role).await {
Ok(v) => v,
Err(e) => {
report.checks.push(error_check(
"owner-role-exists",
format!("could not probe owner role {:?}: {e}", d.owner_role),
));
return false;
}
};
if present {
report.checks.push(RepairCheck {
check: "owner-role-exists".to_string(),
status: RepairStatus::Ok,
detail: format!("owner role {:?} exists", d.owner_role),
ddl: None,
});
return true;
}
match mode {
RepairMode::DryRun => {
report.checks.push(RepairCheck {
check: "owner-role-exists".to_string(),
status: RepairStatus::Drift,
detail: format!(
"owner role {:?} is missing; apply would CREATE it (NOSUPERUSER NOCREATEDB \
NOCREATEROLE NOBYPASSRLS NOREPLICATION) and seal its credential",
d.owner_role
),
ddl: Some(owner_role_ddl(d, "<sealed-owner-password>").join("\n")),
});
false
}
RepairMode::Apply => {
let owner_cred_workload = owner_credential_workload_key(&d.compute, &d.ident);
let owner_pw = match creds.password(&d.project, &owner_cred_workload).await {
Ok(pw) => pw,
Err(e) => {
report.checks.push(error_check(
"owner-role-exists",
format!("could not seal the owner credential: {e}"),
));
return false;
}
};
let stmts = owner_role_ddl(d, &owner_pw);
for stmt in &stmts {
if let Err(e) = maint.run_script(stmt).await {
report.checks.push(error_check(
"owner-role-exists",
format!("creating owner role {:?} failed: {e}", d.owner_role),
));
return false;
}
}
audit_stmts("owner-role-exists", d, &redact_pw(&stmts));
match probe_role_exists(maint, &d.owner_role).await {
Ok(true) => {
report.checks.push(RepairCheck {
check: "owner-role-exists".to_string(),
status: RepairStatus::Repaired,
detail: format!(
"created owner role {:?} (NOSUPERUSER …) and sealed its credential",
d.owner_role
),
ddl: Some(redact_pw(&stmts).join("\n")),
});
true
}
Ok(false) => {
report.checks.push(error_check(
"owner-role-exists",
format!(
"owner role {:?} still absent after CREATE ROLE (unexpected)",
d.owner_role
),
));
false
}
Err(e) => {
report.checks.push(error_check(
"owner-role-exists",
format!("re-probe after creating owner role failed: {e}"),
));
false
}
}
}
}
}
fn owner_role_ddl(d: &Derived, owner_pw: &str) -> Vec<String> {
provision_ddl(
d.kind,
&d.database,
&d.runtime_role,
"unused-runtime-password",
&d.owner_role,
owner_pw,
)
.into_iter()
.filter(|s| {
let upper = s.to_ascii_uppercase();
s.contains(&d.owner_role)
&& upper.contains("ROLE")
&& !upper.contains("CREATE DATABASE")
&& !upper.contains("GRANT ")
&& !upper.contains("REVOKE ")
})
.collect()
}
async fn check_owner_role_attrs(
maint: &Arc<dyn SqlBackend>,
d: &Derived,
mode: RepairMode,
report: &mut RepairReport,
) {
let owner_lit = pg_literal(&d.owner_role);
let rows = match maint
.run_query(&format!(
"SELECT rolsuper, rolcreatedb, rolcreaterole, rolbypassrls, rolreplication \
FROM pg_roles WHERE rolname = {owner_lit};"
))
.await
{
Ok(r) => r,
Err(e) => {
report.checks.push(error_check(
"owner-role-attrs",
format!("could not probe owner-role attributes: {e}"),
));
return;
}
};
let Some(row) = rows.rows.first() else {
report.checks.push(error_check(
"owner-role-attrs",
"owner role vanished between checks (unexpected)",
));
return;
};
let unsafe_attr = row.iter().take(5).any(sql_bool);
let alter = format!(
"ALTER ROLE {} WITH NOSUPERUSER NOCREATEDB NOCREATEROLE NOBYPASSRLS NOREPLICATION;",
quote_ident(d.kind, &d.owner_role)
);
if !unsafe_attr {
report.checks.push(RepairCheck {
check: "owner-role-attrs".to_string(),
status: RepairStatus::Ok,
detail: format!(
"owner role {:?} has the safe (NOSUPERUSER …) attributes",
d.owner_role
),
ddl: None,
});
return;
}
converge_ddl(
maint,
d,
"owner-role-attrs",
&[alter],
mode,
report,
format!(
"owner role {:?} has an unsafe attribute (SUPERUSER/CREATEDB/CREATEROLE/BYPASSRLS/\
REPLICATION); would reset to NOSUPERUSER …",
d.owner_role
),
format!("reset owner role {:?} to the safe attributes", d.owner_role),
)
.await;
}
async fn check_db_owner(
maint: &Arc<dyn SqlBackend>,
d: &Derived,
mode: RepairMode,
report: &mut RepairReport,
) {
let db_lit = pg_literal(&d.database);
let rows = match maint
.run_query(&format!(
"SELECT pg_catalog.pg_get_userbyid(datdba) FROM pg_database WHERE datname = {db_lit};"
))
.await
{
Ok(r) => r,
Err(e) => {
report.checks.push(error_check(
"db-owner",
format!("could not probe the database owner: {e}"),
));
return;
}
};
let actual = rows.rows.first().and_then(|r| r.first()).map(sql_text);
let Some(actual) = actual else {
report.checks.push(error_check(
"db-owner",
"database owner probe returned no row",
));
return;
};
let alter = format!(
"ALTER DATABASE {} OWNER TO {};",
quote_ident(d.kind, &d.database),
quote_ident(d.kind, &d.owner_role)
);
if actual == d.owner_role {
report.checks.push(RepairCheck {
check: "db-owner".to_string(),
status: RepairStatus::Ok,
detail: format!("database {:?} is owned by {:?}", d.database, d.owner_role),
ddl: None,
});
return;
}
converge_ddl(
maint,
d,
"db-owner",
&[alter],
mode,
report,
format!(
"database {:?} is owned by {actual:?}, not the owner role {:?}; would re-own it",
d.database, d.owner_role
),
format!("re-owned database {:?} to {:?}", d.database, d.owner_role),
)
.await;
}
async fn check_object_ownership(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
d: &Derived,
maint_is_superuser: bool,
mode: RepairMode,
report: &mut RepairReport,
) {
let tenant = match build_backend(deploy, creds, d, &d.database).await {
Ok(b) => b,
Err(e) => {
report.checks.push(error_check(
"object-ownership",
format!("could not connect to the tenant database: {e}"),
));
return;
}
};
match probe_current_database(&tenant).await {
Ok(cur) if cur == d.database => {}
Ok(cur) => {
report.checks.push(error_check(
"object-ownership",
format!(
"refusing object re-ownership: connected to {cur:?}, expected the derived \
tenant database {:?}",
d.database
),
));
return;
}
Err(e) => {
report.checks.push(error_check(
"object-ownership",
format!("could not confirm the current database before re-ownership: {e}"),
));
return;
}
}
let owner_lit = pg_literal(&d.owner_role);
let non_owner = match tenant
.run_query(&format!(
"SELECT n.nspname, c.relname, pg_catalog.pg_get_userbyid(c.relowner) AS owner, \
CASE WHEN c.relkind = 'S' THEN 'sequence' ELSE 'table' END AS objkind, \
'' AS args \
FROM pg_catalog.pg_class c \
JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace \
WHERE n.nspname NOT IN ('pg_catalog','information_schema') \
AND n.nspname NOT LIKE 'pg\\_%' ESCAPE '\\' \
AND n.nspname <> 'boatramp_migrations' \
AND c.relkind IN ('r','p','S','v','m') \
AND pg_catalog.pg_get_userbyid(c.relowner) <> {owner_lit} \
UNION ALL \
SELECT n.nspname, p.proname, pg_catalog.pg_get_userbyid(p.proowner) AS owner, \
'function' AS objkind, \
pg_catalog.pg_get_function_identity_arguments(p.oid) AS args \
FROM pg_catalog.pg_proc p \
JOIN pg_catalog.pg_namespace n ON n.oid = p.pronamespace \
WHERE n.nspname NOT IN ('pg_catalog','information_schema') \
AND n.nspname NOT LIKE 'pg\\_%' ESCAPE '\\' \
AND n.nspname <> 'boatramp_migrations' \
AND pg_catalog.pg_get_userbyid(p.proowner) <> {owner_lit};"
))
.await
{
Ok(r) => r,
Err(e) => {
report.checks.push(error_check(
"object-ownership",
format!("could not enumerate object owners: {e}"),
));
return;
}
};
if non_owner.rows.is_empty() {
report.checks.push(RepairCheck {
check: "object-ownership".to_string(),
status: RepairStatus::Ok,
detail: format!(
"all objects in {:?} are owned by {:?}",
d.database, d.owner_role
),
ddl: None,
});
return;
}
let mut ddl: Vec<String> = Vec::new();
let mut runtime_owned = 0usize;
let mut other_owned = 0usize;
let mut skipped_notes: Vec<String> = Vec::new();
for row in &non_owner.rows {
let schema = row.first().map(sql_text).unwrap_or_default();
let name = row.get(1).map(sql_text).unwrap_or_default();
let owner = row.get(2).map(sql_text).unwrap_or_default();
let objkind = row.get(3).map(sql_text).unwrap_or_default();
let args = row.get(4).map(sql_text).unwrap_or_default();
if owner == d.runtime_role {
runtime_owned += 1;
} else {
match targeted_reown_ddl(d, &schema, &name, &objkind, &args) {
Ok(stmt) => {
other_owned += 1;
ddl.push(stmt);
}
Err(reason) => {
tracing::warn!(
target: "boatramp::repair",
derived_db = %d.database,
object = %format!("{schema}.{name}"),
objkind = %objkind,
"object-ownership: skipping targeted re-own (defensive guard): {reason}"
);
skipped_notes.push(format!("{schema:?}.{name:?}: {reason}"));
}
}
}
}
if runtime_owned > 0 {
ddl.insert(
0,
format!(
"REASSIGN OWNED BY {} TO {};",
quote_ident(d.kind, &d.runtime_role),
quote_ident(d.kind, &d.owner_role)
),
);
}
if other_owned > 0 && !maint_is_superuser {
report.checks.push(error_check(
"object-ownership",
format!(
"{other_owned} object(s) are owned by a superuser-loaded identity and re-owning \
them requires a superuser maintenance connection (the current identity is not \
one); refusing — never a GRANT <role> TO <maint> workaround"
),
));
return;
}
if runtime_owned > 0 && !maint_is_superuser {
report.checks.push(error_check(
"object-ownership",
"re-owning runtime-owned objects requires a superuser maintenance connection (the \
current identity is not one); refusing — never a GRANT <role> TO <maint> workaround",
));
return;
}
let skip_note = if skipped_notes.is_empty() {
String::new()
} else {
format!(
" (skipped {} object(s) whose function args failed the defensive guard: {})",
skipped_notes.len(),
skipped_notes.join("; ")
)
};
if ddl.is_empty() {
report.checks.push(skip_check(
"object-ownership",
format!(
"no object re-ownership emitted; every non-owner object was fenced off by the \
defensive guard{skip_note}"
),
));
return;
}
let drift_detail = format!(
"{runtime_owned} runtime-owned + {other_owned} superuser-owned object(s) are not owned by \
{:?}; would re-own them to it{skip_note}",
d.owner_role
);
let repaired_detail = format!(
"re-owned {runtime_owned} runtime-owned + {other_owned} superuser-owned object(s) to \
{:?}{skip_note}",
d.owner_role
);
converge_ddl_on(
&tenant,
d,
"object-ownership",
&ddl,
mode,
report,
drift_detail,
repaired_detail,
)
.await;
}
fn targeted_reown_ddl(
d: &Derived,
schema: &str,
name: &str,
objkind: &str,
args: &str,
) -> Result<String, String> {
let qname = format!(
"{}.{}",
quote_ident(d.kind, schema),
quote_ident(d.kind, name)
);
let owner_id = quote_ident(d.kind, &d.owner_role);
match objkind {
"sequence" => Ok(format!(
"ALTER SEQUENCE IF EXISTS {qname} OWNER TO {owner_id};"
)),
"function" => {
if args.contains(';') {
return Err(format!(
"function identity-args {args:?} contain a ';' (statement separator); refusing \
to emit ALTER FUNCTION — reconcile this object manually"
));
}
if args.matches('\'').count() % 2 != 0 || args.matches('"').count() % 2 != 0 {
return Err(format!(
"function identity-args {args:?} have an unbalanced quote; refusing to emit \
ALTER FUNCTION — reconcile this object manually"
));
}
Ok(format!(
"ALTER FUNCTION {qname}({args}) OWNER TO {owner_id};"
))
}
_ => Ok(format!(
"ALTER TABLE IF EXISTS {qname} OWNER TO {owner_id};"
)),
}
}
async fn check_connect_grants(
maint: &Arc<dyn SqlBackend>,
d: &Derived,
mode: RepairMode,
report: &mut RepairReport,
) {
let db_lit = pg_literal(&d.database);
let public_has = match probe_has_connect(maint, &db_lit, "'public'").await {
Ok(v) => v,
Err(e) => return report_probe_err(report, "connect-grants", e),
};
let owner_has = match probe_has_connect(maint, &db_lit, &pg_literal(&d.owner_role)).await {
Ok(v) => v,
Err(e) => return report_probe_err(report, "connect-grants", e),
};
let runtime_has = match probe_has_connect(maint, &db_lit, &pg_literal(&d.runtime_role)).await {
Ok(v) => v,
Err(e) => return report_probe_err(report, "connect-grants", e),
};
if !public_has && owner_has && runtime_has {
report.checks.push(RepairCheck {
check: "connect-grants".to_string(),
status: RepairStatus::Ok,
detail: format!(
"CONNECT on {:?} is revoked from PUBLIC and granted to owner + runtime",
d.database
),
ddl: None,
});
return;
}
let db_id = quote_ident(d.kind, &d.database);
let owner_id = quote_ident(d.kind, &d.owner_role);
let runtime_id = quote_ident(d.kind, &d.runtime_role);
let ddl = vec![
format!("GRANT CONNECT ON DATABASE {db_id} TO {owner_id};"),
format!("GRANT CONNECT ON DATABASE {db_id} TO {runtime_id};"),
format!("REVOKE CONNECT ON DATABASE {db_id} FROM PUBLIC;"),
];
converge_ddl(
maint,
d,
"connect-grants",
&ddl,
mode,
report,
format!(
"CONNECT lockdown on {:?} has drifted (public_has_connect={public_has}, \
owner_granted={owner_has}, runtime_granted={runtime_has}); would re-issue grants \
then revoke from PUBLIC",
d.database
),
format!("re-issued the CONNECT lockdown on {:?}", d.database),
)
.await;
}
async fn probe_has_connect(
maint: &Arc<dyn SqlBackend>,
db_lit: &str,
grantee_lit: &str,
) -> Result<bool, SqlError> {
let rows = maint
.run_query(&format!(
"SELECT has_database_privilege({grantee_lit}, {db_lit}, 'CONNECT');"
))
.await?;
Ok(rows
.rows
.first()
.and_then(|r| r.first())
.map(sql_bool)
.unwrap_or(false))
}
fn runtime_dml_grants_ok(has_usage: bool, tables_without_select: i64) -> bool {
has_usage && tables_without_select == 0
}
async fn check_runtime_dml_grants(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
d: &Derived,
mode: RepairMode,
report: &mut RepairReport,
) {
let tenant = match build_backend(deploy, creds, d, &d.database).await {
Ok(b) => b,
Err(e) => {
report.checks.push(error_check(
"runtime-grants",
format!("could not connect to the tenant database: {e}"),
));
return;
}
};
let runtime_lit = pg_literal(&d.runtime_role);
let has_usage = match tenant
.run_query(&format!(
"SELECT has_schema_privilege({runtime_lit}, 'public', 'USAGE');"
))
.await
{
Ok(rows) => rows
.rows
.first()
.and_then(|r| r.first())
.map(sql_bool)
.unwrap_or(false),
Err(e) => {
report.checks.push(error_check(
"runtime-grants",
format!("could not probe the runtime role's schema privileges: {e}"),
));
return;
}
};
let tables_without_select = match tenant
.run_query(&format!(
"SELECT count(*)::bigint FROM pg_catalog.pg_tables \
WHERE schemaname = 'public' \
AND NOT has_table_privilege({runtime_lit}, \
format('%I.%I', schemaname, tablename), 'SELECT');"
))
.await
{
Ok(rows) => rows
.rows
.first()
.and_then(|r| r.first())
.map(sql_i64)
.unwrap_or(1),
Err(e) => {
report.checks.push(error_check(
"runtime-grants",
format!("could not probe the runtime role's table privileges: {e}"),
));
return;
}
};
let ddl = grant_app_role_ddl(d.kind, &d.runtime_role, &d.owner_role);
if runtime_dml_grants_ok(has_usage, tables_without_select) {
report.checks.push(RepairCheck {
check: "runtime-grants".to_string(),
status: RepairStatus::Ok,
detail: format!(
"runtime role {:?} has USAGE on public + SELECT on every public table \
(owner-keyed default privileges assumed present)",
d.runtime_role
),
ddl: None,
});
return;
}
let drift_reason = if !has_usage {
format!("runtime role {:?} lacks USAGE on public", d.runtime_role)
} else {
format!(
"runtime role {:?} lacks SELECT on {tables_without_select} public table(s) (e.g. after \
a re-ownership REASSIGN stripped its implicit owner privileges)",
d.runtime_role
)
};
converge_ddl_on(
&tenant,
d,
"runtime-grants",
&ddl,
mode,
report,
format!("{drift_reason}; would (re-)grant DML + set owner-keyed default privileges"),
format!(
"granted runtime role {:?} app DML + owner-keyed default privileges",
d.runtime_role
),
)
.await;
}
async fn check_ledger(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
d: &Derived,
mode: RepairMode,
report: &mut RepairReport,
) {
let tenant = match build_backend(deploy, creds, d, &d.database).await {
Ok(b) => b,
Err(e) => {
report.checks.push(error_check(
"ledger",
format!("could not connect to the tenant database: {e}"),
));
return;
}
};
let schema_lit = pg_literal(LEDGER_SCHEMA);
let table_lit = pg_literal(LEDGER_TABLE);
let schema_owner = match tenant
.run_query(&format!(
"SELECT pg_catalog.pg_get_userbyid(nspowner) FROM pg_catalog.pg_namespace \
WHERE nspname = {schema_lit};"
))
.await
{
Ok(rows) => rows.rows.first().and_then(|r| r.first()).map(sql_text),
Err(e) => {
report.checks.push(error_check(
"ledger",
format!("could not probe the ledger schema: {e}"),
));
return;
}
};
let table_owner = match tenant
.run_query(&format!(
"SELECT tableowner FROM pg_catalog.pg_tables \
WHERE schemaname = {schema_lit} AND tablename = {table_lit};"
))
.await
{
Ok(rows) => rows.rows.first().and_then(|r| r.first()).map(sql_text),
Err(e) => {
report.checks.push(error_check(
"ledger",
format!("could not probe the ledger table: {e}"),
));
return;
}
};
let schema_id = quote_ident(d.kind, LEDGER_SCHEMA);
let ledger_id = format!(
"{}.{}",
quote_ident(d.kind, LEDGER_SCHEMA),
quote_ident(d.kind, LEDGER_TABLE)
);
let owner_id = quote_ident(d.kind, &d.owner_role);
let schema_ok = schema_owner.as_deref() == Some(d.owner_role.as_str());
let table_ok = table_owner.as_deref() == Some(d.owner_role.as_str());
if schema_ok && table_ok {
report.checks.push(RepairCheck {
check: "ledger".to_string(),
status: RepairStatus::Ok,
detail: format!(
"ledger schema + table exist and are owned by {:?}",
d.owner_role
),
ddl: None,
});
return;
}
let ddl = vec![
format!("CREATE SCHEMA IF NOT EXISTS {schema_id};"),
format!(
"CREATE TABLE IF NOT EXISTS {ledger_id} (id text PRIMARY KEY, ordinal integer NOT NULL, \
content_hash text NOT NULL, kind text NOT NULL, applied_at timestamptz NOT NULL \
DEFAULT now(), applied_by text);"
),
format!("ALTER SCHEMA {schema_id} OWNER TO {owner_id};"),
format!("ALTER TABLE {ledger_id} OWNER TO {owner_id};"),
];
converge_ddl_on(
&tenant,
d,
"ledger",
&ddl,
mode,
report,
format!(
"ledger drift (schema_owner={:?}, table_owner={:?}); would scaffold + re-own to {:?}",
schema_owner, table_owner, d.owner_role
),
format!("scaffolded + re-owned the ledger to {:?}", d.owner_role),
)
.await;
}
async fn check_owner_credential_sealed(
creds: &ManagedSqlCredentials,
d: &Derived,
mode: RepairMode,
report: &mut RepairReport,
) {
let owner_cred_workload = owner_credential_workload_key(&d.compute, &d.ident);
let sealed = match creds.is_sealed(&d.project, &owner_cred_workload).await {
Ok(v) => v,
Err(e) => {
report.checks.push(error_check(
"owner-credential-sealed",
format!("could not probe the owner credential: {e}"),
));
return;
}
};
if sealed {
report.checks.push(RepairCheck {
check: "owner-credential-sealed".to_string(),
status: RepairStatus::Ok,
detail: "owner credential is sealed in the control-plane KV".to_string(),
ddl: None,
});
return;
}
match mode {
RepairMode::DryRun => report.checks.push(RepairCheck {
check: "owner-credential-sealed".to_string(),
status: RepairStatus::Drift,
detail: "owner credential is not sealed; apply would generate + seal it \
(a control-plane KV write, not SQL)"
.to_string(),
ddl: None,
}),
RepairMode::Apply => {
match creds.password(&d.project, &owner_cred_workload).await {
Ok(_) => report.checks.push(RepairCheck {
check: "owner-credential-sealed".to_string(),
status: RepairStatus::Repaired,
detail: "generated + sealed the owner credential (control-plane KV write)"
.to_string(),
ddl: None,
}),
Err(e) => report.checks.push(error_check(
"owner-credential-sealed",
format!("could not seal the owner credential: {e}"),
)),
}
}
}
}
async fn check_connectivity(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
d: &Derived,
mode: RepairMode,
report: &mut RepairReport,
) {
let owner_cred_workload = owner_credential_workload_key(&d.compute, &d.ident);
let owner = probe_role_can_connect(
deploy,
creds,
d,
&d.owner_role,
&d.project,
&owner_cred_workload,
mode,
)
.await;
let runtime_cred_workload = crate::tenant_sql::credential_workload_key(&d.compute, &d.ident);
let runtime = probe_role_can_connect(
deploy,
creds,
d,
&d.runtime_role,
&d.project,
&runtime_cred_workload,
mode,
)
.await;
let mut skipped = Vec::new();
let mut problems = Vec::new();
for (label, role, res) in [
("owner", &d.owner_role, owner),
("runtime", &d.runtime_role, runtime),
] {
match res {
ConnectProbe::Connected => {}
ConnectProbe::Unsealed => skipped.push(format!(
"{label} role {role:?} credential not yet sealed — connectivity verified after apply"
)),
ConnectProbe::Refused => {
problems.push(format!("{label} role {role:?} could not connect"))
}
ConnectProbe::Failed(e) => problems.push(format!("{label}-connect probe failed: {e}")),
}
}
if !problems.is_empty() {
report
.checks
.push(error_check("connectivity", problems.join("; ")));
} else if !skipped.is_empty() {
report
.checks
.push(skip_check("connectivity", skipped.join("; ")));
} else {
report.checks.push(RepairCheck {
check: "connectivity".to_string(),
status: RepairStatus::Ok,
detail: format!(
"both the owner ({:?}) and runtime ({:?}) roles connect to {:?}",
d.owner_role, d.runtime_role, d.database
),
ddl: None,
});
}
}
enum ConnectProbe {
Connected,
Refused,
Unsealed,
Failed(String),
}
async fn probe_role_can_connect(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
d: &Derived,
role: &str,
cred_project: &str,
cred_workload: &str,
mode: RepairMode,
) -> ConnectProbe {
use boatramp_storage::sql_compute::ComputeResolvedSqlBackend;
let password = match mode {
RepairMode::DryRun => match creds.get_sealed_password(cred_project, cred_workload).await {
Ok(Some(pw)) => pw,
Ok(None) => return ConnectProbe::Unsealed,
Err(e) => return ConnectProbe::Failed(e),
},
RepairMode::Apply => match creds.password(cred_project, cred_workload).await {
Ok(pw) => pw,
Err(_) => return ConnectProbe::Refused,
},
};
let resolver = Arc::new(crate::managed_sql::DeployEndpointResolver::new(
deploy.clone(),
boatramp_core::project::DEFAULT_PROJECT,
));
let backend = ComputeResolvedSqlBackend::new(
resolver,
&d.compute,
d.kind,
d.database.clone(),
role,
password,
Some(1),
true, Some(std::time::Duration::from_secs(10)),
);
match backend.run_query("SELECT 1;").await {
Ok(_) => ConnectProbe::Connected,
Err(SqlError::Unavailable(m)) => ConnectProbe::Failed(m),
Err(e) => {
let msg = e.to_string().to_ascii_lowercase();
if msg.contains("password") || msg.contains("authentication") || msg.contains("denied")
{
ConnectProbe::Refused
} else {
ConnectProbe::Failed(e.to_string())
}
}
}
}
#[allow(clippy::too_many_arguments)]
async fn converge_ddl(
backend: &Arc<dyn SqlBackend>,
d: &Derived,
check: &str,
ddl: &[String],
mode: RepairMode,
report: &mut RepairReport,
drift_detail: String,
repaired_detail: String,
) {
converge_ddl_on(
backend,
d,
check,
ddl,
mode,
report,
drift_detail,
repaired_detail,
)
.await;
}
#[allow(clippy::too_many_arguments)]
async fn converge_ddl_on(
backend: &Arc<dyn SqlBackend>,
d: &Derived,
check: &str,
ddl: &[String],
mode: RepairMode,
report: &mut RepairReport,
drift_detail: String,
repaired_detail: String,
) {
let joined = ddl.join("\n");
match mode {
RepairMode::DryRun => report.checks.push(RepairCheck {
check: check.to_string(),
status: RepairStatus::Drift,
detail: drift_detail,
ddl: Some(joined),
}),
RepairMode::Apply => {
for stmt in ddl {
if let Err(e) = backend.run_script(stmt).await {
report.checks.push(RepairCheck {
check: check.to_string(),
status: RepairStatus::Error,
detail: format!("converge failed at `{stmt}`: {e}"),
ddl: Some(joined),
});
return;
}
}
audit_stmts(check, d, ddl);
report.checks.push(RepairCheck {
check: check.to_string(),
status: RepairStatus::Repaired,
detail: repaired_detail,
ddl: Some(joined),
});
}
}
}
async fn build_backend(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
d: &Derived,
database: &str,
) -> Result<Arc<dyn SqlBackend>, String> {
let b = shared_admin_backend_for_db(deploy, creds, d.kind, &d.compute, &d.superuser, database)
.await?;
Ok(Arc::new(b) as Arc<dyn SqlBackend>)
}
enum BackendBuildError {
Unsealed,
Other(String),
}
impl std::fmt::Display for BackendBuildError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Unsealed => write!(
f,
"managed credential not yet sealed — connectivity verified after apply"
),
Self::Other(m) => write!(f, "{m}"),
}
}
}
async fn resolve_diagnostic_password(
creds: &ManagedSqlCredentials,
cred_project: &str,
cred_workload: &str,
mode: RepairMode,
) -> Result<String, BackendBuildError> {
match mode {
RepairMode::DryRun => match creds.get_sealed_password(cred_project, cred_workload).await {
Ok(Some(pw)) => Ok(pw),
Ok(None) => Err(BackendBuildError::Unsealed),
Err(e) => Err(BackendBuildError::Other(e)),
},
RepairMode::Apply => creds
.password(cred_project, cred_workload)
.await
.map_err(BackendBuildError::Other),
}
}
#[allow(clippy::too_many_arguments)]
async fn build_compute_backend(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
kind: ExternalSqlKind,
workload: &str,
database: &str,
user: &str,
cred_project: &str,
cred_workload: &str,
endpoint_project: &str,
mode: RepairMode,
) -> Result<Arc<dyn SqlBackend>, BackendBuildError> {
use boatramp_storage::sql_compute::ComputeResolvedSqlBackend;
let password = resolve_diagnostic_password(creds, cred_project, cred_workload, mode).await?;
let resolver = Arc::new(DeployEndpointResolver::new(
deploy.clone(),
endpoint_project.to_string(),
));
Ok(Arc::new(ComputeResolvedSqlBackend::new(
resolver,
workload,
kind,
database.to_string(),
user,
password,
Some(1),
true, Some(std::time::Duration::from_secs(10)),
)) as Arc<dyn SqlBackend>)
}
async fn build_runtime_backend(
deploy: &DeployStore,
creds: &ManagedSqlCredentials,
binding: &ExternalDatabaseConfig,
kind: ExternalSqlKind,
project: &str,
mode: RepairMode,
) -> Result<Arc<dyn SqlBackend>, BackendBuildError> {
let compute_backed = binding.compute.as_deref().is_some_and(|c| !c.is_empty());
if compute_backed {
let compute = binding.compute.as_deref().unwrap_or_default();
let database = binding_db(binding);
let user = binding.user.as_deref().unwrap_or_default();
let (tenant_ident_raw, is_default) = tenant_key(binding.tenant_scope, project, "");
let names = tenant_names(
binding.tenant,
compute,
&database,
&tenant_ident_raw,
is_default,
);
let (cred_project, cred_workload) = match binding.tenant {
TenantIsolation::Single => (
single_credential_project(project, is_default),
names.workload.clone(),
),
TenantIsolation::Shared => (
boatramp_core::project::DEFAULT_PROJECT.to_string(),
compute.to_string(),
),
};
let endpoint_project = match binding.tenant {
TenantIsolation::Single if !is_default => project.to_string(),
_ => boatramp_core::project::DEFAULT_PROJECT.to_string(),
};
build_compute_backend(
deploy,
creds,
kind,
&names.workload,
&names.database,
user,
&cred_project,
&cred_workload,
&endpoint_project,
mode,
)
.await
} else {
build_external_backend(binding, kind).map_err(|e| BackendBuildError::Other(e.to_string()))
}
}
fn build_external_backend(
binding: &ExternalDatabaseConfig,
kind: ExternalSqlKind,
) -> Result<Arc<dyn SqlBackend>, SqlError> {
use boatramp_storage::sql_sqlx::{connect, ExternalSqlOptions};
if binding.url_env.is_empty() {
return Err(SqlError::other(
"external binding has no `url_env` set".to_string(),
));
}
let url = std::env::var(&binding.url_env)
.map_err(|_| SqlError::other(format!("env var {} (url) is unset", binding.url_env)))?;
let timeout = binding
.connect_timeout_secs
.map(std::time::Duration::from_secs);
let opts = ExternalSqlOptions::new(url)
.with_max_connections(Some(1))
.read_only(true)
.with_connect_timeout(timeout);
connect(kind, &opts)
}
async fn check_credential_sealed(
creds: &ManagedSqlCredentials,
cred_project: &str,
cred_workload: &str,
mode: RepairMode,
check: &str,
report: &mut RepairReport,
) {
let sealed = match creds.is_sealed(cred_project, cred_workload).await {
Ok(v) => v,
Err(e) => {
report.checks.push(error_check(
check,
format!("could not probe the sealed credential: {e}"),
));
return;
}
};
if sealed {
report.checks.push(ok_check(
check,
"the sealed server credential is present in the control-plane KV",
));
return;
}
match mode {
RepairMode::DryRun => report.checks.push(RepairCheck {
check: check.to_string(),
status: RepairStatus::Drift,
detail: "the sealed server credential is not present; apply would generate + seal it \
(a control-plane KV write, not SQL)"
.to_string(),
ddl: None,
}),
RepairMode::Apply => match creds.password(cred_project, cred_workload).await {
Ok(_) => report.checks.push(RepairCheck {
check: check.to_string(),
status: RepairStatus::Repaired,
detail: "generated + sealed the server credential (control-plane KV write)"
.to_string(),
ddl: None,
}),
Err(e) => report.checks.push(error_check(
check,
format!("could not seal the server credential: {e}"),
)),
},
}
}
async fn check_pg_ledger_exists(
backend: &Arc<dyn SqlBackend>,
kind: ExternalSqlKind,
mode: RepairMode,
report: &mut RepairReport,
) {
let schema_lit = pg_literal(LEDGER_SCHEMA);
let table_lit = pg_literal(LEDGER_TABLE);
let present = match backend
.run_query(&format!(
"SELECT 1 FROM pg_catalog.pg_tables WHERE schemaname = {schema_lit} \
AND tablename = {table_lit};"
))
.await
{
Ok(rows) => !rows.rows.is_empty(),
Err(e) => {
report.checks.push(error_check(
"ledger",
format!("could not probe the migrate ledger: {e}"),
));
return;
}
};
if present {
report.checks.push(ok_check(
"ledger",
"the migrate ledger (boatramp_migrations.schema_migrations) exists",
));
return;
}
let schema_id = quote_ident(kind, LEDGER_SCHEMA);
let ledger_id = format!(
"{}.{}",
quote_ident(kind, LEDGER_SCHEMA),
quote_ident(kind, LEDGER_TABLE)
);
let ddl = vec![
format!("CREATE SCHEMA IF NOT EXISTS {schema_id};"),
format!(
"CREATE TABLE IF NOT EXISTS {ledger_id} (id text PRIMARY KEY, ordinal integer NOT \
NULL, content_hash text NOT NULL, kind text NOT NULL, applied_at timestamptz NOT \
NULL DEFAULT now(), applied_by text);"
),
];
converge_ledger(backend, &ddl, mode, report).await;
}
async fn check_mysql_ledger_exists(runtime: &Arc<dyn SqlBackend>, report: &mut RepairReport) {
let db_lit = pg_literal(LEDGER_SCHEMA);
match runtime
.run_query(&format!(
"SELECT 1 FROM information_schema.schemata WHERE schema_name = {db_lit};"
))
.await
{
Ok(rows) if !rows.rows.is_empty() => report.checks.push(ok_check(
"ledger",
"the migrate ledger database (boatramp_migrations) exists",
)),
Ok(_) => report.checks.push(RepairCheck {
check: "ledger".to_string(),
status: RepairStatus::Drift,
detail: "the migrate ledger database (boatramp_migrations) does not exist; it is \
created by the DISTINCT DDL identity on first migrate (repair reconciles the \
runtime side; the ledger is created by the migrate path as the DDL identity)"
.to_string(),
ddl: None,
}),
Err(e) => report.checks.push(error_check(
"ledger",
format!("could not probe the migrate ledger database: {e}"),
)),
}
}
async fn check_pg_or_mysql_ledger_exists(
backend: &Arc<dyn SqlBackend>,
kind: ExternalSqlKind,
mode: RepairMode,
report: &mut RepairReport,
) {
match kind {
ExternalSqlKind::Postgres => check_pg_ledger_exists(backend, kind, mode, report).await,
ExternalSqlKind::Mysql => check_mysql_ledger_exists(backend, report).await,
}
}
async fn converge_ledger(
backend: &Arc<dyn SqlBackend>,
ddl: &[String],
mode: RepairMode,
report: &mut RepairReport,
) {
let joined = ddl.join("\n");
match mode {
RepairMode::DryRun => report.checks.push(RepairCheck {
check: "ledger".to_string(),
status: RepairStatus::Drift,
detail:
"the migrate ledger is absent; apply would scaffold it (CREATE … IF NOT EXISTS)"
.to_string(),
ddl: Some(joined),
}),
RepairMode::Apply => {
for stmt in ddl {
if let Err(e) = backend.run_script(stmt).await {
report.checks.push(RepairCheck {
check: "ledger".to_string(),
status: RepairStatus::Error,
detail: format!("scaffolding the ledger failed at `{stmt}`: {e}"),
ddl: Some(joined),
});
return;
}
}
report.checks.push(RepairCheck {
check: "ledger".to_string(),
status: RepairStatus::Repaired,
detail: "scaffolded the migrate ledger (CREATE … IF NOT EXISTS)".to_string(),
ddl: Some(joined),
});
}
}
}
fn check_mysql_ddl_identity(
binding: &ExternalDatabaseConfig,
db: &str,
compute_backed: bool,
report: &mut RepairReport,
) {
if compute_backed && binding.url_env.is_empty() {
report.checks.push(error_check(
"ddl-identity",
format!(
"database {db:?}: compute-backed managed MySQL has no derivable distinct DDL \
identity — migrate is refused fail-closed (boatramp cannot yet auto-mint a \
least-privilege DDL grant); use an external MySQL binding with a distinct \
`migration_url_env`, or await the auto-minted DDL-grant follow-up"
),
));
return;
}
let Some(migration_var) = binding
.migration_url_env
.as_deref()
.filter(|v| !v.is_empty())
else {
report.checks.push(RepairCheck {
check: "ddl-identity".to_string(),
status: RepairStatus::Drift,
detail: format!(
"database {db:?}: MySQL migrations require a DISTINCT DDL identity — set \
`migration_url_env` to an admin/DDL login that is NOT the runtime `user`/`url_env` \
(MySQL has no owner/runtime role split; running DDL as the runtime user is refused)"
),
ddl: None,
});
return;
};
let ddl_url = match std::env::var(migration_var) {
Ok(u) => u,
Err(_) => {
report.checks.push(error_check(
"ddl-identity",
format!(
"env var {migration_var} (the MySQL DDL/migration url for {db:?}) is unset — \
the distinct DDL identity is configured but not reachable"
),
));
return;
}
};
if !binding.url_env.is_empty() {
if let Ok(runtime_url) = std::env::var(&binding.url_env) {
if runtime_url == ddl_url {
report.checks.push(error_check(
"ddl-identity",
format!(
"database {db:?}: `migration_url_env` resolves to the SAME connection as \
the runtime `url_env` — the DDL identity must be DISTINCT (refused)"
),
));
return;
}
#[cfg(feature = "sql-mysql")]
{
let ddl_user = boatramp_storage::sql_sqlx::mysql_dsn_username(&ddl_url);
let runtime_user = boatramp_storage::sql_sqlx::mysql_dsn_username(&runtime_url);
if let (Some(du), Some(ru)) = (&ddl_user, &runtime_user) {
if du == ru {
report.checks.push(error_check(
"ddl-identity",
format!(
"database {db:?}: `migration_url_env` authenticates as the SAME \
MySQL user ({du:?}) as the runtime — the DDL login must be a \
DISTINCT identity (refused)"
),
));
return;
}
}
}
}
}
report.checks.push(ok_check(
"ddl-identity",
format!(
"a distinct DDL identity is configured (`{migration_var}`) and distinct from the runtime"
),
));
}
async fn check_mysql_runtime_grant(
runtime: &Arc<dyn SqlBackend>,
db: &str,
report: &mut RepairReport,
) {
let db_lit = pg_literal(db);
match runtime
.run_query(&format!(
"SELECT 1 FROM information_schema.schema_privileges \
WHERE table_schema = {db_lit} \
AND grantee LIKE CONCAT('''', SUBSTRING_INDEX(CURRENT_USER(), '@', 1), '''@%') \
LIMIT 1;"
))
.await
{
Ok(rows) if !rows.rows.is_empty() => report.checks.push(ok_check(
"runtime-user",
format!("the runtime user holds schema privileges on {db:?}"),
)),
Ok(_) => report.checks.push(RepairCheck {
check: "runtime-user".to_string(),
status: RepairStatus::Drift,
detail: format!(
"the runtime user has no visible schema-level grant on {db:?}; provisioning grants \
`ALL ON {db}.*` — repair reports this (it never issues a GRANT itself)"
),
ddl: None,
}),
Err(e) => report.checks.push(error_check(
"runtime-user",
format!("could not probe the runtime user's schema privileges on {db:?}: {e}"),
)),
}
}
async fn check_mysql_database_exists(
runtime: &Arc<dyn SqlBackend>,
db: &str,
report: &mut RepairReport,
) {
let db_lit = pg_literal(db);
match runtime
.run_query(&format!(
"SELECT 1 FROM information_schema.schemata WHERE schema_name = {db_lit};"
))
.await
{
Ok(rows) if !rows.rows.is_empty() => report.checks.push(ok_check(
"database",
format!("the tenant database {db:?} exists"),
)),
Ok(_) => report.checks.push(RepairCheck {
check: "database".to_string(),
status: RepairStatus::Drift,
detail: format!(
"the tenant database {db:?} does not exist; provision the tenant (repair reconciles \
an existing server, it does not CREATE the database)"
),
ddl: None,
}),
Err(e) => report.checks.push(error_check(
"database",
format!("could not probe the tenant database {db:?}: {e}"),
)),
}
}
#[cfg(feature = "migrate")]
async fn check_libsql_ledger(
sql: &boatramp_storage::LibsqlSql,
mode: RepairMode,
report: &mut RepairReport,
) {
use boatramp_core::sql::SqlBackend;
const LIBSQL_LEDGER: &str = "boatramp_migrations_schema_migrations";
let name_lit = pg_literal(LIBSQL_LEDGER);
let present = match sql
.run_query(&format!(
"SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = {name_lit};"
))
.await
{
Ok(rows) => !rows.rows.is_empty(),
Err(e) => {
report.checks.push(error_check(
"ledger",
format!("could not probe the libsql ledger table: {e}"),
));
return;
}
};
if present {
report.checks.push(ok_check(
"ledger",
format!("the libsql ledger table {LIBSQL_LEDGER:?} exists"),
));
return;
}
let create = format!(
"CREATE TABLE IF NOT EXISTS \"{LIBSQL_LEDGER}\" (id TEXT PRIMARY KEY, ordinal INTEGER NOT \
NULL, content_hash TEXT NOT NULL, kind TEXT NOT NULL, applied_at TEXT NOT NULL DEFAULT \
CURRENT_TIMESTAMP, applied_by TEXT);"
);
match mode {
RepairMode::DryRun => report.checks.push(RepairCheck {
check: "ledger".to_string(),
status: RepairStatus::Drift,
detail: format!(
"the libsql ledger table {LIBSQL_LEDGER:?} is absent; apply would scaffold it \
(CREATE TABLE IF NOT EXISTS)"
),
ddl: Some(create),
}),
RepairMode::Apply => match sql.run_script(&create).await {
Ok(()) => report.checks.push(RepairCheck {
check: "ledger".to_string(),
status: RepairStatus::Repaired,
detail: format!("scaffolded the libsql ledger table {LIBSQL_LEDGER:?}"),
ddl: Some(create),
}),
Err(e) => report.checks.push(error_check(
"ledger",
format!("scaffolding the libsql ledger table failed: {e}"),
)),
},
}
}
fn is_libsql_kind(kind: &str) -> bool {
matches!(
kind.trim().to_ascii_lowercase().as_str(),
"libsql" | "sqlite" | "sqlite3"
)
}
fn audit_stmts(check: &str, d: &Derived, stmts: &[String]) {
for stmt in stmts {
tracing::info!(
target: "boatramp::repair",
project = %d.project,
check = %check,
derived_db = %d.database,
derived_runtime_role = %d.runtime_role,
derived_owner_role = %d.owner_role,
maintenance_user = %d.superuser,
statement = %stmt,
"repair: executed converge statement"
);
}
}
fn redact_pw(stmts: &[String]) -> Vec<String> {
stmts.iter().map(|s| redact_password_literal(s)).collect()
}
fn redact_password_literal(stmt: &str) -> String {
let upper = stmt.to_ascii_uppercase();
let Some(kw) = upper.find("PASSWORD ") else {
return stmt.to_string();
};
let after = kw + "PASSWORD ".len();
let bytes = stmt.as_bytes();
let Some(open_rel) = stmt[after..].find('\'') else {
return stmt.to_string();
};
let open = after + open_rel;
let mut i = open + 1;
while i < bytes.len() {
if bytes[i] == b'\'' {
if i + 1 < bytes.len() && bytes[i + 1] == b'\'' {
i += 2;
continue;
}
break;
}
i += 1;
}
let close = i.min(bytes.len().saturating_sub(1));
format!(
"{}'<redacted>'{}",
&stmt[..open],
&stmt[(close + 1).min(stmt.len())..]
)
}
fn report_probe_err(report: &mut RepairReport, check: &str, e: SqlError) {
report
.checks
.push(error_check(check, format!("probe failed: {e}")));
}
fn classify_backend(binding: &ExternalDatabaseConfig) -> String {
let compute_backed = binding.compute.as_deref().is_some_and(|c| !c.is_empty());
match ExternalSqlKind::parse(&binding.kind) {
Some(ExternalSqlKind::Postgres) if compute_backed => {
let iso = match binding.tenant {
TenantIsolation::Shared => "shared",
TenantIsolation::Single => "single",
};
format!("{iso}-postgres")
}
Some(ExternalSqlKind::Postgres) => "external".to_string(),
Some(ExternalSqlKind::Mysql) => "mysql".to_string(),
None if is_libsql_kind(&binding.kind) => "libsql".to_string(),
None => "external".to_string(),
}
}
fn binding_db(binding: &ExternalDatabaseConfig) -> String {
binding.database.as_deref().unwrap_or_default().to_string()
}
fn maintenance_db(kind: ExternalSqlKind) -> &'static str {
match kind {
ExternalSqlKind::Postgres => "postgres",
ExternalSqlKind::Mysql => "mysql",
}
}
fn base_report(
binding: &ExternalDatabaseConfig,
project: &str,
backend_class: &str,
mode: RepairMode,
) -> RepairReport {
RepairReport {
tenant: format!("{}@{}", binding_db(binding), project),
backend: backend_class.to_string(),
mode: mode.as_str().to_string(),
checks: Vec::new(),
}
}
fn ok_check(check: &str, detail: impl Into<String>) -> RepairCheck {
RepairCheck {
check: check.to_string(),
status: RepairStatus::Ok,
detail: detail.into(),
ddl: None,
}
}
fn skip_check(check: &str, detail: impl Into<String>) -> RepairCheck {
RepairCheck {
check: check.to_string(),
status: RepairStatus::Skipped,
detail: detail.into(),
ddl: None,
}
}
fn error_check(check: &str, detail: impl Into<String>) -> RepairCheck {
RepairCheck {
check: check.to_string(),
status: RepairStatus::Error,
detail: detail.into(),
ddl: None,
}
}
enum LiveState {
Live,
SoftDeletedOnly,
Absent,
}
async fn probe_live_or_soft_deleted(
maint: &Arc<dyn SqlBackend>,
d: &Derived,
) -> Result<LiveState, SqlError> {
let db_lit = pg_literal(&d.database);
let exact = maint
.run_query(&format!(
"SELECT 1 FROM pg_database WHERE datname = {db_lit};"
))
.await?;
if !exact.rows.is_empty() {
return Ok(LiveState::Live);
}
let prefix_lit = pg_literal(&format!("{}__deleted\\_%", d.database));
let sibling = maint
.run_query(&format!(
"SELECT 1 FROM pg_database WHERE datname LIKE {prefix_lit} ESCAPE '\\';"
))
.await?;
if !sibling.rows.is_empty() {
Ok(LiveState::SoftDeletedOnly)
} else {
Ok(LiveState::Absent)
}
}
async fn probe_current_user_superuser(maint: &Arc<dyn SqlBackend>) -> Result<bool, SqlError> {
let rows = maint
.run_query("SELECT rolsuper FROM pg_roles WHERE rolname = current_user;")
.await?;
Ok(rows
.rows
.first()
.and_then(|r| r.first())
.map(sql_bool)
.unwrap_or(false))
}
async fn probe_current_database(backend: &Arc<dyn SqlBackend>) -> Result<String, SqlError> {
let rows = backend.run_query("SELECT current_database();").await?;
Ok(rows
.rows
.first()
.and_then(|r| r.first())
.map(sql_text)
.unwrap_or_default())
}
async fn probe_role_exists(maint: &Arc<dyn SqlBackend>, role: &str) -> Result<bool, SqlError> {
let role_lit = pg_literal(role);
let rows = maint
.run_query(&format!(
"SELECT 1 FROM pg_roles WHERE rolname = {role_lit} AND rolcanlogin;"
))
.await?;
Ok(!rows.rows.is_empty())
}
fn sql_bool(v: &SqlValue) -> bool {
match v {
SqlValue::Boolean(b) => *b,
SqlValue::Integer(n) => *n != 0,
SqlValue::Text(s) => matches!(s.as_str(), "t" | "true" | "TRUE" | "1"),
_ => false,
}
}
fn sql_i64(v: &SqlValue) -> i64 {
match v {
SqlValue::Integer(n) => *n,
SqlValue::Text(s) => s.trim().parse::<i64>().unwrap_or(1),
SqlValue::Real(f) => *f as i64,
_ => 1,
}
}
fn sql_text(v: &SqlValue) -> String {
match v {
SqlValue::Text(s) => s.clone(),
other => format!("{other:?}"),
}
}
fn pg_literal(s: &str) -> String {
format!("'{}'", s.replace('\'', "''"))
}
#[cfg(test)]
mod tests;