#![cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use std::collections::BTreeMap;
use std::sync::Arc;
use async_trait::async_trait;
use boatramp_core::compute::{ApplyDatabase, ApplyDatabaseKind, apply_db_caps};
use boatramp_core::deploy::DeployStore;
use boatramp_core::envelope::KeyEnvelope;
use boatramp_core::kv::KvStore;
use boatramp_core::project::ProjectRef;
use boatramp_core::sql::{DeclareError, ManagedDbDeclare};
use crate::config::{ExternalDatabaseConfig, TenantIsolation, TenantScope};
pub fn derived_workload(project: &str, name: &str) -> String {
boatramp_storage::tenant_provision::derived_managed_db_workload(project, name)
}
pub fn lower(project: &str, db: &ApplyDatabase) -> ExternalDatabaseConfig {
let workload = derived_workload(project, &db.name);
let kind = match db.kind {
ApplyDatabaseKind::Postgres => "postgres",
ApplyDatabaseKind::Mysql => "mysql",
};
let (_, _, volume_size_mib) = db.size.resources();
ExternalDatabaseConfig {
kind: kind.to_string(),
url_env: String::new(),
read_url_env: None,
migration_url_env: None,
image: None,
path: None,
password_env: None,
compute: Some(workload),
database: Some(db.name.clone()),
user: Some(db.name.clone()),
pool_max: db.pool_max.map(|n| n.min(apply_db_caps::MAX_POOL)),
connect_timeout_secs: db
.connect_timeout_secs
.map(|n| n.min(apply_db_caps::MAX_CONNECT_TIMEOUT_SECS)),
startup_grace_secs: db
.startup_grace_secs
.map(|n| n.min(apply_db_caps::MAX_STARTUP_GRACE_SECS)),
volume_size_mib: Some(volume_size_mib),
read_only: db.read_only,
allow_preview: false,
tenant: match db.tenant {
boatramp_core::compute::ApplyDatabaseTenant::Single => TenantIsolation::Single,
boatramp_core::compute::ApplyDatabaseTenant::Shared => TenantIsolation::Shared,
},
tenant_scope: match db.tenant_scope {
boatramp_core::compute::ApplyDatabaseScope::Project => TenantScope::Project,
boatramp_core::compute::ApplyDatabaseScope::Site => TenantScope::Site,
},
rls_session: db.rls_session,
tenant_guc: db.tenant_guc.clone(),
session_guc: db.session_guc.clone(),
tenant_all_marker: db.tenant_all_marker.clone(),
}
}
pub async fn resolve_binding(
static_dbs: &BTreeMap<String, ExternalDatabaseConfig>,
deploy: &DeployStore,
project: &str,
name: &str,
) -> Result<Option<ExternalDatabaseConfig>, DeclareError> {
let in_static = static_dbs.contains_key(name);
let declared = deploy
.get_project_database(ProjectRef::new(project), name)
.await
.map_err(|e| DeclareError::Other(e.to_string()))?;
match (in_static, declared) {
(true, Some(_)) => Err(DeclareError::DaemonConflict(name.to_string())),
(true, None) => Ok(static_dbs.get(name).cloned()),
(false, Some(db)) => Ok(Some(lower(project, &db))),
(false, None) => Ok(None),
}
}
pub struct NodeManagedDbDeclare {
static_dbs: BTreeMap<String, ExternalDatabaseConfig>,
deploy: DeployStore,
kv: Arc<dyn KvStore>,
envelope: Arc<dyn KeyEnvelope>,
quota: DeclareQuota,
}
#[derive(Debug, Clone, Copy)]
pub struct DeclareQuota {
pub max_databases: usize,
pub max_volume_mib: u64,
}
impl Default for DeclareQuota {
fn default() -> Self {
Self {
max_databases: apply_db_caps::DEFAULT_MAX_DECLARED_DATABASES_PER_PROJECT,
max_volume_mib: apply_db_caps::DEFAULT_MAX_DECLARED_VOLUME_MIB,
}
}
}
impl DeclareQuota {
pub fn from_config(max_databases: Option<usize>, max_volume_mib: Option<u64>) -> Self {
let d = Self::default();
Self {
max_databases: max_databases.unwrap_or(d.max_databases),
max_volume_mib: max_volume_mib.unwrap_or(d.max_volume_mib),
}
}
}
impl NodeManagedDbDeclare {
pub fn new(
static_dbs: BTreeMap<String, ExternalDatabaseConfig>,
deploy: DeployStore,
kv: Arc<dyn KvStore>,
envelope: Arc<dyn KeyEnvelope>,
quota: DeclareQuota,
) -> Self {
Self {
static_dbs,
deploy,
kv,
envelope,
quota,
}
}
async fn enforce_quota(
&self,
project: &str,
name: &str,
incoming: &ApplyDatabase,
) -> Result<(), DeclareError> {
let existing = self
.deploy
.list_project_databases(ProjectRef::new(project))
.await
.map_err(|e| DeclareError::Other(e.to_string()))?;
let (_, _, incoming_vol) = incoming.size.resources();
let incoming_vol = u64::from(incoming_vol);
let mut others = 0usize;
let mut others_vol: u64 = 0;
for db in &existing {
if db.name == name {
continue;
}
others += 1;
let (_, _, v) = db.size.resources();
others_vol = others_vol.saturating_add(u64::from(v));
}
let projected_count = others + 1;
if projected_count > self.quota.max_databases {
return Err(DeclareError::QuotaExceeded {
db: name.to_string(),
reason: format!(
"declaring it would bring project {project:?} to {projected_count} managed \
databases, over the ceiling of {}",
self.quota.max_databases
),
});
}
let projected_vol = others_vol.saturating_add(incoming_vol);
if projected_vol > self.quota.max_volume_mib {
return Err(DeclareError::QuotaExceeded {
db: name.to_string(),
reason: format!(
"declaring it would bring project {project:?} to {projected_vol} MiB of \
provisioned managed-database volume, over the ceiling of {} MiB",
self.quota.max_volume_mib
),
});
}
Ok(())
}
async fn provision(
&self,
binding: &ExternalDatabaseConfig,
project: &str,
) -> Result<(), DeclareError> {
if matches!(binding.tenant_scope, TenantScope::Site) {
return Ok(());
}
crate::tenant_sql::provision_tenant(
&self.deploy,
&self.kv,
&self.envelope,
binding,
project,
"",
)
.await
.map_err(DeclareError::Other)
}
}
#[async_trait]
impl ManagedDbDeclare for NodeManagedDbDeclare {
async fn declare(
&self,
project: &str,
name: &str,
db: &ApplyDatabase,
) -> Result<(), DeclareError> {
boatramp_core::project::validate_resource_name("project", project)
.map_err(|e| DeclareError::Other(e.to_string()))?;
boatramp_core::project::validate_resource_name("database", name)
.map_err(|e| DeclareError::Other(e.to_string()))?;
if db.name != name {
return Err(DeclareError::Other(format!(
"declared database name {:?} does not match the path segment {name:?}",
db.name
)));
}
if self.static_dbs.contains_key(name) {
return Err(DeclareError::DaemonConflict(name.to_string()));
}
if let Some(prior) = self
.deploy
.get_project_database(ProjectRef::new(project), name)
.await
.map_err(|e| DeclareError::Other(e.to_string()))?
&& let Some(field) = db.identity_change(&prior)
{
return Err(DeclareError::IdentityChange {
db: name.to_string(),
field,
});
}
self.enforce_quota(project, name, db).await?;
let binding = lower(project, db);
if !binding.is_managed_credential() {
return Err(DeclareError::Other(format!(
"internal: lowered database {name:?} is not the managed-credential path"
)));
}
binding.validate(name).map_err(DeclareError::Other)?;
crate::managed_sql::auto_register_managed_db_workloads(
&self.deploy,
&BTreeMap::from([(name.to_string(), binding.clone())]),
)
.await;
self.deploy
.set_project_database(ProjectRef::new(project), db)
.await
.map_err(|e| DeclareError::Other(e.to_string()))?;
self.provision(&binding, project).await
}
async fn ensure(&self, project: &str, name: &str) -> Result<(), DeclareError> {
boatramp_core::project::validate_resource_name("project", project)
.map_err(|e| DeclareError::Other(e.to_string()))?;
boatramp_core::project::validate_resource_name("database", name)
.map_err(|e| DeclareError::Other(e.to_string()))?;
let binding = match resolve_binding(&self.static_dbs, &self.deploy, project, name).await? {
Some(b) => b,
None => return Err(DeclareError::NotDeclared(name.to_string())),
};
if !binding.is_managed_credential() {
return Err(DeclareError::NotDeclared(name.to_string()));
}
crate::managed_sql::auto_register_managed_db_workloads(
&self.deploy,
&BTreeMap::from([(name.to_string(), binding.clone())]),
)
.await;
self.provision(&binding, project).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use boatramp_core::compute::ApplyDatabaseSize;
fn declared(name: &str) -> ApplyDatabase {
ApplyDatabase {
name: name.to_string(),
kind: ApplyDatabaseKind::Postgres,
version: Some(16),
extensions: vec![],
size: ApplyDatabaseSize::Small,
tenant: boatramp_core::compute::ApplyDatabaseTenant::Shared,
tenant_scope: boatramp_core::compute::ApplyDatabaseScope::Project,
read_only: false,
rls_session: false,
tenant_guc: None,
session_guc: None,
tenant_all_marker: None,
pool_max: Some(1000),
connect_timeout_secs: Some(9999),
startup_grace_secs: Some(99999),
}
}
#[test]
fn lowering_is_the_managed_credential_path() {
let cfg = lower("acme", &declared("app"));
assert!(cfg.is_managed_credential(), "must be managed-credential");
assert!(
cfg.password_env.is_none(),
"no password_env in a declared DB"
);
assert!(cfg.url_env.is_empty(), "no url_env");
assert!(cfg.read_url_env.is_none(), "no read_url_env");
assert!(cfg.migration_url_env.is_none(), "no migration_url_env");
assert!(cfg.image.is_none(), "no arbitrary image");
assert!(cfg.path.is_none(), "no host-fs path");
assert_eq!(
cfg.compute.as_deref(),
Some(derived_workload("acme", "app").as_str())
);
assert!(
cfg.compute
.as_deref()
.is_some_and(|c| c.starts_with("bramp-db-")),
"derived compute keeps the reserved prefix"
);
assert!(cfg.validate("app").is_ok(), "the lowered binding validates");
}
#[test]
fn tuning_knobs_are_capped() {
let cfg = lower("acme", &declared("app"));
assert_eq!(cfg.pool_max, Some(apply_db_caps::MAX_POOL));
assert_eq!(
cfg.connect_timeout_secs,
Some(apply_db_caps::MAX_CONNECT_TIMEOUT_SECS)
);
assert_eq!(
cfg.startup_grace_secs,
Some(apply_db_caps::MAX_STARTUP_GRACE_SECS)
);
}
#[test]
fn derived_workload_is_project_qualified() {
assert_ne!(
derived_workload("acme", "app"),
derived_workload("globex", "app"),
);
assert_ne!(
derived_workload("acme-app", "db"),
derived_workload("acme", "app-db"),
"ambiguous hyphen pair must NOT collide onto one shared server + superuser cred"
);
assert_ne!(derived_workload("a-b", "c"), derived_workload("a", "b-c"));
assert_ne!(
derived_workload("acme-app-db", ""),
derived_workload("", "acme-app-db"),
);
for (p, n) in [("acme-app", "db"), ("acme", "app-db"), ("globex", "app")] {
let w = derived_workload(p, n);
assert!(
w.starts_with("bramp-db-"),
"{w:?} keeps the reserved prefix"
);
assert!(w.len() <= 63, "{w:?} must be <= 63 bytes");
}
}
#[test]
fn size_preset_maps_to_bounded_volume() {
assert_eq!(
lower("p", &db_size(ApplyDatabaseSize::Small)).volume_size_mib,
Some(10 * 1024)
);
assert_eq!(
lower("p", &db_size(ApplyDatabaseSize::Medium)).volume_size_mib,
Some(50 * 1024)
);
assert_eq!(
lower("p", &db_size(ApplyDatabaseSize::Large)).volume_size_mib,
Some(200 * 1024)
);
}
fn db_size(size: ApplyDatabaseSize) -> ApplyDatabase {
let mut d = declared("app");
d.size = size;
d
}
use async_trait::async_trait;
use boatramp_core::envelope::{EnvelopeError, KeyEnvelope};
use boatramp_core::kv::MemoryKv;
use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
struct XorEnvelope;
#[async_trait]
impl KeyEnvelope for XorEnvelope {
async fn wrap(&self, p: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
Ok(p.iter().map(|b| b ^ 0x5a).collect())
}
async fn unwrap(&self, c: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
Ok(c.iter().map(|b| b ^ 0x5a).collect())
}
}
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 single_project(name: &str) -> ApplyDatabase {
let mut d = declared(name);
d.tenant = boatramp_core::compute::ApplyDatabaseTenant::Single;
d.pool_max = None;
d.connect_timeout_secs = None;
d.startup_grace_secs = None;
d
}
fn declare_cap(
static_dbs: BTreeMap<String, ExternalDatabaseConfig>,
) -> (
NodeManagedDbDeclare,
DeployStore,
std::sync::Arc<dyn KvStore>,
) {
declare_cap_quota(static_dbs, DeclareQuota::default())
}
fn declare_cap_quota(
static_dbs: BTreeMap<String, ExternalDatabaseConfig>,
quota: DeclareQuota,
) -> (
NodeManagedDbDeclare,
DeployStore,
std::sync::Arc<dyn KvStore>,
) {
let kv: std::sync::Arc<dyn KvStore> = std::sync::Arc::new(MemoryKv::new());
let deploy = DeployStore::new(std::sync::Arc::new(NullStorage), kv.clone());
let cap = NodeManagedDbDeclare::new(
static_dbs,
deploy.clone(),
kv.clone(),
std::sync::Arc::new(XorEnvelope),
quota,
);
(cap, deploy, kv)
}
#[tokio::test]
async fn managed_db_declaration_scoped_gate() {
let low = lower("acme", &single_project("app"));
assert!(
low.is_managed_credential() && low.password_env.is_none(),
"GATE: lowered binding must be the managed-credential path (no password_env)"
);
assert!(
low.url_env.is_empty() && low.image.is_none() && low.path.is_none(),
"GATE: no url_env / image / host-fs path may appear in a declared binding"
);
let mut static_dbs = BTreeMap::new();
static_dbs.insert("byo".to_string(), ExternalDatabaseConfig::default());
let (cap, _deploy, _kv) = declare_cap(static_dbs);
let refused = cap.declare("acme", "byo", &single_project("byo")).await;
assert!(
matches!(refused, Err(DeclareError::DaemonConflict(ref n)) if n == "byo"),
"GATE: a project manifest may not shadow a daemon-static binding — got {refused:?}"
);
let mut both = BTreeMap::new();
both.insert("byo".to_string(), ExternalDatabaseConfig::default());
let (cap2, deploy2, _kv2) = declare_cap(both);
deploy2
.set_project_database(ProjectRef::new("acme"), &single_project("byo"))
.await
.unwrap();
let merged = resolve_binding(&cap2.static_dbs, &deploy2, "acme", "byo").await;
assert!(
matches!(merged, Err(DeclareError::DaemonConflict(_))),
"GATE: the merge point fails closed when both sources define the name — got {merged:?}"
);
let (cap3, deploy3, kv3) = declare_cap(BTreeMap::new());
cap3.declare("acme", "app", &single_project("app"))
.await
.expect("GATE: a clean declare succeeds");
let stored = deploy3
.get_project_database(ProjectRef::new("acme"), "app")
.await
.unwrap();
assert!(stored.is_some(), "GATE: the declaration is persisted");
let cred_keys = kv3.list_prefix("managed-sql-cred/acme/").await.unwrap();
assert!(
!cred_keys.is_empty(),
"GATE: the eager provision minted + sealed a managed credential under the caller's project"
);
let sealed = kv3.get(&cred_keys[0]).await.unwrap().unwrap();
assert!(
sealed.iter().all(|b| *b != 0) && sealed.len() == 64,
"GATE: the credential is sealed at rest (64-byte sealed blob)"
);
let mut changed = single_project("app");
changed.tenant = boatramp_core::compute::ApplyDatabaseTenant::Shared;
let refused_change = cap3.declare("acme", "app", &changed).await;
assert!(
matches!(
refused_change,
Err(DeclareError::IdentityChange { field, .. }) if field == "tenant"
),
"GATE: changing `tenant` on an existing declared DB is refused — got {refused_change:?}"
);
assert_eq!(
lower("acme", &single_project("app")).compute.as_deref(),
Some(derived_workload("acme", "app").as_str()),
);
assert!(
lower("acme", &single_project("app"))
.compute
.as_deref()
.is_some_and(|c| c.starts_with("bramp-db-")),
"GATE: the derived server workload keeps the reserved `bramp-db-` prefix"
);
assert_ne!(
lower("globex", &single_project("app")).compute,
lower("acme", &single_project("app")).compute,
"GATE: a declaration provisions onto its OWN project's derived server only"
);
assert_ne!(
derived_workload("acme-app", "db"),
derived_workload("acme", "app-db"),
"GATE: the ambiguous `(project, name)` pair must NOT collide onto one server + superuser credential"
);
assert_ne!(
lower("acme-app", &single_project("db")).compute,
lower("acme", &single_project("app-db")).compute,
"GATE: two distinct projects must not lower onto the same shared server"
);
assert!(
!kv3.list_prefix("managed-sql-cred/acme/")
.await
.unwrap()
.is_empty(),
"GATE: nothing in the declare API removes the credential (inert on removal)"
);
assert!(
deploy3
.get_project_database(ProjectRef::new("acme"), "app")
.await
.unwrap()
.is_some(),
"GATE: the declaration record survives (removal is an explicit imperative verb)"
);
println!("MANAGED-DB DECLARATION SCOPED OK");
}
#[tokio::test]
async fn per_project_count_ceiling_is_fail_closed() {
let quota = DeclareQuota {
max_databases: 2,
max_volume_mib: u64::MAX, };
let (cap, deploy, _kv) = declare_cap_quota(BTreeMap::new(), quota);
cap.declare("acme", "a", &single_project("a"))
.await
.unwrap();
cap.declare("acme", "b", &single_project("b"))
.await
.unwrap();
let refused = cap.declare("acme", "c", &single_project("c")).await;
assert!(
matches!(refused, Err(DeclareError::QuotaExceeded { ref db, .. }) if db == "c"),
"declaring past the count ceiling must be refused fail-closed — got {refused:?}"
);
assert!(
deploy
.get_project_database(ProjectRef::new("acme"), "c")
.await
.unwrap()
.is_none(),
"the over-ceiling declaration must not be persisted"
);
cap.declare("acme", "a", &single_project("a"))
.await
.expect("re-declaring an existing DB is within the ceiling (no new slot)");
cap.declare("globex", "x", &single_project("x"))
.await
.expect("a different project has its own count budget");
}
#[tokio::test]
async fn per_project_volume_ceiling_is_fail_closed() {
let quota = DeclareQuota {
max_databases: usize::MAX, max_volume_mib: 25 * 1024, };
let (cap, deploy, _kv) = declare_cap_quota(BTreeMap::new(), quota);
cap.declare("acme", "a", &single_project("a"))
.await
.unwrap();
cap.declare("acme", "b", &single_project("b"))
.await
.unwrap();
let refused = cap.declare("acme", "c", &single_project("c")).await;
assert!(
matches!(refused, Err(DeclareError::QuotaExceeded { ref db, .. }) if db == "c"),
"declaring past the aggregate-volume ceiling must be refused — got {refused:?}"
);
assert!(
deploy
.get_project_database(ProjectRef::new("acme"), "c")
.await
.unwrap()
.is_none(),
"the over-ceiling declaration must not be persisted"
);
let (cap2, _deploy2, _) = declare_cap_quota(BTreeMap::new(), quota);
let mut medium = single_project("big");
medium.size = ApplyDatabaseSize::Medium;
let refused_medium = cap2.declare("zeta", "big", &medium).await;
assert!(
matches!(refused_medium, Err(DeclareError::QuotaExceeded { .. })),
"a single over-ceiling volume must be refused — got {refused_medium:?}"
);
}
#[test]
fn declare_quota_from_config_folds_defaults() {
let d = DeclareQuota::from_config(None, None);
assert_eq!(
d.max_databases,
apply_db_caps::DEFAULT_MAX_DECLARED_DATABASES_PER_PROJECT
);
assert_eq!(
d.max_volume_mib,
apply_db_caps::DEFAULT_MAX_DECLARED_VOLUME_MIB
);
let o = DeclareQuota::from_config(Some(4), Some(1234));
assert_eq!(o.max_databases, 4);
assert_eq!(o.max_volume_mib, 1234);
let z = DeclareQuota::from_config(Some(0), Some(0));
assert_eq!(z.max_databases, 0);
assert_eq!(z.max_volume_mib, 0);
}
}