#![cfg(target_os = "linux")]
#![cfg(feature = "sql-postgres")]
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use boatramp_container::ContainerBackend;
use boatramp_core::compute::{
Artifact, ComputeBackend, LaunchRequest, ObservedInstance, PrivilegeDirective, ReplicaPhase,
};
use boatramp_core::deploy::DeployStore;
use boatramp_core::envelope::{EnvelopeError, KeyEnvelope};
use boatramp_core::kv::{KvStore, MemoryKv};
use boatramp_core::project::ProjectRef;
use boatramp_core::sql::{SqlBackend, SqlValue};
use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
use boatramp_node::config::{ExternalDatabaseConfig, TenantIsolation, TenantScope};
use boatramp_node::managed_sql::ManagedSqlCredentials;
use boatramp_node::tenant_sql::{provision_tenant, NodeTenantSqlResolver};
use boatramp_storage::sql_sqlx::PerTenantSqlResolver;
use bytes::Bytes;
use futures::StreamExt;
const COMPUTE: &str = "pg";
const DATABASE: &str = "appdb";
const APP_USER: &str = "app";
const TENANTS: &[&str] = &["acme", "globex"];
const SHARED_VOLUME_BUG_NAME: &str = "data";
struct FileBlob(Vec<u8>);
#[async_trait]
impl Storage for FileBlob {
async fn get(&self, _key: &str) -> Result<GetObject, StorageError> {
let bytes = Bytes::from(self.0.clone());
let size = self.0.len() as u64;
let body: ByteStream = futures::stream::once(async move { Ok(bytes) }).boxed();
Ok(GetObject {
meta: ObjectMeta {
key: String::new(),
size: Some(size),
content_type: None,
etag: None,
},
body,
})
}
async fn get_range(&self, _: &str, _: u64, _: Option<u64>) -> Result<GetObject, StorageError> {
Err(StorageError::unsupported("range"))
}
async fn put(&self, _: &str, _: ByteStream, _: PutMeta) -> Result<ObjectMeta, StorageError> {
Err(StorageError::unsupported("put"))
}
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())
}
}
struct RevEnvelope;
#[async_trait]
impl KeyEnvelope for RevEnvelope {
async fn wrap(&self, p: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
Ok(p.iter().rev().copied().collect())
}
async fn unwrap(&self, w: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
Ok(w.iter().rev().copied().collect())
}
}
fn single_project_binding() -> ExternalDatabaseConfig {
ExternalDatabaseConfig {
kind: "postgres".into(),
compute: Some(COMPUTE.into()),
database: Some(DATABASE.into()),
user: Some(APP_USER.into()),
tenant: TenantIsolation::Single,
tenant_scope: TenantScope::Project,
connect_timeout_secs: Some(10),
..Default::default()
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[ignore = "needs Linux + root + a bridge (privileged live seam); pulls pgvector, boots 2 containers"]
async fn single_mode_isolates_volumes_and_each_authenticates() {
let Some(bin) = std::env::var_os("BOATRAMP_BIN") else {
eprintln!(
"single_volume_live: set BOATRAMP_BIN (root + a br-boatramp bridge) to run; \
skipping (never fails on a dev box)"
);
return;
};
let bin = PathBuf::from(bin);
let bridge = std::env::var("CONTAINER_BRIDGE").unwrap_or_else(|_| "br-boatramp".into());
let subnet = std::env::var("CONTAINER_SUBNET").unwrap_or_else(|_| "10.0.0.0/24".into());
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let storage: Arc<dyn Storage> = Arc::new(FileBlob(Vec::new()));
let deploy = DeployStore::new(storage, kv.clone());
let envelope: Arc<dyn KeyEnvelope> = Arc::new(RevEnvelope);
let creds = ManagedSqlCredentials::new(kv.clone(), envelope.clone());
let data_dir = std::env::temp_dir().join(format!("boatramp-single-vol-{}", std::process::id()));
let backend = Arc::new(
ContainerBackend::new(
Arc::new(FileBlob(Vec::new())),
data_dir.clone(),
bridge,
&subnet,
bin,
)
.expect("backend"),
);
let outcome = run_assertions(&deploy, &kv, &envelope, &creds, &backend, &data_dir).await;
for handle in &outcome.handles {
let _ = backend.stop(handle).await;
}
let _ = std::fs::remove_dir_all(&data_dir);
outcome
.result
.expect("single-mode volume isolation assertions");
println!("SINGLE VOLUME ISOLATION OK: acme + globex on distinct volumes, each authenticates");
}
struct Outcome {
result: Result<(), String>,
handles: Vec<boatramp_core::compute::InstanceHandle>,
}
async fn run_assertions(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
creds: &ManagedSqlCredentials,
backend: &Arc<ContainerBackend>,
data_dir: &std::path::Path,
) -> Outcome {
let mut handles = Vec::new();
let mut volume_dirs: Vec<(String, PathBuf)> = Vec::new();
let binding = single_project_binding();
for tenant in TENANTS {
match provision_and_launch_tenant(
deploy, kv, envelope, creds, backend, data_dir, &binding, tenant,
)
.await
{
Ok((handle, vol_dir)) => {
handles.push(handle);
volume_dirs.push((tenant.to_string(), vol_dir));
}
Err(e) => {
return Outcome {
result: Err(format!("tenant {tenant}: {e}")),
handles,
};
}
}
}
let result = (|| {
let [(a_name, a_dir), (b_name, b_dir)] = volume_dirs.as_slice() else {
return Err(format!(
"expected {} tenants, got {}",
TENANTS.len(),
volume_dirs.len()
));
};
if a_dir == b_dir {
return Err(format!(
"CROSS-TENANT VOLUME COLLISION: {a_name} and {b_name} share the volume dir {a_dir:?}"
));
}
eprintln!("distinct volumes: {a_name} -> {a_dir:?} != {b_name} -> {b_dir:?} OK");
Ok(())
})();
Outcome { result, handles }
}
#[allow(clippy::too_many_arguments)]
async fn provision_and_launch_tenant(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
creds: &ManagedSqlCredentials,
backend: &Arc<ContainerBackend>,
data_dir: &std::path::Path,
binding: &ExternalDatabaseConfig,
tenant: &str,
) -> Result<(boatramp_core::compute::InstanceHandle, PathBuf), String> {
provision_tenant(deploy, kv, envelope, binding, tenant, "-")
.await
.map_err(|e| format!("provision: {e}"))?;
let proj = ProjectRef::new(tenant);
let workload = deploy
.get_compute_workload(proj, &single_workload_name(tenant))
.await
.map_err(|e| format!("read workload: {e}"))?
.ok_or_else(|| "provision did not register the per-tenant Single workload".to_string())?;
let mut spec = deploy
.get_compute_spec(&workload.active)
.await
.map_err(|e| format!("read spec: {e}"))?
.ok_or_else(|| "workload's active spec is missing".to_string())?;
let vol = spec
.volumes
.first()
.ok_or_else(|| "managed Single spec must carry a data volume".to_string())?;
if vol.name != workload.name {
return Err(format!(
"volume name {:?} is not keyed to the per-tenant workload {:?} (the v0.3.11 fix)",
vol.name, workload.name
));
}
if vol.name == SHARED_VOLUME_BUG_NAME {
return Err(format!(
"volume name is the shared {SHARED_VOLUME_BUG_NAME:?} — the pre-v0.3.11 bug where \
every tenant mounts the SAME PGDATA"
));
}
let vol_dir = data_dir.join("compute").join("volumes").join(&vol.name);
eprintln!(
"[{tenant}] workload={:?} volume-name={:?} -> {vol_dir:?} (not {SHARED_VOLUME_BUG_NAME:?}) OK",
workload.name, vol.name
);
let tenant_pw = creds
.password(tenant, &workload.name)
.await
.map_err(|e| format!("resolve per-tenant credential: {e}"))?;
PrivilegeDirective::Rootless { uid: 999, gid: 999 }.apply(&mut spec);
spec.env
.insert("POSTGRES_USER".to_string(), APP_USER.to_string());
spec.env
.insert("POSTGRES_PASSWORD".to_string(), tenant_pw.clone());
spec.env
.insert("POSTGRES_DB".to_string(), DATABASE.to_string());
let artifact = backend
.materialize(&spec)
.await
.map_err(|e| format!("materialize (pull) image: {e}"))?;
if !matches!(artifact, Artifact::Rootfs { .. }) {
return Err("expected a rootfs artifact from the image pull".to_string());
}
let _ = deploy
.put_compute_spec(&spec)
.await
.map_err(|e| format!("store staged spec: {e}"))?;
let req = LaunchRequest {
project: "default".into(),
workload: workload.name.clone(),
replica: 0,
spec,
artifact,
};
let inst = backend
.launch(&req)
.await
.map_err(|e| format!("launch container: {e}"))?;
let handle = inst.handle.clone();
let host = inst.endpoint.host.clone();
let port = inst.endpoint.port;
eprintln!("[{tenant}] == dedicated pgvector launched == endpoint={host}:{port}");
if !vol_dir.is_dir() {
return Err(format!(
"the per-tenant volume backing dir {vol_dir:?} was not created on disk"
));
}
let observed = ObservedInstance {
handle: inst.handle.clone(),
node: 0,
backend: backend.id().to_string(),
endpoint: inst.endpoint.clone(),
region: None,
healthy: true,
started_at: None,
phase: ReplicaPhase::Running,
snapshot: None,
};
deploy
.set_replica_state(ProjectRef::new(tenant), &observed)
.await
.map_err(|e| format!("publish replica state: {e}"))?;
let backend_for_tenant = resolve_tenant_backend(deploy, kv, envelope, binding, tenant)
.await
.map_err(|e| format!("resolve shipped backend: {e}"))?;
let current_user = wait_for_authenticated_query(&backend_for_tenant).await?;
if current_user != APP_USER {
return Err(format!(
"authenticated, but current_user was {current_user:?}, expected {APP_USER:?}"
));
}
eprintln!(
"[{tenant}] resolved-backend SELECT current_user = {current_user:?} (authenticated) OK"
);
Ok((handle, vol_dir))
}
fn single_workload_name(tenant: &str) -> String {
format!(
"{COMPUTE}-{}",
boatramp_storage::tenant_provision::sanitize_ident(tenant)
)
}
async fn resolve_tenant_backend(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
binding: &ExternalDatabaseConfig,
tenant: &str,
) -> Result<Arc<dyn SqlBackend>, String> {
let resolver =
NodeTenantSqlResolver::new(deploy.clone(), kv.clone(), envelope.clone(), binding)
.ok_or_else(|| {
"resolver should build for a compute-backed managed binding".to_string()
})?;
resolver
.resolve(tenant, "-")
.await
.map_err(|e| format!("resolve: {e}"))
}
async fn wait_for_authenticated_query(backend: &Arc<dyn SqlBackend>) -> Result<String, String> {
let mut last = String::new();
for _ in 0..120 {
match backend.run_query("SELECT current_user").await {
Ok(rows) => match rows.rows.first().and_then(|r| r.first()) {
Some(SqlValue::Text(s)) => return Ok(s.clone()),
other => {
last = format!("current_user was not text: {other:?}");
}
},
Err(e) => last = e.to_string(),
}
tokio::time::sleep(Duration::from_millis(500)).await;
}
Err(format!(
"the resolved per-tenant backend never authenticated + answered `SELECT current_user` \
within the budget (last: {last:?}) — the pre-v0.3.11 symptom was \"password \
authentication failed\" because the Single container reused another workload's volume \
and skipped initdb"
))
}