Skip to main content

boatramp_node/
node.rs

1//! The node-graph assembly: given a built store (blobs + KV), a configured
2//! [`Auth`](boatramp_server::Auth), and resolved
3//! [`ServerOptions`](boatramp_server::ServerOptions), wire the deploy store,
4//! handler runtime, compute reconcile loop, and domain-verify reconcile loop
5//! into a [`RunningNode`] ready to hand to a transport (`serve_with` & friends).
6//!
7//! This is the headline extraction of `PLAN-node-library`: the binary's
8//! `serve::run` used to inline this wiring, so no embedder or in-process test
9//! could exercise the same graph the `boatramp serve` binary runs. `run` now
10//! resolves the *environment* (args -> backends -> store, signal handlers,
11//! migration, auth) and calls [`assemble`]; the cluster path keeps its own inline
12//! copy until a later step converges it here.
13
14use std::path::Path;
15use std::sync::Arc;
16
17use boatramp_core::deploy::DeployStore;
18use boatramp_core::kv::KvStore;
19use boatramp_core::Storage;
20
21use crate::config::ServerConfig;
22use crate::error::{Error, Result};
23
24/// How often the compute reconcile loop converges desired vs actual workloads.
25/// Defaults to 30s; override with `BOATRAMP_COMPUTE_RECONCILE_TICK_MS` (milliseconds)
26/// so compute-backed tests can converge in a fraction of a second instead of
27/// waiting a full tick for the launch/scale reconcile.
28pub fn compute_reconcile_tick() -> std::time::Duration {
29    std::env::var("BOATRAMP_COMPUTE_RECONCILE_TICK_MS")
30        .ok()
31        .and_then(|s| s.parse::<u64>().ok())
32        .filter(|&ms| ms > 0)
33        .map(std::time::Duration::from_millis)
34        .unwrap_or(std::time::Duration::from_secs(30))
35}
36/// How often the domain-verify reconcile loop re-checks pending challenges.
37pub const DOMAIN_VERIFY_RECONCILE_TICK: std::time::Duration = std::time::Duration::from_secs(60);
38/// How long a compute workload may be idle before scale-to-zero sleeps it.
39pub const COMPUTE_IDLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(300);
40
41/// The built store + resolved config handed to [`assemble`]. Owns the blob/KV
42/// backends and the auth/options the caller already resolved; borrows the parsed
43/// config and data directory.
44pub struct NodeInput<'a> {
45    /// The full parsed server config (the handler + compute sections are read here).
46    pub config: &'a ServerConfig,
47    /// The node data directory (per-site SQL, handler state).
48    pub data_dir: &'a Path,
49    /// The object store built by [`crate::blobs::build_blobs`].
50    pub storage: Arc<dyn Storage>,
51    /// The metadata KV, already cache-fronted, built by [`crate::backends::build_kv`].
52    pub kv: Arc<dyn KvStore>,
53    /// The control-plane auth built by [`crate::auth::configure_auth`].
54    pub auth: boatramp_server::Auth,
55    /// Server options, already carrying the resolved posture, daemon runtime, and
56    /// (post-`configure_auth`/`configure_oidc`) issuer / OIDC verifier.
57    pub options: boatramp_server::ServerOptions,
58    /// The public HTTP serve bind address, if known — used (under
59    /// `allow_guest_self_egress`) to let a handler guest's `wasi:http` reach this
60    /// instance's own front door over loopback. `None` (an in-process embedder with no
61    /// listener) disables self-egress.
62    pub serve_addr: Option<std::net::SocketAddr>,
63    /// The cloud blob-change watch provider (FA-5b2), if the backend is a cloud one.
64    pub watch_provider: Option<Arc<dyn boatramp_core::blob_provision::WatchProvider>>,
65    /// The provisioning tier for the watch provider.
66    pub provision_tier: boatramp_core::blob_notify::ProvisionTier,
67    /// The `wasi:messaging` substrate override for the handler runtime. `None` uses
68    /// the single-node default (`LogMessaging` over the same backends); the cluster
69    /// path passes its Raft-backed coordinator.
70    pub messaging: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
71    /// The single leader gate for cron firing + the compute / domain-verify reconcile
72    /// loops. Single-node passes an always-true gate (there is one node); the cluster
73    /// passes its Raft `is_leader` check so a single node drives each sweep.
74    pub is_leader: boatramp_server::CronLeaderGate,
75    /// This node's compute scheduler id (`0` single-node; the cluster node id in a
76    /// fleet, so replicas are tagged to the right node).
77    pub node_id: u64,
78    /// The binary the re-exec'd compute workers run as — the container backend's
79    /// `__sandbox` jailer and the microVM backends' `__vmm-run`/`__vz-run` VM hosts.
80    /// `None` uses this process's own executable (`current_exe`), which is what
81    /// `boatramp serve` wants (the child *is* boatramp). An **embedding harness**
82    /// whose own binary doesn't implement those subcommands should point this at a
83    /// built `boatramp` binary, so it can drive the real container/microVM backends
84    /// in-process (only the per-workload worker re-execs; the serving plane stays
85    /// embedded). The docker backend needs neither — it talks to a daemon.
86    pub worker_exe: Option<std::path::PathBuf>,
87}
88
89/// A fully wired node: the deploy store, handler runtime, auth, and options a
90/// transport consumes, plus the detached reconcile loops kept alive for the
91/// node's serving life. Destructure it and hold `reconcile` across the serve
92/// await so the loops outlive assembly.
93pub struct RunningNode {
94    /// The deploy store (blob + KV) the router serves from.
95    pub deploy: DeployStore,
96    /// The handler runtime for wasm handlers (a disabled build ⇒ a no-op runtime).
97    pub handlers: boatramp_server::HandlerRuntime,
98    /// The control-plane auth.
99    pub auth: boatramp_server::Auth,
100    /// The resolved server options.
101    pub options: boatramp_server::ServerOptions,
102    /// The detached reconcile loops (compute + domain-verify). Tokio `JoinHandle`s
103    /// do not abort on drop, so the loops run for the process life regardless; the
104    /// handles are retained so an embedder can join/abort them on shutdown.
105    pub reconcile: Vec<tokio::task::JoinHandle<()>>,
106}
107
108/// The instance's own serve socket(s) a guest self-call may reach, given the bind `addr` and
109/// whether the posture (`allow_guest_self_egress`) permits it. A wildcard bind
110/// (`0.0.0.0`/`::`) is reachable over loopback, so it normalizes to `127.0.0.1` **and** `::1`
111/// on the serve port; a specific bind is reachable at itself. Empty when disabled or no
112/// listener.
113fn self_egress_addrs(
114    addr: Option<std::net::SocketAddr>,
115    enabled: bool,
116) -> Vec<std::net::SocketAddr> {
117    use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
118    let Some(addr) = addr.filter(|_| enabled) else {
119        return Vec::new();
120    };
121    if addr.ip().is_unspecified() {
122        let port = addr.port();
123        vec![
124            SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), port),
125            SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), port),
126        ]
127    } else {
128        vec![addr]
129    }
130}
131
132/// Wire [`NodeInput`] into a [`RunningNode`]: build the handler runtime, the
133/// deploy store (materializing the reserved `default` project), the compute
134/// backends + reconcile loop, and the domain-verify reconcile loop.
135///
136/// The caller has already built the store and configured auth/OIDC on `options`;
137/// this is the pure node-graph wiring, identical to what `boatramp serve` runs.
138pub async fn assemble(input: NodeInput<'_>) -> Result<RunningNode> {
139    let NodeInput {
140        config,
141        data_dir,
142        storage,
143        kv,
144        auth,
145        options,
146        serve_addr,
147        watch_provider,
148        provision_tier,
149        messaging,
150        is_leader,
151        node_id,
152        worker_exe,
153    } = input;
154    // Copy out the posture scalars up front so `options` can be moved into the
155    // returned `RunningNode` without a lingering borrow.
156    let max_handler_blob_bytes = options.posture.max_handler_blob_bytes;
157    let max_component_bytes = options.posture.max_component_bytes;
158    let allow_guest_private_egress = options.posture.allow_guest_private_egress;
159    // The instance's own serve socket(s) a guest self-call may reach, when the posture allows
160    // it: a wildcard bind (`0.0.0.0`/`::`) is reachable on loopback, so normalize to
161    // `127.0.0.1`/`::1`; a specific bind is itself.
162    let self_egress_addrs = self_egress_addrs(serve_addr, options.posture.allow_guest_self_egress);
163    let allow_shared_kernel = options.posture.allow_shared_kernel_compute;
164    let domain_verify_allow_private = options.posture.domain_verify_allow_private;
165
166    // The deploy store the router serves from — built up front so the handler
167    // runtime's managed compute-backed `sql` binding can resolve DB endpoints from
168    // the same store the reconcile writes.
169    let compute_storage = storage.clone();
170    let deploy = DeployStore::new(storage, kv.clone());
171    // The `[secrets]` envelope (local KEK / Vault) that seals a managed SQL
172    // credential at rest. `None` ⇒ no wrapping (a managed DB then fails closed).
173    let secrets_envelope = build_secrets_envelope(config.secrets.as_ref(), data_dir)?;
174
175    // The handler runtime reuses the same blob/KV backends (per-site prefixed)
176    // for its wasi:blobstore/keyvalue bindings; the sql binding is selected by
177    // `[handlers.bindings.sql]` (default: per-site libsql files under <data-dir>).
178    let handlers = crate::handlers::build_handler_runtime(
179        kv.clone(),
180        compute_storage.clone(),
181        data_dir,
182        config.handlers.as_ref(),
183        messaging,
184        max_handler_blob_bytes,
185        max_component_bytes,
186        allow_guest_private_egress,
187        self_egress_addrs,
188        &deploy,
189        secrets_envelope.clone(),
190    )
191    .await?;
192    // Leader-gate cron firing (cluster: only the Raft leader fires; single-node: an
193    // always-true gate, equivalent to the unset default). The same gate drives the
194    // reconcile loops below, so all three converge on one leader per fleet. Only the
195    // handler runtime has a scheduler, so this is a no-op without the `handlers` feature.
196    #[cfg(feature = "handlers")]
197    handlers.set_cron_leader_gate(is_leader.clone());
198    // FA-5b2: on a cloud backend, wire the blob-change notification provisioner +
199    // its tier so adding a `blob` trigger provisions (and removing it retracts).
200    #[cfg(feature = "handlers")]
201    if let Some(provider) = watch_provider {
202        handlers.set_watch_provider(provider);
203        handlers.set_provision_tier(provision_tier);
204    }
205    #[cfg(not(feature = "handlers"))]
206    let _ = (watch_provider, provision_tier);
207
208    // Materialize the reserved `default` project so `project ls` / `project show
209    // default` reflect it on a fresh install, not only after a migration. Best
210    // effort: the reader backstop keeps listings correct even if this write can't
211    // land, so a transient failure must never block serving.
212    match deploy.ensure_default_project().await {
213        Ok(true) => tracing::info!("materialized the reserved `default` project record"),
214        Ok(false) => {}
215        Err(e) => tracing::warn!(
216            error = %e,
217            "could not materialize the `default` project record; readers use the synthesized default"
218        ),
219    }
220    // Wire the function-to-function invoke resolver now the deploy store exists,
221    // so a function granted `invoke` can call a sibling in-process (FI).
222    #[cfg(feature = "handlers")]
223    handlers.set_invoker(deploy.clone());
224
225    // Compute reconcile loop. Single-node is always the "leader". Backends are
226    // built from the `[compute]` config + capability detection; a no-op when none
227    // are registered. Detached for the server's life.
228    let (compute_backends, compute_node) = crate::compute::build_compute(
229        config.compute.as_ref(),
230        compute_storage,
231        data_dir,
232        node_id,
233        !allow_shared_kernel,
234        options.daemon_runtime.clone(),
235        worker_exe.as_deref(),
236    )
237    .await;
238    // Activate the compute sql-shim (PLAN-compute-bindings): bind its listener +
239    // build the resolver when a sql provider and `compute.sql_shim_url` are both present.
240    #[cfg(feature = "handlers")]
241    let sql_resolver = boatramp_server::sql_shim::spawn_sql_shim(
242        handlers.sql_backends(),
243        config.compute.as_ref().and_then(|c| c.sql_shim_url.clone()),
244    )
245    .await;
246    #[cfg(not(feature = "handlers"))]
247    let sql_resolver: Option<Arc<dyn boatramp_core::compute::ComputeBindingResolver>> = None;
248
249    // Managed compute-backed SQL (PLAN-managed-compute-sql P2-b): if the handler
250    // `sql` config declares any managed database, inject its `POSTGRES_*`/`MYSQL_*`
251    // server-init env into the DB workload at launch from the sealed credential.
252    // Reaching here with a managed DB implies an envelope (build_handler_runtime
253    // fails closed otherwise), so the credential store always has one to seal with.
254    // Keep a clone of the secrets envelope for the operator-SQL capability below
255    // (the managed_db_resolver match moves the original).
256    #[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
257    let operator_envelope = secrets_envelope.clone();
258    #[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
259    let managed_db_resolver: Option<Arc<dyn boatramp_core::compute::ManagedDbEnvResolver>> = match (
260        config
261            .handlers
262            .as_ref()
263            .and_then(|h| h.bindings.sql.as_ref()),
264        secrets_envelope,
265    ) {
266        (Some(sql), Some(envelope)) if !sql.databases.is_empty() => {
267            let creds = crate::managed_sql::ManagedSqlCredentials::new(kv.clone(), envelope);
268            let privilege = config
269                .compute
270                .as_ref()
271                .map(|c| c.managed_db_privilege)
272                .unwrap_or_default();
273            let env =
274                crate::managed_sql::ManagedDbEnv::from_config(&sql.databases, creds, privilege);
275            (!env.is_empty()).then(|| Arc::new(env) as Arc<_>)
276        }
277        _ => None,
278    };
279    #[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
280    let managed_db_resolver: Option<Arc<dyn boatramp_core::compute::ManagedDbEnvResolver>> = None;
281
282    // Turnkey managed DB: auto-register the compute workload backing each managed
283    // co-located database that has none yet, so declaring the `databases` binding is
284    // enough to boot the DB (no separate `compute set` / apply). Non-clobbering and
285    // idempotent; runs before the reconcile loop so its first tick can launch it.
286    #[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
287    if let Some(sql) = config
288        .handlers
289        .as_ref()
290        .and_then(|h| h.bindings.sql.as_ref())
291        .filter(|sql| !sql.databases.is_empty())
292    {
293        crate::managed_sql::auto_register_managed_db_workloads(&deploy, &sql.databases).await;
294    }
295
296    // Operator SQL capability (managed-DB migrations/queries via the sealed
297    // credential, resolved server-side) — backs `POST /api/sql/{db}/{exec,query}`.
298    #[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
299    let operator_sql: Option<Arc<dyn boatramp_core::sql::OperatorSql>> = config
300        .handlers
301        .as_ref()
302        .and_then(|h| h.bindings.sql.as_ref())
303        .filter(|sql| !sql.databases.is_empty())
304        .map(|sql| {
305            Arc::new(crate::managed_sql::NodeOperatorSql::new(
306                sql.databases.clone(),
307                kv.clone(),
308                operator_envelope,
309                deploy.clone(),
310            )) as Arc<_>
311        });
312    #[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
313    let operator_sql: Option<Arc<dyn boatramp_core::sql::OperatorSql>> = None;
314
315    // Operator compute-exec capability (run a command inside a running workload) —
316    // backs `POST /api/compute/{name}/exec`, gated by the `allow_compute_exec`
317    // posture. Clone the backend registry before the reconcile loop consumes it.
318    let compute_exec: Option<Arc<dyn boatramp_core::compute::ComputeExec>> = Some(Arc::new(
319        crate::compute::NodeComputeExec::new(compute_backends.clone(), deploy.clone()),
320    ) as Arc<_>);
321
322    let compute_reconcile = boatramp_server::spawn_compute_reconcile(
323        deploy.clone(),
324        compute_backends,
325        vec![compute_node],
326        boatramp_core::compute::BackendPolicy::from_shared_kernel_allowed(allow_shared_kernel),
327        is_leader.clone(),
328        compute_reconcile_tick(),
329        COMPUTE_IDLE_TIMEOUT,
330        sql_resolver,
331        managed_db_resolver,
332    );
333
334    // Domain-verify auto-complete: periodically re-check every site's pending
335    // ownership challenges and attach any that now pass — a published token (e.g.
336    // via `domain add --provider`) converges without a manual `domain verify`.
337    let dv_reconcile = boatramp_server::spawn_domain_verify_reconcile(
338        deploy.clone(),
339        domain_verify_allow_private,
340        is_leader,
341        DOMAIN_VERIFY_RECONCILE_TICK,
342    );
343
344    // Wire the operator capabilities onto the options the router is built from.
345    let mut options = options;
346    options.operator_sql = operator_sql;
347    options.compute_exec = compute_exec;
348
349    Ok(RunningNode {
350        deploy,
351        handlers,
352        auth,
353        options,
354        reconcile: vec![compute_reconcile, dv_reconcile],
355    })
356}
357
358/// Build the `[secrets]` envelope (secrets-at-rest wrapping) from `boatramp.cfg`'s
359/// `[secrets]` section: `local` (a machine-local AES-256-GCM KEK) or `vault` (Vault
360/// Transit). `None`/empty ⇒ no wrapping. The Vault token is read from the
361/// environment (`token_env`), never a file. This seals a managed SQL credential at
362/// rest; a managed database fails closed without it.
363fn build_secrets_envelope(
364    secrets: Option<&crate::config::SecretsConfig>,
365    data_dir: &Path,
366) -> Result<Option<Arc<dyn boatramp_core::envelope::KeyEnvelope>>> {
367    use boatramp_server::envelope::{build_envelope, EnvelopeSpec};
368    let Some(cfg) = secrets else {
369        return Ok(None);
370    };
371    let spec = match cfg.envelope.as_str() {
372        "" => EnvelopeSpec::None,
373        "local" => EnvelopeSpec::Local {
374            kek_file: cfg
375                .kek_file
376                .clone()
377                .unwrap_or_else(|| data_dir.join("secrets/kek")),
378        },
379        "vault" => {
380            let v = cfg.vault.as_ref().ok_or_else(|| {
381                Error::Envelope(
382                    "secrets.envelope = \"vault\" needs a [secrets.vault] section".into(),
383                )
384            })?;
385            let token = std::env::var(&v.token_env).map_err(|_| {
386                Error::Envelope(format!("Vault token env `{}` is not set", v.token_env))
387            })?;
388            EnvelopeSpec::Vault {
389                addr: v.addr.clone(),
390                key: v.key.clone(),
391                token,
392            }
393        }
394        other => {
395            return Err(Error::Envelope(format!(
396                "unknown secrets.envelope {other:?} (want \"local\" or \"vault\")"
397            )))
398        }
399    };
400    build_envelope(spec).map_err(|e| Error::Envelope(e.to_string()))
401}
402
403#[cfg(all(test, feature = "fs"))]
404mod tests {
405    use super::*;
406    use boatramp_core::kv::MemoryKv;
407    use boatramp_core::security::SecurityProfile;
408
409    /// The headline in-process fidelity check (PLAN-node-library N2b.3): `assemble`
410    /// over a temp `FsStorage` + `MemoryKv` produces a `RunningNode` whose deploy
411    /// store is live (the reserved `default` project was materialized during
412    /// assembly) and whose router — the exact one `boatramp serve` builds — answers
413    /// `/healthz`. No listener is bound: the request is driven through the router
414    /// via `tower::oneshot`, so the whole assembly runs in-process.
415    #[tokio::test]
416    async fn assemble_produces_a_serving_node_over_a_temp_store() {
417        use axum::body::Body;
418        use axum::http::{Request, StatusCode};
419        use tower::ServiceExt;
420
421        let tmp = tempfile::tempdir().unwrap();
422        let storage: Arc<dyn Storage> = Arc::new(boatramp_storage::FsStorage::new(tmp.path()));
423        let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
424        let config = ServerConfig::default();
425        let options = boatramp_server::ServerOptions {
426            // The strict `multi-tenant` posture, as an unconfigured `serve` resolves.
427            posture: SecurityProfile::MultiTenant.preset(),
428            ..Default::default()
429        };
430
431        let node = assemble(NodeInput {
432            config: &config,
433            data_dir: tmp.path(),
434            storage,
435            kv,
436            auth: boatramp_server::Auth::disabled(),
437            options,
438            serve_addr: None,
439            watch_provider: None,
440            provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
441            messaging: None,
442            is_leader: Arc::new(|| true),
443            node_id: 0,
444            worker_exe: None,
445        })
446        .await
447        .expect("assemble a node over a temp store");
448
449        // The deploy store is live: `assemble` already materialized the reserved
450        // `default` project, so a second ensure reports "already present" (`false`).
451        assert!(
452            !node
453                .deploy
454                .ensure_default_project()
455                .await
456                .expect("read the default project"),
457            "assemble should have materialized the default project"
458        );
459
460        // The assembled router (the same wiring `serve` binds) answers /healthz.
461        let router =
462            boatramp_server::router_with(node.deploy, node.auth, node.handlers, node.options);
463        let response = router
464            .oneshot(
465                Request::builder()
466                    .uri("/healthz")
467                    .body(Body::empty())
468                    .unwrap(),
469            )
470            .await
471            .expect("route /healthz");
472        assert_eq!(response.status(), StatusCode::OK);
473    }
474}