use std::net::SocketAddr;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use boatramp_core::cache_coherence::Changelog;
use boatramp_core::deploy::DeployStore;
use boatramp_core::kv::{CachedKv, KvStore};
use boatramp_core::migrate;
use boatramp_node::backends::{BlobBackend, KvBackend};
use clap::ValueEnum;
use crate::config::ServerConfig;
pub(crate) use boatramp_node::backends::build_kv as build_control_plane_kv;
use boatramp_node::blobs::{build_blobs, BlobArgs};
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[cfg(not(feature = "cluster"))]
#[error(
"[cluster] config is present but this build has no cluster support; \
rebuild with `--features cluster`"
)]
NoClusterSupport,
#[cfg(not(feature = "tls"))]
#[error("this build has no TLS support; rebuild with `--features tls`")]
NoTlsSupport,
#[cfg(not(feature = "acme-dns"))]
#[error("this build has no ACME DNS-01 support; rebuild with `--features acme-dns`")]
NoAcmeDnsSupport,
#[cfg(feature = "cluster")]
#[error("invalid auth root private key: {0}")]
AuthPrivKey(String),
#[cfg(any(feature = "cluster", feature = "tls"))]
#[error(transparent)]
RpkTls(#[from] boatramp_rpktls::RpkError),
#[cfg(feature = "cluster")]
#[error(
"refusing to serve the peer mesh on {0} (non-loopback) with an empty trust \
set: found with --cluster-init or join with --cluster-join <ticket>"
)]
MeshUnconfigured(std::net::SocketAddr),
#[cfg(feature = "cluster")]
#[error("cluster startup: {0}")]
ClusterStartup(String),
#[cfg(all(feature = "cluster", feature = "acme-dns"))]
#[error("secrets envelope: {0}")]
Envelope(String),
#[cfg(feature = "oidc")]
#[error("OIDC setup failed: {0}")]
OidcSetup(String),
#[cfg(feature = "oidc")]
#[error(
"OIDC is enabled without an audience, but the security posture requires one \
(set --oidc-audience, or relax `oidc_require_audience`)"
)]
OidcAudienceRequired,
#[cfg(feature = "tls")]
#[error("--tls-cert is required for --tls custom")]
TlsCertRequired,
#[cfg(feature = "tls")]
#[error("--tls-key is required for --tls custom")]
TlsKeyRequired,
#[cfg(feature = "http3")]
#[error("no certificates in {0}")]
NoCert(String),
#[cfg(feature = "http3")]
#[error("no private key in {0}")]
NoPrivateKey(String),
#[cfg(feature = "tls")]
#[error("at least one --acme-domain is required for --tls acme")]
NoAcmeDomain,
#[cfg(feature = "acme-dns")]
#[error("at least one --acme-domain is required for --tls acme-dns")]
NoAcmeDomainDns,
#[cfg(feature = "acme-dns")]
#[error("unknown --acme-dns-provider {0:?} (expected manual | cloudflare | route53 | oci)")]
UnknownDnsProvider(String),
#[cfg(all(feature = "cluster", feature = "acme-dns"))]
#[error(
"no certificates available yet — awaiting the cluster leader to issue (retry shortly)"
)]
NoCertsYet,
#[error(transparent)]
Assembly(#[from] boatramp_node::Error),
#[error(transparent)]
Security(#[from] boatramp_core::security::SecurityError),
#[error(transparent)]
Serve(#[from] boatramp_server::ServeError),
#[cfg(any(feature = "tls", feature = "cluster"))]
#[error(transparent)]
Io(#[from] std::io::Error),
#[cfg(feature = "acme-dns")]
#[error(transparent)]
AcmeDns(#[from] crate::acme_dns::Error),
#[cfg(feature = "http3")]
#[error(transparent)]
Http3(#[from] boatramp_server::Http3Error),
#[cfg(feature = "tls")]
#[error(transparent)]
Rustls(#[from] rustls::Error),
#[cfg(all(feature = "cluster", feature = "acme-dns"))]
#[error(transparent)]
ClusterTls(#[from] crate::cluster_tls::Error),
#[cfg(feature = "cluster")]
#[error(transparent)]
Bootstrap(Box<boatramp_cluster::node::BootstrapError>),
#[cfg(any(feature = "slatedb", feature = "cloudflare-kv"))]
#[error(transparent)]
Kv(#[from] boatramp_core::kv::KvError),
#[cfg(feature = "handlers")]
#[error(transparent)]
Handler(#[from] boatramp_handlers::HandlerError),
#[error(
"the control-plane store is not migrated to the project-scoped (0.2.0) layout; \
run `boatramp migrate` first, or start `serve --auto-migrate`"
)]
UnmigratedStore,
#[error("store migration failed: {0}")]
Migrate(String),
}
impl From<boatramp_core::migrate::MigrateError> for Error {
fn from(e: boatramp_core::migrate::MigrateError) -> Self {
Self::Migrate(e.to_string())
}
}
type Result<T> = std::result::Result<T, Error>;
#[cfg(feature = "cluster")]
impl From<boatramp_cluster::node::BootstrapError> for Error {
fn from(e: boatramp_cluster::node::BootstrapError) -> Self {
Self::Bootstrap(Box::new(e))
}
}
#[cfg(feature = "cluster")]
const _: () = assert!(std::mem::size_of::<Error>() <= 128);
#[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)]
enum TlsMode {
Off,
Custom,
Acme,
AcmeDns,
Rpk,
}
#[derive(Debug, clap::Args)]
pub struct ServeArgs {
#[arg(long, env = "BOATRAMP_ADDR")]
addr: Option<SocketAddr>,
#[arg(long, env = "BOATRAMP_DATA_DIR")]
data_dir: Option<PathBuf>,
#[arg(long, value_enum, default_value_t = BlobBackend::Fs)]
blobs: BlobBackend,
#[arg(long, value_enum, default_value_t = KvBackend::Slatedb)]
kv: KvBackend,
#[arg(long)]
auto_migrate: bool,
#[arg(long, env = "BOATRAMP_S3_BUCKET")]
s3_bucket: Option<String>,
#[arg(long, env = "BOATRAMP_S3_ENDPOINT")]
s3_endpoint: Option<String>,
#[arg(long, env = "BOATRAMP_S3_REGION")]
s3_region: Option<String>,
#[arg(long, env = "BOATRAMP_S3_PATH_STYLE")]
s3_path_style: bool,
#[arg(long, env = "BOATRAMP_GCS_BUCKET")]
gcs_bucket: Option<String>,
#[arg(long, env = "BOATRAMP_GCS_ENDPOINT")]
gcs_endpoint: Option<String>,
#[arg(long, env = "BOATRAMP_GCS_ANONYMOUS")]
gcs_anonymous: bool,
#[arg(long, env = "BOATRAMP_AZURE_ACCOUNT")]
azure_account: Option<String>,
#[arg(long, env = "BOATRAMP_AZURE_CONTAINER")]
azure_container: Option<String>,
#[arg(long, env = "BOATRAMP_AZURE_ACCESS_KEY")]
azure_access_key: Option<String>,
#[arg(long, env = "BOATRAMP_AZURE_EMULATOR")]
azure_emulator: bool,
#[arg(long, default_value_t = 256)]
cache_entries: usize,
#[arg(long, env = "BOATRAMP_AUTH_ROOT_PRIVATE_KEY")]
auth_root_private_key: Option<String>,
#[arg(long, env = "BOATRAMP_AUTH_ROOT_PUBLIC_KEY")]
auth_root_public_key: Option<String>,
#[arg(long, env = "BOATRAMP_BOOTSTRAP_SECRET")]
bootstrap_secret: Option<String>,
#[arg(long, value_enum, default_value_t = TlsMode::Off)]
tls: TlsMode,
#[arg(long, requires = "tls_key")]
tls_cert: Option<PathBuf>,
#[arg(long, requires = "tls_cert")]
tls_key: Option<PathBuf>,
#[arg(long = "acme-domain")]
acme_domain: Vec<String>,
#[arg(long, default_value = "https://acme-v02.api.letsencrypt.org/directory")]
acme_directory: String,
#[arg(long)]
acme_contact: Option<String>,
#[arg(long)]
acme_ca_cert: Option<PathBuf>,
#[arg(long, default_value = "./data/acme")]
acme_cache: PathBuf,
#[arg(long, default_value = "manual")]
acme_dns_provider: String,
#[arg(long)]
acme_wildcard_preview: bool,
#[arg(long, env = "BOATRAMP_MAX_UPLOAD_BYTES")]
max_upload_bytes: Option<u64>,
#[arg(long, env = "BOATRAMP_UPLOAD_IDLE_TIMEOUT")]
upload_idle_timeout_secs: Option<u64>,
#[arg(long, env = "BOATRAMP_MAX_CONCURRENT_UPLOADS")]
max_concurrent_uploads: Option<usize>,
#[arg(long, env = "BOATRAMP_HTTP_REDIRECT_ADDR")]
http_redirect_addr: Option<SocketAddr>,
#[arg(long, env = "BOATRAMP_DEFAULT_SITE")]
default_site: Option<String>,
#[arg(long, env = "BOATRAMP_POP_ORIGIN")]
pop_origin: Option<String>,
#[arg(long, env = "BOATRAMP_CLUSTER_RATE_LIMIT")]
cluster_rate_limit: bool,
#[arg(long, env = "BOATRAMP_CLUSTER_INIT")]
cluster_init: bool,
#[arg(long, env = "BOATRAMP_CLUSTER_ADVERTISE_ADDR")]
cluster_advertise_addr: Option<String>,
#[arg(long, env = "BOATRAMP_CLUSTER_JOIN")]
cluster_join: Option<String>,
#[arg(long, env = "BOATRAMP_SHARED_CACHE_COHERENCE")]
shared_cache_coherence: bool,
#[arg(long, env = "BOATRAMP_PROTECT_PREVIEWS")]
protect_previews: bool,
#[cfg(feature = "http3")]
#[arg(long)]
http3: bool,
#[cfg(feature = "oidc")]
#[arg(long, env = "BOATRAMP_OIDC_ISSUER")]
oidc_issuer: Option<String>,
#[cfg(feature = "oidc")]
#[arg(long, env = "BOATRAMP_OIDC_AUDIENCE")]
oidc_audience: Option<String>,
#[cfg(feature = "oidc")]
#[arg(long, env = "BOATRAMP_OIDC_SCOPE_CLAIM")]
oidc_scope_claim: Option<String>,
}
impl ServeArgs {
fn server_limits(
&self,
serve_cfg: &crate::config::ServeConfig,
posture: &boatramp_core::security::SecurityPosture,
) -> boatramp_server::ServerLimits {
boatramp_server::ServerLimits {
max_upload_bytes: self
.max_upload_bytes
.or(serve_cfg.max_upload_bytes)
.or_else(|| (posture.max_upload_bytes != 0).then_some(posture.max_upload_bytes)),
upload_idle_timeout: self
.upload_idle_timeout_secs
.or(serve_cfg.upload_idle_timeout_secs)
.map(std::time::Duration::from_secs),
max_concurrent_uploads: self
.max_concurrent_uploads
.or(serve_cfg.max_concurrent_uploads),
}
}
}
pub async fn run(args: ServeArgs, config: &ServerConfig) -> Result<()> {
let serve_cfg = config.serve.clone().unwrap_or_default();
let posture = config.security.clone().unwrap_or_default().resolve()?;
let mut options = boatramp_server::ServerOptions {
limits: args.server_limits(&serve_cfg, &posture),
default_site: args.default_site.clone().or(serve_cfg.default_site.clone()),
pop_origin: args.pop_origin.clone().or(serve_cfg.pop_origin.clone()),
protect_previews: args.protect_previews || serve_cfg.protect_previews,
posture,
served_over_tls: !matches!(args.tls, TlsMode::Off),
bootstrap_secret: args
.bootstrap_secret
.clone()
.or(serve_cfg.bootstrap_secret.clone()),
..Default::default()
};
let cluster_rate_limit = args.cluster_rate_limit || serve_cfg.cluster_rate_limit;
let addr = args
.addr
.or(serve_cfg.addr)
.unwrap_or_else(|| "127.0.0.1:8080".parse().expect("valid default addr"));
options.implicit_routing = options.posture.allow_implicit_routing || addr.ip().is_loopback();
if let Some(console) = serve_cfg.console.as_ref().filter(|c| c.enabled) {
#[cfg(feature = "console")]
{
options.console = Some(boatramp_server::console::ConsoleMount::resolve(
console.host.clone(),
console.path.clone(),
));
}
#[cfg(not(feature = "console"))]
{
let _ = console;
tracing::warn!(
"[serve.console] enabled but this build lacks the `console` feature — \
the console is not served"
);
}
}
let data_dir = args
.data_dir
.clone()
.or(serve_cfg.data_dir)
.unwrap_or_else(|| PathBuf::from("./data"));
let notify_tier = serve_cfg.blob_notify_tier;
let notify_account = serve_cfg.blob_notify_account_id.clone();
let blob_args = BlobArgs {
blobs: args.blobs,
s3_bucket: args.s3_bucket.clone(),
s3_endpoint: args.s3_endpoint.clone(),
s3_region: args.s3_region.clone(),
s3_path_style: args.s3_path_style,
gcs_bucket: args.gcs_bucket.clone(),
gcs_endpoint: args.gcs_endpoint.clone(),
gcs_anonymous: args.gcs_anonymous,
azure_account: args.azure_account.clone(),
azure_container: args.azure_container.clone(),
azure_access_key: args.azure_access_key.clone(),
azure_emulator: args.azure_emulator,
};
let built_blobs = build_blobs(&blob_args, &data_dir, notify_tier, notify_account).await?;
let storage = built_blobs.storage.clone();
#[cfg(feature = "cluster")]
if config.cluster.is_some() || args.cluster_init || args.cluster_join.is_some() {
let cluster_cfg = config.cluster.clone().unwrap_or_else(|| {
crate::config::ClusterConfig {
listen: std::net::SocketAddr::new(addr.ip(), DEFAULT_MESH_PORT),
root_pubkeys: Vec::new(),
seeds: Vec::new(),
join_token: None,
store_dir: None,
mesh: None,
}
});
return run_cluster(
args,
config,
cluster_cfg,
addr,
data_dir,
built_blobs,
options,
)
.await;
}
#[cfg(not(feature = "cluster"))]
if config.cluster.is_some() {
return Err(Error::NoClusterSupport);
}
let kv_backend = boatramp_node::backends::build_kv(args.kv, &data_dir).await?;
let shared_coherence = args.shared_cache_coherence || serve_cfg.shared_cache_coherence;
let changelog = shared_coherence
.then(|| Arc::new(Changelog::new(kv_backend.clone(), CHANGELOG_RETENTION_SECS)));
let mut cached = CachedKv::new(kv_backend.clone(), args.cache_entries);
if let Some(changelog) = &changelog {
cached = cached.with_publisher(changelog.clone());
}
let kv: Arc<dyn KvStore> = Arc::new(cached);
match migrate::status(kv.as_ref()).await? {
migrate::Status::Ready => {}
migrate::Status::Dual => tracing::warn!(
"control-plane store is in the 2-dual soak window; \
run `boatramp migrate --finalize` to reclaim the old-layout keys"
),
migrate::Status::NeedsMigration => {
if args.auto_migrate {
tracing::warn!(
"control-plane store is on the pre-0.2.0 layout; running a one-shot \
project re-keying migration (--auto-migrate)"
);
let report =
migrate::migrate(kv.as_ref(), migrate::MigrateOptions::one_shot()).await?;
tracing::info!(
rekeyed = report.total_rekeyed(),
owner_entries = report.owner_entries,
"control-plane store migrated to the project-scoped layout"
);
} else {
return Err(Error::UnmigratedStore);
}
}
}
let kv_handle = kv.clone();
if cluster_rate_limit {
options.cluster_rate_limit_kv = Some(kv_backend.clone());
}
let daemon_runtime = Arc::new(boatramp_server::DaemonRuntime::new(
boatramp_server::config_baseline(&options),
));
options.daemon_runtime = Some(daemon_runtime.clone());
spawn_sighup_reload(kv.clone(), Some(daemon_runtime.clone()));
if let Some(changelog) = changelog {
spawn_cache_poller(changelog, kv.clone(), Some(daemon_runtime.clone()));
}
let auth = boatramp_node::auth::configure_auth(
serve_cfg.signer.as_ref(),
args.auth_root_private_key
.clone()
.or(serve_cfg.auth_root_private_key.clone()),
args.auth_root_public_key
.clone()
.or(serve_cfg.auth_root_public_key.clone()),
&mut options,
kv.clone(),
)
.await?;
configure_oidc(&args, &mut options).await?;
boatramp_node::auth::enforce_auth_bind(addr, &auth, &options.posture)?;
let boatramp_node::RunningNode {
deploy,
handlers,
auth,
options,
reconcile: _reconcile,
} = boatramp_node::assemble(boatramp_node::NodeInput {
config,
data_dir: data_dir.as_path(),
storage,
kv,
auth,
options,
watch_provider: built_blobs.watch_provider.clone(),
provision_tier: built_blobs.provision_tier,
messaging: None,
is_leader: Arc::new(|| true),
node_id: 0,
})
.await?;
tracing::info!(
blobs = ?args.blobs, kv = ?args.kv, tls = ?args.tls,
auth = !auth.is_disabled(), "starting boatramp"
);
#[cfg(feature = "tls")]
if !matches!(args.tls, TlsMode::Off) {
if let Some(redirect_addr) = args.http_redirect_addr.or(serve_cfg.http_redirect_addr) {
spawn_http_redirect(redirect_addr, deploy.clone(), posture);
}
}
let serve_result = match args.tls {
TlsMode::Off => boatramp_server::serve_with(addr, deploy, auth, handlers, options)
.await
.map_err(Error::Serve),
TlsMode::Custom => serve_custom(&args, addr, deploy, auth, handlers, options).await,
TlsMode::Acme => serve_acme(&args, addr, deploy, auth, handlers, options).await,
TlsMode::AcmeDns => serve_acme_dns(&args, addr, deploy, auth, handlers, options).await,
TlsMode::Rpk => serve_rpk(&args, addr, deploy, auth, handlers, options, &data_dir).await,
};
if let Err(e) = kv_handle.flush().await {
tracing::warn!(error = %e, "metadata store flush on shutdown failed");
}
serve_result
}
const CHANGELOG_RETENTION_SECS: u64 = 60;
fn spawn_cache_poller(
changelog: Arc<Changelog>,
cache: Arc<dyn KvStore>,
daemon: Option<Arc<boatramp_server::DaemonRuntime>>,
) {
use std::time::Duration;
tokio::spawn(async move {
let poll = Duration::from_secs(1);
let flush_every = Duration::from_secs(300);
let mut cursor = changelog.current_cursor().await;
let mut since_trim = Duration::ZERO;
let mut since_flush = Duration::ZERO;
loop {
tokio::time::sleep(poll).await;
let changed = changelog.poll(&mut cursor).await;
if !changed.is_empty() {
cache.invalidate_keys(&changed);
if let Some(daemon) = &daemon {
if changed.iter().any(|k| k.starts_with("daemon/")) {
daemon.notify_reload();
}
}
}
since_trim += poll;
if since_trim >= Duration::from_secs(30) {
changelog.trim().await;
since_trim = Duration::ZERO;
}
since_flush += poll;
if since_flush >= flush_every {
cache.invalidate_cache();
cursor = changelog.current_cursor().await;
since_flush = Duration::ZERO;
}
}
});
}
#[cfg(unix)]
fn spawn_sighup_reload(kv: Arc<dyn KvStore>, daemon: Option<Arc<boatramp_server::DaemonRuntime>>) {
tokio::spawn(async move {
let mut hup = match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::hangup()) {
Ok(sig) => sig,
Err(err) => {
tracing::warn!(%err, "could not install SIGHUP handler");
return;
}
};
while hup.recv().await.is_some() {
kv.invalidate_cache();
if let Some(daemon) = &daemon {
daemon.notify_reload();
}
tracing::info!("SIGHUP: invalidated config cache (next reads reload from the store)");
}
});
}
#[cfg(not(unix))]
fn spawn_sighup_reload(
_kv: Arc<dyn KvStore>,
_daemon: Option<Arc<boatramp_server::DaemonRuntime>>,
) {
}
#[cfg(feature = "tls")]
fn spawn_http_redirect(
addr: SocketAddr,
deploy: DeployStore,
posture: boatramp_core::security::SecurityPosture,
) {
tokio::spawn(async move {
tracing::info!(%addr, "serving HTTP→HTTPS redirect listener");
let service = boatramp_server::http_redirect_router(deploy, posture).into_make_service();
if let Err(err) = axum_server::bind(addr).serve(service).await {
tracing::error!(%addr, %err, "HTTP redirect listener failed");
}
});
}
#[cfg(feature = "cluster")]
const MESH_ROTATION_PROPAGATION: std::time::Duration = std::time::Duration::from_secs(2);
#[cfg(feature = "cluster")]
const DEFAULT_MESH_PORT: u16 = 7000;
#[cfg(feature = "cluster")]
fn parse_rotation_interval(spec: &str) -> Option<std::time::Duration> {
let spec = spec.trim();
let split = spec.find(|c: char| !c.is_ascii_digit())?;
let (num, unit) = spec.split_at(split);
let n: u64 = num.parse().ok()?;
let secs = match unit {
"s" => n,
"m" => n.checked_mul(60)?,
"h" => n.checked_mul(3600)?,
"d" => n.checked_mul(86_400)?,
_ => return None,
};
(secs > 0).then(|| std::time::Duration::from_secs(secs))
}
#[cfg(feature = "cluster")]
const JOIN_PROOF_MAX_SKEW_SECS: u64 = 300;
#[cfg(feature = "cluster")]
const MEMBER_ASSERTION_TTL_SECS: u64 = 300;
#[cfg(feature = "cluster")]
struct ClusterMeshControl {
node: Arc<boatramp_cluster::node::ClusterNode>,
issuer: Option<Arc<dyn boatramp_core::cose::Signer>>,
}
#[cfg(feature = "cluster")]
#[async_trait::async_trait]
impl boatramp_server::MeshControl for ClusterMeshControl {
async fn admit(
&self,
mesh_pubkey_hex: &str,
jti: &str,
possession_proof: &[u8],
proof_iat: u64,
now: u64,
advertise_addr: Option<&str>,
) -> std::result::Result<boatramp_server::JoinOutcome, String> {
use boatramp_server::JoinOutcome;
let fresh = proof_iat <= now.saturating_add(JOIN_PROOF_MAX_SKEW_SECS)
&& now <= proof_iat.saturating_add(JOIN_PROOF_MAX_SKEW_SECS);
if !fresh {
return Ok(JoinOutcome::ProofInvalid);
}
let Ok(spki) = boatramp_cluster::mesh::parse_public_key(mesh_pubkey_hex) else {
return Ok(JoinOutcome::ProofInvalid);
};
let challenge = boatramp_core::cose::join_challenge(jti, mesh_pubkey_hex, proof_iat);
if !boatramp_rpktls::verify_signature(&spki, &challenge, possession_proof) {
return Ok(JoinOutcome::ProofInvalid);
}
use boatramp_cluster::raft::AdmitOutcome;
match self
.node
.admit(mesh_pubkey_hex, jti, advertise_addr)
.await
.map_err(|e| e.to_string())?
{
AdmitOutcome::Admitted => {}
AdmitOutcome::Spent => return Ok(JoinOutcome::TokenSpent),
AdmitOutcome::Revoked => return Ok(JoinOutcome::Revoked),
}
let Some(issuer) = self.issuer.as_ref() else {
return Err("cluster node has no root signing key to vouch for members".to_string());
};
let mut members = Vec::new();
for (node_id, pubkey) in self.node.trusted_member_keys().await {
let assertion = boatramp_core::cose::mint_member_assertion(
node_id,
&pubkey,
MEMBER_ASSERTION_TTL_SECS,
now,
issuer.as_ref(),
)
.await
.map_err(|e| e.to_string())?;
members.push(assertion);
}
let addrs = self.node.peer_addrs();
Ok(JoinOutcome::Admitted { members, addrs })
}
async fn rotate_key(&self) -> std::result::Result<String, String> {
let new_pub = self
.node
.rotate_key(MESH_ROTATION_PROPAGATION)
.await
.map_err(|e| e.to_string())?;
Ok(new_pub.iter().map(|b| format!("{b:02x}")).collect())
}
async fn revoke(&self, node: u64) -> std::result::Result<(), String> {
self.node.revoke(node).await.map_err(|e| e.to_string())
}
async fn members(&self) -> std::result::Result<Vec<boatramp_server::MeshMember>, String> {
let addrs = self.node.peer_addrs();
Ok(self
.node
.members()
.into_iter()
.map(|m| boatramp_server::MeshMember {
node: m.node,
voter: m.voter,
caught_up: m.caught_up,
leader: m.leader,
addr: addrs.get(&m.node).cloned(),
})
.collect())
}
async fn promote(&self, node: u64) -> std::result::Result<(), String> {
self.node.promote(node).await.map_err(|e| e.to_string())
}
}
#[cfg(feature = "cluster")]
struct MeshWriteAuthz {
public: boatramp_core::cose::TokenPublicKey,
}
#[cfg(feature = "cluster")]
impl boatramp_cluster::http::ClientWriteAuthz for MeshWriteAuthz {
fn authorize(&self, capability: Option<&str>) -> bool {
let Some(token) = capability else {
return false;
};
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let Ok(verified) = boatramp_core::cose::verify(token, &self.public, now) else {
return false;
};
verified.roles.iter().any(|r| r.name == "cluster-write")
}
}
#[cfg(feature = "cluster")]
#[allow(clippy::type_complexity)]
async fn build_mesh_write_gate(
args: &ServeArgs,
config: &ServerConfig,
mesh_cfg: &crate::config::MeshConfig,
) -> Result<(
Option<String>,
Option<Arc<dyn boatramp_cluster::http::ClientWriteAuthz>>,
)> {
use boatramp_core::authz::GrantedRole;
use boatramp_core::cose::{self, Claims, LocalSigner, Signer};
if !mesh_cfg.gate_client_writes.unwrap_or(false) {
return Ok((None, None));
}
let priv_hex = args.auth_root_private_key.clone().or_else(|| {
config
.serve
.as_ref()
.and_then(|s| s.auth_root_private_key.clone())
});
let Some(priv_hex) = priv_hex else {
tracing::warn!(
"cluster.mesh.gate_client_writes is set but no token root private key is \
configured — mesh client-write gating is disabled"
);
return Ok((None, None));
};
let signer =
LocalSigner::from_private_hex(&priv_hex).map_err(|e| Error::AuthPrivKey(e.to_string()))?;
let claims = Claims {
roles: vec![GrantedRole::global("cluster-write")],
kind: cose::KIND_CLUSTER_WRITE.to_string(),
ttl_secs: None,
now_unix: 0,
};
let capability = cose::mint(&claims, &signer)
.await
.map_err(|e| Error::AuthPrivKey(format!("minting cluster-write capability: {e}")))?;
let authz: Arc<dyn boatramp_cluster::http::ClientWriteAuthz> = Arc::new(MeshWriteAuthz {
public: signer.public_key(),
});
Ok((Some(capability), Some(authz)))
}
#[cfg(all(feature = "cluster", feature = "acme-dns"))]
fn build_cert_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(feature = "cluster")]
struct DaemonConfigObserver(Arc<boatramp_server::DaemonRuntime>);
#[cfg(feature = "cluster")]
impl boatramp_cluster::raft::ApplyObserver for DaemonConfigObserver {
fn on_apply(&self, muts: &[boatramp_core::kv::WriteOp]) {
use boatramp_core::kv::WriteOp;
let touched = muts.iter().any(|m| match m {
WriteOp::Put(k, _) | WriteOp::Delete(k) => k.starts_with("daemon/"),
});
if touched {
self.0.notify_reload();
}
}
fn on_reset(&self, data: &std::collections::BTreeMap<String, Vec<u8>>) {
if data.keys().any(|k| k.starts_with("daemon/")) {
self.0.notify_reload();
}
}
}
#[cfg(feature = "cluster")]
#[allow(clippy::too_many_arguments)]
async fn run_cluster(
args: ServeArgs,
config: &ServerConfig,
mut cluster_cfg: crate::config::ClusterConfig,
addr: SocketAddr,
data_dir: PathBuf,
built_blobs: boatramp_node::blobs::BuiltBlobs,
mut options: boatramp_server::ServerOptions,
) -> Result<()> {
use boatramp_cluster::node::{build_node, ClusterParams};
let storage = built_blobs.storage.clone();
let store_dir = cluster_cfg
.store_dir
.clone()
.unwrap_or_else(|| data_dir.join("raft"));
let store_dir_existed = store_dir.exists();
let durable_kv: Arc<dyn KvStore> = Arc::new(
boatramp_storage::SlateKv::open_local_with_flush(
store_dir,
boatramp_node::backends::CONTROL_PLANE_FLUSH,
)
.await?,
);
let has_committed_state = boatramp_cluster::persist::has_committed_trust(&durable_kv)
.await
.map_err(|e| Error::ClusterStartup(e.to_string()))?;
let durable_kv_handle = durable_kv.clone();
use boatramp_cluster::mesh::{self, MeshIdentity, MeshTls, TrustSet};
let mut peers = std::collections::BTreeMap::new();
let mut genesis_trust = std::collections::BTreeMap::new();
let mesh_cfg = cluster_cfg.mesh.clone().unwrap_or_default();
let key_file = mesh_cfg
.key_file
.clone()
.unwrap_or_else(|| data_dir.join("mesh/identity.key"));
let identity = MeshIdentity::load_or_generate(&key_file)?;
let node_id = boatramp_cluster::raft::derive_node_id(identity.public_key());
tracing::info!(
node_id,
pubkey = %identity.public_key_hex(),
"cluster: mesh identity"
);
if let Some(blob) = args.cluster_join.as_deref() {
let ticket = crate::join::JoinTicket::decode(blob)
.map_err(|e| Error::ClusterStartup(e.to_string()))?;
cluster_cfg.seeds = ticket.seeds;
cluster_cfg.root_pubkeys = ticket.root_pubkeys;
cluster_cfg.join_token = Some(ticket.token);
}
let seeds_present = !cluster_cfg.seeds.is_empty();
let is_operator_founder =
std::env::var("BOATRAMP_POD_NAME").is_ok_and(|name| name.rsplit('-').next() == Some("0"));
let init_requested = args.cluster_init || is_operator_founder;
let action = crate::join::decide_startup(&crate::join::StartupInputs {
has_committed_state,
ever_member: store_dir_existed,
seeds_present,
init_requested,
});
let is_resume = matches!(action, crate::join::StartupAction::Resume);
let self_advertise = args
.cluster_advertise_addr
.clone()
.unwrap_or_else(|| format!("https://{}", cluster_cfg.listen));
let mut do_bootstrap = false;
match action {
crate::join::StartupAction::FailClosed(reason) => {
return Err(Error::ClusterStartup(reason));
}
crate::join::StartupAction::Found => {
peers.insert(node_id, self_advertise.clone());
genesis_trust.insert(node_id, identity.public_key().to_vec());
do_bootstrap = true;
tracing::info!(node_id, "cluster: founding a new cluster (init)");
}
crate::join::StartupAction::Join => {
let roots = if cluster_cfg.root_pubkeys.is_empty() {
config
.serve
.as_ref()
.and_then(|s| s.auth_root_public_key.clone())
.into_iter()
.collect()
} else {
cluster_cfg.root_pubkeys.clone()
};
let token = cluster_cfg
.join_token
.as_deref()
.and_then(|s| crate::join::resolve_join_token(s).transpose())
.transpose()
.map_err(|e| Error::ClusterStartup(e.to_string()))?
.ok_or_else(|| {
Error::ClusterStartup(
"joining requires [cluster].join_token (env:/path:/inline)".into(),
)
})?;
let ticket = crate::join::JoinTicket {
seeds: cluster_cfg.seeds.clone(),
root_pubkeys: roots,
token,
};
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let adopted = crate::join::join_cluster(&ticket, &identity, Some(&self_advertise), now)
.await
.map_err(|e| Error::ClusterStartup(e.to_string()))?;
tracing::info!(
node_id,
members = adopted.len(),
"cluster: joined via seeds"
);
for m in adopted {
if let Ok(spki) = mesh::parse_public_key(&m.mesh_pubkey_hex) {
genesis_trust.insert(m.node_id, spki);
}
if let Some(addr) = m.mesh_addr {
peers.insert(m.node_id, addr);
}
}
}
crate::join::StartupAction::Resume => {
tracing::info!(node_id, "cluster: resuming from durable state");
}
}
if !cluster_cfg.listen.ip().is_loopback() && genesis_trust.is_empty() && !is_resume {
return Err(Error::MeshUnconfigured(cluster_cfg.listen));
}
let mesh_tls = Arc::new(MeshTls::new(
Arc::new(identity),
TrustSet::from_map(genesis_trust),
));
let (write_capability, write_authz) = build_mesh_write_gate(&args, config, &mesh_cfg).await?;
if write_authz.is_some() {
tracing::info!("cluster: mesh client-write gating enabled");
}
let daemon_runtime = Arc::new(boatramp_server::DaemonRuntime::new(
boatramp_server::config_baseline(&options),
));
options.daemon_runtime = Some(daemon_runtime.clone());
let daemon_observer: Arc<dyn boatramp_cluster::raft::ApplyObserver> =
Arc::new(DaemonConfigObserver(daemon_runtime));
let node = Arc::new(
build_node(ClusterParams {
node_id,
peers,
voters: std::collections::BTreeSet::new(),
durable_kv,
storage: storage.clone(),
mesh: mesh_tls.clone(),
cluster_write_capability: write_capability,
extra_observers: vec![daemon_observer],
})
.await?,
);
if is_resume && !cluster_cfg.listen.ip().is_loopback() && mesh_tls.trust().snapshot().is_empty()
{
return Err(Error::MeshUnconfigured(cluster_cfg.listen));
}
let mesh_router = match write_authz {
Some(authz) => node.router.clone().layer(axum::Extension::<
boatramp_cluster::http::WriteAuthz,
>(Some(authz))),
None => node.router.clone(),
};
let mesh_config =
axum_server::tls_rustls::RustlsConfig::from_config(Arc::new(mesh_tls.server()?));
let listen = cluster_cfg.listen;
tracing::info!(
node_id, %listen,
"cluster: serving peer mesh (mutual TLS)"
);
tokio::spawn(async move {
if let Err(err) = axum_server::bind_rustls(listen, mesh_config)
.serve(mesh_router.into_make_service())
.await
{
tracing::error!(%err, "cluster: peer mesh server exited");
}
});
if do_bootstrap {
node.bootstrap().await?;
node.advertise_addr(&self_advertise).await?;
tracing::info!("cluster: bootstrapped membership");
}
if let Some(interval) = mesh_cfg
.key_rotation
.as_deref()
.and_then(parse_rotation_interval)
{
let rotate_node = node.clone();
let stagger = std::time::Duration::from_secs(node_id % 60);
tokio::spawn(async move {
loop {
tokio::time::sleep(interval + stagger).await;
match rotate_node.rotate_key(MESH_ROTATION_PROPAGATION).await {
Ok(pubkey) => tracing::info!(
pubkey = %pubkey.iter().map(|b| format!("{b:02x}")).collect::<String>(),
"cluster: rotated mesh key on schedule"
),
Err(err) => {
tracing::error!(%err, "cluster: scheduled mesh key rotation failed");
}
}
}
});
}
let kv: Arc<dyn KvStore> = node.kv.clone();
match migrate::status(kv.as_ref()).await? {
migrate::Status::Ready => {}
migrate::Status::Dual => tracing::warn!(
"control-plane store is in the dual soak window; \
run `boatramp migrate --finalize` to reclaim the old-layout keys"
),
migrate::Status::NeedsMigration => {
if args.auto_migrate {
tracing::warn!(
"control-plane store is below the current schema version; running a \
one-shot migration through the cluster (--auto-migrate)"
);
let report =
migrate::migrate(kv.as_ref(), migrate::MigrateOptions::one_shot()).await?;
tracing::info!(
rekeyed = report.total_rekeyed(),
owner_entries = report.owner_entries,
"control-plane store migrated to the project-scoped layout"
);
} else {
return Err(Error::UnmigratedStore);
}
}
}
if args.cluster_rate_limit || config.serve.as_ref().is_some_and(|s| s.cluster_rate_limit) {
options.cluster_rate_limit_kv = Some(kv.clone());
}
spawn_sighup_reload(kv.clone(), options.daemon_runtime.clone());
let cluster_serve_cfg = config.serve.clone().unwrap_or_default();
let auth = boatramp_node::auth::configure_auth(
cluster_serve_cfg.signer.as_ref(),
args.auth_root_private_key
.clone()
.or(cluster_serve_cfg.auth_root_private_key.clone()),
args.auth_root_public_key
.clone()
.or(cluster_serve_cfg.auth_root_public_key),
&mut options,
kv.clone(),
)
.await?;
configure_oidc(&args, &mut options).await?;
boatramp_node::auth::enforce_auth_bind(addr, &auth, &options.posture)?;
options.mesh_control = Some(Arc::new(ClusterMeshControl {
node: node.clone(),
issuer: options.issuer.clone(),
}));
let leader_raft = node.raft.clone();
let leader_node_id = node.node_id;
let is_leader: boatramp_server::CronLeaderGate =
Arc::new(move || boatramp_cluster::raft::is_leader(&leader_raft, leader_node_id));
let boatramp_node::RunningNode {
deploy,
handlers,
auth,
options,
reconcile: _reconcile,
} = boatramp_node::assemble(boatramp_node::NodeInput {
config,
data_dir: data_dir.as_path(),
storage,
kv,
auth,
options,
watch_provider: built_blobs.watch_provider.clone(),
provision_tier: built_blobs.provision_tier,
messaging: Some(node.messaging.clone()),
is_leader,
node_id: node.node_id,
})
.await?;
tracing::info!(tls = ?args.tls, "cluster: serving public traffic");
#[cfg(feature = "tls")]
if !matches!(args.tls, TlsMode::Off) {
let redirect = args
.http_redirect_addr
.or_else(|| config.serve.as_ref().and_then(|s| s.http_redirect_addr));
if let Some(redirect_addr) = redirect {
spawn_http_redirect(redirect_addr, deploy.clone(), options.posture);
}
}
let serve_result = match args.tls {
TlsMode::Off => boatramp_server::serve_with(addr, deploy, auth, handlers, options)
.await
.map_err(Error::Serve),
TlsMode::Custom => serve_custom(&args, addr, deploy, auth, handlers, options).await,
TlsMode::Acme => serve_acme(&args, addr, deploy, auth, handlers, options).await,
TlsMode::Rpk => serve_rpk(&args, addr, deploy, auth, handlers, options, &data_dir).await,
#[cfg(feature = "acme-dns")]
TlsMode::AcmeDns => {
let cert_store: Arc<dyn boatramp_core::cert::CertStore> =
match build_cert_envelope(config.secrets.as_ref(), &data_dir)? {
Some(envelope) => Arc::new(boatramp_core::cert::KvCertStore::with_envelope(
node.kv.clone(),
envelope,
)),
None => Arc::new(boatramp_core::cert::KvCertStore::new(node.kv.clone())),
};
let cert_raft = node.raft.clone();
let cert_node_id = node.node_id;
serve_cluster_acme_dns(
&args,
addr,
deploy,
auth,
handlers,
options,
cert_store,
move || boatramp_cluster::raft::is_leader(&cert_raft, cert_node_id),
)
.await
}
#[cfg(not(feature = "acme-dns"))]
TlsMode::AcmeDns => serve_acme_dns(&args, addr, deploy, auth, handlers, options).await,
};
if let Err(e) = durable_kv_handle.flush().await {
tracing::warn!(error = %e, "cluster: durable Raft store flush on shutdown failed");
} else {
tracing::info!("cluster: durable Raft store flushed on shutdown");
}
serve_result
}
#[cfg(all(feature = "cluster", feature = "acme-dns"))]
#[allow(clippy::too_many_arguments)]
async fn serve_cluster_acme_dns(
args: &ServeArgs,
addr: SocketAddr,
deploy: DeployStore,
auth: boatramp_server::Auth,
handlers: boatramp_server::HandlerRuntime,
options: boatramp_server::ServerOptions,
cert_store: Arc<dyn boatramp_core::cert::CertStore>,
is_leader: impl Fn() -> bool + Send + Sync + Clone + 'static,
) -> Result<()> {
use boatramp_acme::acme::CertRequest;
use std::time::Duration;
if args.acme_domain.is_empty() {
return Err(Error::NoAcmeDomainDns);
}
install_crypto_provider();
let kind = parse_dns_provider(&args.acme_dns_provider)?;
let provider: Arc<dyn boatramp_acme::dns::DnsProvider> =
crate::acme_dns::build_provider(kind).await?.into();
let base = CertRequest {
directory_url: args.acme_directory.clone(),
contact_email: args.acme_contact.clone(),
domains: Vec::new(),
dns_ttl: 60,
propagation_delay: Duration::from_secs(15),
timeout: Duration::from_secs(120),
};
let domains = crate::acme_dns::server_domains(&args.acme_domain, args.acme_wildcard_preview);
let cache = args.acme_cache.clone();
let entries =
cluster_refresh_certs(&cert_store, &domains, is_leader(), &provider, &base, &cache).await?;
if entries.is_empty() {
return Err(Error::NoCertsYet);
}
#[cfg(feature = "http3")]
let (config, h3_endpoint) = if args.http3 {
let (tcp, h3) = crate::acme_dns::build_server_configs(entries)?;
let endpoint =
boatramp_server::http3_endpoint(addr, boatramp_server::quinn_server_config(h3)?)?;
(tcp, Some(endpoint))
} else {
(crate::acme_dns::build_server_config(entries)?, None)
};
#[cfg(not(feature = "http3"))]
let config = crate::acme_dns::build_server_config(entries)?;
let tls = axum_server::tls_rustls::RustlsConfig::from_config(Arc::new(config));
{
let (tls, cert_store, provider, base, cache, domains, is_leader) = (
tls.clone(),
cert_store.clone(),
provider.clone(),
base.clone(),
cache.clone(),
domains.clone(),
is_leader.clone(),
);
#[cfg(feature = "http3")]
let h3_renew = h3_endpoint.clone();
tokio::spawn(async move {
loop {
tokio::time::sleep(Duration::from_secs(6 * 3600)).await;
match cluster_refresh_certs(
&cert_store,
&domains,
is_leader(),
&provider,
&base,
&cache,
)
.await
{
Ok(entries) if !entries.is_empty() => {
#[cfg(feature = "http3")]
if let Some(endpoint) = &h3_renew {
match crate::acme_dns::build_server_configs(entries) {
Ok((tcp, h3)) => {
tls.reload_from_config(Arc::new(tcp));
match boatramp_server::quinn_server_config(h3) {
Ok(qc) => endpoint.set_server_config(Some(qc)),
Err(err) => {
tracing::error!(%err, "cluster acme-dns: rebuilding h3 config failed");
}
}
}
Err(err) => {
tracing::error!(%err, "cluster acme-dns: rebuilding TLS config failed");
}
}
} else {
match crate::acme_dns::build_server_config(entries) {
Ok(config) => tls.reload_from_config(Arc::new(config)),
Err(err) => {
tracing::error!(%err, "cluster acme-dns: rebuilding TLS config failed");
}
}
}
#[cfg(not(feature = "http3"))]
match crate::acme_dns::build_server_config(entries) {
Ok(config) => tls.reload_from_config(Arc::new(config)),
Err(err) => {
tracing::error!(%err, "cluster acme-dns: rebuilding TLS config failed");
}
}
}
Ok(_) => {} Err(err) => tracing::error!(%err, "cluster acme-dns: renewal failed"),
}
}
});
}
tracing::info!(%addr, domains = ?domains, "cluster: serving HTTPS (cluster-managed ACME DNS-01)");
#[cfg(feature = "handlers")]
let _scheduler = handlers.spawn_scheduler(deploy.clone());
let handle = spawn_tls_shutdown();
let app = boatramp_server::router_with(deploy, auth, handlers, options);
#[cfg(feature = "http3")]
let app = if let Some(endpoint) = h3_endpoint {
let app_h3 = app.clone();
tokio::spawn(async move {
if let Err(err) = boatramp_server::serve_http3_endpoint(endpoint, app_h3).await {
tracing::error!(%err, "cluster acme-dns: HTTP/3 listener failed");
}
});
boatramp_server::advertise_http3(app, addr.port())
} else {
app
};
axum_server::bind_rustls(addr, tls)
.handle(handle)
.serve(app.into_make_service_with_connect_info::<SocketAddr>())
.await?;
Ok(())
}
#[cfg(all(feature = "cluster", feature = "acme-dns"))]
async fn cluster_refresh_certs(
cert_store: &Arc<dyn boatramp_core::cert::CertStore>,
domains: &[String],
is_leader: bool,
provider: &Arc<dyn boatramp_acme::dns::DnsProvider>,
base: &boatramp_acme::acme::CertRequest,
cache: &Path,
) -> Result<Vec<(String, boatramp_acme::acme::IssuedCert)>> {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let entries = crate::cluster_tls::refresh_entries(
cert_store.as_ref(),
domains,
is_leader,
now,
|domain| {
let (provider, base, cache) = (provider.clone(), base.clone(), cache.to_path_buf());
async move {
let issued =
crate::acme_dns::obtain_or_load(&domain, &base, provider.as_ref(), &cache)
.await?;
Ok::<_, crate::acme_dns::Error>(crate::cluster_tls::issued_to_stored(&issued, now))
}
},
)
.await?;
Ok(entries)
}
#[cfg(feature = "acme-dns")]
fn parse_dns_provider(value: &str) -> Result<crate::acme_dns::DnsProviderKind> {
use clap::ValueEnum;
crate::acme_dns::DnsProviderKind::from_str(value, true)
.map_err(|_| Error::UnknownDnsProvider(value.to_string()))
}
#[cfg(feature = "acme-dns")]
async fn serve_acme_dns(
args: &ServeArgs,
addr: SocketAddr,
deploy: DeployStore,
auth: boatramp_server::Auth,
handlers: boatramp_server::HandlerRuntime,
options: boatramp_server::ServerOptions,
) -> Result<()> {
use std::time::Duration;
use boatramp_acme::acme::CertRequest;
if args.acme_domain.is_empty() {
return Err(Error::NoAcmeDomainDns);
}
install_crypto_provider();
let kind = parse_dns_provider(&args.acme_dns_provider)?;
let provider = crate::acme_dns::build_provider(kind).await?;
let base = CertRequest {
directory_url: args.acme_directory.clone(),
contact_email: args.acme_contact.clone(),
domains: Vec::new(),
dns_ttl: 60,
propagation_delay: Duration::from_secs(15),
timeout: Duration::from_secs(120),
};
let domains = crate::acme_dns::server_domains(&args.acme_domain, args.acme_wildcard_preview);
let entries = obtain_all(&domains, &base, provider.as_ref(), &args.acme_cache).await?;
#[cfg(feature = "http3")]
let (config, h3_endpoint) = if args.http3 {
let (tcp, h3) = crate::acme_dns::build_server_configs(entries)?;
let endpoint =
boatramp_server::http3_endpoint(addr, boatramp_server::quinn_server_config(h3)?)?;
(tcp, Some(endpoint))
} else {
(crate::acme_dns::build_server_config(entries)?, None)
};
#[cfg(not(feature = "http3"))]
let config = crate::acme_dns::build_server_config(entries)?;
let tls = axum_server::tls_rustls::RustlsConfig::from_config(Arc::new(config));
{
let (tls, base, cache) = (tls.clone(), base.clone(), args.acme_cache.clone());
let domains = domains.clone();
let provider = crate::acme_dns::build_provider(kind).await?;
#[cfg(feature = "http3")]
let h3_renew = h3_endpoint.clone();
tokio::spawn(async move {
loop {
tokio::time::sleep(Duration::from_secs(6 * 3600)).await;
match obtain_all(&domains, &base, provider.as_ref(), &cache).await {
Ok(entries) => {
#[cfg(feature = "http3")]
if let Some(endpoint) = &h3_renew {
match crate::acme_dns::build_server_configs(entries) {
Ok((tcp, h3)) => {
tls.reload_from_config(Arc::new(tcp));
match boatramp_server::quinn_server_config(h3) {
Ok(qc) => endpoint.set_server_config(Some(qc)),
Err(err) => {
tracing::error!(%err, "acme-dns: rebuilding h3 config failed");
}
}
}
Err(err) => {
tracing::error!(%err, "acme-dns: rebuilding TLS config failed");
}
}
} else {
match crate::acme_dns::build_server_config(entries) {
Ok(config) => tls.reload_from_config(Arc::new(config)),
Err(err) => {
tracing::error!(%err, "acme-dns: rebuilding TLS config failed");
}
}
}
#[cfg(not(feature = "http3"))]
match crate::acme_dns::build_server_config(entries) {
Ok(config) => tls.reload_from_config(Arc::new(config)),
Err(err) => {
tracing::error!(%err, "acme-dns: rebuilding TLS config failed");
}
}
}
Err(err) => tracing::error!(%err, "acme-dns: renewal failed"),
}
}
});
}
tracing::info!(%addr, domains = ?domains, "serving HTTPS (ACME DNS-01)");
#[cfg(feature = "handlers")]
let _scheduler = handlers.spawn_scheduler(deploy.clone());
let handle = spawn_tls_shutdown();
let app = boatramp_server::router_with(deploy, auth, handlers, options);
#[cfg(feature = "http3")]
let app = if let Some(endpoint) = h3_endpoint {
let app_h3 = app.clone();
tokio::spawn(async move {
if let Err(err) = boatramp_server::serve_http3_endpoint(endpoint, app_h3).await {
tracing::error!(%err, "acme-dns: HTTP/3 listener failed");
}
});
boatramp_server::advertise_http3(app, addr.port())
} else {
app
};
axum_server::bind_rustls(addr, tls)
.handle(handle)
.serve(app.into_make_service_with_connect_info::<SocketAddr>())
.await?;
Ok(())
}
#[cfg(feature = "acme-dns")]
async fn obtain_all(
domains: &[String],
base: &boatramp_acme::acme::CertRequest,
provider: &dyn boatramp_acme::dns::DnsProvider,
cache: &Path,
) -> Result<Vec<(String, boatramp_acme::acme::IssuedCert)>> {
let mut entries = Vec::with_capacity(domains.len());
for domain in domains {
let cert = crate::acme_dns::obtain_or_load(domain, base, provider, cache).await?;
entries.push((domain.clone(), cert));
}
Ok(entries)
}
#[cfg(not(feature = "acme-dns"))]
async fn serve_acme_dns(
_args: &ServeArgs,
_addr: SocketAddr,
_deploy: DeployStore,
_auth: boatramp_server::Auth,
_handlers: boatramp_server::HandlerRuntime,
_options: boatramp_server::ServerOptions,
) -> Result<()> {
Err(Error::NoAcmeDnsSupport)
}
#[cfg(feature = "oidc")]
async fn configure_oidc(
args: &ServeArgs,
options: &mut boatramp_server::ServerOptions,
) -> Result<()> {
let Some(issuer) = args.oidc_issuer.clone() else {
return Ok(());
};
if options.posture.oidc_require_audience && args.oidc_audience.is_none() {
return Err(Error::OidcAudienceRequired);
}
let mut config = boatramp_server::OidcConfig::new(issuer);
config.audience = args.oidc_audience.clone();
if let Some(claim) = args.oidc_scope_claim.clone() {
config.scope_claim = claim;
}
let http = reqwest::Client::new();
let verifier = Arc::new(
boatramp_server::OidcVerifier::from_discovery(&http, &config)
.await
.map_err(|err| Error::OidcSetup(err.to_string()))?,
);
tracing::info!(issuer = %config.issuer, "OIDC → token exchange enabled");
{
let verifier = verifier.clone();
tokio::spawn(async move {
loop {
tokio::time::sleep(std::time::Duration::from_secs(3600)).await;
if let Err(err) = verifier.refresh().await {
tracing::warn!(%err, "OIDC JWKS refresh failed (keeping current keys)");
}
}
});
}
options.oidc_verifier = Some(verifier);
Ok(())
}
#[cfg(not(feature = "oidc"))]
async fn configure_oidc(
_args: &ServeArgs,
_options: &mut boatramp_server::ServerOptions,
) -> Result<()> {
Ok(())
}
#[cfg(feature = "tls")]
async fn serve_custom(
args: &ServeArgs,
addr: SocketAddr,
deploy: DeployStore,
auth: boatramp_server::Auth,
handlers: boatramp_server::HandlerRuntime,
options: boatramp_server::ServerOptions,
) -> Result<()> {
install_crypto_provider();
let cert = args.tls_cert.clone().ok_or(Error::TlsCertRequired)?;
let key = args.tls_key.clone().ok_or(Error::TlsKeyRequired)?;
let config = axum_server::tls_rustls::RustlsConfig::from_pem_file(&cert, &key).await?;
tracing::info!(%addr, "serving HTTPS (custom certificate)");
#[cfg(feature = "handlers")]
let _scheduler = handlers.spawn_scheduler(deploy.clone());
let handle = spawn_tls_shutdown();
let app = boatramp_server::router_with(deploy, auth, handlers, options);
#[cfg(feature = "http3")]
let app = if args.http3 {
let (certs, key) = load_cert_chain_and_key(&cert, &key)?;
let app_h3 = app.clone();
tokio::spawn(async move {
if let Err(err) = boatramp_server::serve_http3(addr, certs, key, app_h3).await {
tracing::error!(%err, "HTTP/3 listener failed");
}
});
boatramp_server::advertise_http3(app, addr.port())
} else {
app
};
axum_server::bind_rustls(addr, config)
.handle(handle)
.serve(app.into_make_service_with_connect_info::<SocketAddr>())
.await?;
Ok(())
}
#[cfg(feature = "tls")]
async fn serve_rpk(
_args: &ServeArgs,
addr: SocketAddr,
deploy: DeployStore,
auth: boatramp_server::Auth,
handlers: boatramp_server::HandlerRuntime,
mut options: boatramp_server::ServerOptions,
data_dir: &Path,
) -> Result<()> {
install_crypto_provider();
let key_file = data_dir.join("controlplane-tls.key");
let identity = boatramp_rpktls::RpkIdentity::load_or_generate(&key_file)?;
let fingerprint = identity.public_key_hex();
if let Some(signer) = options.issuer.clone() {
const ATTESTATION_TTL_SECS: u64 = 365 * 24 * 60 * 60;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
match boatramp_core::cose::mint_attestation(
&fingerprint,
ATTESTATION_TTL_SECS,
now,
signer.as_ref(),
)
.await
{
Ok(att) => options.bootstrap_attestation = Some(att),
Err(err) => {
tracing::warn!(%err, "could not mint the bootstrap-TLS attestation; --root-pubkey pinning unavailable");
}
}
}
let rpk =
boatramp_rpktls::RpkTls::new(Arc::new(identity), boatramp_rpktls::TrustSet::default());
let config = axum_server::tls_rustls::RustlsConfig::from_config(Arc::new(rpk.server_auth()?));
tracing::info!(%addr, pubkey = %fingerprint, "serving HTTPS (RPK bootstrap TLS)");
println!(
"control-plane RPK TLS identity — pin the client with:\n --server-pubkey {fingerprint}"
);
#[cfg(feature = "handlers")]
let _scheduler = handlers.spawn_scheduler(deploy.clone());
let handle = spawn_tls_shutdown();
let app = boatramp_server::router_with(deploy, auth, handlers, options);
axum_server::bind_rustls(addr, config)
.handle(handle)
.serve(app.into_make_service_with_connect_info::<SocketAddr>())
.await?;
Ok(())
}
#[cfg(feature = "http3")]
fn load_cert_chain_and_key(
cert: &Path,
key: &Path,
) -> Result<(
Vec<rustls::pki_types::CertificateDer<'static>>,
rustls::pki_types::PrivateKeyDer<'static>,
)> {
let cert_pem = std::fs::read(cert)?;
let certs =
rustls_pemfile::certs(&mut &cert_pem[..]).collect::<std::result::Result<Vec<_>, _>>()?;
if certs.is_empty() {
return Err(Error::NoCert(cert.display().to_string()));
}
let key_pem = std::fs::read(key)?;
let key = rustls_pemfile::private_key(&mut &key_pem[..])?
.ok_or_else(|| Error::NoPrivateKey(key.display().to_string()))?;
Ok((certs, key))
}
#[cfg(feature = "tls")]
async fn serve_acme(
args: &ServeArgs,
addr: SocketAddr,
deploy: DeployStore,
auth: boatramp_server::Auth,
handlers: boatramp_server::HandlerRuntime,
options: boatramp_server::ServerOptions,
) -> Result<()> {
use futures::StreamExt;
use rustls_acme::{caches::DirCache, AcmeConfig};
if args.acme_domain.is_empty() {
return Err(Error::NoAcmeDomain);
}
install_crypto_provider();
let mut config = AcmeConfig::new(args.acme_domain.clone())
.cache(DirCache::new(args.acme_cache.clone()))
.directory(args.acme_directory.clone());
if let Some(contact) = &args.acme_contact {
config = config.contact_push(format!("mailto:{contact}"));
}
if let Some(ca) = &args.acme_ca_cert {
config = config.client_tls_config(acme_client_config(ca)?);
}
let mut state = config.state();
let acceptor = state.axum_acceptor(state.default_rustls_config());
tokio::spawn(async move {
loop {
match state.next().await {
Some(Ok(event)) => tracing::info!("acme: {event:?}"),
Some(Err(err)) => tracing::error!("acme error: {err}"),
None => break,
}
}
});
tracing::info!(%addr, domains = ?args.acme_domain, "serving HTTPS (ACME)");
#[cfg(feature = "handlers")]
let _scheduler = handlers.spawn_scheduler(deploy.clone());
let handle = spawn_tls_shutdown();
axum_server::bind(addr)
.handle(handle)
.acceptor(acceptor)
.serve(
boatramp_server::router_with(deploy, auth, handlers, options)
.into_make_service_with_connect_info::<SocketAddr>(),
)
.await?;
Ok(())
}
#[cfg(feature = "tls")]
fn install_crypto_provider() {
let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
}
#[cfg(feature = "tls")]
fn spawn_tls_shutdown() -> axum_server::Handle<SocketAddr> {
let handle = axum_server::Handle::new();
let trigger = handle.clone();
tokio::spawn(async move {
boatramp_server::shutdown_signal().await;
trigger.graceful_shutdown(Some(std::time::Duration::from_secs(10)));
});
handle
}
#[cfg(feature = "tls")]
fn acme_client_config(ca_path: &std::path::Path) -> Result<Arc<rustls::ClientConfig>> {
let pem = std::fs::read(ca_path)?;
let mut roots = rustls::RootCertStore::empty();
for cert in rustls_pemfile::certs(&mut &pem[..]) {
roots.add(cert?)?;
}
let config = rustls::ClientConfig::builder()
.with_root_certificates(roots)
.with_no_client_auth();
Ok(Arc::new(config))
}
#[cfg(not(feature = "tls"))]
async fn serve_custom(
_args: &ServeArgs,
_addr: SocketAddr,
_deploy: DeployStore,
_auth: boatramp_server::Auth,
_handlers: boatramp_server::HandlerRuntime,
_options: boatramp_server::ServerOptions,
) -> Result<()> {
Err(Error::NoTlsSupport)
}
#[cfg(not(feature = "tls"))]
async fn serve_acme(
_args: &ServeArgs,
_addr: SocketAddr,
_deploy: DeployStore,
_auth: boatramp_server::Auth,
_handlers: boatramp_server::HandlerRuntime,
_options: boatramp_server::ServerOptions,
) -> Result<()> {
Err(Error::NoTlsSupport)
}
#[cfg(not(feature = "tls"))]
async fn serve_rpk(
_args: &ServeArgs,
_addr: SocketAddr,
_deploy: DeployStore,
_auth: boatramp_server::Auth,
_handlers: boatramp_server::HandlerRuntime,
_options: boatramp_server::ServerOptions,
_data_dir: &Path,
) -> Result<()> {
Err(Error::NoTlsSupport)
}
#[cfg(test)]
mod tests {
#[cfg(feature = "cluster")]
use super::*;
#[cfg(feature = "cluster")]
#[tokio::test]
async fn mesh_write_authz_accepts_only_a_cluster_write_capability() {
use boatramp_cluster::http::ClientWriteAuthz;
use boatramp_core::authz::GrantedRole;
use boatramp_core::cose::{self, Claims, LocalSigner, Signer, TokenAlg};
async fn cap(signer: &dyn Signer, role: &str) -> String {
let claims = Claims {
roles: vec![GrantedRole::global(role)],
kind: cose::KIND_CLUSTER_WRITE.to_string(),
ttl_secs: None,
now_unix: 0,
};
cose::mint(&claims, signer).await.unwrap()
}
let signer = LocalSigner::generate(TokenAlg::Es256);
let authz = MeshWriteAuthz {
public: signer.public_key(),
};
assert!(
authz.authorize(Some(&cap(&signer, "cluster-write").await)),
"a real cluster-write capability"
);
assert!(!authz.authorize(None), "no capability");
assert!(
!authz.authorize(Some(&cap(&signer, "admin").await)),
"wrong role"
);
assert!(!authz.authorize(Some("not-a-token")), "garbage");
let other = LocalSigner::generate(TokenAlg::Es256);
assert!(
!authz.authorize(Some(&cap(&other, "cluster-write").await)),
"foreign root key"
);
}
#[cfg(feature = "cluster")]
#[test]
fn rotation_interval_parses_units_and_rejects_junk() {
use std::time::Duration;
assert_eq!(
parse_rotation_interval("30d"),
Some(Duration::from_secs(30 * 86_400))
);
assert_eq!(
parse_rotation_interval("12h"),
Some(Duration::from_secs(12 * 3600))
);
assert_eq!(
parse_rotation_interval("90m"),
Some(Duration::from_secs(90 * 60))
);
assert_eq!(
parse_rotation_interval(" 45s "),
Some(Duration::from_secs(45))
);
assert_eq!(parse_rotation_interval("30"), None);
assert_eq!(parse_rotation_interval("5w"), None);
assert_eq!(parse_rotation_interval("0d"), None);
assert_eq!(parse_rotation_interval(""), None);
}
}