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::Result;
pub const COMPUTE_RECONCILE_TICK: std::time::Duration = 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 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 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<()>>,
}
pub async fn assemble(input: NodeInput<'_>) -> Result<RunningNode> {
let NodeInput {
config,
data_dir,
storage,
kv,
auth,
options,
watch_provider,
provision_tier,
messaging,
is_leader,
node_id,
} = input;
let max_handler_blob_bytes = options.posture.max_handler_blob_bytes;
let max_component_bytes = options.posture.max_component_bytes;
let allow_shared_kernel = options.posture.allow_shared_kernel_compute;
let domain_verify_allow_private = options.posture.domain_verify_allow_private;
let handlers = crate::handlers::build_handler_runtime(
kv.clone(),
storage.clone(),
data_dir,
config.handlers.as_ref(),
messaging,
max_handler_blob_bytes,
max_component_bytes,
)?;
#[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);
let compute_storage = storage.clone();
let deploy = DeployStore::new(storage, kv);
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(),
)
.await;
#[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;
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,
);
let dv_reconcile = boatramp_server::spawn_domain_verify_reconcile(
deploy.clone(),
domain_verify_allow_private,
is_leader,
DOMAIN_VERIFY_RECONCILE_TICK,
);
Ok(RunningNode {
deploy,
handlers,
auth,
options,
reconcile: vec![compute_reconcile, dv_reconcile],
})
}
#[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,
watch_provider: None,
provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
messaging: None,
is_leader: Arc::new(|| true),
node_id: 0,
})
.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);
}
}