use std::path::Path;
use std::sync::Arc;
use boatramp_core::deploy::DeployStore;
use boatramp_core::kv::KvStore;
use boatramp_core::Storage;
use crate::config::ServerConfig;
use crate::error::{Error, Result};
pub fn compute_reconcile_tick() -> std::time::Duration {
std::env::var("BOATRAMP_COMPUTE_RECONCILE_TICK_MS")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|&ms| ms > 0)
.map(std::time::Duration::from_millis)
.unwrap_or(std::time::Duration::from_secs(30))
}
pub const DOMAIN_VERIFY_RECONCILE_TICK: std::time::Duration = std::time::Duration::from_secs(60);
pub const COMPUTE_IDLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(300);
pub struct NodeInput<'a> {
pub config: &'a ServerConfig,
pub data_dir: &'a Path,
pub storage: Arc<dyn Storage>,
pub kv: Arc<dyn KvStore>,
pub auth: boatramp_server::Auth,
pub options: boatramp_server::ServerOptions,
pub serve_addr: Option<std::net::SocketAddr>,
pub watch_provider: Option<Arc<dyn boatramp_core::blob_provision::WatchProvider>>,
pub provision_tier: boatramp_core::blob_notify::ProvisionTier,
pub messaging: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
pub is_leader: boatramp_server::CronLeaderGate,
pub node_id: u64,
pub worker_exe: Option<std::path::PathBuf>,
}
pub struct RunningNode {
pub deploy: DeployStore,
pub handlers: boatramp_server::HandlerRuntime,
pub auth: boatramp_server::Auth,
pub options: boatramp_server::ServerOptions,
pub reconcile: Vec<tokio::task::JoinHandle<()>>,
}
fn self_egress_addrs(
addr: Option<std::net::SocketAddr>,
enabled: bool,
) -> Vec<std::net::SocketAddr> {
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
let Some(addr) = addr.filter(|_| enabled) else {
return Vec::new();
};
if addr.ip().is_unspecified() {
let port = addr.port();
vec![
SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), port),
SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), port),
]
} else {
vec![addr]
}
}
pub async fn assemble(input: NodeInput<'_>) -> Result<RunningNode> {
let NodeInput {
config,
data_dir,
storage,
kv,
auth,
options,
serve_addr,
watch_provider,
provision_tier,
messaging,
is_leader,
node_id,
worker_exe,
} = input;
let max_handler_blob_bytes = options.posture.max_handler_blob_bytes;
let max_component_bytes = options.posture.max_component_bytes;
let allow_guest_private_egress = options.posture.allow_guest_private_egress;
let allow_env_secret_refs = options.posture.allow_env_secret_refs;
let allow_guest_email = options.posture.allow_guest_email;
let self_egress_addrs = self_egress_addrs(serve_addr, options.posture.allow_guest_self_egress);
let allow_shared_kernel = options.posture.allow_shared_kernel_compute;
let domain_verify_allow_private = options.posture.domain_verify_allow_private;
let compute_storage = storage.clone();
let deploy = DeployStore::new(storage, kv.clone());
let secrets_envelope = build_secrets_envelope(config.secrets.as_ref(), data_dir)?;
let secret_store = secrets_envelope.clone().map(|envelope| {
Arc::new(boatramp_core::secret_store::SecretStore::new(
kv.clone(),
envelope,
))
});
let email_profile_store = secrets_envelope.clone().map(|envelope| {
Arc::new(boatramp_core::email_config::EmailProfileStore::new(
kv.clone(),
envelope,
))
});
let handlers = crate::handlers::build_handler_runtime(
kv.clone(),
compute_storage.clone(),
data_dir,
config.handlers.as_ref(),
messaging,
max_handler_blob_bytes,
max_component_bytes,
allow_guest_private_egress,
self_egress_addrs,
allow_env_secret_refs,
allow_guest_email,
options.posture.require_tenancy_declaration,
options.posture.allow_cross_tenant_db,
&deploy,
secrets_envelope.clone(),
)
.await?;
#[cfg(feature = "handlers")]
if let Some(issuer) = options.issuer.clone() {
handlers.set_session_signer(issuer);
}
#[cfg(feature = "capability")]
if options.posture.allow_guest_mint_capability {
handlers.set_capability_minting(options.posture.max_guest_capability_ttl_secs);
}
#[cfg(feature = "handlers")]
{
let base = &options.posture;
let overrides: std::collections::BTreeMap<
String,
boatramp_core::security::ResolvedProjectTenancy,
> = config
.security
.as_ref()
.map(|s| {
s.projects
.iter()
.map(|(project, ovr)| (project.clone(), base.project_tenancy(ovr)))
.collect()
})
.unwrap_or_default();
handlers.set_project_tenancy_overrides(overrides);
}
#[cfg(feature = "admin")]
{
use boatramp_handlers::AdminSurface;
let p = &options.posture;
let mut surfaces = std::collections::BTreeSet::new();
if p.allow_guest_admin_domains {
surfaces.insert(AdminSurface::Domains);
}
if p.allow_guest_admin_email {
surfaces.insert(AdminSurface::Email);
}
if p.allow_guest_admin_site {
surfaces.insert(AdminSurface::Site);
}
if p.allow_guest_admin_secrets {
surfaces.insert(AdminSurface::Secrets);
}
if !surfaces.is_empty() {
let controller = Arc::new(boatramp_server::ServerAdminController::with_server_probe(
deploy.clone(),
email_profile_store.clone(),
secret_store.clone(),
p.domain_verify_allow_private,
));
handlers.set_admin(controller, surfaces);
}
}
#[cfg(feature = "handlers")]
handlers.set_cron_leader_gate(is_leader.clone());
#[cfg(feature = "handlers")]
if let Some(provider) = watch_provider {
handlers.set_watch_provider(provider);
handlers.set_provision_tier(provision_tier);
}
#[cfg(not(feature = "handlers"))]
let _ = (watch_provider, provision_tier);
match deploy.ensure_default_project().await {
Ok(true) => tracing::info!("materialized the reserved `default` project record"),
Ok(false) => {}
Err(e) => tracing::warn!(
error = %e,
"could not materialize the `default` project record; readers use the synthesized default"
),
}
#[cfg(feature = "handlers")]
handlers.set_invoker(deploy.clone());
let (compute_backends, compute_node) = crate::compute::build_compute(
config.compute.as_ref(),
compute_storage,
data_dir,
node_id,
!allow_shared_kernel,
options.daemon_runtime.clone(),
worker_exe.as_deref(),
)
.await;
crate::compute::adopt_running_replica_ips(&deploy, &compute_backends).await;
let internal_dns =
crate::compute::spawn_internal_dns(config.compute.as_ref(), &compute_backends, &deploy);
#[cfg(feature = "handlers")]
let sql_resolver = boatramp_server::sql_shim::spawn_sql_shim(
handlers.sql_backends(),
config.compute.as_ref().and_then(|c| c.sql_shim_url.clone()),
)
.await;
#[cfg(not(feature = "handlers"))]
let sql_resolver: Option<Arc<dyn boatramp_core::compute::ComputeBindingResolver>> = None;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
let operator_envelope = secrets_envelope.clone();
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
let deprovision_envelope = secrets_envelope.clone();
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
let reaper_envelope = secrets_envelope.clone();
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
let managed_db_resolver: Option<Arc<dyn boatramp_core::compute::ManagedDbEnvResolver>> = match (
config
.handlers
.as_ref()
.and_then(|h| h.bindings.sql.as_ref()),
secrets_envelope,
) {
(Some(sql), Some(envelope)) if !sql.databases.is_empty() => {
let creds = crate::managed_sql::ManagedSqlCredentials::new(kv.clone(), envelope);
let privilege = config
.compute
.as_ref()
.map(|c| c.managed_db_privilege)
.unwrap_or_default();
let env =
crate::managed_sql::ManagedDbEnv::from_config(&sql.databases, creds, privilege);
(!env.is_empty()).then(|| Arc::new(env) as Arc<_>)
}
_ => None,
};
#[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
let managed_db_resolver: Option<Arc<dyn boatramp_core::compute::ManagedDbEnvResolver>> = None;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
if let Some(sql) = config
.handlers
.as_ref()
.and_then(|h| h.bindings.sql.as_ref())
.filter(|sql| !sql.databases.is_empty())
{
crate::managed_sql::auto_register_managed_db_workloads(&deploy, &sql.databases).await;
}
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
let operator_sql: Option<Arc<dyn boatramp_core::sql::OperatorSql>> = config
.handlers
.as_ref()
.and_then(|h| h.bindings.sql.as_ref())
.filter(|sql| !sql.databases.is_empty())
.map(|sql| {
Arc::new(crate::managed_sql::NodeOperatorSql::new(
sql.databases.clone(),
kv.clone(),
operator_envelope,
deploy.clone(),
)) as Arc<_>
});
#[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
let operator_sql: Option<Arc<dyn boatramp_core::sql::OperatorSql>> = None;
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
let deprovision_grace_secs = config
.handlers
.as_ref()
.and_then(|h| h.bindings.sql.as_ref())
.and_then(|sql| sql.deprovision_grace_secs)
.unwrap_or(crate::tenant_sql::DEFAULT_DEPROVISION_GRACE_SECS);
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
let tenant_deprovisioner: Option<Arc<dyn boatramp_core::sql::TenantDeprovisioner>> = config
.handlers
.as_ref()
.and_then(|h| h.bindings.sql.as_ref())
.filter(|sql| !sql.databases.is_empty())
.zip(deprovision_envelope)
.map(|(sql, envelope)| {
Arc::new(crate::tenant_sql::NodeTenantDeprovisioner::new(
deploy.clone(),
kv.clone(),
envelope,
sql.databases.clone(),
deprovision_grace_secs,
)) as Arc<_>
});
#[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
let tenant_deprovisioner: Option<Arc<dyn boatramp_core::sql::TenantDeprovisioner>> = None;
let compute_exec: Option<Arc<dyn boatramp_core::compute::ComputeExec>> = Some(Arc::new(
crate::compute::NodeComputeExec::new(compute_backends.clone(), deploy.clone()),
) as Arc<_>);
let compute_volumes: Option<Arc<dyn boatramp_core::compute::ComputeVolumes>> = Some(Arc::new(
crate::compute::NodeComputeVolumes::new(compute_backends.clone(), deploy.clone()),
)
as Arc<_>);
let compute_control: Option<Arc<dyn boatramp_core::compute::ComputeControl>> = Some(Arc::new(
crate::compute::NodeComputeControl::new(compute_backends.clone(), deploy.clone()),
)
as Arc<_>);
let compute_reconcile = boatramp_server::spawn_compute_reconcile(
deploy.clone(),
compute_backends,
vec![compute_node],
boatramp_core::compute::BackendPolicy::from_shared_kernel_allowed(allow_shared_kernel),
is_leader.clone(),
compute_reconcile_tick(),
COMPUTE_IDLE_TIMEOUT,
sql_resolver,
managed_db_resolver,
);
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
let tombstone_reaper: Option<tokio::task::JoinHandle<()>> = config
.handlers
.as_ref()
.and_then(|h| h.bindings.sql.as_ref())
.filter(|sql| !sql.databases.is_empty())
.zip(reaper_envelope)
.map(|(_sql, envelope)| {
crate::tenant_sql::spawn_tenant_tombstone_reaper(
deploy.clone(),
kv.clone(),
envelope,
is_leader.clone(),
crate::tenant_sql::TOMBSTONE_REAPER_TICK,
)
});
#[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
let tombstone_reaper: Option<tokio::task::JoinHandle<()>> = None;
let dv_reconcile = boatramp_server::spawn_domain_verify_reconcile(
deploy.clone(),
domain_verify_allow_private,
is_leader,
DOMAIN_VERIFY_RECONCILE_TICK,
);
let mut options = options;
options.operator_sql = operator_sql;
options.tenant_deprovisioner = tenant_deprovisioner;
options.compute_exec = compute_exec;
options.compute_volumes = compute_volumes;
options.compute_control = compute_control;
options.secret_store = secret_store;
options.email_profile_store = email_profile_store;
let mut reconcile = vec![compute_reconcile, dv_reconcile];
if let Some(reaper) = tombstone_reaper {
reconcile.push(reaper);
}
if let Some(dns) = internal_dns {
reconcile.push(dns);
}
Ok(RunningNode {
deploy,
handlers,
auth,
options,
reconcile,
})
}
fn build_secrets_envelope(
secrets: Option<&crate::config::SecretsConfig>,
data_dir: &Path,
) -> Result<Option<Arc<dyn boatramp_core::envelope::KeyEnvelope>>> {
use boatramp_server::envelope::{build_envelope, EnvelopeSpec};
let Some(cfg) = secrets else {
return Ok(None);
};
let spec = match cfg.envelope.as_str() {
"" => EnvelopeSpec::None,
"local" => EnvelopeSpec::Local {
kek_file: cfg
.kek_file
.clone()
.unwrap_or_else(|| data_dir.join("secrets/kek")),
},
"vault" => {
let v = cfg.vault.as_ref().ok_or_else(|| {
Error::Envelope(
"secrets.envelope = \"vault\" needs a [secrets.vault] section".into(),
)
})?;
let token = std::env::var(&v.token_env).map_err(|_| {
Error::Envelope(format!("Vault token env `{}` is not set", v.token_env))
})?;
EnvelopeSpec::Vault {
addr: v.addr.clone(),
key: v.key.clone(),
token,
}
}
other => {
return Err(Error::Envelope(format!(
"unknown secrets.envelope {other:?} (want \"local\" or \"vault\")"
)))
}
};
build_envelope(spec).map_err(|e| Error::Envelope(e.to_string()))
}
#[cfg(all(test, feature = "fs"))]
mod tests {
use super::*;
use boatramp_core::kv::MemoryKv;
use boatramp_core::security::SecurityProfile;
#[tokio::test]
async fn assemble_produces_a_serving_node_over_a_temp_store() {
use axum::body::Body;
use axum::http::{Request, StatusCode};
use tower::ServiceExt;
let tmp = tempfile::tempdir().unwrap();
let storage: Arc<dyn Storage> = Arc::new(boatramp_storage::FsStorage::new(tmp.path()));
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let config = ServerConfig::default();
let options = boatramp_server::ServerOptions {
posture: SecurityProfile::MultiTenant.preset(),
..Default::default()
};
let node = assemble(NodeInput {
config: &config,
data_dir: tmp.path(),
storage,
kv,
auth: boatramp_server::Auth::disabled(),
options,
serve_addr: None,
watch_provider: None,
provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
messaging: None,
is_leader: Arc::new(|| true),
node_id: 0,
worker_exe: None,
})
.await
.expect("assemble a node over a temp store");
assert!(
!node
.deploy
.ensure_default_project()
.await
.expect("read the default project"),
"assemble should have materialized the default project"
);
let router =
boatramp_server::router_with(node.deploy, node.auth, node.handlers, node.options);
let response = router
.oneshot(
Request::builder()
.uri("/healthz")
.body(Body::empty())
.unwrap(),
)
.await
.expect("route /healthz");
assert_eq!(response.status(), StatusCode::OK);
}
}