#![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::{
managed_db_spec, Artifact, ComputeBackend, ComputeWorkload, LaunchRequest, ManagedDbEngine,
ObservedInstance, PlacementConstraints, PrivilegeDirective, ReplicaPhase,
};
use boatramp_core::deploy::DeployStore;
use boatramp_core::envelope::{EnvelopeError, KeyEnvelope};
use boatramp_core::kv::{KvStore, MemoryKv};
use boatramp_core::project::{ProjectRef, DEFAULT_PROJECT};
use boatramp_core::sql::{
reject_reserved_session_writes, SqlBackend, SqlError, SqlTransaction, 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::NodeTenantSqlResolver;
use boatramp_storage::sql_sqlx::PerTenantSqlResolver;
use bytes::Bytes;
use futures::StreamExt;
const COMPUTE: &str = "pg";
const DATABASE: &str = "appdb";
const SUPERUSER: &str = "postgres";
const TENANT_PROJECT: &str = "acme";
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 rls_project_binding() -> ExternalDatabaseConfig {
ExternalDatabaseConfig {
kind: "postgres".into(),
compute: Some(COMPUTE.into()),
database: Some(DATABASE.into()),
user: Some(SUPERUSER.into()),
tenant: TenantIsolation::Shared,
tenant_scope: TenantScope::Project,
rls_session: true,
connect_timeout_secs: Some(10),
..Default::default()
}
}
fn wait_for_pg(conn: &str) -> (bool, String) {
let mut last = String::new();
for _ in 0..60 {
if let Ok(out) = std::process::Command::new("psql")
.args([conn, "-tAc", "select 1"])
.output()
{
last = format!(
"{}{}",
String::from_utf8_lossy(&out.stdout),
String::from_utf8_lossy(&out.stderr)
);
if out.status.success() && last.trim() == "1" {
return (true, last);
}
}
std::thread::sleep(Duration::from_millis(500));
}
(false, last)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[ignore = "needs Linux + root + a bridge + psql (privileged live seam); pulls pgvector"]
async fn rls_session_guard_holds_live_postgres() {
let Some(bin) = std::env::var_os("BOATRAMP_BIN") else {
eprintln!(
"rls_session_live: set BOATRAMP_BIN (and have `psql` on PATH, root, and a \
br-boatramp bridge) to run"
);
return;
};
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 superuser_pw = creds
.password(DEFAULT_PROJECT, COMPUTE)
.await
.expect("seal superuser credential");
let data_dir =
std::env::temp_dir().join(format!("boatramp-rls-session-{}", std::process::id()));
let backend = ContainerBackend::new(
Arc::new(FileBlob(Vec::new())),
data_dir.clone(),
bridge,
&subnet,
PathBuf::from(bin),
)
.expect("backend");
let mut spec = managed_db_spec(
ManagedDbEngine::Postgres,
Some("pgvector/pgvector:pg16"),
512,
);
PrivilegeDirective::Rootless { uid: 999, gid: 999 }.apply(&mut spec);
spec.env
.insert("POSTGRES_USER".to_string(), SUPERUSER.to_string());
spec.env
.insert("POSTGRES_PASSWORD".to_string(), superuser_pw.clone());
spec.env
.insert("POSTGRES_DB".to_string(), DATABASE.to_string());
let artifact = backend
.materialize(&spec)
.await
.expect("materialize (pull) pgvector image");
assert!(matches!(artifact, Artifact::Rootfs { .. }));
let spec_id = deploy
.put_compute_spec(&spec)
.await
.expect("store compute spec");
deploy
.set_compute_workload(
ProjectRef::DEFAULT,
&ComputeWorkload {
version: 1,
name: COMPUTE.to_string(),
active: spec_id,
replicas: 1,
placement: PlacementConstraints::default(),
},
)
.await
.expect("register compute workload");
let req = LaunchRequest {
project: "default".into(),
workload: COMPUTE.to_string(),
replica: 0,
spec,
artifact,
};
let inst = backend
.launch(&req)
.await
.expect("launch pgvector container");
let handle = inst.handle.clone();
let host = inst.endpoint.host.clone();
let port = inst.endpoint.port;
eprintln!("== shared pgvector launched == endpoint={host}:{port}");
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::DEFAULT, &observed)
.await
.expect("publish replica state");
let outcome = run_assertions(&deploy, &kv, &envelope, &superuser_pw, &host, port).await;
let _ = backend.stop(&handle).await;
let _ = std::fs::remove_dir_all(&data_dir);
outcome.expect("rls_session guard assertions");
println!(
"RLS SESSION GUARD OK (postgres): guest spoof refused, boatramp.project stayed {TENANT_PROJECT}"
);
}
const PG_SPOOFS: &[&str] = &[
"SELECT set_config('boatramp.project','victim',false)",
"SET boatramp.project='victim'",
"DO $$ BEGIN PERFORM set_config('boatramp.project','victim',false); END $$;",
];
const PG_UNGUARDED_OVERWRITE: &str = "SELECT set_config('boatramp.project','victim',true)";
async fn guest_entry_point_execute(
backend: &Arc<dyn SqlBackend>,
tx: &mut Box<dyn SqlTransaction>,
statement: &str,
) -> Result<(), SqlError> {
if backend.injects_session_context() {
reject_reserved_session_writes(statement)?;
}
tx.execute(statement, &[]).await.map(|_| ())
}
async fn read_boatramp_project(tx: &mut Box<dyn SqlTransaction>) -> Result<String, SqlError> {
let rows = tx
.query("SELECT current_setting('boatramp.project')", &[])
.await?;
match rows.rows.first().and_then(|r| r.first()) {
Some(SqlValue::Text(s)) => Ok(s.clone()),
other => Err(SqlError::other(format!(
"current_setting('boatramp.project') was not text: {other:?}"
))),
}
}
async fn run_assertions(
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
envelope: &Arc<dyn KeyEnvelope>,
superuser_pw: &str,
host: &str,
port: u16,
) -> Result<(), String> {
let super_conn = format!("postgresql://{SUPERUSER}:{superuser_pw}@{host}:{port}/{DATABASE}");
let (up, last) = wait_for_pg(&super_conn);
if !up {
return Err(format!(
"shared Postgres should answer `select 1` as the superuser (last: {last:?})"
));
}
let binding = rls_project_binding();
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()
})?;
let backend: Arc<dyn SqlBackend> = resolver
.resolve(TENANT_PROJECT, "-")
.await
.map_err(|e| format!("resolve tenant {TENANT_PROJECT}: {e}"))?;
if !backend.injects_session_context() {
return Err(
"the resolved rls_session backend must report injects_session_context() == true \
(else the guest guard would be inert and this gate vacuous)"
.to_string(),
);
}
let mut tx = backend
.begin()
.await
.map_err(|e| format!("begin injecting transaction: {e}"))?;
let injected = read_boatramp_project(&mut tx)
.await
.map_err(|e| format!("read injected boatramp.project: {e}"))?;
if injected != TENANT_PROJECT {
return Err(format!(
"injection not wired: boatramp.project read back {injected:?}, expected {TENANT_PROJECT:?}"
));
}
eprintln!("injection: boatramp.project == {injected:?} OK");
for spoof in PG_SPOOFS {
let refused = guest_entry_point_execute(&backend, &mut tx, spoof).await;
if refused.is_ok() {
return Err(format!(
"GUARD BYPASS: guest spoof was ALLOWED through the entry point: {spoof:?}"
));
}
let after = read_boatramp_project(&mut tx)
.await
.map_err(|e| format!("re-read boatramp.project after refused spoof {spoof:?}: {e}"))?;
if after != TENANT_PROJECT {
return Err(format!(
"RLS SPOOF LANDED: after refused spoof {spoof:?}, boatramp.project = {after:?} \
(expected it to stay {TENANT_PROJECT:?})"
));
}
eprintln!("spoof refused + value unchanged: {spoof:?} OK");
}
{
let mut ctl = backend
.begin()
.await
.map_err(|e| format!("begin control transaction: {e}"))?;
let base = read_boatramp_project(&mut ctl)
.await
.map_err(|e| format!("control read (pre): {e}"))?;
if base != TENANT_PROJECT {
return Err(format!(
"control precondition: expected {TENANT_PROJECT:?}, got {base:?}"
));
}
ctl.execute(PG_UNGUARDED_OVERWRITE, &[])
.await
.map_err(|e| format!("control unguarded overwrite: {e}"))?;
let after = read_boatramp_project(&mut ctl)
.await
.map_err(|e| format!("control read (post): {e}"))?;
if after != "victim" {
return Err(format!(
"control invalid: an UNGUARDED set_config should change boatramp.project to \
\"victim\", got {after:?} — the spoofs would be inert, making the \
value-unchanged assertion meaningless"
));
}
let _ = ctl.rollback().await;
eprintln!("control: an UNGUARDED overwrite DOES change boatramp.project to \"victim\" OK");
}
let _ = tx.rollback().await;
Ok(())
}