use std::net::{IpAddr, SocketAddr};
use std::path::PathBuf;
use std::time::Duration;
use crate::Server;
use tracing::info;
#[cfg(feature = "ssr")]
use tracing::warn;
const DEFAULT_SERVE_PORT: u16 = 3000;
#[derive(clap::Args, Debug)]
pub struct ServeArgs {
workload: Option<PathBuf>,
#[arg(long)]
bundle: Option<PathBuf>,
#[arg(long)]
listen: Option<SocketAddr>,
#[arg(long)]
idle_ttl: Option<u64>,
#[arg(long, default_value_t = DEFAULT_SERVE_PORT)]
port: u16,
#[arg(long, default_value = "0.0.0.0")]
host: IpAddr,
#[arg(long)]
revalidate: bool,
#[arg(long, default_value = "mesofact.config.toml")]
publish_config: PathBuf,
#[arg(long, env = "MESOFACT_MIRROR_KEY")]
mirror_key: Option<String>,
#[arg(long, env = "MESOFACT_MIRROR_KEY_FILE")]
mirror_key_file: Option<PathBuf>,
#[arg(long = "allow-route", value_name = "ROUTE")]
allow_route: Vec<String>,
#[arg(long, conflicts_with_all = ["workload", "publish_config", "allow_route"])]
tenants: Option<PathBuf>,
#[arg(
long,
env = "MESOFACT_TRUST_EDGE_AUTH",
num_args = 0..=1,
default_value_t = false,
default_missing_value = "true",
value_parser = parse_truthy,
)]
trust_edge_auth: bool,
#[arg(
long,
env = "MESOFACT_ROUTE_HEADERS",
default_value = "",
value_parser = parse_route_headers,
)]
route_headers: crate::RouteHeaderTable,
#[arg(
long,
env = "MESOFACT_POLICY_DELEGATED",
value_delimiter = ',',
value_parser = parse_policy_field,
)]
policy_delegated: Vec<mesofact_core::RoutePolicy>,
}
fn parse_truthy(raw: &str) -> Result<bool, String> {
match raw.trim().to_ascii_lowercase().as_str() {
"" | "0" | "false" | "no" | "off" => Ok(false),
"1" | "true" | "yes" | "on" => Ok(true),
other => Err(format!(
"expected a boolean (1/true/yes/on or 0/false/no/off), got {other:?}"
)),
}
}
fn parse_route_headers(raw: &str) -> Result<crate::RouteHeaderTable, String> {
crate::RouteHeaderTable::parse(raw).map_err(|e| format!("{e:#}"))
}
fn parse_policy_field(raw: &str) -> Result<mesofact_core::RoutePolicy, String> {
mesofact_core::RoutePolicy::parse(raw.trim()).ok_or_else(|| {
format!(
"unknown route policy {raw:?} — known policies are {}",
mesofact_core::RoutePolicy::ALL
.iter()
.map(|p| p.field())
.collect::<Vec<_>>()
.join(", "),
)
})
}
fn with_declared_route_headers(server: Server, args: &ServeArgs) -> Server {
if !args.route_headers.is_empty() {
info!(
rules = args.route_headers.len(),
"applying the domain manifest's per-route response headers (MESOFACT_ROUTE_HEADERS)",
);
}
server.with_route_headers(args.route_headers.clone())
}
impl ServeArgs {
fn bind_addr(&self) -> SocketAddr {
self.listen
.unwrap_or_else(|| SocketAddr::new(self.host, self.port))
}
}
pub async fn run(args: ServeArgs) -> anyhow::Result<()> {
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env().unwrap_or_else(|_| {
tracing_subscriber::EnvFilter::new("mesofact=info,mesofact_dev=info,tower_http=info")
}),
)
.init();
if let Some(bundle) = args.bundle.as_ref() {
let bundle_abs = bundle.canonicalize().unwrap_or_else(|_| bundle.clone());
let server = Server::from_bundle(&bundle_abs)?;
let app = bundle_abs.join("app");
assert_declared_policy_is_enforced(&app, &serve_policy_support(&args))?;
let server = attach_bundle_ssr(server, &bundle_abs).await?;
let server = with_declared_route_headers(server, &args);
let server = with_declared_cache_policy(server, &app)?;
let idle_ttl = args.idle_ttl.filter(|s| *s > 0).map(Duration::from_secs);
let listener = match socket_activation_listener()? {
Some(l) => {
info!(bundle = %bundle_abs.display(), "mesofact serve: serving bundle on inherited LISTEN_FDS socket");
l
}
None => {
let addr = args.bind_addr();
info!(%addr, bundle = %bundle_abs.display(), "mesofact-serve listening (bundle, static v0)");
tokio::net::TcpListener::bind(addr).await?
}
};
return server.serve_on_listener(listener, idle_ttl).await;
}
run_workload_modes(args).await
}
fn serve_policy_support(args: &ServeArgs) -> mesofact_core::PolicySupport {
use mesofact_core::RoutePolicy;
let mut support =
mesofact_core::PolicySupport::new("mesofact serve").enforces(RoutePolicy::CachePolicy);
#[cfg(feature = "ssr")]
{
support = support.enforces(RoutePolicy::Resilience);
}
if args.trust_edge_auth {
support = support.delegate(RoutePolicy::Requires);
}
for policy in &args.policy_delegated {
support = support.delegate(*policy);
}
support
}
fn assert_declared_policy_is_enforced(
workload: &std::path::Path,
support: &mesofact_core::PolicySupport,
) -> anyhow::Result<()> {
let Some(raw) = crate::read_manifest_bytes(workload).map_err(|e| {
anyhow::anyhow!(
"refusing to start: cannot read the route manifest under {} to check for declared \
policy this binary does not enforce: {e}",
workload.display(),
)
})?
else {
return Ok(());
};
mesofact_core::check_manifest(&raw, support).map_err(|e| anyhow::anyhow!("{e}"))?;
let delegated: Vec<&str> = support.delegated().map(|p| p.field()).collect();
if !delegated.is_empty() {
info!(
policies = ?delegated,
enforced_here = ?support.enforced().map(|p| p.field()).collect::<Vec<_>>(),
"route policies asserted to be enforced IN FRONT of this process, not by it",
);
}
Ok(())
}
fn with_declared_cache_policy(
server: Server,
workload: &std::path::Path,
) -> anyhow::Result<Server> {
let table = crate::declared_cache_policy(workload).map_err(|e| {
anyhow::anyhow!(
"refusing to start: cannot read the route manifest under {} to derive the declared \
cache policy: {e}",
workload.display(),
)
})?;
if !table.is_empty() {
info!(
routes = table.len(),
"enforcing the manifest's declared cache_policy as response Cache-Control/Vary",
);
}
Ok(server.with_cache_policy(table))
}
#[cfg(feature = "ssr")]
async fn attach_bundle_ssr(server: crate::Server, bundle: &std::path::Path) -> anyhow::Result<crate::Server> {
use crate::{ssr, SsrSpawnOptions};
let app = bundle.join("app");
let opts = SsrSpawnOptions::new(app.clone(), app.join("dist"), app.join(".mesofact-serve"))
.with_env(std::env::vars().collect());
Ok(match ssr::spawn(opts).await? {
Some(child) => {
info!(prefixes = ?child.prefixes(), "mesofact serve --bundle: ssr runtime attached");
server.with_ssr(child)
}
None => server,
})
}
#[cfg(not(feature = "ssr"))]
async fn attach_bundle_ssr(server: crate::Server, bundle: &std::path::Path) -> anyhow::Result<crate::Server> {
refuse_unservable_ssr_routes(&bundle.join("app"))?;
Ok(server)
}
#[cfg(not(feature = "ssr"))]
fn refuse_unservable_ssr_routes(workload: &std::path::Path) -> anyhow::Result<()> {
let ssr_routes = crate::routes_declaring_ssr(workload).map_err(|e| {
anyhow::anyhow!(
"refusing to start: cannot read the route manifest under {} to check for \
mode:\"ssr\" routes: {e}",
workload.display(),
)
})?;
if ssr_routes.is_empty() {
return Ok(());
}
anyhow::bail!(
"refusing to start: {} route(s) declare mode:\"ssr\" ({}) but this `mesofact serve` \
binary was built without the `ssr` feature (no V8) — it can only serve them as a 404. \
Serve this workload with an ssr-enabled build (the shipped stock runtime is built \
`--features deploy`), or drop the ssr route(s) if it is meant to be static.",
ssr_routes.len(),
ssr_routes.join(", "),
);
}
#[cfg(unix)]
fn socket_activation_listener() -> anyhow::Result<Option<tokio::net::TcpListener>> {
use std::os::fd::FromRawFd;
let n_fds: i32 = std::env::var("LISTEN_FDS")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(0);
if n_fds < 1 {
return Ok(None);
}
if let Ok(pid) = std::env::var("LISTEN_PID") {
if pid.parse::<u32>().ok() != Some(std::process::id()) {
return Ok(None);
}
}
const SD_LISTEN_FDS_START: i32 = 3;
let std_listener = unsafe { std::net::TcpListener::from_raw_fd(SD_LISTEN_FDS_START) };
std_listener
.set_nonblocking(true)
.map_err(|e| anyhow::anyhow!("set_nonblocking on inherited LISTEN_FDS socket: {e}"))?;
let listener = tokio::net::TcpListener::from_std(std_listener)
.map_err(|e| anyhow::anyhow!("adopting inherited LISTEN_FDS socket: {e}"))?;
Ok(Some(listener))
}
#[cfg(not(unix))]
fn socket_activation_listener() -> anyhow::Result<Option<tokio::net::TcpListener>> {
Ok(None)
}
fn resolve_mirror_key(direct: Option<String>, file: Option<PathBuf>) -> anyhow::Result<Option<String>> {
match direct {
Some(key) => Ok(Some(key)),
None => match file {
Some(path) => {
let contents = std::fs::read_to_string(&path)
.map_err(|e| anyhow::anyhow!("reading --mirror-key-file {}: {e}", path.display()))?;
Ok(Some(contents.trim().to_string()))
}
None => Ok(None),
},
}
}
#[cfg(feature = "ssr")]
async fn run_workload_modes(args: ServeArgs) -> anyhow::Result<()> {
use crate::{revalidate, ssr, tenants, SsrSpawnOptions};
let addr = args.bind_addr();
if let Some(tenants_dir) = args.tenants.as_ref() {
if !args.revalidate {
warn!("--tenants implies the revalidate receiver; running multi-tenant receiver");
}
if args.mirror_key.is_some() {
warn!(
"--mirror-key / MESOFACT_MIRROR_KEY is ignored in --tenants mode — \
each tenant's bearer comes from its own mirror_key_env",
);
}
let files = tenants::load_tenants(tenants_dir)?;
let resolved = tenants::resolve_tenants(files, |name| std::env::var(name).ok());
let registry = tenants::TenantRegistry::new(resolved);
if registry.is_empty() {
anyhow::bail!(
"--tenants {} contains no tenants/<id>.toml files — a receiver with an \
empty registry rejects every poke",
tenants_dir.display()
);
}
registry.validate()?;
info!(tenants = registry.len(), dir = %tenants_dir.display(), "multi-tenant revalidate receiver");
return tenants::serve(registry, addr.ip(), addr.port()).await;
}
if args.revalidate {
let workload = args
.workload
.clone()
.ok_or_else(|| anyhow::anyhow!("--revalidate needs a <workload> dir (or use --tenants)"))?;
let workload_abs = workload.canonicalize().unwrap_or(workload);
let mirror_key = resolve_mirror_key(args.mirror_key, args.mirror_key_file)?;
return revalidate::serve(
revalidate::RevalidateConfig {
workload: workload_abs,
publish_config: args.publish_config,
mirror_key,
routes: args.allow_route,
},
addr.ip(),
addr.port(),
)
.await;
}
let workload = args
.workload
.clone()
.ok_or_else(|| anyhow::anyhow!("a <workload> dir is required for the SSR host (or --bundle for static serving)"))?;
let server = Server::from_workload(&workload)?;
let workload_abs = workload.canonicalize().unwrap_or(workload);
assert_declared_policy_is_enforced(&workload_abs, &serve_policy_support(&args))?;
let opts = SsrSpawnOptions::new(
workload_abs.clone(),
workload_abs.join("dist"),
workload_abs.join(".mesofact-serve"),
)
.with_env(std::env::vars().collect());
let server = match ssr::spawn(opts).await? {
Some(child) => {
info!(prefixes = ?child.prefixes(), "mesofact-serve ssr runtime attached");
server.with_ssr(child)
}
None => {
warn!("mesofact-serve: no SSR routes (or no manifest) — serving static only");
server
}
};
let server = with_declared_route_headers(server, &args);
let server = with_declared_cache_policy(server, &workload_abs)?;
let addr = args.bind_addr();
info!(%addr, workload = %workload_abs.display(), "mesofact-serve listening");
server.serve_on(addr).await
}
#[cfg(not(feature = "ssr"))]
async fn run_workload_modes(args: ServeArgs) -> anyhow::Result<()> {
if args.revalidate || args.tenants.is_some() {
anyhow::bail!(
"--revalidate / --tenants need the `ssr` build feature (V8); this is a static-only build"
);
}
let workload = args
.workload
.clone()
.ok_or_else(|| anyhow::anyhow!("a <workload> dir is required (or --bundle to serve a W272 bundle)"))?;
let server = Server::from_workload(&workload)?;
assert_declared_policy_is_enforced(&workload, &serve_policy_support(&args))?;
refuse_unservable_ssr_routes(&workload)?;
let server = with_declared_route_headers(server, &args);
let server = with_declared_cache_policy(server, &workload)?;
let addr = args.bind_addr();
info!(%addr, workload = %workload.display(), "mesofact-serve listening (static only, no ssr)");
server.serve_on(addr).await
}
#[cfg(test)]
mod tests {
use super::*;
fn workload_with_manifest(json: &str) -> tempfile::TempDir {
let dir = tempfile::tempdir().unwrap();
let dist = dir.path().join("dist");
std::fs::create_dir_all(&dist).unwrap();
std::fs::write(dist.join("manifest.json"), json).unwrap();
dir
}
const AUTHED: &str = r#"{"routes":[{"route":"/","mode":"ssr","requires":["user"]}]}"#;
#[derive(clap::Parser)]
struct Harness {
#[command(flatten)]
args: ServeArgs,
}
fn args_from(extra: &[&str]) -> ServeArgs {
let mut argv = vec!["mesofact-serve"];
argv.extend_from_slice(extra);
Harness::parse_from(argv).args
}
#[test]
fn resolve_mirror_key_prefers_the_direct_value_over_the_file() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("bearer");
std::fs::write(&path, "from-file\n").unwrap();
let resolved =
resolve_mirror_key(Some("from-flag".to_string()), Some(path)).unwrap();
assert_eq!(resolved.as_deref(), Some("from-flag"));
}
#[test]
fn resolve_mirror_key_falls_back_to_the_file_and_trims_it() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("bearer");
std::fs::write(&path, "sentinel-bearer\n").unwrap();
let resolved = resolve_mirror_key(None, Some(path)).unwrap();
assert_eq!(resolved.as_deref(), Some("sentinel-bearer"));
}
#[test]
fn resolve_mirror_key_is_none_when_both_are_unset() {
assert!(resolve_mirror_key(None, None).unwrap().is_none());
}
#[test]
fn resolve_mirror_key_names_the_path_when_the_file_is_missing() {
let err = resolve_mirror_key(None, Some(PathBuf::from("/nonexistent/mirror-key")))
.unwrap_err();
assert!(err.to_string().contains("/nonexistent/mirror-key"));
}
#[test]
fn mirror_key_file_flag_reaches_serve_args() {
let args = args_from(&["--mirror-key-file", "/run/yah/secrets/mesofact/x/mirror-key"]);
assert_eq!(
args.mirror_key_file.as_deref(),
Some(std::path::Path::new(
"/run/yah/secrets/mesofact/x/mirror-key"
))
);
assert!(args.mirror_key.is_none());
}
fn support(extra: &[&str]) -> mesofact_core::PolicySupport {
serve_policy_support(&args_from(extra))
}
#[test]
fn a_declared_authed_route_refuses_to_start_without_an_asserted_edge() {
let dir = workload_with_manifest(AUTHED);
let err = assert_declared_policy_is_enforced(dir.path(), &support(&[]))
.unwrap_err()
.to_string();
assert!(err.contains("refusing to start"), "{err}");
assert!(
err.contains('/'),
"the message must name the route the operator has to look at: {err}"
);
assert!(
err.contains("--policy-delegated"),
"and the remedy, or the operator has a refusal with no next move: {err}"
);
}
#[test]
fn an_asserted_edge_allows_the_declared_authed_route() {
let dir = workload_with_manifest(AUTHED);
assert!(assert_declared_policy_is_enforced(dir.path(), &support(&["--trust-edge-auth"])).is_ok());
assert!(assert_declared_policy_is_enforced(
dir.path(),
&support(&["--policy-delegated", "requires"]),
)
.is_ok());
}
#[test]
fn a_workload_declaring_no_authed_route_starts_unchanged() {
let dir = workload_with_manifest(r#"{"routes":[{"route":"/","mode":"static"}]}"#);
assert!(assert_declared_policy_is_enforced(dir.path(), &support(&[])).is_ok());
let empty = tempfile::tempdir().unwrap();
assert!(assert_declared_policy_is_enforced(empty.path(), &support(&[])).is_ok());
}
#[test]
fn a_policy_this_binary_does_not_implement_refuses_to_start() {
let dir = workload_with_manifest(
r#"{"routes":[{"route":"/busy","mode":"ssr","cache_policy":{"ttl":0},"concurrency":4}]}"#,
);
let err = assert_declared_policy_is_enforced(dir.path(), &support(&[]))
.unwrap_err()
.to_string();
assert!(err.contains("/busy") && err.contains("concurrency"), "{err}");
assert!(assert_declared_policy_is_enforced(
dir.path(),
&support(&["--policy-delegated", "concurrency"]),
)
.is_ok());
}
#[test]
fn a_policy_field_this_binary_predates_refuses_to_start() {
let dir = workload_with_manifest(
r#"{"routes":[{"route":"/x","mode":"ssr","cache_policy":{"ttl":0},"rate_limit":{"rps":10}}]}"#,
);
let err = assert_declared_policy_is_enforced(dir.path(), &support(&["--trust-edge-auth"]))
.unwrap_err()
.to_string();
assert!(err.contains("rate_limit") && err.contains("not delegatable"), "{err}");
}
#[test]
fn a_declared_cache_policy_is_enforced_rather_than_refused() {
let dir = workload_with_manifest(
r#"{"routes":[{"route":"/issues","mode":"static","render_entrypoint":"e.js","cache_policy":{"ttl":3600,"swr":86400}}]}"#,
);
assert!(assert_declared_policy_is_enforced(dir.path(), &support(&[])).is_ok());
let table = crate::declared_cache_policy(dir.path()).unwrap();
assert_eq!(table.len(), 1, "the advertised enforcement must be real");
}
#[test]
fn an_unknown_delegated_policy_is_rejected_not_ignored() {
assert!(parse_policy_field("cache-policy").is_err());
assert!(parse_policy_field("requires").is_ok());
}
#[test]
fn the_edge_assertion_accepts_the_truthy_forms_an_operator_will_type() {
for yes in ["1", "true", "TRUE", "yes", "on", " true "] {
assert_eq!(parse_truthy(yes), Ok(true), "{yes:?}");
}
for no in ["", "0", "false", "no", "off"] {
assert_eq!(parse_truthy(no), Ok(false), "{no:?}");
}
assert_eq!(parse_truthy(""), Ok(false));
assert!(parse_truthy("maybe").is_err());
}
#[derive(clap::Parser)]
struct TestCli {
#[command(flatten)]
serve: ServeArgs,
}
use clap::Parser as _;
fn parse_args(value: &str) -> Result<TestCli, clap::Error> {
TestCli::try_parse_from(["mesofact-serve", "--bundle", "b", "--route-headers", value])
}
#[test]
fn a_declared_table_reaches_the_server() {
let cli = parse_args(
r#"[{"path":"/app/*","headers":{"Cross-Origin-Opener-Policy":"same-origin"}},{"path":"/*","headers":{"X-Tier":"marketing"}}]"#,
)
.expect("a well-formed table must parse");
assert_eq!(cli.serve.route_headers.len(), 2);
}
#[test]
fn an_unset_table_starts_normally() {
let bare = TestCli::try_parse_from(["mesofact-serve", "--bundle", "b"]).unwrap();
assert!(bare.serve.route_headers.is_empty());
assert!(parse_args("").unwrap().serve.route_headers.is_empty());
assert!(parse_args(" ").unwrap().serve.route_headers.is_empty());
}
#[test]
fn a_malformed_table_refuses_the_start() {
for bad in [
"{not json",
"[",
r#"{"path":"/*","headers":{}}"#,
r#"[{"path":"/*"}]"#,
r#"[{"path":"/*","headers":{"Bad Name":"1"}}]"#,
] {
match parse_args(bad) {
Ok(cli) => panic!(
"malformed table {bad:?} started anyway with {} rule(s) — that is the \
serve-anyway posture this must not have",
cli.serve.route_headers.len(),
),
Err(e) => {
let rendered = e.to_string();
assert!(
rendered.contains("route header table")
|| rendered.contains("not a valid HTTP header name")
|| rendered.contains("missing field"),
"the refusal must name the problem: {rendered}",
);
}
}
}
}
#[test]
fn an_unreadable_manifest_refuses_to_start() {
let dir = workload_with_manifest("{ not json");
let err = assert_declared_policy_is_enforced(dir.path(), &support(&[]))
.unwrap_err()
.to_string();
assert!(err.contains("refusing to start"), "{err}");
}
#[cfg(not(feature = "ssr"))]
#[tokio::test]
async fn a_bundle_declaring_ssr_refuses_to_start_on_a_static_only_build() {
let bundle_dir = tempfile::tempdir().unwrap();
let bundle = bundle_dir.path();
let app = bundle.join("app");
std::fs::create_dir_all(app.join("dist")).unwrap();
std::fs::write(
app.join("dist").join("manifest.json"),
r#"{"routes":[{"route":"/live","mode":"ssr"},{"route":"/","mode":"static"}]}"#,
)
.unwrap();
let server = crate::Server::from_workload(&app).unwrap();
let err = match attach_bundle_ssr(server, bundle).await {
Ok(_) => panic!("expected a refusal — bundle declares mode:\"ssr\""),
Err(e) => e.to_string(),
};
assert!(err.contains("refusing to start"), "{err}");
assert!(err.contains("/live"), "must name the route: {err}");
assert!(err.contains("ssr"), "and say why: {err}");
}
#[cfg(not(feature = "ssr"))]
#[tokio::test]
async fn a_static_only_bundle_starts_unchanged_on_a_static_only_build() {
let bundle_dir = tempfile::tempdir().unwrap();
let bundle = bundle_dir.path();
let app = bundle.join("app");
std::fs::create_dir_all(app.join("dist")).unwrap();
std::fs::write(
app.join("dist").join("manifest.json"),
r#"{"routes":[{"route":"/","mode":"static"}]}"#,
)
.unwrap();
let server = crate::Server::from_workload(&app).unwrap();
assert!(attach_bundle_ssr(server, bundle).await.is_ok());
}
}