#[cfg(feature = "handlers")]
use crate::error::Error;
use crate::error::Result;
use boatramp_core::deploy::DeployStore;
use boatramp_core::envelope::KeyEnvelope;
use boatramp_core::kv::KvStore;
use std::path::Path;
use std::sync::Arc;
#[cfg(feature = "handlers")]
const DEFAULT_ASYNC_TIMEOUT_MS: u64 = 15 * 60 * 1000;
#[cfg(feature = "handlers")]
const DEFAULT_ASYNC_CONCURRENCY: usize = 8;
#[cfg(feature = "handlers")]
const DEFAULT_STREAMING_TIMEOUT_MS: u64 = 15 * 60 * 1000;
#[cfg(feature = "handlers")]
const DEFAULT_STREAMING_CONCURRENCY: usize = 64;
#[cfg(feature = "handlers")]
fn sync_concurrency_for_cores(cores: usize) -> usize {
cores.saturating_mul(4).max(64)
}
#[cfg(feature = "handlers")]
fn default_sync_concurrency() -> usize {
let cores = std::thread::available_parallelism()
.map(std::num::NonZeroUsize::get)
.unwrap_or(1);
sync_concurrency_for_cores(cores)
}
#[cfg(feature = "handlers")]
fn resolve_sync_concurrency(configured: Option<usize>) -> usize {
match configured {
None | Some(0) => default_sync_concurrency(),
Some(n) => n,
}
}
#[cfg(all(test, feature = "handlers"))]
mod sync_concurrency_default_tests {
use super::sync_concurrency_for_cores;
#[test]
fn sync_lane_default_is_additive_then_scales_with_cores() {
for cores in [0usize, 1, 2, 4, 8, 16] {
assert_eq!(
sync_concurrency_for_cores(cores),
64,
"≤16 vCPU must stay at the historical 64 (cores={cores})"
);
}
assert_eq!(
sync_concurrency_for_cores(17),
68,
"just past the crossover scales"
);
assert_eq!(sync_concurrency_for_cores(32), 128);
assert_eq!(sync_concurrency_for_cores(64), 256);
}
#[test]
fn sync_max_concurrency_zero_is_not_disable() {
use super::{default_sync_concurrency, resolve_sync_concurrency};
let default = default_sync_concurrency();
assert!(default >= 64, "host default is floored at 64");
assert_eq!(
resolve_sync_concurrency(Some(0)),
default,
"0 is invalid → host default, NOT 1 and NOT disable"
);
assert_eq!(
resolve_sync_concurrency(None),
default,
"unset → host default"
);
assert_eq!(
resolve_sync_concurrency(Some(1)),
1,
"a positive value is honored verbatim"
);
assert_eq!(resolve_sync_concurrency(Some(200)), 200);
}
}
#[cfg(feature = "handlers")]
#[allow(clippy::too_many_arguments)]
pub async fn build_handler_runtime(
kv: Arc<dyn KvStore>,
storage: Arc<dyn boatramp_core::Storage>,
data_dir: &Path,
handlers_cfg: Option<&crate::config::HandlersConfig>,
messaging_override: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
max_blob_bytes: u64,
max_component_bytes: u64,
allow_guest_private_egress: bool,
self_egress_addrs: Vec<std::net::SocketAddr>,
guest_egress_extra_roots: Vec<rustls::pki_types::CertificateDer<'static>>,
allow_env_secret_refs: bool,
allow_guest_email: bool,
require_tenancy_declaration: bool,
allow_cross_tenant_db: bool,
deploy: &DeployStore,
secrets_envelope: Option<Arc<dyn KeyEnvelope>>,
) -> Result<boatramp_server::HandlerRuntime> {
#[cfg(not(feature = "email"))]
let _ = allow_guest_email;
const DEFAULT_BLOB_MARKER_CACHE: usize = 4096;
let marker_cache_entries = handlers_cfg
.and_then(|h| h.blob_marker_cache_entries)
.unwrap_or(DEFAULT_BLOB_MARKER_CACHE);
let storage: Arc<dyn boatramp_core::Storage> = Arc::new(
boatramp_storage::MarkerHeadCache::new(storage, marker_cache_entries),
);
let defaults = boatramp_handlers::Limits::default();
let mib = |mb: usize| mb.saturating_mul(1024 * 1024);
let sync_limits = boatramp_handlers::Limits {
timeout_ms: handlers_cfg
.and_then(|h| h.sync_max_timeout_ms)
.unwrap_or(defaults.timeout_ms),
memory_bytes: handlers_cfg
.and_then(|h| h.sync_max_memory_mb)
.map(mib)
.unwrap_or(defaults.memory_bytes),
max_concurrency: {
let configured = handlers_cfg.and_then(|h| h.sync_max_concurrency);
if configured == Some(0) {
tracing::warn!(
default = default_sync_concurrency(),
"[handlers] sync_max_concurrency = 0 is invalid: the global sync lane is a \
required node-wide safety ceiling and cannot be disabled (unlike \
serve_concurrency, where 0 means pass-through). Using the host-scaled default \
instead — set a positive value to override."
);
}
resolve_sync_concurrency(configured)
},
..defaults
};
let async_limits = boatramp_handlers::Limits {
timeout_ms: handlers_cfg
.and_then(|h| h.async_max_timeout_ms)
.unwrap_or(DEFAULT_ASYNC_TIMEOUT_MS),
max_concurrency: handlers_cfg
.and_then(|h| h.async_max_concurrency)
.unwrap_or(DEFAULT_ASYNC_CONCURRENCY),
fuel: handlers_cfg.and_then(|h| h.async_max_fuel),
memory_bytes: handlers_cfg
.and_then(|h| h.async_max_memory_mb)
.map(mib)
.unwrap_or(defaults.memory_bytes),
..sync_limits
};
let streaming_limits = boatramp_handlers::Limits {
timeout_ms: handlers_cfg
.and_then(|h| h.streaming_max_timeout_ms)
.unwrap_or(DEFAULT_STREAMING_TIMEOUT_MS),
max_concurrency: handlers_cfg
.and_then(|h| h.streaming_max_concurrency)
.unwrap_or(DEFAULT_STREAMING_CONCURRENCY),
fuel: handlers_cfg.and_then(|h| h.streaming_max_fuel),
memory_bytes: handlers_cfg
.and_then(|h| h.streaming_max_memory_mb)
.map(mib)
.unwrap_or(defaults.memory_bytes),
..sync_limits
};
let outbound_timeout = handlers_cfg
.and_then(|h| h.outbound_timeout_ms)
.map(std::time::Duration::from_millis);
const MAX_INSTANCE_CACHE: usize = 16_384;
let requested_cache = handlers_cfg
.and_then(|h| h.instance_cache_size)
.unwrap_or(64);
let instance_cache_size = requested_cache.clamp(1, MAX_INSTANCE_CACHE);
if instance_cache_size != requested_cache {
tracing::warn!(
requested = requested_cache,
applied = instance_cache_size,
"[handlers] instance_cache_size clamped to [1, 16384]"
);
}
let compile_cache_dir = Some(data_dir);
let engine = if handlers_cfg.is_some_and(|h| h.pooling) {
boatramp_handlers::HandlerEngine::with_pooling_lanes(
sync_limits,
async_limits,
streaming_limits,
instance_cache_size,
compile_cache_dir,
)?
} else {
boatramp_handlers::HandlerEngine::new_with_cache_dir(
sync_limits,
instance_cache_size,
compile_cache_dir,
)?
.with_async_limits(async_limits)
.with_streaming_limits(streaming_limits)
}
.with_outbound_timeout(outbound_timeout)
.with_private_egress(allow_guest_private_egress)
.with_self_egress(self_egress_addrs)
.with_guest_egress_extra_roots(guest_egress_extra_roots);
let sql = build_sql_backends(
handlers_cfg.and_then(|h| h.bindings.sql.as_ref()),
data_dir,
deploy,
&kv,
secrets_envelope.as_ref(),
&boatramp_core::env::SystemEnv,
)
.await?;
let max_unflushed = handlers_cfg
.and_then(|h| h.messaging_max_unflushed_msgs)
.unwrap_or(0);
if max_unflushed > 0 {
tracing::warn!(
max_unflushed,
"messaging: RELAXED publish durability ENABLED (messaging_max_unflushed_msgs={max_unflushed}) \
— publish() acks before the WAL flush; up to {max_unflushed} acknowledged-but-unflushed \
messages are lost on a process crash / OOM / SIGKILL / power loss. Control-plane \
durability is UNAFFECTED. Set 0 (the default) for stronger-than-JetStream durability."
);
}
let messaging: Arc<dyn boatramp_core::messaging::Messaging> = messaging_override
.unwrap_or_else(|| {
Arc::new(
boatramp_core::messaging::LogMessaging::new(storage.clone(), kv.clone())
.with_max_unflushed(max_unflushed),
)
});
let kv_for_secrets = kv.clone();
#[cfg(feature = "email")]
let kv_for_email = kv.clone();
#[cfg(feature = "email")]
let messaging_for_email = messaging.clone();
#[cfg(feature = "email")]
let email_envelope = secrets_envelope.clone();
let runtime =
boatramp_server::HandlerRuntime::new(engine, kv, storage, Some(sql), Some(messaging));
runtime.set_max_blob_bytes(max_blob_bytes);
runtime.set_max_component_bytes(max_component_bytes);
if let Some(cap) = handlers_cfg.and_then(|h| h.serve_concurrency) {
runtime.set_serve_concurrency(cap);
}
match runtime.serve_concurrency() {
0 => tracing::info!("serve-admission gate DISABLED ([handlers] serve_concurrency = 0)"),
cap => tracing::info!(
serve_concurrency = cap,
"serve-admission gate active: ≤ {cap} concurrent serves per component"
),
}
tracing::info!(
sync_concurrency = sync_limits.max_concurrency,
"sync lane: ≤ {} concurrent live serves node-wide",
sync_limits.max_concurrency
);
{
let defaults = boatramp_server::DeliveryConfig::default();
let safetynet = handlers_cfg
.and_then(|h| h.messaging_safetynet_interval_ms)
.map(std::time::Duration::from_millis)
.unwrap_or(defaults.safetynet_interval);
let rebuild = handlers_cfg
.and_then(|h| h.messaging_readyset_rebuild_interval_ms)
.map(std::time::Duration::from_millis)
.unwrap_or(defaults.rebuild_interval);
runtime.set_delivery_config(boatramp_server::DeliveryConfig {
safetynet_interval: safetynet,
rebuild_interval: rebuild,
});
}
{
const JWKS_HARD_TTL_CEILING: u64 = 86_400;
let jwks_tier = |refresh_secs: u64,
hard_secs: u64,
tier: &str|
-> boatramp_server::JwksTier {
let refresh = refresh_secs.clamp(60, JWKS_HARD_TTL_CEILING / 2);
let mut hard = hard_secs.clamp(60, JWKS_HARD_TTL_CEILING);
if hard_secs > JWKS_HARD_TTL_CEILING {
tracing::warn!(
tier,
requested_hard_ttl_secs = hard_secs,
applied_hard_ttl_secs = hard,
"[handlers] jwks hard_ttl exceeds the 24 h sanity ceiling — clamped down (a \
longer revocation window is almost certainly a misconfiguration)"
);
}
if hard <= refresh {
let adjusted = refresh.saturating_mul(2).min(JWKS_HARD_TTL_CEILING);
tracing::warn!(
tier,
refresh_secs = refresh,
requested_hard_ttl_secs = hard,
applied_hard_ttl_secs = adjusted,
"[handlers] jwks hard_ttl must exceed refresh — clamped up"
);
hard = adjusted;
}
boatramp_server::JwksTier {
refresh: std::time::Duration::from_secs(refresh),
hard_ttl: std::time::Duration::from_secs(hard),
}
};
let own = jwks_tier(
handlers_cfg
.and_then(|h| h.jwks_refresh_secs)
.unwrap_or(120),
handlers_cfg
.and_then(|h| h.jwks_hard_ttl_secs)
.unwrap_or(600),
"own",
);
let foreign = jwks_tier(
handlers_cfg
.and_then(|h| h.jwks_foreign_refresh_secs)
.unwrap_or(60),
handlers_cfg
.and_then(|h| h.jwks_foreign_hard_ttl_secs)
.unwrap_or(300),
"foreign",
);
let prewarm = handlers_cfg.and_then(|h| h.jwks_prewarm).unwrap_or(true);
boatramp_server::set_jwks_config(boatramp_server::JwksConfig {
own,
foreign,
prewarm,
});
}
runtime.set_allow_env_secret_refs(allow_env_secret_refs);
runtime.set_tenancy_posture(require_tenancy_declaration, allow_cross_tenant_db);
if let Some(envelope) = secrets_envelope {
runtime.set_secret_store(Arc::new(boatramp_core::secret_store::SecretStore::new(
kv_for_secrets,
envelope,
)));
}
#[cfg(feature = "email")]
if allow_guest_email && let Some(envelope) = email_envelope {
let store = Arc::new(boatramp_core::email_config::EmailProfileStore::new(
kv_for_email,
envelope,
));
let backend = Arc::new(boatramp_handlers::LettreBackend::new(
allow_guest_private_egress,
));
let spool = boatramp_server::NodeEmailSpool::spawn(
backend,
Some(messaging_for_email),
store.clone(),
);
runtime.set_email_profile_store(store);
runtime.set_email_spool(spool);
}
Ok(runtime)
}
#[cfg(feature = "handlers")]
async fn build_sql_backends(
cfg: Option<&crate::config::SqlBindingConfig>,
data_dir: &Path,
deploy: &DeployStore,
kv: &Arc<dyn KvStore>,
secrets_envelope: Option<&Arc<dyn KeyEnvelope>>,
env_source: &dyn boatramp_core::env::EnvSource,
) -> Result<Arc<dyn boatramp_core::sql::SqlBackends>> {
let resolve_env = |var: &Option<String>| -> Result<Option<String>> {
match var {
Some(var) => Ok(Some(
env_source
.get(var)
.ok_or_else(|| Error::SqlEnvUnset(var.clone()))?,
)),
None => Ok(None),
}
};
let backend = match cfg.and_then(|c| c.url.as_ref()) {
Some(url) => {
let cfg = cfg.expect("url implies cfg");
let admin_url = cfg.admin_url.as_ref().ok_or(Error::SqlAdminUrlRequired)?;
let token = resolve_env(&cfg.token_env)?.unwrap_or_default();
let admin_token = resolve_env(&cfg.admin_token_env)?;
let backends = boatramp_storage::LibsqlSqlBackends::remote(
url.clone(),
admin_url.clone(),
token,
admin_token,
);
match &cfg.replica_url {
Some(replica_url) => backends.with_read_replica(replica_url.clone()),
None => backends,
}
}
None => {
let dir = cfg
.and_then(|c| c.dir.clone())
.unwrap_or_else(|| data_dir.join("handlers-sql"));
boatramp_storage::LibsqlSqlBackends::local(dir)
}
};
let preview_mode = match cfg.and_then(|c| c.preview_mode.as_deref()) {
None | Some("empty") => boatramp_core::sql::PreviewSqlMode::Empty,
Some("branch") => boatramp_core::sql::PreviewSqlMode::Branch,
Some("shared") => boatramp_core::sql::PreviewSqlMode::Shared,
Some(other) => return Err(Error::UnknownPreviewMode(other.to_string())),
};
let preview_init = match cfg.and_then(|c| c.preview_init.as_ref()) {
Some(path) => {
Some(
std::fs::read_to_string(path).map_err(|err| Error::PreviewInitRead {
path: path.clone(),
source: err,
})?,
)
}
None => None,
};
let default: Arc<dyn boatramp_core::sql::SqlBackends> =
Arc::new(backend.with_preview_policy(preview_mode, preview_init));
let databases = cfg.map(|c| &c.databases);
if databases.is_none_or(std::collections::BTreeMap::is_empty) {
return Ok(default);
}
let databases = databases.expect("checked non-empty above");
#[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
{
use boatramp_core::sql::SqlBackend;
use boatramp_storage::sql_sqlx::{
CompositeSqlBackends, ExternalSqlKind, ExternalSqlOptions, connect,
};
let timeout = |db: &crate::config::ExternalDatabaseConfig| {
db.connect_timeout_secs.map(std::time::Duration::from_secs)
};
let mut composite = CompositeSqlBackends::new(default);
for (name, db) in databases {
let kind = ExternalSqlKind::parse(&db.kind).ok_or_else(|| Error::SqlExternalKind {
name: name.clone(),
kind: db.kind.clone(),
})?;
if db.compute.as_deref().is_some_and(|c| !c.is_empty()) {
if let Some(var) = db.password_env.as_deref().filter(|v| !v.is_empty()) {
let workload = db.compute.as_deref().expect("compute checked above");
let password = env_source
.get(var)
.ok_or_else(|| Error::SqlEnvUnset(var.into()))?;
let resolver = Arc::new(crate::managed_sql::DeployEndpointResolver::new(
deploy.clone(),
boatramp_core::project::DEFAULT_PROJECT,
));
let external: Arc<dyn SqlBackend> = Arc::new(
boatramp_storage::sql_compute::ComputeResolvedSqlBackend::new(
resolver,
workload,
kind,
db.database.clone().unwrap_or_default(),
db.user.clone().unwrap_or_default(),
password,
db.pool_max,
db.read_only,
timeout(db),
),
);
composite = composite.with_external(name.clone(), external, db.allow_preview);
continue;
}
let envelope = secrets_envelope
.cloned()
.ok_or_else(|| Error::SqlManagedNeedsSecrets(name.clone()))?;
let resolver = crate::tenant_sql::NodeTenantSqlResolver::new(
deploy.clone(),
kv.clone(),
envelope,
db,
)
.expect("a compute-backed managed binding builds a per-tenant resolver");
let site_scoped = resolver.site_scoped();
composite = composite.with_per_tenant(
name.clone(),
Arc::new(resolver),
site_scoped,
db.allow_preview,
);
} else {
if db.url_env.trim().is_empty() {
return Err(Error::SqlExternalUrlEnvMissing(name.clone()));
}
let url = env_source
.get(&db.url_env)
.ok_or_else(|| Error::SqlEnvUnset(db.url_env.clone()))?;
let read_url = match &db.read_url_env {
Some(var) => Some(
env_source
.get(var)
.ok_or_else(|| Error::SqlEnvUnset(var.clone()))?,
),
None => None,
};
let opts = ExternalSqlOptions::new(url)
.with_read_url(read_url)
.with_max_connections(db.pool_max)
.read_only(db.read_only)
.with_connect_timeout(timeout(db));
let external: Arc<dyn SqlBackend> =
connect(kind, &opts).map_err(|source| Error::SqlExternalConnect {
name: name.clone(),
source,
})?;
composite = composite.with_external(name.clone(), external, db.allow_preview);
}
}
Ok(Arc::new(composite))
}
#[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
{
let _ = (deploy, kv, secrets_envelope);
let name = databases.keys().next().cloned().unwrap_or_default();
Err(Error::SqlExternalUnavailable(name))
}
}
#[cfg(not(feature = "handlers"))]
#[allow(clippy::too_many_arguments)]
pub async fn build_handler_runtime(
_kv: Arc<dyn KvStore>,
_storage: Arc<dyn boatramp_core::Storage>,
_data_dir: &Path,
_handlers_cfg: Option<&crate::config::HandlersConfig>,
_messaging_override: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
_max_blob_bytes: u64,
_max_component_bytes: u64,
_allow_guest_private_egress: bool,
_self_egress_addrs: Vec<std::net::SocketAddr>,
_guest_egress_extra_roots: Vec<rustls::pki_types::CertificateDer<'static>>,
_allow_env_secret_refs: bool,
_allow_guest_email: bool,
_require_tenancy_declaration: bool,
_allow_cross_tenant_db: bool,
_deploy: &DeployStore,
_secrets_envelope: Option<Arc<dyn KeyEnvelope>>,
) -> Result<boatramp_server::HandlerRuntime> {
Ok(boatramp_server::HandlerRuntime::disabled())
}
#[cfg(all(test, any(feature = "sql-postgres", feature = "sql-mysql")))]
mod tests {
use super::*;
use std::result::Result;
use async_trait::async_trait;
use boatramp_core::envelope::EnvelopeError;
use boatramp_core::kv::MemoryKv;
use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
struct TestEnvelope;
#[async_trait]
impl KeyEnvelope for TestEnvelope {
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())
}
}
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 managed_sql_cfg() -> crate::config::SqlBindingConfig {
let mut databases = std::collections::BTreeMap::new();
databases.insert(
"analytics".to_string(),
crate::config::ExternalDatabaseConfig {
kind: "postgres".into(),
compute: Some("pg".into()),
database: Some("analytics".into()),
user: Some("app".into()),
..Default::default()
},
);
crate::config::SqlBindingConfig {
databases,
..Default::default()
}
}
#[tokio::test]
async fn managed_sql_fails_closed_without_secrets() {
let tmp = tempfile::tempdir().unwrap();
let deploy = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let cfg = managed_sql_cfg();
match build_sql_backends(
Some(&cfg),
tmp.path(),
&deploy,
&kv,
None,
&boatramp_core::env::SystemEnv,
)
.await
{
Err(Error::SqlManagedNeedsSecrets(name)) => assert_eq!(name, "analytics"),
Ok(_) => panic!("a managed DB without [secrets] must fail closed, got Ok"),
Err(other) => panic!("expected SqlManagedNeedsSecrets, got: {other}"),
}
}
#[tokio::test]
async fn managed_sql_builds_lazily_and_seals_the_credential() {
let tmp = tempfile::tempdir().unwrap();
let deploy = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let envelope: Arc<dyn KeyEnvelope> = Arc::new(TestEnvelope);
let cfg = managed_sql_cfg();
let backends = build_sql_backends(
Some(&cfg),
tmp.path(),
&deploy,
&kv,
Some(&envelope),
&boatramp_core::env::SystemEnv,
)
.await
.expect("managed sql builds without a live DB (lazy connect)");
assert!(
kv.get("managed-sql-cred/default/pg")
.await
.unwrap()
.is_none(),
"nothing sealed at build — per-tenant credentials are minted on first open"
);
let _ = backends
.database("default", "blog", "analytics")
.await
.unwrap();
let sealed = kv
.get("managed-sql-cred/default/pg")
.await
.unwrap()
.expect("credential sealed on first resolve under the default project");
assert_ne!(sealed.len(), 0);
}
}