use super::*;
#[cfg(feature = "handlers")]
#[allow(clippy::too_many_arguments)]
pub(super) async fn dispatch_handler(
runtime: &HandlerRuntime,
deploy: &DeployStore,
manifest: &Manifest,
site: &str,
request_path: &str,
site_config: Option<&SiteConfig>,
handler: &boatramp_core::config::HandlerConfig,
mut request: Request,
client_ip: IpAddr,
preview: Option<&str>,
) -> Response {
let Some(inner) = runtime.inner.as_ref() else {
return not_found();
};
let scope = match preview {
Some(id) => format!("{site}/_preview/{id}"),
None => site.to_string(),
};
set_forwarded_headers(&mut request, client_ip);
rewrite_request_uri(&mut request, request_path);
let Some(site_handlers) = site_config
.and_then(|c| c.handlers.as_ref())
.filter(|h| h.enabled)
else {
return not_found();
};
let Some(entry) = manifest.files.get(&handler.component) else {
tracing::warn!(site, component = %handler.component, "handler component missing from deployment");
return handler_unavailable();
};
let wasm = match read_blob_fully(deploy, &entry.hash).await {
Ok(bytes) => bytes,
Err(response) => return response,
};
let bindings = build_bindings(
inner,
site,
&scope,
preview,
&handler.imports,
site_handlers,
&handler.env,
)
.await;
let _site_permit = match acquire_site_permit(inner, &scope, site_handlers) {
Ok(permit) => permit,
Err(()) => {
return (
StatusCode::SERVICE_UNAVAILABLE,
"site handler concurrency limit reached\n",
)
.into_response()
}
};
let limits = effective_limits(site_handlers, handler);
let start = std::time::Instant::now();
let result = inner
.engine
.serve_with_limits(&entry.hash, &wasm, request, bindings, limits)
.await;
inner.metrics.observe(
site,
metrics::Trigger::Http,
&handler.route,
&entry.hash,
metrics::Outcome::from_result(&result),
start.elapsed(),
);
match result {
Ok(response) => {
let (parts, body) = response.into_parts();
axum::http::Response::from_parts(parts, axum::body::Body::new(body))
}
Err(err) => {
tracing::warn!(site, route = %handler.route, %err, "handler invocation failed");
handler_error_response(&err)
}
}
}
#[cfg(feature = "handlers")]
pub(super) fn set_forwarded_headers(request: &mut Request, client_ip: IpAddr) {
let headers = request.headers_mut();
if let Ok(value) = HeaderValue::from_str(&client_ip.to_string()) {
headers.insert(HeaderName::from_static("x-forwarded-for"), value);
}
if let Some(host) = headers.get(header::HOST).cloned() {
headers.insert(HeaderName::from_static("x-forwarded-host"), host);
}
if !headers.contains_key("x-forwarded-proto") {
headers.insert(
HeaderName::from_static("x-forwarded-proto"),
HeaderValue::from_static("http"),
);
}
}
#[cfg(feature = "handlers")]
fn rewrite_request_uri(request: &mut Request, request_path: &str) {
let authority = request
.headers()
.get(header::HOST)
.and_then(|value| value.to_str().ok())
.filter(|host| !host.is_empty())
.unwrap_or("localhost")
.to_string();
let path_and_query = match request.uri().query() {
Some(query) => format!("{request_path}?{query}"),
None => request_path.to_string(),
};
if let Ok(uri) = format!("http://{authority}{path_and_query}").parse() {
*request.uri_mut() = uri;
}
}
#[cfg(feature = "handlers")]
#[allow(clippy::too_many_arguments)]
pub(super) async fn precheck_component(
deploy: &DeployStore,
manifest: &Manifest,
site_handlers: &boatramp_core::config::HandlersSiteConfig,
inner: &HandlerRuntimeInner,
max_component: u64,
imports: &[String],
component: &str,
label: &str,
) -> Result<(), String> {
for import in imports {
if !site_handlers.allow_imports.iter().any(|a| a == import) {
return Err(format!(
"{label} requests import {import:?} the site does not allow"
));
}
if import == "sql" && inner.sql.is_none() {
return Err(format!(
"{label} requests `sql` but this server has no SQL backend configured"
));
}
if import == "wasi:messaging" && inner.messaging.is_none() {
return Err(format!(
"{label} requests `wasi:messaging` but this server has no messaging backend"
));
}
}
let entry = manifest
.files
.get(component)
.ok_or_else(|| format!("{label} component {component:?} missing from deployment"))?;
if max_component != 0 && entry.size > max_component {
return Err(format!(
"{label} component {component:?} is {} bytes, over the {max_component}-byte limit",
entry.size
));
}
let wasm = read_blob_bytes(deploy, &entry.hash)
.await
.map_err(|err| format!("reading {label} component: {err}"))?;
inner
.engine
.precompile(&entry.hash, &wasm)
.map_err(|err| format!("{label} failed to compile: {err}"))?;
Ok(())
}
#[cfg(feature = "handlers")]
pub(super) async fn read_blob_bytes(
deploy: &DeployStore,
hash: &str,
) -> Result<Vec<u8>, DeployError> {
let object = deploy.open_blob(hash).await?;
let mut body = object.body;
let mut buf = Vec::new();
while let Some(chunk) = body.next().await {
buf.extend_from_slice(&chunk?);
}
Ok(buf)
}
#[cfg(feature = "handlers")]
pub(super) async fn read_blob_fully(deploy: &DeployStore, hash: &str) -> Result<Vec<u8>, Response> {
read_blob_bytes(deploy, hash)
.await
.map_err(deploy_error_response)
}
#[cfg(feature = "handlers")]
#[allow(clippy::too_many_arguments)]
pub(super) async fn build_bindings(
inner: &HandlerRuntimeInner,
site: &str,
scope: &str,
preview: Option<&str>,
imports: &[String],
site_handlers: &boatramp_core::config::HandlersSiteConfig,
deploy_env: &std::collections::BTreeMap<String, String>,
) -> boatramp_handlers::Bindings {
let granted = |name: &str| {
imports.iter().any(|i| i == name) && site_handlers.allow_imports.iter().any(|a| a == name)
};
let mut bindings = boatramp_handlers::Bindings::new(scope);
if granted("wasi:keyvalue") {
bindings = bindings.with_keyvalue(scope, inner.kv.clone());
}
if granted("wasi:blobstore") {
let max_blob = inner.max_blob_bytes.get().copied().unwrap_or(0);
bindings = bindings.with_blobstore(scope, inner.storage.clone(), max_blob);
}
if granted("sql") {
if let Some(provider) = &inner.sql {
let opened = match preview {
Some(id) => provider.preview_database(site, "", id).await,
None => provider.database(site, "").await,
};
match opened {
Ok(backend) => bindings = bindings.with_sql("", backend),
Err(err) => tracing::warn!(site, %err, "opening site SQL database failed"),
}
}
}
if granted("wasi:messaging") {
if let Some(messaging) = &inner.messaging {
bindings = bindings.with_messaging(format!("{scope}/"), messaging.clone());
}
}
inner.logs.configure(site, site_handlers.max_log_rate);
bindings = bindings.with_logging(site.to_string(), inner.logs.clone());
bindings = bindings.with_env(resolve_env(site, deploy_env, site_handlers));
bindings
}
#[cfg(feature = "handlers")]
pub(super) fn resolve_env(
site: &str,
deploy_env: &std::collections::BTreeMap<String, String>,
site_handlers: &boatramp_core::config::HandlersSiteConfig,
) -> Vec<(String, String)> {
let mut env: Vec<(String, String)> = deploy_env
.iter()
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
for (guest_name, host_ref) in &site_handlers.secrets {
match std::env::var(host_ref) {
Ok(value) => {
env.retain(|(k, _)| k != guest_name);
env.push((guest_name.clone(), value));
}
Err(_) => tracing::warn!(
site,
secret = %guest_name,
"site secret references env var {host_ref}, which is not set; not injected"
),
}
}
env
}
#[cfg(feature = "handlers")]
#[allow(clippy::too_many_arguments)]
pub(super) async fn dispatch_consumer_batch(
engine: &boatramp_handlers::HandlerEngine,
messaging: &dyn boatramp_core::messaging::Messaging,
metrics: &metrics::Metrics,
site: &str,
namespaced_topic: &str,
scope_prefix: &str,
component_hash: &str,
component: &[u8],
bindings: &boatramp_handlers::Bindings,
limits: boatramp_handlers::Limits,
lease: Duration,
max_attempts: u32,
batch: usize,
) -> usize {
let claimed = match messaging
.claim(namespaced_topic, lease, batch, max_attempts)
.await
{
Ok(claimed) => claimed,
Err(err) => {
tracing::warn!(topic = namespaced_topic, %err, "messaging claim failed");
return 0;
}
};
let mut acked = 0;
for msg in claimed {
let guest_topic = msg.topic.strip_prefix(scope_prefix).unwrap_or(&msg.topic);
let start = std::time::Instant::now();
let result = engine
.dispatch_message(
component_hash,
component,
guest_topic,
&msg.payload,
bindings.clone(),
limits,
)
.await;
metrics.observe(
site,
metrics::Trigger::Consumer,
guest_topic,
component_hash,
metrics::Outcome::from_result(&result),
start.elapsed(),
);
match result {
Ok(()) => match messaging.ack(&msg).await {
Ok(()) => acked += 1,
Err(err) => tracing::warn!(id = msg.id, %err, "messaging ack failed"),
},
Err(err) => {
tracing::warn!(
id = msg.id,
attempts = msg.attempts,
%err,
"consumer failed; redelivering (dead-letters after max attempts)"
);
let _ = messaging.nack(&msg).await;
}
}
}
acked
}