use std::sync::Arc;
use async_trait::async_trait;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use crate::config::ManagedDbPrivilege;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use boatramp_core::compute::{ManagedDbEnvResolver, PrivilegeDirective, ReplicaPhase};
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use boatramp_core::deploy::DeployStore;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use boatramp_core::envelope::KeyEnvelope;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use boatramp_core::kv::KvStore;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use boatramp_core::project::ProjectRef;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use boatramp_core::sql::SqlError;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
use boatramp_core::sql::{
AppliedMigration, LedgerOrigin, MigrateDdl, MigrateDdlError, MigrationError, MigrationStep,
MigrationSubstrate, SubstrateStepOutcome,
};
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use boatramp_storage::sql_compute::{ComputeEndpointResolver, ReplicaDiag};
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use boatramp_storage::ExternalSqlKind;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use std::collections::HashMap;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub fn managed_db_server_env(
kind: ExternalSqlKind,
database: &str,
user: &str,
password: &str,
) -> Vec<(String, String)> {
match kind {
ExternalSqlKind::Postgres => vec![
("POSTGRES_USER".into(), user.into()),
("POSTGRES_PASSWORD".into(), password.into()),
("POSTGRES_DB".into(), database.into()),
],
ExternalSqlKind::Mysql => vec![
("MYSQL_USER".into(), user.into()),
("MYSQL_PASSWORD".into(), password.into()),
("MYSQL_DATABASE".into(), database.into()),
("MYSQL_ROOT_PASSWORD".into(), password.into()),
],
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub struct ManagedSqlCredentials {
kv: Arc<dyn KvStore>,
envelope: Arc<dyn KeyEnvelope>,
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
impl ManagedSqlCredentials {
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub fn new(kv: Arc<dyn KvStore>, envelope: Arc<dyn KeyEnvelope>) -> Self {
Self { kv, envelope }
}
fn key(project: &str, workload: &str) -> String {
format!("managed-sql-cred/{project}/{workload}")
}
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub async fn password(&self, project: &str, workload: &str) -> Result<String, String> {
let key = Self::key(project, workload);
if let Some(sealed) = self.kv.get(&key).await.map_err(|e| e.to_string())? {
let plain = self
.envelope
.unwrap(&sealed)
.await
.map_err(|e| e.to_string())?;
return String::from_utf8(plain).map_err(|_| {
format!("managed sql credential for {workload:?} is not valid UTF-8")
});
}
let mut bytes = [0u8; 32];
getrandom::getrandom(&mut bytes).map_err(|e| format!("rng: {e}"))?;
let password: String = bytes.iter().map(|b| format!("{b:02x}")).collect();
let sealed = self
.envelope
.wrap(password.as_bytes())
.await
.map_err(|e| e.to_string())?;
self.kv.put(&key, sealed).await.map_err(|e| e.to_string())?;
Ok(password)
}
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub async fn is_sealed(&self, project: &str, workload: &str) -> Result<bool, String> {
let key = Self::key(project, workload);
Ok(self
.kv
.get(&key)
.await
.map_err(|e| e.to_string())?
.is_some())
}
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub async fn get_sealed_password(
&self,
project: &str,
workload: &str,
) -> Result<Option<String>, String> {
let key = Self::key(project, workload);
let Some(sealed) = self.kv.get(&key).await.map_err(|e| e.to_string())? else {
return Ok(None);
};
let plain = self
.envelope
.unwrap(&sealed)
.await
.map_err(|e| e.to_string())?;
String::from_utf8(plain)
.map(Some)
.map_err(|_| format!("managed sql credential for {workload:?} is not valid UTF-8"))
}
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub async fn delete(&self, project: &str, workload: &str) -> Result<(), String> {
self.kv
.delete(&Self::key(project, workload))
.await
.map_err(|e| e.to_string())
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
struct ManagedDbSpec {
kind: ExternalSqlKind,
database: String,
user: String,
tenant: crate::config::TenantIsolation,
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub struct ManagedDbEnv {
dbs: HashMap<String, ManagedDbSpec>,
creds: ManagedSqlCredentials,
privilege: ManagedDbPrivilege,
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
fn managed_db_default_ids(_kind: ExternalSqlKind) -> (u32, u32) {
(999, 999)
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
fn managed_db_caps() -> Vec<String> {
["CHOWN", "DAC_OVERRIDE", "FOWNER", "SETUID", "SETGID"]
.iter()
.map(|s| (*s).to_string())
.collect()
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
impl ManagedDbEnv {
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub fn from_config(
databases: &std::collections::BTreeMap<String, crate::config::ExternalDatabaseConfig>,
creds: ManagedSqlCredentials,
privilege: ManagedDbPrivilege,
) -> Self {
let mut dbs = HashMap::new();
for db in databases.values() {
if !db.is_managed_credential() {
continue;
}
let (Some(workload), Some(kind), Some(database), Some(user)) = (
db.compute.clone(),
ExternalSqlKind::parse(&db.kind),
db.database.clone(),
db.user.clone(),
) else {
continue;
};
dbs.insert(
workload,
ManagedDbSpec {
kind,
database,
user,
tenant: db.tenant,
},
);
}
Self {
dbs,
creds,
privilege,
}
}
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub fn is_empty(&self) -> bool {
self.dbs.is_empty()
}
fn resolve_spec(&self, workload: &str) -> Option<&ManagedDbSpec> {
if let Some(spec) = self.dbs.get(workload) {
return Some(spec);
}
self.dbs
.iter()
.filter(|(base, spec)| {
matches!(spec.tenant, crate::config::TenantIsolation::Single)
&& workload
.strip_prefix(base.as_str())
.is_some_and(|rest| rest.starts_with('-') && rest.len() > 1)
})
.max_by_key(|(base, _)| base.len())
.map(|(_, spec)| spec)
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[async_trait]
impl ManagedDbEnvResolver for ManagedDbEnv {
async fn managed_db_env(&self, project: &str, workload: &str) -> Vec<(String, String)> {
let Some(db) = self.resolve_spec(workload) else {
return Vec::new();
};
match self.creds.password(project, workload).await {
Ok(password) => managed_db_server_env(db.kind, &db.database, &db.user, &password),
Err(e) => {
tracing::error!(
%workload,
error = %e,
"managed sql: could not resolve the sealed credential; DB launched without managed env"
);
Vec::new()
}
}
}
fn managed_db_privilege(&self, _project: &str, workload: &str) -> Option<PrivilegeDirective> {
let db = self.resolve_spec(workload)?;
Some(match self.privilege {
ManagedDbPrivilege::Rootless => {
let (uid, gid) = managed_db_default_ids(db.kind);
PrivilegeDirective::Rootless { uid, gid }
}
ManagedDbPrivilege::Caps => PrivilegeDirective::Caps(managed_db_caps()),
})
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
pub async fn auto_register_managed_db_workloads(
deploy: &DeployStore,
databases: &std::collections::BTreeMap<String, crate::config::ExternalDatabaseConfig>,
) {
use crate::config::TenantIsolation;
for db in databases.values() {
if !db.is_managed_credential() {
continue;
}
let Some(workload) = db.compute.as_deref().filter(|c| !c.is_empty()) else {
continue;
};
if ExternalSqlKind::parse(&db.kind).is_none() {
continue;
}
match db.tenant {
TenantIsolation::Shared => {
register_shared_server(deploy, db, workload).await;
}
TenantIsolation::Single => {}
}
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
async fn register_shared_server(
deploy: &DeployStore,
db: &crate::config::ExternalDatabaseConfig,
workload: &str,
) {
use boatramp_core::compute::{
managed_db_spec, ComputeWorkload, ManagedDbEngine, PlacementConstraints,
};
const DEFAULT_VOLUME_MIB: u32 = 10 * 1024;
let engine = match ExternalSqlKind::parse(&db.kind) {
Some(ExternalSqlKind::Postgres) => ManagedDbEngine::Postgres,
Some(ExternalSqlKind::Mysql) => ManagedDbEngine::Mysql,
None => return,
};
match deploy
.get_compute_workload(ProjectRef::DEFAULT, workload)
.await
{
Ok(Some(_)) => return,
Ok(None) => {}
Err(e) => {
tracing::warn!(%workload, error = %e, "managed sql: could not check for an existing compute workload; skipping auto-register");
return;
}
}
let image = db.image.as_deref();
let mut spec = managed_db_spec(
engine,
image,
db.volume_size_mib.unwrap_or(DEFAULT_VOLUME_MIB),
);
if let Some(grace) = db.startup_grace_secs {
spec.startup_grace_secs = grace;
}
let spec_id = match deploy.put_compute_spec(&spec).await {
Ok(id) => id,
Err(e) => {
tracing::warn!(%workload, error = %e, "managed sql: could not store the auto-registered compute spec");
return;
}
};
let wl = ComputeWorkload {
version: 1,
name: workload.to_string(),
active: spec_id,
replicas: 1,
placement: PlacementConstraints::default(),
};
match deploy.set_compute_workload(ProjectRef::DEFAULT, &wl).await {
Ok(()) => tracing::info!(
%workload,
image = %image.unwrap_or_else(|| engine.default_image()),
"managed sql: auto-registered the shared co-located database compute workload"
),
Err(e) => {
tracing::warn!(%workload, error = %e, "managed sql: could not register the auto-registered compute workload")
}
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub struct DeployEndpointResolver {
deploy: DeployStore,
project: String,
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
impl DeployEndpointResolver {
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub fn new(deploy: DeployStore, project: impl Into<String>) -> Self {
Self {
deploy,
project: project.into(),
}
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[async_trait]
impl ComputeEndpointResolver for DeployEndpointResolver {
async fn endpoints(&self, workload: &str) -> Result<Vec<(String, u16)>, SqlError> {
let states = self
.deploy
.list_replica_states(ProjectRef::new(&self.project), workload)
.await
.map_err(SqlError::other)?;
Ok(states
.into_iter()
.filter(|s| s.phase == ReplicaPhase::Running && s.healthy)
.map(|s| (s.endpoint.host, s.endpoint.port))
.collect())
}
async fn replica_diagnostics(&self, workload: &str) -> Vec<ReplicaDiag> {
let states = match self
.deploy
.list_replica_states(ProjectRef::new(&self.project), workload)
.await
{
Ok(s) => s,
Err(_) => return Vec::new(),
};
states
.into_iter()
.map(|s| ReplicaDiag {
endpoint: format!("{}:{}", s.endpoint.host, s.endpoint.port),
healthy: s.healthy,
phase: format!("{:?}", s.phase),
})
.collect()
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
pub struct NodeOperatorSql {
databases: std::collections::BTreeMap<String, crate::config::ExternalDatabaseConfig>,
kv: Arc<dyn KvStore>,
envelope: Option<Arc<dyn KeyEnvelope>>,
deploy: DeployStore,
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
impl NodeOperatorSql {
pub fn new(
databases: std::collections::BTreeMap<String, crate::config::ExternalDatabaseConfig>,
kv: Arc<dyn KvStore>,
envelope: Option<Arc<dyn KeyEnvelope>>,
deploy: DeployStore,
) -> Self {
Self {
databases,
kv,
envelope,
deploy,
}
}
pub(crate) async fn backend_for(
&self,
project: &str,
db: &str,
) -> Result<Arc<dyn boatramp_core::sql::SqlBackend>, SqlError> {
self.connect_for(project, db, false).await
}
pub(crate) async fn owner_backend_for(
&self,
project: &str,
db: &str,
) -> Result<Arc<dyn boatramp_core::sql::SqlBackend>, SqlError> {
self.connect_for(project, db, true).await
}
pub(crate) fn engine_kind(&self, db: &str) -> Option<ExternalSqlKind> {
self.databases
.get(db)
.and_then(|cfg| ExternalSqlKind::parse(&cfg.kind))
}
async fn mysql_ddl_backend_for(
&self,
db: &str,
) -> Result<Arc<dyn boatramp_core::sql::SqlBackend>, SqlError> {
use boatramp_storage::sql_sqlx::{connect, ExternalSqlOptions};
let cfg = self
.databases
.get(db)
.ok_or_else(|| SqlError::other(format!("no database named {db:?}")))?;
let compute_backed = cfg.compute.as_deref().is_some_and(|c| !c.is_empty());
if compute_backed && cfg.url_env.is_empty() {
return Err(SqlError::other(format!(
"database {db:?}: compute-backed managed MySQL migration is not supported this \
release: boatramp cannot yet auto-derive a distinct least-privilege DDL identity \
for a managed MySQL database; use an external MySQL binding with a distinct \
`migration_url_env`, or await the auto-minted DDL-grant follow-up"
)));
}
let migration_var = cfg
.migration_url_env
.as_deref()
.filter(|v| !v.is_empty())
.ok_or_else(|| {
SqlError::other(format!(
"database {db:?}: MySQL schema migrations require a distinct DDL identity — set \
`migration_url_env` to an admin/DDL connection URL that is NOT the runtime \
`user`/`url_env` (MySQL has no owner/runtime role split; running DDL as the \
runtime tenant user is refused fail-closed)"
))
})?;
let ddl_url = std::env::var(migration_var).map_err(|_| {
SqlError::other(format!(
"env var {migration_var} (migration/DDL url for {db:?}) is unset"
))
})?;
if !cfg.url_env.is_empty() {
if let Ok(runtime_url) = std::env::var(&cfg.url_env) {
if runtime_url == ddl_url {
return Err(SqlError::other(format!(
"database {db:?}: `migration_url_env` resolves to the SAME connection as the \
runtime `url_env` — the MySQL DDL identity must be distinct from the runtime \
user (refused fail-closed)"
)));
}
#[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(ddl_user), Some(runtime_user)) = (&ddl_user, &runtime_user) {
if ddl_user == runtime_user {
return Err(SqlError::other(format!(
"database {db:?}: `migration_url_env` authenticates as the SAME \
MySQL user ({ddl_user:?}) as the runtime `url_env` — the DDL login \
must be a DISTINCT identity from the runtime tenant login (refused \
fail-closed)"
)));
}
}
}
}
}
let timeout = cfg.connect_timeout_secs.map(std::time::Duration::from_secs);
let opts = ExternalSqlOptions::new(ddl_url)
.with_max_connections(cfg.pool_max)
.read_only(false)
.with_connect_timeout(timeout);
connect(ExternalSqlKind::Mysql, &opts)
}
async fn connect_for(
&self,
project: &str,
db: &str,
owner: bool,
) -> Result<Arc<dyn boatramp_core::sql::SqlBackend>, SqlError> {
use boatramp_storage::sql_compute::ComputeResolvedSqlBackend;
use boatramp_storage::sql_sqlx::{connect, ExternalSqlOptions};
boatramp_core::project::validate_resource_name("database", db)
.map_err(|err| SqlError::other(err.to_string()))?;
let cfg = self
.databases
.get(db)
.ok_or_else(|| SqlError::other(format!("no database named {db:?}")))?;
let kind = ExternalSqlKind::parse(&cfg.kind).ok_or_else(|| {
SqlError::other(format!("database {db:?}: unknown engine {:?}", cfg.kind))
})?;
let timeout = cfg.connect_timeout_secs.map(std::time::Duration::from_secs);
if cfg.compute.as_deref().is_some_and(|c| !c.is_empty()) {
let target = if owner {
owner_target(cfg, project, db)?
} else {
operator_target(cfg, project, db)?
};
let password = match cfg.password_env.as_deref().filter(|v| !v.is_empty()) {
Some(var) => std::env::var(var)
.map_err(|_| SqlError::other(format!("env var {var} (password) is unset")))?,
None => {
let envelope = self.envelope.clone().ok_or_else(|| {
SqlError::other(format!(
"managed database {db:?} needs a [secrets] envelope to unseal its credential"
))
})?;
ManagedSqlCredentials::new(self.kv.clone(), envelope)
.password(&target.cred_project, &target.cred_workload)
.await
.map_err(SqlError::other)?
}
};
let resolver = Arc::new(DeployEndpointResolver::new(
self.deploy.clone(),
target.endpoint_project,
));
Ok(Arc::new(ComputeResolvedSqlBackend::new(
resolver,
target.workload,
kind,
target.database,
target.user,
password,
cfg.pool_max,
cfg.read_only,
timeout,
)))
} else {
let url = std::env::var(&cfg.url_env)
.map_err(|_| SqlError::other(format!("env var {} (url) is unset", cfg.url_env)))?;
let read_url =
match &cfg.read_url_env {
Some(var) => Some(std::env::var(var).map_err(|_| {
SqlError::other(format!("env var {var} (read url) is unset"))
})?),
None => None,
};
let opts = ExternalSqlOptions::new(url)
.with_read_url(read_url)
.with_max_connections(cfg.pool_max)
.read_only(cfg.read_only)
.with_connect_timeout(timeout);
connect(kind, &opts)
}
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[derive(Debug)]
pub(crate) struct OperatorTarget {
pub workload: String,
pub database: String,
pub user: String,
pub endpoint_project: String,
pub cred_project: String,
pub cred_workload: String,
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
pub(crate) fn operator_target(
cfg: &crate::config::ExternalDatabaseConfig,
project: &str,
db: &str,
) -> Result<OperatorTarget, SqlError> {
use crate::config::{TenantIsolation, TenantScope};
use crate::tenant_sql::{single_credential_project, tenant_key, tenant_names};
use boatramp_core::project::DEFAULT_PROJECT;
let compute = cfg
.compute
.as_deref()
.filter(|c| !c.is_empty())
.ok_or_else(|| SqlError::other(format!("database {db:?}: not compute-backed")))?;
let database = cfg.database.as_deref().unwrap_or_default();
let user = cfg.user.as_deref().unwrap_or_default();
if matches!(cfg.tenant_scope, TenantScope::Site) {
return Err(SqlError::other(format!(
"database {db:?} is a site-scoped managed database; operator sql exec/query is \
project-level and cannot target a specific site's database"
)));
}
let (tenant_ident_raw, is_default) = tenant_key(cfg.tenant_scope, project, "");
let names = tenant_names(cfg.tenant, compute, database, &tenant_ident_raw, is_default);
let (cred_project, cred_workload) = match cfg.tenant {
TenantIsolation::Single => (
single_credential_project(project, is_default),
names.workload.clone(),
),
TenantIsolation::Shared => (DEFAULT_PROJECT.to_string(), compute.to_string()),
};
let endpoint_project = match cfg.tenant {
TenantIsolation::Single if !is_default => project.to_string(),
_ => DEFAULT_PROJECT.to_string(),
};
Ok(OperatorTarget {
workload: names.workload,
database: names.database,
user: user.to_string(),
endpoint_project,
cred_project,
cred_workload,
})
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
pub(crate) fn owner_target(
cfg: &crate::config::ExternalDatabaseConfig,
project: &str,
db: &str,
) -> Result<OperatorTarget, SqlError> {
use crate::config::{TenantIsolation, TenantScope};
use crate::tenant_sql::{owner_credential_workload_key, tenant_key, tenant_names};
let mut target = operator_target(cfg, project, db)?;
if matches!(cfg.tenant, TenantIsolation::Shared)
&& !matches!(cfg.tenant_scope, TenantScope::Site)
{
let compute = cfg.compute.as_deref().unwrap_or_default();
let database = cfg.database.as_deref().unwrap_or_default();
let (tenant_ident_raw, is_default) = tenant_key(cfg.tenant_scope, project, "");
if !is_default {
let names = tenant_names(cfg.tenant, compute, database, &tenant_ident_raw, is_default);
let ident = boatramp_storage::tenant_provision::sanitize_ident(&tenant_ident_raw);
target.user = names.owner_role;
target.cred_project = project.to_string();
target.cred_workload = owner_credential_workload_key(compute, &ident);
}
}
Ok(target)
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[async_trait]
impl boatramp_core::sql::OperatorSql for NodeOperatorSql {
async fn exec_script(&self, project: &str, db: &str, script: &str) -> Result<(), SqlError> {
self.backend_for(project, db)
.await?
.run_script(script)
.await
}
async fn query(
&self,
project: &str,
db: &str,
sql: &str,
) -> Result<boatramp_core::sql::SqlRows, SqlError> {
self.backend_for(project, db).await?.run_query(sql).await
}
async fn ping(
&self,
project: &str,
db: &str,
) -> Result<Vec<boatramp_core::sql::SqlPingReplica>, SqlError> {
use boatramp_core::sql::SqlPingReplica;
use std::time::Duration;
let cfg = self
.databases
.get(db)
.ok_or_else(|| SqlError::other(format!("no database named {db:?}")))?;
if cfg.compute.as_deref().is_none_or(|c| c.is_empty()) {
return Err(SqlError::other(format!(
"database {db:?} is not compute-backed; `sql ping` probes managed co-located \
replicas only"
)));
}
let target = operator_target(cfg, project, db)?;
let resolver = DeployEndpointResolver::new(self.deploy.clone(), target.endpoint_project);
let diags = resolver.replica_diagnostics(&target.workload).await;
let mut out = Vec::with_capacity(diags.len());
for d in diags {
let reachable = match d.endpoint.parse::<std::net::SocketAddr>() {
Ok(addr) => matches!(
tokio::time::timeout(
Duration::from_secs(2),
tokio::net::TcpStream::connect(addr),
)
.await,
Ok(Ok(_))
),
Err(_) => false,
};
out.push(SqlPingReplica {
endpoint: d.endpoint,
healthy: d.healthy,
phase: d.phase,
tcp_reachable: reachable,
});
}
Ok(out)
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
pub struct NodeMigrationRunner {
op: Arc<NodeOperatorSql>,
trusted_extensions: std::collections::BTreeSet<String>,
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
const LEDGER_SCHEMA: &str = "boatramp_migrations";
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
const LEDGER_TABLE: &str = "schema_migrations";
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
fn sql_quote_literal(s: &str) -> String {
format!("'{}'", s.replace('\'', "''"))
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
fn valid_migration_id(id: &str) -> bool {
!id.is_empty()
&& id
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'))
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
fn mentions_create_extension(script: &str, kind: ExternalSqlKind) -> bool {
boatramp_core::sql::script_has_create_extension_in(script, guard_dialect(kind))
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
fn guard_dialect(kind: ExternalSqlKind) -> boatramp_core::sql::GuardDialect {
match kind {
ExternalSqlKind::Mysql => boatramp_core::sql::GuardDialect::Mysql,
ExternalSqlKind::Postgres => boatramp_core::sql::GuardDialect::Postgres,
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
fn mentions_txn_control(script: &str, kind: ExternalSqlKind) -> bool {
boatramp_core::sql::script_has_txn_control_in(script, guard_dialect(kind))
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
fn mentions_ledger_schema(script: &str, kind: ExternalSqlKind) -> bool {
boatramp_core::sql::script_references_word_in(script, LEDGER_SCHEMA, guard_dialect(kind))
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
pub(crate) struct OwnerDdl {
owner: Arc<dyn boatramp_core::sql::SqlBackend>,
kind: ExternalSqlKind,
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
impl OwnerDdl {
fn guard(&self, script: &str) -> Result<(), MigrateDdlError> {
if mentions_ledger_schema(script, self.kind) {
return Err(MigrateDdlError::LedgerProtected);
}
if mentions_txn_control(script, self.kind) {
return Err(MigrateDdlError::TxnControl);
}
Ok(())
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[async_trait]
impl MigrateDdl for OwnerDdl {
async fn exec(&self, script: &str) -> Result<(), MigrateDdlError> {
self.guard(script)?;
self.owner
.run_script(script)
.await
.map_err(|e| MigrateDdlError::Sql(sanitize_migration_error(&e.to_string())))
}
async fn exec_batch(&self, scripts: Vec<String>) -> Result<(), MigrateDdlError> {
for script in &scripts {
self.exec(script).await?;
}
Ok(())
}
async fn query(&self, sql: &str) -> Result<boatramp_core::sql::SqlRows, MigrateDdlError> {
self.guard(sql)?;
let rows = self
.owner
.run_query(sql)
.await
.map_err(|e| MigrateDdlError::Sql(sanitize_migration_error(&e.to_string())))?;
if rows.rows.len() > MIGRATE_QUERY_MAX_ROWS {
return Err(MigrateDdlError::Sql(format!(
"query returned {} rows (cap {MIGRATE_QUERY_MAX_ROWS}); add a LIMIT — a migration \
verification query should read a bounded set",
rows.rows.len()
)));
}
Ok(rows)
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
const MIGRATE_QUERY_MAX_ROWS: usize = 100_000;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
impl NodeMigrationRunner {
pub fn new(
op: Arc<NodeOperatorSql>,
trusted_extensions: std::collections::BTreeSet<String>,
) -> Self {
Self {
op,
trusted_extensions,
}
}
fn ledger(kind: ExternalSqlKind) -> String {
use boatramp_storage::tenant_provision::quote_ident;
format!(
"{}.{}",
quote_ident(kind, LEDGER_SCHEMA),
quote_ident(kind, LEDGER_TABLE)
)
}
fn ledger_insert(
kind: ExternalSqlKind,
step: &MigrationStep,
ordinal: usize,
effective_hash: &str,
origin: LedgerOrigin,
) -> String {
format!(
"INSERT INTO {ledger} (id, ordinal, content_hash, kind, applied_by) \
VALUES ({id}, {ord}, {hash}, {step_kind}, {origin});",
ledger = Self::ledger(kind),
id = sql_quote_literal(&step.id),
ord = ordinal,
hash = sql_quote_literal(effective_hash),
step_kind = sql_quote_literal(step.kind()),
origin = sql_quote_literal(origin.as_str()),
)
}
async fn ensure_ledger(
&self,
kind: ExternalSqlKind,
owner: &Arc<dyn boatramp_core::sql::SqlBackend>,
) -> Result<(), SqlError> {
use boatramp_storage::tenant_provision::quote_ident;
match kind {
ExternalSqlKind::Postgres => {
let schema = quote_ident(kind, LEDGER_SCHEMA);
owner
.run_script(&format!("CREATE SCHEMA IF NOT EXISTS {schema};"))
.await?;
owner
.run_script(&format!(
"CREATE TABLE IF NOT EXISTS {ledger} (\
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);",
ledger = Self::ledger(kind)
))
.await?;
}
ExternalSqlKind::Mysql => {
let db = quote_ident(kind, LEDGER_SCHEMA);
owner
.run_script(&format!("CREATE DATABASE IF NOT EXISTS {db};"))
.await?;
owner
.run_script(&format!(
"CREATE TABLE IF NOT EXISTS {ledger} (\
id VARCHAR(255) PRIMARY KEY, \
ordinal INT NOT NULL, \
content_hash VARCHAR(255) NOT NULL, \
kind VARCHAR(32) NOT NULL, \
applied_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, \
applied_by VARCHAR(32)) ENGINE=InnoDB;",
ledger = Self::ledger(kind)
))
.await?;
}
}
Ok(())
}
async fn read_applied(
&self,
kind: ExternalSqlKind,
owner: &Arc<dyn boatramp_core::sql::SqlBackend>,
) -> Result<Vec<AppliedMigration>, SqlError> {
use boatramp_core::sql::SqlValue;
let applied_at = match kind {
ExternalSqlKind::Postgres => "applied_at::text",
ExternalSqlKind::Mysql => "CAST(applied_at AS CHAR)",
};
let rows = owner
.run_query(&format!(
"SELECT id, ordinal, content_hash, kind, {applied_at}, applied_by \
FROM {ledger} ORDER BY ordinal;",
ledger = Self::ledger(kind)
))
.await?;
let text = |v: &SqlValue| match v {
SqlValue::Text(s) => s.clone(),
other => format!("{other:?}"),
};
let int = |v: &SqlValue| match v {
SqlValue::Integer(n) => *n,
_ => 0,
};
Ok(rows
.rows
.iter()
.map(|r| {
let origin = match r.get(5) {
Some(SqlValue::Text(s)) if s == LedgerOrigin::Baseline.as_str() => {
LedgerOrigin::Baseline.as_str().to_string()
}
_ => LedgerOrigin::Apply.as_str().to_string(),
};
AppliedMigration {
id: r.first().map(&text).unwrap_or_default(),
ordinal: r.get(1).map(&int).unwrap_or_default(),
content_hash: r.get(2).map(&text).unwrap_or_default(),
kind: r.get(3).map(&text).unwrap_or_default(),
applied_at: r.get(4).map(&text).unwrap_or_default(),
origin,
}
})
.collect())
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
impl NodeMigrationRunner {
fn engine_gate(&self, db: &str) -> Result<ExternalSqlKind, MigrationError> {
match self.op.engine_kind(db) {
Some(kind @ (ExternalSqlKind::Postgres | ExternalSqlKind::Mysql)) => Ok(kind),
None => Err(MigrationError::NotConfigured),
}
}
async fn connect_ddl(
&self,
project: &str,
db: &str,
) -> Result<Arc<dyn boatramp_core::sql::SqlBackend>, MigrationError> {
let kind = self.engine_gate(db)?;
let result = match kind {
ExternalSqlKind::Postgres => self.op.owner_backend_for(project, db).await,
ExternalSqlKind::Mysql => self.op.mysql_ddl_backend_for(db).await,
};
result.map_err(|e| match e {
SqlError::Unavailable(m) => MigrationError::Unavailable(m),
other => MigrationError::Sql(other),
})
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[async_trait]
impl MigrationSubstrate for NodeMigrationRunner {
async fn preflight(
&self,
project: &str,
db: &str,
) -> Result<Vec<AppliedMigration>, MigrationError> {
let kind = self.engine_gate(db)?;
let owner = self.connect_ddl(project, db).await?;
self.ensure_ledger(kind, &owner).await?;
Ok(self.read_applied(kind, &owner).await?)
}
async fn apply_substrate_step(
&self,
project: &str,
db: &str,
step: &MigrationStep,
ordinal: usize,
effective_hash: &str,
) -> Result<SubstrateStepOutcome, MigrationError> {
use boatramp_core::sql::MigrationAction;
if !valid_migration_id(&step.id) {
return Ok(SubstrateStepOutcome::Failed(
"invalid migration id (allowed: A-Za-z0-9._-)".to_string(),
));
}
let kind = self.engine_gate(db)?;
let owner = self.connect_ddl(project, db).await?;
let ledger_insert =
Self::ledger_insert(kind, step, ordinal, effective_hash, LedgerOrigin::Apply);
let outcome: Result<(), String> = match &step.action {
MigrationAction::Sql {
script,
no_transaction,
} => {
if mentions_create_extension(script, kind) {
let hint = match kind {
ExternalSqlKind::Postgres => "use an extension step",
ExternalSqlKind::Mysql => "MySQL has no CREATE EXTENSION",
};
Err(format!("a sql step may not CREATE EXTENSION — {hint}"))
} else if mentions_ledger_schema(script, kind) {
Err(
"a sql step may not reference the host-owned migration-ledger schema"
.to_string(),
)
} else if mentions_txn_control(script, kind)
&& (kind == ExternalSqlKind::Mysql || !*no_transaction)
{
let why = match kind {
ExternalSqlKind::Postgres => {
"a transactional sql step may not contain its own BEGIN/COMMIT/ROLLBACK \
(it would desync the atomic wrapper) — use a no_transaction step to \
manage the transaction yourself"
}
ExternalSqlKind::Mysql => {
"a MySQL sql step may not contain its own BEGIN/COMMIT/ROLLBACK — DDL \
implicitly commits on MySQL, so there is no transaction to control (a \
multi-DDL step is applied per-statement, not atomically)"
}
};
Err(why.to_string())
} else {
match kind {
ExternalSqlKind::Postgres if !*no_transaction => {
let batch = format!("BEGIN;\n{script};\n{ledger_insert}\nCOMMIT;");
owner.run_script(&batch).await.map_err(|e| e.to_string())
}
_ => match owner.run_script(script).await {
Ok(()) => owner
.run_script(&ledger_insert)
.await
.map_err(|e| {
mysql_partial_apply_note(kind, &step.id, &e.to_string())
}),
Err(e) => Err(mysql_multi_ddl_note(kind, &step.id, &e.to_string())),
},
}
}
}
MigrationAction::Extension { name } => match kind {
ExternalSqlKind::Mysql => Err(format!(
"extension step {name:?} is not supported on MySQL (MySQL has no CREATE \
EXTENSION) — install any plugin operator-side and use a plain sql step"
)),
ExternalSqlKind::Postgres => {
if !self.trusted_extensions.contains(name) {
Err(format!(
"extension {name:?} is not on the operator trusted-extension allowlist"
))
} else {
use boatramp_storage::tenant_provision::quote_ident;
let create = format!(
"CREATE EXTENSION IF NOT EXISTS {};",
quote_ident(kind, name)
);
match self.op.backend_for(project, db).await {
Ok(su) => match su.run_script(&create).await {
Ok(()) => owner
.run_script(&ledger_insert)
.await
.map_err(|e| e.to_string()),
Err(e) => Err(e.to_string()),
},
Err(e) => Err(e.to_string()),
}
}
}
},
MigrationAction::Function { .. } => {
return Err(MigrationError::Other(
"internal: a function step must be invoked by the orchestrator, not the \
substrate"
.to_string(),
))
}
};
Ok(match outcome {
Ok(()) => SubstrateStepOutcome::Applied,
Err(error) => SubstrateStepOutcome::Failed(sanitize_migration_error(&error)),
})
}
async fn record(
&self,
project: &str,
db: &str,
step: &MigrationStep,
ordinal: usize,
effective_hash: &str,
origin: LedgerOrigin,
) -> Result<(), MigrationError> {
if !valid_migration_id(&step.id) {
return Err(MigrationError::Other(format!(
"invalid migration id {:?} (allowed: A-Za-z0-9._-)",
step.id
)));
}
let kind = self.engine_gate(db)?;
let owner = self.connect_ddl(project, db).await?;
self.ensure_ledger(kind, &owner).await?;
owner
.run_script(&Self::ledger_insert(
kind,
step,
ordinal,
effective_hash,
origin,
))
.await?;
Ok(())
}
async fn owner_ddl(
&self,
project: &str,
db: &str,
) -> Result<Arc<dyn MigrateDdl>, MigrationError> {
let kind = self.engine_gate(db)?;
let owner = self.connect_ddl(project, db).await?;
Ok(Arc::new(OwnerDdl { owner, kind }))
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
fn mysql_multi_ddl_note(kind: ExternalSqlKind, id: &str, err: &str) -> String {
match kind {
ExternalSqlKind::Mysql => format!(
"step {id:?} failed mid-way and may be PARTIALLY APPLIED (MySQL commits each DDL \
statement implicitly; earlier statements were not rolled back). Re-apply re-runs the \
whole step — make it idempotent. Underlying error: {err}"
),
ExternalSqlKind::Postgres => err.to_string(),
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
fn mysql_partial_apply_note(kind: ExternalSqlKind, id: &str, err: &str) -> String {
match kind {
ExternalSqlKind::Mysql => format!(
"step {id:?} DDL applied but its ledger row could not be recorded — the step is \
PARTIALLY APPLIED (schema changed, unrecorded). Re-apply re-runs the whole step \
(make it idempotent). Underlying error: {err}"
),
ExternalSqlKind::Postgres => err.to_string(),
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
fn sanitize_migration_error(e: &str) -> String {
let first = e.lines().next().unwrap_or(e);
if first.chars().count() > 300 {
let truncated: String = first.chars().take(300).collect();
format!("{truncated}…")
} else {
first.to_string()
}
}
#[cfg(feature = "migrate")]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub struct LibsqlMigrationRunner {
databases: std::collections::BTreeMap<String, crate::config::ExternalDatabaseConfig>,
}
#[cfg(feature = "migrate")]
pub(crate) fn kind_is_libsql(kind: &str) -> bool {
matches!(
kind.trim().to_ascii_lowercase().as_str(),
"libsql" | "sqlite" | "sqlite3"
)
}
#[cfg(feature = "migrate")]
const LIBSQL_LEDGER_WORD: &str = concat!("boatramp_migrations", "_", "schema_migrations");
#[cfg(feature = "migrate")]
fn libsql_ledger() -> String {
format!("\"{LEDGER_SCHEMA}_{LEDGER_TABLE}\"")
}
#[cfg(feature = "migrate")]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
impl LibsqlMigrationRunner {
pub fn new(
databases: std::collections::BTreeMap<String, crate::config::ExternalDatabaseConfig>,
) -> Self {
Self { databases }
}
pub(crate) fn is_libsql(&self, db: &str) -> bool {
self.databases
.get(db)
.is_some_and(|cfg| kind_is_libsql(&cfg.kind))
}
fn engine_gate(&self, db: &str) -> Result<(), MigrationError> {
if self.is_libsql(db) {
Ok(())
} else {
Err(MigrationError::NotConfigured)
}
}
async fn open(&self, db: &str) -> Result<boatramp_storage::LibsqlSql, MigrationError> {
let cfg = self
.databases
.get(db)
.ok_or(MigrationError::NotConfigured)?;
let Some(path) = cfg.path.as_deref().filter(|p| !p.as_os_str().is_empty()) else {
return Err(MigrationError::Other(format!(
"libsql database {db:?}: schema migrations need a single-node `path` (the embedded \
file); a remote-sqld `libsql` binding is not a migration target"
)));
};
if let Some(parent) = path.parent() {
if !parent.as_os_str().is_empty() {
std::fs::create_dir_all(parent)
.map_err(|e| MigrationError::Other(format!("libsql database {db:?}: {e}")))?;
}
}
boatramp_storage::LibsqlSql::open_local(path)
.await
.map_err(MigrationError::Sql)
}
fn ledger_insert(
step: &MigrationStep,
ordinal: usize,
effective_hash: &str,
origin: LedgerOrigin,
) -> String {
format!(
"INSERT INTO {ledger} (id, ordinal, content_hash, kind, applied_by) \
VALUES ({id}, {ord}, {hash}, {step_kind}, {origin});",
ledger = libsql_ledger(),
id = sql_quote_literal(&step.id),
ord = ordinal,
hash = sql_quote_literal(effective_hash),
step_kind = sql_quote_literal(step.kind()),
origin = sql_quote_literal(origin.as_str()),
)
}
async fn ensure_ledger(&self, sql: &boatramp_storage::LibsqlSql) -> Result<(), MigrationError> {
use boatramp_core::sql::SqlBackend;
sql.run_script(&format!(
"CREATE TABLE IF NOT EXISTS {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);",
ledger = libsql_ledger()
))
.await
.map_err(MigrationError::Sql)
}
async fn read_applied(
&self,
sql: &boatramp_storage::LibsqlSql,
) -> Result<Vec<AppliedMigration>, MigrationError> {
use boatramp_core::sql::{SqlBackend, SqlValue};
let rows = sql
.run_query(&format!(
"SELECT id, ordinal, content_hash, kind, applied_at, applied_by \
FROM {ledger} ORDER BY ordinal;",
ledger = libsql_ledger()
))
.await
.map_err(MigrationError::Sql)?;
let text = |v: &SqlValue| match v {
SqlValue::Text(s) => s.clone(),
other => format!("{other:?}"),
};
let int = |v: &SqlValue| match v {
SqlValue::Integer(n) => *n,
_ => 0,
};
Ok(rows
.rows
.iter()
.map(|r| {
let origin = match r.get(5) {
Some(SqlValue::Text(s)) if s == LedgerOrigin::Baseline.as_str() => {
LedgerOrigin::Baseline.as_str().to_string()
}
_ => LedgerOrigin::Apply.as_str().to_string(),
};
AppliedMigration {
id: r.first().map(&text).unwrap_or_default(),
ordinal: r.get(1).map(&int).unwrap_or_default(),
content_hash: r.get(2).map(&text).unwrap_or_default(),
kind: r.get(3).map(&text).unwrap_or_default(),
applied_at: r.get(4).map(&text).unwrap_or_default(),
origin,
}
})
.collect())
}
}
#[cfg(feature = "migrate")]
#[async_trait]
impl MigrationSubstrate for LibsqlMigrationRunner {
async fn preflight(
&self,
_project: &str,
db: &str,
) -> Result<Vec<AppliedMigration>, MigrationError> {
self.engine_gate(db)?;
let sql = self.open(db).await?;
self.ensure_ledger(&sql).await?;
self.read_applied(&sql).await
}
async fn apply_substrate_step(
&self,
_project: &str,
db: &str,
step: &MigrationStep,
ordinal: usize,
effective_hash: &str,
) -> Result<SubstrateStepOutcome, MigrationError> {
use boatramp_core::sql::{MigrationAction, SqlBackend};
if !valid_migration_id(&step.id) {
return Ok(SubstrateStepOutcome::Failed(
"invalid migration id (allowed: A-Za-z0-9._-)".to_string(),
));
}
self.engine_gate(db)?;
let sql = self.open(db).await?;
let ledger_insert = Self::ledger_insert(step, ordinal, effective_hash, LedgerOrigin::Apply);
let dialect = boatramp_core::sql::GuardDialect::Sqlite;
let outcome: Result<(), String> = match &step.action {
MigrationAction::Sql {
script,
no_transaction,
} => {
if boatramp_core::sql::script_has_create_extension_in(script, dialect) {
Err(
"a sql step may not CREATE EXTENSION — SQLite loadable extensions are \
host-controlled, never enabled by a migration step"
.to_string(),
)
} else if boatramp_core::sql::script_references_word_in(
script,
LIBSQL_LEDGER_WORD,
dialect,
) {
Err(
"a sql step may not reference the host-owned migration-ledger table"
.to_string(),
)
} else if boatramp_core::sql::script_has_txn_control_in(script, dialect)
&& !*no_transaction
{
Err("a transactional sql step may not contain its own BEGIN/COMMIT/ROLLBACK (it \
would desync the atomic wrapper) — use a no_transaction step to manage the \
transaction yourself"
.to_string())
} else if *no_transaction {
match sql.run_script(script).await {
Ok(()) => sql
.run_script(&ledger_insert)
.await
.map_err(|e| e.to_string()),
Err(e) => Err(e.to_string()),
}
} else {
sql.run_migration_txn(script, &ledger_insert)
.await
.map_err(|e| e.to_string())
}
}
MigrationAction::Extension { name } => Err(format!(
"extension step {name:?} is not supported on libsql/SQLite — SQLite loadable \
extensions are native host-controlled files, never enabled by a migration step; \
install any extension operator-side and use a plain sql step"
)),
MigrationAction::Function { .. } => {
return Err(MigrationError::Other(
"internal: a function step must be invoked by the orchestrator, not the \
substrate"
.to_string(),
))
}
};
Ok(match outcome {
Ok(()) => SubstrateStepOutcome::Applied,
Err(error) => SubstrateStepOutcome::Failed(sanitize_migration_error(&error)),
})
}
async fn record(
&self,
_project: &str,
db: &str,
step: &MigrationStep,
ordinal: usize,
effective_hash: &str,
origin: LedgerOrigin,
) -> Result<(), MigrationError> {
use boatramp_core::sql::SqlBackend;
if !valid_migration_id(&step.id) {
return Err(MigrationError::Other(format!(
"invalid migration id {:?} (allowed: A-Za-z0-9._-)",
step.id
)));
}
self.engine_gate(db)?;
let sql = self.open(db).await?;
self.ensure_ledger(&sql).await?;
sql.run_script(&Self::ledger_insert(step, ordinal, effective_hash, origin))
.await
.map_err(MigrationError::Sql)
}
async fn owner_ddl(
&self,
_project: &str,
db: &str,
) -> Result<Arc<dyn MigrateDdl>, MigrationError> {
self.engine_gate(db)?;
let sql = self.open(db).await?;
Ok(Arc::new(LibsqlDdl { sql }))
}
}
#[cfg(feature = "migrate")]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub(crate) struct LibsqlDdl {
sql: boatramp_storage::LibsqlSql,
}
#[cfg(feature = "migrate")]
impl LibsqlDdl {
fn guard(&self, script: &str) -> Result<(), MigrateDdlError> {
let dialect = boatramp_core::sql::GuardDialect::Sqlite;
if boatramp_core::sql::script_references_word_in(script, LIBSQL_LEDGER_WORD, dialect) {
return Err(MigrateDdlError::LedgerProtected);
}
if boatramp_core::sql::script_has_txn_control_in(script, dialect) {
return Err(MigrateDdlError::TxnControl);
}
Ok(())
}
}
#[cfg(feature = "migrate")]
#[async_trait]
impl MigrateDdl for LibsqlDdl {
async fn exec(&self, script: &str) -> Result<(), MigrateDdlError> {
use boatramp_core::sql::SqlBackend;
self.guard(script)?;
self.sql
.run_script(script)
.await
.map_err(|e| MigrateDdlError::Sql(sanitize_migration_error(&e.to_string())))
}
async fn exec_batch(&self, scripts: Vec<String>) -> Result<(), MigrateDdlError> {
for script in &scripts {
self.exec(script).await?;
}
Ok(())
}
async fn query(&self, sql: &str) -> Result<boatramp_core::sql::SqlRows, MigrateDdlError> {
use boatramp_core::sql::SqlBackend;
self.guard(sql)?;
let rows = self
.sql
.run_query(sql)
.await
.map_err(|e| MigrateDdlError::Sql(sanitize_migration_error(&e.to_string())))?;
if rows.rows.len() > MIGRATE_QUERY_MAX_ROWS {
return Err(MigrateDdlError::Sql(format!(
"query returned {} rows (cap {MIGRATE_QUERY_MAX_ROWS}); add a LIMIT — a migration \
verification query should read a bounded set",
rows.rows.len()
)));
}
Ok(rows)
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub struct DispatchMigrationRunner {
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
sqlx: NodeMigrationRunner,
#[cfg(feature = "migrate")]
libsql: LibsqlMigrationRunner,
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
impl DispatchMigrationRunner {
pub fn new(
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))] op: Arc<NodeOperatorSql>,
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
trusted_extensions: std::collections::BTreeSet<String>,
#[cfg(feature = "migrate")] databases: std::collections::BTreeMap<
String,
crate::config::ExternalDatabaseConfig,
>,
) -> Self {
Self {
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
sqlx: NodeMigrationRunner::new(op, trusted_extensions),
#[cfg(feature = "migrate")]
libsql: LibsqlMigrationRunner::new(databases),
}
}
#[cfg(feature = "migrate")]
fn routes_to_libsql(&self, db: &str) -> bool {
self.libsql.is_libsql(db)
}
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql", feature = "migrate"))]
#[async_trait]
impl MigrationSubstrate for DispatchMigrationRunner {
async fn preflight(
&self,
project: &str,
db: &str,
) -> Result<Vec<AppliedMigration>, MigrationError> {
#[cfg(feature = "migrate")]
if self.routes_to_libsql(db) {
return self.libsql.preflight(project, db).await;
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
{
return self.sqlx.preflight(project, db).await;
}
#[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
{
let _ = (project, db);
Err(MigrationError::NotConfigured)
}
}
async fn apply_substrate_step(
&self,
project: &str,
db: &str,
step: &MigrationStep,
ordinal: usize,
effective_hash: &str,
) -> Result<SubstrateStepOutcome, MigrationError> {
#[cfg(feature = "migrate")]
if self.routes_to_libsql(db) {
return self
.libsql
.apply_substrate_step(project, db, step, ordinal, effective_hash)
.await;
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
{
return self
.sqlx
.apply_substrate_step(project, db, step, ordinal, effective_hash)
.await;
}
#[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
{
let _ = (project, db, step, ordinal, effective_hash);
Err(MigrationError::NotConfigured)
}
}
async fn record(
&self,
project: &str,
db: &str,
step: &MigrationStep,
ordinal: usize,
effective_hash: &str,
origin: LedgerOrigin,
) -> Result<(), MigrationError> {
#[cfg(feature = "migrate")]
if self.routes_to_libsql(db) {
return self
.libsql
.record(project, db, step, ordinal, effective_hash, origin)
.await;
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
{
return self
.sqlx
.record(project, db, step, ordinal, effective_hash, origin)
.await;
}
#[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
{
let _ = (project, db, step, ordinal, effective_hash, origin);
Err(MigrationError::NotConfigured)
}
}
async fn owner_ddl(
&self,
project: &str,
db: &str,
) -> Result<Arc<dyn MigrateDdl>, MigrationError> {
#[cfg(feature = "migrate")]
if self.routes_to_libsql(db) {
return self.libsql.owner_ddl(project, db).await;
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
{
return self.sqlx.owner_ddl(project, db).await;
}
#[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
{
let _ = (project, db);
Err(MigrationError::NotConfigured)
}
}
}
#[cfg(all(test, any(feature = "sql-postgres", feature = "sql-mysql")))]
mod tests {
use super::*;
use async_trait::async_trait;
use boatramp_core::envelope::EnvelopeError;
use boatramp_core::kv::MemoryKv;
struct ReverseEnvelope;
#[async_trait]
impl KeyEnvelope for ReverseEnvelope {
async fn wrap(&self, plaintext: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
Ok(plaintext.iter().rev().copied().collect())
}
async fn unwrap(&self, wrapped: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
Ok(wrapped.iter().rev().copied().collect())
}
}
#[tokio::test]
async fn password_is_generated_once_sealed_and_stable() {
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let creds = ManagedSqlCredentials::new(kv.clone(), Arc::new(ReverseEnvelope));
let pw = creds.password("default", "pg").await.unwrap();
assert_eq!(pw.len(), 64, "32 random bytes, hex-encoded");
assert_eq!(creds.password("default", "pg").await.unwrap(), pw);
let raw = kv
.get("managed-sql-cred/default/pg")
.await
.unwrap()
.unwrap();
assert_ne!(
raw,
pw.as_bytes(),
"the stored blob is sealed, not the password"
);
assert_eq!(
raw.iter().rev().copied().collect::<Vec<u8>>(),
pw.as_bytes()
);
let after_restart = ManagedSqlCredentials::new(kv, Arc::new(ReverseEnvelope));
assert_eq!(after_restart.password("default", "pg").await.unwrap(), pw);
assert_ne!(creds.password("default", "other").await.unwrap(), pw);
}
#[tokio::test]
async fn is_sealed_is_read_only_and_accurate() {
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let creds = ManagedSqlCredentials::new(kv.clone(), Arc::new(ReverseEnvelope));
assert!(!creds.is_sealed("acme", "pg/tenant/owner").await.unwrap());
assert!(
kv.get("managed-sql-cred/acme/pg/tenant/owner")
.await
.unwrap()
.is_none(),
"is_sealed must not create the credential (dry-run purity)"
);
let _ = creds.password("acme", "pg/tenant/owner").await.unwrap();
assert!(creds.is_sealed("acme", "pg/tenant/owner").await.unwrap());
assert!(!creds.is_sealed("acme", "pg/tenant").await.unwrap());
}
#[test]
fn server_env_recipe_per_engine() {
let pg = managed_db_server_env(ExternalSqlKind::Postgres, "analytics", "app", "pw");
assert_eq!(
pg,
vec![
("POSTGRES_USER".into(), "app".into()),
("POSTGRES_PASSWORD".into(), "pw".into()),
("POSTGRES_DB".into(), "analytics".into()),
]
);
let my = managed_db_server_env(ExternalSqlKind::Mysql, "shop", "app", "pw");
assert!(my.contains(&("MYSQL_USER".into(), "app".into())));
assert!(my.contains(&("MYSQL_DATABASE".into(), "shop".into())));
assert!(my.contains(&("MYSQL_ROOT_PASSWORD".into(), "pw".into())));
}
use crate::config::ExternalDatabaseConfig;
use std::collections::BTreeMap;
fn db(
kind: &str,
compute: Option<&str>,
url_env: &str,
pw_env: Option<&str>,
) -> ExternalDatabaseConfig {
ExternalDatabaseConfig {
kind: kind.into(),
url_env: url_env.into(),
compute: compute.map(Into::into),
database: compute.map(|_| "analytics".into()),
user: compute.map(|_| "app".into()),
password_env: pw_env.map(Into::into),
..Default::default()
}
}
#[tokio::test]
async fn managed_db_env_only_covers_managed_workloads() {
let mut dbs = BTreeMap::new();
dbs.insert(
"analytics".to_string(),
db("postgres", Some("pg"), "", None),
);
dbs.insert(
"byo".to_string(),
db("postgres", Some("pg2"), "", Some("PG2_PW")),
);
dbs.insert("ext".to_string(), db("mysql", None, "MYSQL_URL", None));
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let creds = ManagedSqlCredentials::new(kv, Arc::new(ReverseEnvelope));
let env = ManagedDbEnv::from_config(&dbs, creds, ManagedDbPrivilege::default());
assert!(!env.is_empty());
assert_eq!(
env.managed_db_privilege("default", "pg"),
Some(PrivilegeDirective::Rootless { uid: 999, gid: 999 })
);
assert_eq!(env.managed_db_privilege("default", "nope"), None);
let pg = env.managed_db_env("default", "pg").await;
assert!(pg.contains(&("POSTGRES_USER".into(), "app".into())));
assert!(pg.contains(&("POSTGRES_DB".into(), "analytics".into())));
let password = pg
.iter()
.find(|(k, _)| k == "POSTGRES_PASSWORD")
.map(|(_, v)| v.clone())
.expect("password present");
assert_eq!(password.len(), 64, "managed 32-byte hex password");
let pg2 = env.managed_db_env("default", "pg").await;
assert_eq!(pg, pg2);
assert!(env.managed_db_env("default", "pg2").await.is_empty());
assert!(env.managed_db_env("default", "nope").await.is_empty());
}
#[tokio::test]
async fn resolve_spec_prefers_the_longest_matching_single_base() {
use crate::config::TenantIsolation;
let mut pg = db("postgres", Some("pg"), "", None);
pg.tenant = TenantIsolation::Single;
pg.database = Some("appdb".into());
pg.user = Some("app".into());
let mut pg_metrics = db("postgres", Some("pg-metrics"), "", None);
pg_metrics.tenant = TenantIsolation::Single;
pg_metrics.database = Some("metricsdb".into());
pg_metrics.user = Some("metrics".into());
let mut dbs = BTreeMap::new();
dbs.insert("analytics".to_string(), pg);
dbs.insert("metrics".to_string(), pg_metrics);
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let creds = ManagedSqlCredentials::new(kv, Arc::new(ReverseEnvelope));
let env = ManagedDbEnv::from_config(&dbs, creds, ManagedDbPrivilege::default());
let e = env.managed_db_env("acme", "pg-metrics-acme").await;
assert!(
e.contains(&("POSTGRES_DB".into(), "metricsdb".into())),
"longest base (`pg-metrics`) must win over `pg`: {e:?}"
);
assert!(e.contains(&("POSTGRES_USER".into(), "metrics".into())));
let e = env.managed_db_env("acme", "pg-acme").await;
assert!(e.contains(&("POSTGRES_DB".into(), "appdb".into())));
assert!(e.contains(&("POSTGRES_USER".into(), "app".into())));
assert_eq!(
env.managed_db_privilege("acme", "pg-metrics-acme"),
Some(PrivilegeDirective::Rootless { uid: 999, gid: 999 })
);
}
use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
struct NullStorage;
#[async_trait]
impl Storage for NullStorage {
async fn get(&self, _: &str) -> Result<GetObject, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn get_range(
&self,
_: &str,
_: u64,
_: Option<u64>,
) -> Result<GetObject, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn put(
&self,
_: &str,
_: ByteStream,
_: PutMeta,
) -> Result<ObjectMeta, StorageError> {
Err(StorageError::unsupported("null"))
}
async fn head(&self, _: &str) -> Result<ObjectMeta, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn delete(&self, _: &str) -> Result<(), StorageError> {
Ok(())
}
async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
Ok(Vec::new())
}
}
fn replica(
workload: &str,
replica: u32,
host: &str,
port: u16,
healthy: bool,
phase: ReplicaPhase,
) -> boatramp_core::compute::ObservedInstance {
use boatramp_core::compute::{Endpoint, InstanceHandle, Scheme};
boatramp_core::compute::ObservedInstance {
handle: InstanceHandle {
project: "default".into(),
workload: workload.into(),
replica,
backend_ref: String::new(),
},
node: 0,
backend: "fake".into(),
endpoint: Endpoint {
scheme: Scheme::Http,
host: host.into(),
port,
},
region: None,
healthy,
started_at: None,
phase,
snapshot: None,
}
}
#[tokio::test]
async fn endpoint_resolver_returns_only_healthy_running_replicas() {
let deploy = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
let p = ProjectRef::DEFAULT;
deploy
.set_replica_state(
p,
&replica("pg", 0, "10.0.0.1", 5432, true, ReplicaPhase::Running),
)
.await
.unwrap();
deploy
.set_replica_state(
p,
&replica("pg", 1, "10.0.0.2", 5432, true, ReplicaPhase::Running),
)
.await
.unwrap();
deploy
.set_replica_state(
p,
&replica("pg", 2, "10.0.0.3", 5432, false, ReplicaPhase::Running),
)
.await
.unwrap();
deploy
.set_replica_state(
p,
&replica("pg", 3, "10.0.0.4", 5432, false, ReplicaPhase::Zero),
)
.await
.unwrap();
let resolver = DeployEndpointResolver::new(deploy, "default");
let mut eps = resolver.endpoints("pg").await.unwrap();
eps.sort();
assert_eq!(
eps,
vec![
("10.0.0.1".to_string(), 5432),
("10.0.0.2".to_string(), 5432)
],
"only the healthy running replicas, unhealthy + Zero filtered out"
);
assert!(resolver.endpoints("absent").await.unwrap().is_empty());
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[tokio::test]
async fn auto_register_shared_registers_one_server_idempotently_and_never_clobbers() {
use crate::config::TenantIsolation;
use boatramp_core::compute::{ComputeWorkload, PlacementConstraints};
let deploy = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
let p = ProjectRef::DEFAULT;
let mut shared = db("postgres", Some("pg"), "", None);
shared.tenant = TenantIsolation::Shared;
let mut dbs = BTreeMap::new();
dbs.insert("analytics".to_string(), shared);
dbs.insert(
"byo".to_string(),
db("postgres", Some("byopg"), "", Some("PW")),
);
dbs.insert("ext".to_string(), db("mysql", None, "MYSQL_URL", None));
auto_register_managed_db_workloads(&deploy, &dbs).await;
let wl = deploy
.get_compute_workload(p, "pg")
.await
.unwrap()
.expect("shared server workload `pg` auto-registered");
assert_eq!(wl.replicas, 1);
assert!(!wl.active.is_empty(), "an active spec hash was stored");
assert!(
deploy
.get_compute_workload(p, "byopg")
.await
.unwrap()
.is_none(),
"a BYO-credential DB is not auto-registered"
);
auto_register_managed_db_workloads(&deploy, &dbs).await;
let wl2 = deploy.get_compute_workload(p, "pg").await.unwrap().unwrap();
assert_eq!(wl2.active, wl.active, "re-run is a no-op");
let operator = ComputeWorkload {
version: 1,
name: "pg".to_string(),
active: "operatorspec".to_string(),
replicas: 3,
placement: PlacementConstraints::default(),
};
deploy.set_compute_workload(p, &operator).await.unwrap();
auto_register_managed_db_workloads(&deploy, &dbs).await;
let after = deploy.get_compute_workload(p, "pg").await.unwrap().unwrap();
assert_eq!(
after.replicas, 3,
"auto-register must not overwrite the operator's workload"
);
assert_eq!(after.active, "operatorspec");
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
async fn seed_project_with_site(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
project: &str,
site: &str,
) {
deploy
.put_project(&boatramp_core::project::Project {
version: 1,
name: project.to_string(),
created_at: 0,
meta: Default::default(),
config: Default::default(),
secrets_ref: None,
})
.await
.expect("seed the project pointer");
let key = format!("project/{project}/current/{site}");
kv.put(&key, b"deadbeef".to_vec())
.await
.expect("seed a current site deployment pointer");
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[tokio::test]
async fn auto_register_single_registers_nothing_at_boot() {
use crate::config::{TenantIsolation, TenantScope};
use crate::tenant_sql::tenant_key;
use boatramp_storage::tenant_provision::sanitize_ident;
let store_kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let deploy = DeployStore::new(Arc::new(NullStorage), store_kv.clone());
seed_project_with_site(&deploy, &store_kv, "construens", "app").await;
let mut single = db("postgres", Some("pg"), "", None);
single.tenant = TenantIsolation::Single; single.tenant_scope = TenantScope::Project;
let mut dbs = BTreeMap::new();
dbs.insert("main".to_string(), single);
auto_register_managed_db_workloads(&deploy, &dbs).await;
let (raw, _is_default) = tenant_key(TenantScope::Project, "construens", "");
let derived = format!("pg-{}", sanitize_ident(&raw));
assert!(
deploy
.get_compute_workload(ProjectRef::new("construens"), &derived)
.await
.unwrap()
.is_none(),
"a Single binding must NOT boot-warm a per-tenant `pg-<ident>` (that is the lazy \
resolve's job on first `sql` use)"
);
assert!(
deploy
.get_compute_workload(ProjectRef::DEFAULT, "pg")
.await
.unwrap()
.is_none(),
"a Single binding must NOT register a tenant-blind bare `pg`/`default`"
);
assert!(
deploy
.list_compute_workloads_all()
.await
.unwrap()
.is_empty(),
"a Single binding registers no managed workload at boot"
);
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[tokio::test]
async fn auto_register_single_static_only_project_gets_no_db_at_boot() {
use crate::config::{TenantIsolation, TenantScope};
let store_kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let deploy = DeployStore::new(Arc::new(NullStorage), store_kv.clone());
seed_project_with_site(&deploy, &store_kv, "default", "www").await;
let mut single = db("postgres", Some("pg"), "", None);
single.tenant = TenantIsolation::Single;
single.tenant_scope = TenantScope::Project;
let mut dbs = BTreeMap::new();
dbs.insert("main".to_string(), single);
auto_register_managed_db_workloads(&deploy, &dbs).await;
assert!(
deploy
.list_compute_workloads_all()
.await
.unwrap()
.is_empty(),
"a static-only project must not get a spurious managed `pg` at boot"
);
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[test]
fn operator_target_single_targets_per_tenant_workload_and_cred() {
use crate::config::{TenantIsolation, TenantScope};
use crate::tenant_sql::tenant_key;
use boatramp_storage::tenant_provision::sanitize_ident;
let mut single = db("postgres", Some("pg"), "", None);
single.tenant = TenantIsolation::Single;
single.tenant_scope = TenantScope::Project;
single.database = Some("appdb".into());
single.user = Some("app".into());
let (raw, is_default) = tenant_key(TenantScope::Project, "construens", "");
assert!(!is_default);
let ident = sanitize_ident(&raw);
let derived = format!("pg-{ident}");
let t = operator_target(&single, "construens", "main").unwrap();
assert_eq!(
t.workload, derived,
"targets the per-tenant workload, not bare `pg`"
);
assert_eq!(t.database, "appdb");
assert_eq!(t.user, "app");
assert_eq!(
t.endpoint_project, "construens",
"a Single per-tenant workload's replicas live under its project"
);
assert_eq!(t.cred_project, "construens");
assert_eq!(t.cred_workload, derived);
assert_ne!(
t.cred_workload, "pg",
"never the bare tenant-blind workload"
);
let d = operator_target(&single, "default", "main").unwrap();
assert_eq!(d.workload, "pg");
assert_eq!(d.cred_project, "default");
assert_eq!(d.cred_workload, "pg");
assert_eq!(d.endpoint_project, "default");
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[test]
fn operator_target_shared_uses_superuser_cred_against_tenant_db() {
use crate::config::{TenantIsolation, TenantScope};
use crate::tenant_sql::tenant_key;
use boatramp_core::project::DEFAULT_PROJECT;
use boatramp_storage::tenant_provision::{sanitize_ident, tenant_db_name};
let mut shared = db("postgres", Some("pg"), "", None);
shared.tenant = TenantIsolation::Shared;
shared.tenant_scope = TenantScope::Project;
shared.database = Some("appdb".into());
shared.user = Some("postgres".into());
let (raw, _) = tenant_key(TenantScope::Project, "construens", "");
let ident = sanitize_ident(&raw);
let t = operator_target(&shared, "construens", "main").unwrap();
assert_eq!(t.workload, "pg");
assert_eq!(t.database, tenant_db_name("appdb", &ident));
assert_eq!(t.user, "postgres");
assert_eq!(t.cred_project, DEFAULT_PROJECT);
assert_eq!(t.cred_workload, "pg");
assert_eq!(t.endpoint_project, DEFAULT_PROJECT);
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[test]
fn operator_target_site_scoped_errors_clearly() {
use crate::config::{TenantIsolation, TenantScope};
let mut site = db("postgres", Some("pg"), "", None);
site.tenant = TenantIsolation::Single;
site.tenant_scope = TenantScope::Site;
let err = operator_target(&site, "construens", "main").unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("site-scoped"),
"the error explains a site-scoped DB needs a site: {msg}"
);
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
fn op_for(kind: &str, migration_url_env: Option<&str>) -> Arc<NodeOperatorSql> {
use crate::config::{TenantIsolation, TenantScope};
let mut databases = BTreeMap::new();
databases.insert(
"main".to_string(),
ExternalDatabaseConfig {
kind: kind.to_string(),
url_env: "RUNTIME_URL_ENV".to_string(),
migration_url_env: migration_url_env.map(Into::into),
pool_max: Some(2),
read_only: false,
connect_timeout_secs: Some(5),
tenant: TenantIsolation::Shared,
tenant_scope: TenantScope::Project,
..Default::default()
},
);
Arc::new(NodeOperatorSql::new(
databases,
Arc::new(MemoryKv::new()),
None,
DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new())),
))
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
fn runner_over(op: Arc<NodeOperatorSql>) -> NodeMigrationRunner {
NodeMigrationRunner::new(op, std::collections::BTreeSet::new())
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[test]
fn engine_gate_admits_postgres_and_mysql() {
let pg = runner_over(op_for("postgres", None));
assert!(matches!(
pg.engine_gate("main"),
Ok(ExternalSqlKind::Postgres)
));
let my = runner_over(op_for("mysql", Some("X")));
assert!(matches!(my.engine_gate("main"), Ok(ExternalSqlKind::Mysql)));
assert!(matches!(
pg.engine_gate("absent"),
Err(MigrationError::NotConfigured)
));
}
#[cfg(feature = "sql-mysql")]
#[tokio::test]
async fn mysql_refuses_without_a_distinct_ddl_user() {
let sub = runner_over(op_for("mysql", None));
let err = sub.preflight("default", "main").await.unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("migration_url_env") && msg.contains("distinct"),
"MySQL migrate must refuse without a distinct DDL identity: {msg}"
);
}
#[cfg(feature = "sql-mysql")]
#[tokio::test]
async fn mysql_refuses_when_ddl_url_equals_runtime_url() {
std::env::set_var("RUNTIME_URL_ENV", "mysql://app:pw@localhost:3306/appdb");
std::env::set_var(
"MIGRATE_URL_ENV_SAME",
"mysql://app:pw@localhost:3306/appdb",
);
let sub = runner_over(op_for("mysql", Some("MIGRATE_URL_ENV_SAME")));
let err = sub.preflight("default", "main").await.unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("SAME connection") || msg.contains("distinct"),
"a DDL url identical to the runtime url must be refused: {msg}"
);
std::env::remove_var("RUNTIME_URL_ENV");
std::env::remove_var("MIGRATE_URL_ENV_SAME");
}
#[cfg(feature = "sql-mysql")]
#[tokio::test]
async fn mysql_ddl_url_env_unset_is_a_clear_error() {
std::env::remove_var("MIGRATE_URL_ENV_MISSING");
let sub = runner_over(op_for("mysql", Some("MIGRATE_URL_ENV_MISSING")));
let err = sub.preflight("default", "main").await.unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("MIGRATE_URL_ENV_MISSING") && msg.contains("unset"),
"an unset migration url env must be reported clearly: {msg}"
);
}
#[cfg(feature = "sql-mysql")]
fn mysql_runner_with_env(runtime_env: &str, migrate_env: &str) -> NodeMigrationRunner {
use crate::config::{TenantIsolation, TenantScope};
let mut databases = BTreeMap::new();
databases.insert(
"main".to_string(),
ExternalDatabaseConfig {
kind: "mysql".to_string(),
url_env: runtime_env.to_string(),
migration_url_env: Some(migrate_env.to_string()),
pool_max: Some(2),
read_only: false,
connect_timeout_secs: Some(5),
tenant: TenantIsolation::Shared,
tenant_scope: TenantScope::Project,
..Default::default()
},
);
runner_over(Arc::new(NodeOperatorSql::new(
databases,
Arc::new(MemoryKv::new()),
None,
DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new())),
)))
}
#[cfg(feature = "sql-mysql")]
#[tokio::test]
async fn mysql_refuses_same_username_even_when_dsn_strings_differ() {
let same_user_pairs = [
(
"mysql://app:pw@host:3306/db",
"mysql://app:pw@host:3306/db?charset=utf8",
),
("mysql://app:pw@host:3306/db", "mysql://app:pw@host/db"),
("mysql://app:pw@host:3306/db", "mysql://app:pw@host:3306/"),
(
"mysql://app:pw@host/db",
"mysql://app:other@host:3306/otherdb",
),
];
for (i, (runtime, ddl)) in same_user_pairs.iter().enumerate() {
let rvar = format!("HIGH1_RT_{i}");
let mvar = format!("HIGH1_DDL_{i}");
std::env::set_var(&rvar, runtime);
std::env::set_var(&mvar, ddl);
let sub = mysql_runner_with_env(&rvar, &mvar);
let err = sub.preflight("default", "main").await.unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("SAME MySQL user") || msg.contains("SAME connection"),
"same-username DSNs {runtime:?} vs {ddl:?} must be refused: {msg}"
);
std::env::remove_var(&rvar);
std::env::remove_var(&mvar);
}
}
#[cfg(feature = "sql-mysql")]
#[tokio::test]
async fn mysql_allows_a_distinct_ddl_username() {
std::env::set_var("HIGH1_RT_OK", "mysql://app:pw@127.0.0.1:1/db");
std::env::set_var("HIGH1_DDL_OK", "mysql://root_migrate:pw@127.0.0.1:1/db");
let sub = mysql_runner_with_env("HIGH1_RT_OK", "HIGH1_DDL_OK");
let err = sub.preflight("default", "main").await.unwrap_err();
let msg = err.to_string();
assert!(
!msg.contains("SAME MySQL user") && !msg.contains("SAME connection"),
"a distinct DDL username must pass the distinctness barrier: {msg}"
);
std::env::remove_var("HIGH1_RT_OK");
std::env::remove_var("HIGH1_DDL_OK");
}
#[cfg(feature = "sql-mysql")]
#[tokio::test]
async fn mysql_refuses_compute_backed_managed_migration() {
use crate::config::{TenantIsolation, TenantScope};
let mut databases = BTreeMap::new();
let mut cfg = db("mysql", Some("my-compute"), "", None);
cfg.tenant = TenantIsolation::Shared;
cfg.tenant_scope = TenantScope::Project;
cfg.migration_url_env = Some("SOME_DDL_URL".to_string());
databases.insert("main".to_string(), cfg);
let sub = runner_over(Arc::new(NodeOperatorSql::new(
databases,
Arc::new(MemoryKv::new()),
None,
DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new())),
)));
let err = sub.preflight("default", "main").await.unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("compute-backed managed MySQL migration is not supported"),
"compute-backed managed MySQL migrate must be refused fail-closed: {msg}"
);
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[test]
fn ledger_name_is_dialect_quoted() {
assert_eq!(
NodeMigrationRunner::ledger(ExternalSqlKind::Postgres),
"\"boatramp_migrations\".\"schema_migrations\""
);
assert_eq!(
NodeMigrationRunner::ledger(ExternalSqlKind::Mysql),
"`boatramp_migrations`.`schema_migrations`"
);
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[test]
fn ledger_insert_quotes_all_values() {
use boatramp_core::sql::{MigrationAction, MigrationStep};
let step = MigrationStep {
id: "0001_init".to_string(),
action: MigrationAction::Sql {
script: "CREATE TABLE t (id int)".to_string(),
no_transaction: false,
},
};
let sql = NodeMigrationRunner::ledger_insert(
ExternalSqlKind::Mysql,
&step,
0,
"abc123",
LedgerOrigin::Apply,
);
assert!(sql.contains("`boatramp_migrations`.`schema_migrations`"));
assert!(sql.contains("'0001_init'"));
assert!(sql.contains("'abc123'"));
assert!(sql.contains("'sql'"));
assert!(sql.contains("'apply'"));
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[test]
fn guard_dialect_follows_engine() {
use boatramp_core::sql::GuardDialect;
assert_eq!(guard_dialect(ExternalSqlKind::Mysql), GuardDialect::Mysql);
assert_eq!(
guard_dialect(ExternalSqlKind::Postgres),
GuardDialect::Postgres
);
assert!(!mentions_txn_control(
"CREATE TABLE t (id int) # COMMIT",
ExternalSqlKind::Mysql
));
assert!(mentions_txn_control(
"DROP TABLE t; COMMIT",
ExternalSqlKind::Mysql
));
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
#[test]
fn mysql_partial_apply_notes_are_marked_and_postgres_passthrough() {
let my_mid = mysql_multi_ddl_note(ExternalSqlKind::Mysql, "0002_x", "boom");
assert!(my_mid.contains("PARTIALLY APPLIED") && my_mid.contains("0002_x"));
assert!(my_mid.contains("boom"));
let my_ledger = mysql_partial_apply_note(ExternalSqlKind::Mysql, "0002_x", "ledger down");
assert!(my_ledger.contains("PARTIALLY APPLIED") && my_ledger.contains("unrecorded"));
assert_eq!(
mysql_multi_ddl_note(ExternalSqlKind::Postgres, "0002_x", "boom"),
"boom"
);
assert_eq!(
mysql_partial_apply_note(ExternalSqlKind::Postgres, "0002_x", "boom"),
"boom"
);
}
#[cfg(feature = "migrate")]
#[test]
fn dispatch_routes_by_engine_kind() {
let mut databases = BTreeMap::new();
databases.insert("lite".to_string(), db("libsql", None, "", None));
databases.get_mut("lite").unwrap().path = Some("/tmp/lite.db".into());
databases.insert("pg".to_string(), db("postgres", None, "PG_URL", None));
let op = Arc::new(NodeOperatorSql::new(
databases.clone(),
Arc::new(MemoryKv::new()),
None,
DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new())),
));
let dispatch =
DispatchMigrationRunner::new(op, std::collections::BTreeSet::new(), databases);
assert!(
dispatch.routes_to_libsql("lite"),
"a libsql binding routes to the libsql substrate"
);
assert!(
!dispatch.routes_to_libsql("pg"),
"a postgres binding does NOT route to the libsql substrate"
);
assert!(
!dispatch.routes_to_libsql("absent"),
"an unknown name does not route to libsql (the sqlx path then returns NotConfigured)"
);
}
}
#[cfg(all(test, feature = "migrate"))]
mod libsql_migrate_tests {
use super::*;
use boatramp_core::sql::{
GuardDialect, LedgerOrigin, MigrateDdl, MigrateDdlError, MigrationAction, MigrationStep,
MigrationSubstrate, SubstrateStepOutcome,
};
use std::collections::BTreeMap;
fn sql_step(id: &str, script: &str) -> MigrationStep {
MigrationStep {
id: id.to_string(),
action: MigrationAction::Sql {
script: script.to_string(),
no_transaction: false,
},
}
}
fn runner_and_db(tag: &str) -> (LibsqlMigrationRunner, String) {
let dir = std::env::temp_dir().join(format!(
"boatramp-migrate-libsql-{tag}-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("app.db");
let mut databases = BTreeMap::new();
databases.insert(
"app".to_string(),
crate::config::ExternalDatabaseConfig {
kind: "libsql".to_string(),
path: Some(path),
..Default::default()
},
);
(LibsqlMigrationRunner::new(databases), "app".to_string())
}
#[test]
fn libsql_ledger_name_is_a_single_prefixed_table() {
assert_eq!(libsql_ledger(), "\"boatramp_migrations_schema_migrations\"");
}
#[test]
fn engine_gate_admits_only_libsql_kinds() {
assert!(kind_is_libsql("libsql"));
assert!(kind_is_libsql("LibSQL"));
assert!(kind_is_libsql("sqlite"));
assert!(kind_is_libsql("sqlite3"));
assert!(!kind_is_libsql("postgres"));
assert!(!kind_is_libsql("mysql"));
let mut databases = BTreeMap::new();
databases.insert(
"app".to_string(),
crate::config::ExternalDatabaseConfig {
kind: "libsql".to_string(),
path: Some("/tmp/x.db".into()),
..Default::default()
},
);
databases.insert(
"pg".to_string(),
crate::config::ExternalDatabaseConfig {
kind: "postgres".to_string(),
url_env: "X".into(),
..Default::default()
},
);
let sub = LibsqlMigrationRunner::new(databases);
assert!(sub.is_libsql("app"));
assert!(!sub.is_libsql("pg"));
assert!(!sub.is_libsql("absent"));
assert!(matches!(
sub.engine_gate("pg"),
Err(MigrationError::NotConfigured)
));
assert!(matches!(
sub.engine_gate("absent"),
Err(MigrationError::NotConfigured)
));
assert!(sub.engine_gate("app").is_ok());
}
#[test]
fn libsql_ledger_insert_quotes_all_values() {
let step = sql_step("0001_init", "CREATE TABLE t (id integer)");
let sql = LibsqlMigrationRunner::ledger_insert(&step, 0, "abc123", LedgerOrigin::Apply);
assert!(sql.contains("\"boatramp_migrations_schema_migrations\""));
assert!(sql.contains("'0001_init'"));
assert!(sql.contains("'abc123'"));
assert!(sql.contains("'sql'"));
assert!(sql.contains("'apply'"));
}
#[test]
fn libsql_guards_use_the_sqlite_dialect() {
assert!(boatramp_core::sql::script_has_txn_control_in(
"DROP TABLE t; COMMIT",
GuardDialect::Sqlite
));
assert!(!boatramp_core::sql::script_has_txn_control_in(
"CREATE TABLE t (id int) -- COMMIT",
GuardDialect::Sqlite
));
assert!(boatramp_core::sql::script_references_word_in(
"SELECT * FROM [boatramp_migrations_schema_migrations]",
LIBSQL_LEDGER_WORD,
GuardDialect::Sqlite
));
assert!(boatramp_core::sql::script_references_word_in(
"DROP TABLE `boatramp_migrations_schema_migrations`",
LIBSQL_LEDGER_WORD,
GuardDialect::Sqlite
));
let dir =
std::env::temp_dir().join(format!("boatramp-libsql-guard-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let sql = futures_lite_block_on(boatramp_storage::LibsqlSql::open_local(dir.join("g.db")))
.unwrap();
let ddl = LibsqlDdl { sql };
assert!(matches!(
ddl.guard("SELECT * FROM `boatramp_migrations_schema_migrations`"),
Err(MigrateDdlError::LedgerProtected)
));
assert!(matches!(
ddl.guard("BEGIN; CREATE TABLE x(i int); COMMIT"),
Err(MigrateDdlError::TxnControl)
));
assert!(ddl.guard("CREATE TABLE x (i integer)").is_ok());
}
fn futures_lite_block_on<F: std::future::Future>(f: F) -> F::Output {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(f)
}
#[tokio::test]
#[ignore = "opens a real embedded libsql file; run unignored on the host toolchain (static-musl segfaults in libsql)"]
async fn owner_ddl_is_host_mediated_guest_holds_no_credential() {
let (sub, db) = runner_and_db("nocred");
let ddl: Arc<dyn MigrateDdl> = sub.owner_ddl("default", &db).await.unwrap();
ddl.exec("CREATE TABLE proof (n integer)").await.unwrap();
ddl.exec("INSERT INTO proof (n) VALUES (1)").await.unwrap();
let rows = ddl.query("SELECT n FROM proof").await.unwrap();
assert_eq!(rows.rows.len(), 1);
assert!(matches!(
ddl.exec("SELECT * FROM boatramp_migrations_schema_migrations")
.await
.unwrap_err(),
MigrateDdlError::LedgerProtected
));
}
#[tokio::test]
#[ignore = "opens a real embedded libsql file; run unignored on the host toolchain (static-musl segfaults in libsql)"]
async fn sql_step_rolls_back_cleanly_on_failure() {
let (sub, db) = runner_and_db("rollback");
assert!(sub.preflight("default", &db).await.unwrap().is_empty());
let bad = sql_step(
"0001_atomic",
"CREATE TABLE widget (id integer primary key); \
CREATE TABLE widget (id integer primary key)", );
let eff = bad.content_hash();
match sub
.apply_substrate_step("default", &db, &bad, 0, &eff)
.await
.unwrap()
{
SubstrateStepOutcome::Failed(_) => {}
other => panic!("expected a failed step, got {other:?}"),
}
let applied = sub.preflight("default", &db).await.unwrap();
assert!(!applied.iter().any(|a| a.id == "0001_atomic"));
let good = sql_step(
"0001_atomic",
"CREATE TABLE widget (id integer primary key)",
);
assert!(matches!(
sub.apply_substrate_step("default", &db, &good, 0, &good.content_hash())
.await
.unwrap(),
SubstrateStepOutcome::Applied
));
let applied = sub.preflight("default", &db).await.unwrap();
assert_eq!(applied.len(), 1);
assert_eq!(applied[0].id, "0001_atomic");
assert_eq!(applied[0].origin, "apply");
assert_eq!(
applied[0].content_hash,
good.content_hash(),
"the ledger records the effective hash"
);
}
#[tokio::test]
#[ignore = "opens a real embedded libsql file; run unignored on the host toolchain (static-musl segfaults in libsql)"]
async fn extension_step_is_refused_on_libsql() {
let (sub, db) = runner_and_db("ext");
let step = MigrationStep {
id: "0001_ext".to_string(),
action: MigrationAction::Extension {
name: "spellfix".to_string(),
},
};
match sub
.apply_substrate_step("default", &db, &step, 0, &step.content_hash())
.await
.unwrap()
{
SubstrateStepOutcome::Failed(msg) => {
assert!(
msg.contains("libsql") || msg.contains("SQLite"),
"extension refused on libsql: {msg}"
);
}
other => panic!("expected extension refused, got {other:?}"),
}
let raw = sql_step("0001_rawext", "CREATE EXTENSION IF NOT EXISTS whatever");
assert!(matches!(
sub.apply_substrate_step("default", &db, &raw, 0, &raw.content_hash())
.await
.unwrap(),
SubstrateStepOutcome::Failed(_)
));
}
#[tokio::test]
#[ignore = "opens a real embedded libsql file; run unignored on the host toolchain (static-musl segfaults in libsql)"]
async fn txn_control_and_ledger_reference_are_refused_under_sqlite_dialect() {
let (sub, db) = runner_and_db("guards");
let txn = sql_step("0001_txn", "BEGIN; CREATE TABLE x (i int); COMMIT");
assert!(matches!(
sub.apply_substrate_step("default", &db, &txn, 0, &txn.content_hash())
.await
.unwrap(),
SubstrateStepOutcome::Failed(_)
));
let ledger = sql_step(
"0001_led",
"INSERT INTO boatramp_migrations_schema_migrations (id) VALUES ('x')",
);
assert!(matches!(
sub.apply_substrate_step("default", &db, &ledger, 0, &ledger.content_hash())
.await
.unwrap(),
SubstrateStepOutcome::Failed(_)
));
}
#[tokio::test]
#[ignore = "opens a real embedded libsql file; run unignored on the host toolchain (static-musl segfaults in libsql)"]
async fn baseline_records_without_running() {
let (sub, db) = runner_and_db("baseline");
let baselined = sql_step("0001_baselined", "CREATE TABLE never_run (id integer)");
sub.record(
"default",
&db,
&baselined,
0,
&baselined.content_hash(),
LedgerOrigin::Baseline,
)
.await
.unwrap();
let applied = sub.preflight("default", &db).await.unwrap();
let row = applied
.iter()
.find(|a| a.id == "0001_baselined")
.expect("baselined row present");
assert_eq!(row.origin, "baseline");
let good = sql_step("0002_after", "CREATE TABLE never_run (id integer)");
assert!(matches!(
sub.apply_substrate_step("default", &db, &good, 1, &good.content_hash())
.await
.unwrap(),
SubstrateStepOutcome::Applied
));
}
}