use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
use axum::Router;
use axum::extract::{DefaultBodyLimit, State};
use axum::http::{HeaderMap, StatusCode, header};
use axum::middleware::from_fn_with_state;
use axum::response::IntoResponse;
use axum::routing::get;
use serde_json::json;
use tracing::error;
use crate::admin::require_admin;
use crate::config::Config;
use crate::config::ConfigError;
use crate::events::EventBus;
use crate::http::{
Json, MAX_BODY_BYTES, ScopeState, cors_layer, scope_layer, security_headers_layer,
};
use crate::module::{HARNESS_API, Module, ModuleContext, harness_api_mismatch};
use crate::ports::Dispatcher;
use crate::ports::{Clock, Database, Port, Ports, Statement, SystemClock, warn_undeclared_ports};
use crate::problem::Problem;
use crate::sidecar::SidecarMount;
use crate::surface::{RenderedSurface, SurfaceDocument, SurfaceSource, UiContext, UiMount};
use crate::template::{Template, TemplateRegistry};
use crate::venture::Venture;
pub trait Runtime: Send + Sync + 'static {
fn provides(&self) -> Vec<Port>;
}
pub struct Harness {
venture: Arc<Venture>,
modules: Vec<Arc<dyn Module>>,
templates: Arc<TemplateRegistry>,
events: EventBus,
runtime: Option<Arc<dyn Runtime>>,
well_known: Option<Router>,
surface: Arc<SurfaceVariants>,
ui: Option<Arc<dyn UiMount>>,
}
struct SurfaceVariants {
document: Arc<SurfaceDocument>,
full: RenderedSurface,
public: RenderedSurface,
}
impl SurfaceVariants {
fn compose(
venture: &Venture,
modules: &[Arc<dyn Module>],
ui: Option<&Arc<dyn UiMount>>,
) -> Self {
let mut document = SurfaceDocument::compose(venture, modules);
document.ui = ui.and_then(|ui| ui.describe());
Self {
full: RenderedSurface::render(&document),
public: RenderedSurface::render(&document.public()),
document: Arc::new(document),
}
}
}
impl std::fmt::Debug for Harness {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let modules: Vec<&str> = self.modules.iter().map(|m| m.name()).collect();
f.debug_struct("Harness")
.field("venture", &self.venture.name)
.field("modules", &modules)
.finish_non_exhaustive()
}
}
impl Harness {
pub fn builder() -> HarnessBuilder {
HarnessBuilder::default()
}
pub fn venture(&self) -> &Arc<Venture> {
&self.venture
}
pub fn modules(&self) -> &[Arc<dyn Module>] {
&self.modules
}
pub fn templates(&self) -> &Arc<TemplateRegistry> {
&self.templates
}
pub fn events(&self) -> &EventBus {
&self.events
}
pub fn module_context(&self, module: &dyn Module, ports: &Ports) -> ModuleContext {
ModuleContext {
config: Arc::clone(&ports.config),
ports: ports.view_for(module),
events: self.events.clone(),
templates: Arc::clone(&self.templates),
venture: Arc::clone(&self.venture),
ui_mounted: self.ui.is_some(),
}
}
fn sidecar_mounts(&self, ports: &Ports) -> Vec<SidecarMount> {
let mounts = match crate::sidecar::SidecarMounts::from_config(ports.config.as_ref()) {
Ok(mounts) => mounts,
Err(errors) => {
for error in errors {
tracing::error!(error, "ignoring the sidecar mount table");
}
crate::sidecar::SidecarMounts::default()
}
};
let module_names: Vec<&str> = self.modules.iter().map(|m| m.name()).collect();
for collision in mounts.collisions(&module_names) {
tracing::error!(error = collision, "ignoring the colliding sidecar mount");
}
mounts
.iter()
.filter(|m| !module_names.contains(&m.name.as_str()))
.cloned()
.collect()
}
pub fn runtime(&self) -> Option<&Arc<dyn Runtime>> {
self.runtime.as_ref()
}
#[must_use]
pub fn surface(&self) -> &SurfaceDocument {
&self.surface.document
}
pub fn router(&self, ports: Ports) -> Router {
let mut api = Router::new();
for module in &self.modules {
let ctx = self.module_context(module.as_ref(), &ports);
api = api.nest(&format!("/v1/{}", module.name()), module.router(ctx));
}
let mounted = self.sidecar_mounts(&ports);
for mount in &mounted {
api = api.nest(
&format!("/v1/{}", mount.name),
crate::sidecar::router(mount.clone(), ports.dispatcher.clone()),
);
}
let api = api
.layer(axum::middleware::from_fn(security_headers_layer))
.layer(DefaultBodyLimit::max(MAX_BODY_BYTES));
let surface_source: Arc<dyn SurfaceSource> = Arc::new(MergedSurface {
base: Arc::clone(&self.surface),
mounts: mounted,
dispatcher: ports.dispatcher.clone(),
});
let ui = self.ui.as_ref().map(|ui| {
ui.router(UiContext {
surface: Arc::clone(&surface_source),
api: api.clone(),
config: Arc::clone(&ports.config),
venture: Arc::clone(&self.venture),
captcha_configured: ports.captcha.is_some(),
signer: ports.signer.clone(),
rate_limiter: ports.rate_limiter.clone(),
})
.layer(DefaultBodyLimit::max(MAX_BODY_BYTES))
});
let Ports {
config,
db,
mailer,
captcha,
rate_limiter: _,
signer: _,
kv: _,
blob: _,
http: _,
clock,
id_gen,
defer,
dispatcher: _,
} = ports;
let health_state = HealthState {
venture: Arc::clone(&self.venture),
modules: self.modules.clone(),
harness_build: config
.get("HARNESS_BUILD")
.filter(|build| !build.is_empty()),
mailer_configured: mailer.is_some(),
captcha_configured: captcha.is_some(),
};
let scope_state = ScopeState {
defer: defer.unwrap_or_else(|| Arc::new(crate::ports::NoopDefer)),
id_gen: id_gen.unwrap_or_else(|| Arc::new(crate::ports::UlidIdGen)),
};
let ready_state = ReadyState {
db,
clock: clock.unwrap_or_else(|| Arc::new(SystemClock)),
};
let surface_state = SurfaceState {
config,
source: surface_source,
};
let root = Router::new()
.route("/__health", get(health_handler))
.with_state(health_state)
.route("/__ready", get(ready_handler))
.with_state(ready_state)
.route("/__surface", get(surface_handler))
.with_state(surface_state)
.merge(api);
let root = match &self.well_known {
Some(well_known) => root.nest(
"/.well-known",
well_known
.clone()
.layer(DefaultBodyLimit::max(MAX_BODY_BYTES)),
),
None => root,
};
let root = match ui {
Some(ui) => root.nest("/ui", ui),
None => root,
};
root.layer(from_fn_with_state(scope_state, scope_layer))
.layer(cors_layer(&self.venture.cors_origins))
}
}
#[derive(Clone)]
struct HealthState {
venture: Arc<Venture>,
modules: Vec<Arc<dyn Module>>,
harness_build: Option<String>,
mailer_configured: bool,
captcha_configured: bool,
}
async fn health_handler(State(state): State<HealthState>) -> impl IntoResponse {
let modules: Vec<serde_json::Value> = state
.modules
.iter()
.map(|module| {
json!({
"name": module.name(),
"version": module.version(),
"emits": module.emits(),
})
})
.collect();
Json(json!({
"venture": state.venture.name,
"env": state.venture.env.as_str(),
"harness_api": HARNESS_API,
"harness_build": state.harness_build,
"mailer": if state.mailer_configured { "configured" } else { "not_configured" },
"captcha": if state.captcha_configured { "configured" } else { "absent" },
"modules": modules,
}))
}
#[derive(Clone)]
struct ReadyState {
db: Option<Arc<dyn Database>>,
clock: Arc<dyn Clock>,
}
async fn ready_handler(State(state): State<ReadyState>) -> impl IntoResponse {
let Some(db) = state.db else {
return Problem::not_ready("database port is not configured").into_response();
};
let stmt = Statement::new("SELECT 1");
let query = async move { db.query(&stmt).await };
match crate::ports::timeout(&*state.clock, query, Duration::from_secs(2)).await {
Some(Ok(_rows)) => Json(json!({ "ok": true })).into_response(),
Some(Err(err)) => {
error!(error = %err, "readiness probe query failed");
Problem::not_ready("database query failed").into_response()
}
None => Problem::not_ready("database did not answer within 2 s").into_response(),
}
}
#[derive(Clone)]
struct SurfaceState {
config: Arc<dyn Config>,
source: Arc<dyn SurfaceSource>,
}
struct MergedSurface {
base: Arc<SurfaceVariants>,
mounts: Vec<SidecarMount>,
dispatcher: Option<Arc<dyn Dispatcher>>,
}
impl MergedSurface {
async fn sidecar_modules(&self) -> Vec<crate::surface::ModuleSurface> {
let mut extra = Vec::new();
let Some(dispatcher) = &self.dispatcher else {
return extra;
};
for mount in &self.mounts {
if !dispatcher.has(&mount.binding) {
continue;
}
let request = axum::http::Request::builder()
.method(axum::http::Method::GET)
.uri("/__surface")
.header(header::ACCEPT, "application/json")
.body(bytes::Bytes::new())
.expect("static request builds");
let answer = match dispatcher.dispatch(&mount.binding, request).await {
Ok(response) if response.status().is_success() => response,
Ok(response) => {
tracing::warn!(module = mount.name, status = %response.status(), "sidecar surface not available");
continue;
}
Err(err) => {
tracing::warn!(module = mount.name, error = %err, "sidecar surface fetch failed");
continue;
}
};
match serde_json::from_slice::<SurfaceDocument>(answer.body()) {
Ok(document) => extra.extend(
document
.modules
.into_iter()
.filter(|m| m.name == mount.name)
.map(|m| crate::surface::ModuleSurface {
name: m.name,
version: m.version,
surface: m.surface.public(),
}),
),
Err(err) => {
tracing::warn!(module = mount.name, error = %err, "sidecar surface is not a surface document");
}
}
}
extra
}
}
#[async_trait::async_trait]
impl SurfaceSource for MergedSurface {
async fn current(&self) -> Arc<SurfaceDocument> {
if self.mounts.is_empty() {
return Arc::clone(&self.base.document);
}
let mut document = (*self.base.document).clone();
document.modules.extend(self.sidecar_modules().await);
Arc::new(document)
}
fn built(&self) -> Arc<SurfaceDocument> {
Arc::clone(&self.base.document)
}
fn rendered(&self, admin: bool) -> Option<&RenderedSurface> {
Some(if admin {
&self.base.full
} else {
&self.base.public
})
}
}
async fn surface_handler(
State(state): State<SurfaceState>,
headers: HeaderMap,
) -> impl IntoResponse {
let admin = require_admin(&*state.config, &headers).is_ok();
let current = state.source.current().await;
let built = state.source.built();
let prerendered = Arc::ptr_eq(¤t, &built).then(|| state.source.rendered(admin));
let fresh;
let rendered: &RenderedSurface = if let Some(rendered) = prerendered.flatten() {
rendered
} else {
fresh = if admin {
RenderedSurface::render(¤t)
} else {
RenderedSurface::render(¤t.public())
};
&fresh
};
let matches = headers
.get(header::IF_NONE_MATCH)
.and_then(|value| value.to_str().ok())
.is_some_and(|value| {
value
.split(',')
.map(str::trim)
.any(|tag| tag == "*" || tag == rendered.etag)
});
let mut response = if matches {
StatusCode::NOT_MODIFIED.into_response()
} else {
(
[(header::CONTENT_TYPE, "application/json")],
rendered.json.clone(),
)
.into_response()
};
let response_headers = response.headers_mut();
response_headers.insert(
header::ETAG,
header::HeaderValue::from_str(&rendered.etag).expect("hex etag is a valid header"),
);
response_headers.insert(
header::CACHE_CONTROL,
header::HeaderValue::from_static("no-cache"),
);
response_headers.insert(
header::VARY,
header::HeaderValue::from_static("Authorization"),
);
response
}
#[derive(Default)]
pub struct HarnessBuilder {
venture: Option<Venture>,
modules: Vec<Arc<dyn Module>>,
provides: Vec<Port>,
runtime: Option<Arc<dyn Runtime>>,
module_templates: Vec<(String, Box<dyn Template>)>,
overrides: Vec<(String, Box<dyn Template>)>,
ui: Option<Arc<dyn UiMount>>,
}
impl HarnessBuilder {
#[must_use]
pub fn venture(mut self, venture: Venture) -> Self {
self.venture = Some(venture);
self
}
#[must_use]
pub fn module(mut self, module: impl Module) -> Self {
self.modules.push(Arc::new(module));
self
}
#[must_use]
pub fn module_arc(mut self, module: Arc<dyn Module>) -> Self {
self.modules.push(module);
self
}
#[must_use]
pub fn runtime(mut self, runtime: impl Runtime) -> Self {
let runtime: Arc<dyn Runtime> = Arc::new(runtime);
self.provides = runtime.provides();
self.runtime = Some(runtime);
self
}
#[must_use]
pub fn templates(
mut self,
templates: impl IntoIterator<Item = (String, Box<dyn Template>)>,
) -> Self {
self.module_templates.extend(templates);
self
}
#[must_use]
pub fn ui(mut self, ui: impl UiMount) -> Self {
self.ui = Some(Arc::new(ui));
self
}
#[must_use]
pub fn template(mut self, id: impl Into<String>, template: Box<dyn Template>) -> Self {
self.overrides.push((id.into(), template));
self
}
pub fn build(self) -> Result<Harness, ConfigError> {
let mut errors = ConfigError::default();
let venture = if let Some(venture) = self.venture {
venture.validate(&mut errors);
venture
} else {
errors.push("missing venture: call .venture(Venture::new(..)) before .build()");
Venture::new("invalid", "invalid.invalid")
};
let mut names: HashMap<&'static str, usize> = HashMap::new();
let mut tables: HashMap<&'static str, &'static str> = HashMap::new();
for module in &self.modules {
if module.harness_api() != HARNESS_API {
errors.push(harness_api_mismatch(module.as_ref()));
}
let name = module.name();
if name.is_empty() || !is_module_name(name) {
errors.push(format!(
"module name `{name}` must be kebab-case ([a-z0-9]+ separated by '-')"
));
}
match names.get(name) {
Some(_) => errors.push(format!("duplicate module name `{name}`")),
None => {
names.insert(name, 1);
}
}
for port in module.requires().iter().chain(module.optional()) {
if !Port::ALL.contains(port) {
errors.push(format!(
"module `{name}` declares unknown port {}",
port.name()
));
}
}
for port in module.requires() {
if module.optional().contains(port) {
errors.push(format!(
"module `{name}` lists port {} in both requires() and optional()",
port.name()
));
}
}
module.surface().validate(name, &mut errors);
for table in module.tables() {
match tables.get(table) {
Some(owner) => errors.push(format!(
"duplicate table `{table}` claimed by modules `{owner}` and `{name}`"
)),
None => {
tables.insert(table, name);
}
}
}
}
let well_known = collect_well_known(&self.modules, &mut errors);
for module in &self.modules {
for port in module.requires() {
if !self.provides.contains(port) {
errors.push(format!(
"module `{}` requires port {} which the runtime does not provide",
module.name(),
port.name()
));
}
}
warn_undeclared_ports(module.as_ref(), &self.provides);
}
check_template_ids(
self.overrides.iter().chain(self.module_templates.iter()),
&names,
&mut errors,
);
let surface = Arc::new(SurfaceVariants::compose(
&venture,
&self.modules,
self.ui.as_ref(),
));
if let Some(ui) = &self.ui {
ui.validate(&surface.document, &mut errors);
}
errors.into_result()?;
let mut registry = TemplateRegistry::new();
registry.register_all(self.module_templates);
registry.register_all(self.overrides);
let mut events = EventBus::new();
for module in &self.modules {
for (name, handler) in module.events() {
events = events.on(name, handler);
}
}
Ok(Harness {
venture: Arc::new(venture),
modules: self.modules,
templates: Arc::new(registry),
events,
runtime: self.runtime,
well_known,
surface,
ui: self.ui,
})
}
}
fn check_template_ids<'a>(
ids: impl Iterator<Item = &'a (String, Box<dyn Template>)>,
names: &HashMap<&'static str, usize>,
errors: &mut ConfigError,
) {
for (id, _) in ids {
let Some(module_name) = id.split('/').next() else {
continue;
};
if !names.contains_key(module_name) {
errors.push(format!(
"template `{id}` names module `{module_name}` which is not registered"
));
}
}
}
fn is_module_name(name: &str) -> bool {
name.split('-').all(|part| {
!part.is_empty()
&& part
.chars()
.all(|c| c.is_ascii_lowercase() || c.is_ascii_digit())
})
}
fn collect_well_known(modules: &[Arc<dyn Module>], errors: &mut ConfigError) -> Option<Router> {
let mut providers: Vec<&'static str> = Vec::new();
let mut well_known = None;
for module in modules {
if let Some(router) = module.well_known() {
providers.push(module.name());
well_known = Some(router);
}
}
if providers.len() > 1 {
let listed = providers
.iter()
.map(|name| format!("`{name}`"))
.collect::<Vec<_>>()
.join(", ");
errors.push(format!(
"modules {listed} all provide a well-known router; at most one module may occupy /.well-known"
));
}
well_known
}