use std::collections::HashMap;
use std::sync::Arc;
use async_trait::async_trait;
use boatramp_core::compute::{ManagedDbEnvResolver, PrivilegeDirective, ReplicaPhase};
use crate::config::ManagedDbPrivilege;
use boatramp_core::deploy::DeployStore;
use boatramp_core::envelope::KeyEnvelope;
use boatramp_core::kv::KvStore;
use boatramp_core::project::ProjectRef;
use boatramp_core::sql::SqlError;
use boatramp_storage::sql_compute::{ComputeEndpointResolver, ReplicaDiag};
use boatramp_storage::ExternalSqlKind;
#[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_attr(not(feature = "handlers"), allow(dead_code))]
pub struct ManagedSqlCredentials {
kv: Arc<dyn KvStore>,
envelope: Arc<dyn KeyEnvelope>,
}
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 delete(&self, project: &str, workload: &str) -> Result<(), String> {
self.kv
.delete(&Self::key(project, workload))
.await
.map_err(|e| e.to_string())
}
}
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
struct ManagedDbSpec {
kind: ExternalSqlKind,
database: String,
user: String,
tenant: crate::config::TenantIsolation,
}
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
pub struct ManagedDbEnv {
dbs: HashMap<String, ManagedDbSpec>,
creds: ManagedSqlCredentials,
privilege: ManagedDbPrivilege,
}
#[cfg_attr(not(feature = "handlers"), allow(dead_code))]
fn managed_db_default_ids(_kind: ExternalSqlKind) -> (u32, u32) {
(999, 999)
}
#[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()
}
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)
}
}
#[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_attr(not(feature = "handlers"), allow(dead_code))]
pub struct DeployEndpointResolver {
deploy: DeployStore,
project: String,
}
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(),
}
}
}
#[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,
}
}
async fn backend_for(
&self,
project: &str,
db: &str,
) -> Result<Arc<dyn boatramp_core::sql::SqlBackend>, SqlError> {
use boatramp_storage::sql_compute::ComputeResolvedSqlBackend;
use boatramp_storage::sql_sqlx::{connect, ExternalSqlOptions};
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 = 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"))]
#[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(test)]
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);
}
#[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}"
);
}
}