use super::*;
use boatramp_core::project::ProjectRef;
#[cfg(feature = "handlers")]
const INVOKE_AUTHORITY: &str = "function.invoke";
#[cfg(feature = "handlers")]
const MAX_INVOKE_ATTEMPTS: u32 = 5;
#[cfg(feature = "handlers")]
const MAX_ASYNC_BODY_BYTES: usize = 16 * 1024 * 1024;
#[cfg(feature = "handlers")]
#[derive(serde::Deserialize)]
pub(super) struct InvokeQuery {
#[serde(default)]
mode: Option<String>,
#[serde(default)]
version: Option<String>,
}
#[cfg(feature = "handlers")]
pub(super) fn new_invocation_id() -> String {
let mut bytes = [0u8; 16];
if getrandom::getrandom(&mut bytes).is_err() {
tracing::error!("getrandom failed generating an invocation id");
}
hex::encode(bytes)
}
#[cfg(feature = "handlers")]
pub(super) async fn invoke_function(
State(deploy): State<DeployStore>,
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path(name): Path<String>,
axum::extract::Query(query): axum::extract::Query<InvokeQuery>,
request: Request,
) -> Response {
let Some(inner) = handlers.inner.as_ref() else {
return handler_unavailable();
};
let function = match deploy.get_function(project.as_ref(), &name).await {
Ok(Some(f)) => f,
Ok(None) => {
return (StatusCode::NOT_FOUND, format!("no function {name:?}\n")).into_response()
}
Err(err) => return deploy_error_response(err),
};
let reference = query.version.as_deref().unwrap_or(&function.active);
let Some(component) = function.resolve(reference).map(str::to_owned) else {
return (
StatusCode::NOT_FOUND,
format!("no version {reference:?} in function {name:?}\n"),
)
.into_response();
};
let is_async = query.mode.as_deref() == Some("async");
let idem_key = request
.headers()
.get("idempotency-key")
.and_then(|v| v.to_str().ok())
.map(str::to_owned);
if let Some(key) = &idem_key {
match deploy.get_idempotency(project.as_ref(), &name, key).await {
Ok(Some(id)) => {
if let Ok(Some(inv)) = deploy.get_invocation(project.as_ref(), &name, &id).await {
return replay_invocation(&inv);
}
}
Ok(None) => {}
Err(err) => return deploy_error_response(err),
}
}
if let Err(response) = admit_by_quota(inner, &deploy, project.as_ref(), &function).await {
return response;
}
if is_async {
enqueue_invocation(
&deploy,
project.as_ref(),
&function,
&component,
request,
idem_key,
)
.await
} else {
execute_sync(
inner,
&deploy,
project.as_ref(),
&function,
&component,
request,
idem_key,
)
.await
}
}
#[cfg(feature = "handlers")]
async fn execute_sync(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
component: &str,
request: Request,
idem_key: Option<String>,
) -> Response {
let (response, duration_ms) =
execute_function(inner, deploy, project, function, component, request, 0).await;
let Some(key) = idem_key else {
let sample = boatramp_core::function::MeteringSample {
success: response.status().as_u16() < 500,
duration_ms,
bytes_in: 0,
bytes_out: 0,
};
record_metering(inner, deploy, project, &function.name, &sample).await;
return response;
};
let (status, content_type, body) = capture_response(response).await;
let sample = boatramp_core::function::MeteringSample {
success: status.as_u16() < 500,
duration_ms,
bytes_in: 0,
bytes_out: body.len() as u64,
};
record_metering(inner, deploy, project, &function.name, &sample).await;
let now = now_unix();
let id = new_invocation_id();
let inv = boatramp_core::function::Invocation {
id: id.clone(),
function: function.name.clone(),
version: component.to_string(),
mode: boatramp_core::function::InvokeMode::Sync,
status: boatramp_core::function::InvocationStatus::Succeeded,
idempotency_key: Some(key.clone()),
attempts: 1,
request_b64: None,
request_content_type: None,
result: Some(boatramp_core::function::InvocationResult {
status: status.as_u16(),
content_type: content_type.clone(),
body_b64: b64_encode(&body),
}),
created: now,
updated: now,
};
if let Err(err) = deploy.put_invocation(project, &inv).await {
return deploy_error_response(err);
}
if let Err(err) = deploy
.put_idempotency(project, &function.name, &key, &id)
.await
{
return deploy_error_response(err);
}
rebuild_response(status, content_type.as_deref(), body)
}
#[cfg(feature = "handlers")]
async fn enqueue_invocation(
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
component: &str,
request: Request,
idem_key: Option<String>,
) -> Response {
let content_type = request
.headers()
.get(header::CONTENT_TYPE)
.and_then(|v| v.to_str().ok())
.map(str::to_owned);
let body = match axum::body::to_bytes(request.into_body(), MAX_ASYNC_BODY_BYTES).await {
Ok(bytes) => bytes,
Err(_) => {
return (
StatusCode::PAYLOAD_TOO_LARGE,
"async invoke body exceeds the buffer cap\n",
)
.into_response()
}
};
let now = now_unix();
let id = new_invocation_id();
let inv = boatramp_core::function::Invocation {
id: id.clone(),
function: function.name.clone(),
version: component.to_string(),
mode: boatramp_core::function::InvokeMode::Async,
status: boatramp_core::function::InvocationStatus::Queued,
idempotency_key: idem_key.clone(),
attempts: 0,
request_b64: (!body.is_empty()).then(|| b64_encode(&body)),
request_content_type: content_type,
result: None,
created: now,
updated: now,
};
if let Err(err) = deploy.put_invocation(project, &inv).await {
return deploy_error_response(err);
}
if let Some(key) = &idem_key {
if let Err(err) = deploy
.put_idempotency(project, &function.name, key, &id)
.await
{
return deploy_error_response(err);
}
}
(StatusCode::ACCEPTED, Json(inv)).into_response()
}
#[cfg(feature = "handlers")]
pub(super) async fn get_invocation_record(
State(deploy): State<DeployStore>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path((name, id)): Path<(String, String)>,
) -> Response {
match deploy.get_invocation(project.as_ref(), &name, &id).await {
Ok(Some(inv)) => Json(inv).into_response(),
Ok(None) => (
StatusCode::NOT_FOUND,
format!("no invocation {id:?} for function {name:?}\n"),
)
.into_response(),
Err(err) => deploy_error_response(err),
}
}
#[cfg(feature = "handlers")]
fn replay_invocation(inv: &boatramp_core::function::Invocation) -> Response {
match &inv.result {
Some(result) => {
let body = b64_decode(&result.body_b64);
rebuild_response(
StatusCode::from_u16(result.status).unwrap_or(StatusCode::OK),
result.content_type.as_deref(),
body,
)
}
None => (StatusCode::ACCEPTED, Json(inv.clone())).into_response(),
}
}
#[cfg(feature = "handlers")]
pub(super) async fn execute_function(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
component: &str,
request: Request,
depth: u32,
) -> (Response, u64) {
let permit_key = project.qualified(&function.name);
let _permit = match acquire_function_permit(inner, &permit_key, &function.config.quota) {
Ok(permit) => permit,
Err(()) => {
return (
(
StatusCode::SERVICE_UNAVAILABLE,
"function concurrency limit reached\n",
)
.into_response(),
0,
)
}
};
let wasm = match read_blob_fully(deploy, component).await {
Ok(bytes) => bytes,
Err(response) => return (response, 0),
};
let fn_ident = format!("fn/{}", function.name);
let scope = project.qualified(&fn_ident);
let bindings =
build_function_bindings(inner, project, &scope, &fn_ident, &function.config, depth).await;
let limits = function_limits(function.config.limits.as_ref());
let request = prepare_invoke_request(request);
let start = std::time::Instant::now();
let result = inner
.engine
.serve_with_limits(component, &wasm, request, bindings, limits)
.await;
let elapsed = start.elapsed();
inner.metrics.observe(
&function.name,
metrics::Trigger::Invoke,
"invoke",
component,
metrics::Outcome::from_result(&result),
elapsed,
);
let response = 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!(function = %function.name, %err, "function invocation failed");
handler_error_response(&err)
}
};
(response, elapsed.as_millis() as u64)
}
#[cfg(feature = "handlers")]
async fn build_function_bindings(
inner: &HandlerRuntimeInner,
project: ProjectRef<'_>,
scope: &str,
sql_site: &str,
config: &boatramp_core::function::FunctionConfig,
depth: u32,
) -> boatramp_handlers::Bindings {
let granted = |name: &str| config.imports.iter().any(|i| i == 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 let Some(provider) = &inner.sql {
let mut names: Vec<&str> = Vec::new();
if granted("sql") {
names.push(""); }
for imp in &config.imports {
if let Some(name) = imp.strip_prefix("sql:") {
if !name.is_empty() && name != "*" {
names.push(name);
}
}
}
for name in names {
match provider.database(project.as_str(), sql_site, name).await {
Ok(backend) => bindings = bindings.with_sql(name, backend),
Err(err) => {
tracing::warn!(scope, database = name, %err, "opening function SQL database failed");
}
}
}
}
if granted("wasi:messaging") {
if let Some(messaging) = &inner.messaging {
bindings = bindings.with_messaging(format!("{scope}/"), messaging.clone());
}
}
if granted("invoke") && !config.invoke_targets.is_empty() {
if let Some(invoker) = inner.invoker.get() {
bindings = bindings.with_invoke(
invoker.scoped(project),
config.invoke_targets.clone(),
depth,
);
}
}
if granted("graphql") {
if let Some(runner) = inner.federation_runner.get() {
bindings = bindings.with_graphql(runner.scoped(project), depth);
}
}
inner.logs.configure(scope, None);
bindings = bindings.with_logging(scope.to_string(), None, inner.logs.clone());
let env: Vec<(String, String)> = config
.env
.iter()
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
bindings.with_env(env)
}
#[cfg(feature = "handlers")]
fn function_limits(
limits: Option<&boatramp_core::config::HandlerLimits>,
) -> boatramp_handlers::Limits {
let mut l = boatramp_handlers::Limits::default();
if let Some(hl) = limits {
if let Some(mb) = hl.memory_mb {
l.memory_bytes = (mb as usize).saturating_mul(1024 * 1024);
}
if let Some(ms) = hl.timeout_ms {
l.timeout_ms = ms as u64;
}
if let Some(fuel) = hl.fuel {
l.fuel = Some(fuel);
}
}
l
}
#[cfg(feature = "handlers")]
fn prepare_invoke_request(mut request: Request) -> Request {
let already_internal = request
.uri()
.authority()
.is_some_and(|a| a.host() == INVOKE_AUTHORITY);
if !already_internal {
if let Ok(uri) = format!("http://{INVOKE_AUTHORITY}/").parse() {
*request.uri_mut() = uri;
}
}
request
.headers_mut()
.insert(header::HOST, HeaderValue::from_static(INVOKE_AUTHORITY));
request
}
#[cfg(feature = "handlers")]
pub(crate) struct FunctionInvoker {
deploy: DeployStore,
runtime: std::sync::Weak<HandlerRuntimeInner>,
project: String,
}
#[cfg(feature = "handlers")]
impl FunctionInvoker {
pub(crate) fn new(deploy: DeployStore, runtime: std::sync::Weak<HandlerRuntimeInner>) -> Self {
Self {
deploy,
runtime,
project: ProjectRef::DEFAULT.as_str().to_string(),
}
}
pub(crate) fn scoped(&self, project: ProjectRef<'_>) -> Arc<dyn boatramp_handlers::Invoker> {
Arc::new(Self {
deploy: self.deploy.clone(),
runtime: self.runtime.clone(),
project: project.as_str().to_string(),
})
}
}
#[cfg(feature = "handlers")]
#[async_trait::async_trait]
impl boatramp_handlers::Invoker for FunctionInvoker {
async fn invoke(
&self,
target: &str,
request: boatramp_handlers::InvokeRequest,
depth: u32,
) -> Result<boatramp_handlers::InvokeResponse, boatramp_handlers::InvokeError> {
use boatramp_handlers::InvokeError;
let Some(inner) = self.runtime.upgrade() else {
return Err(InvokeError::Failed(
"handler runtime is shutting down".into(),
));
};
let project = ProjectRef::new(&self.project);
let function = match self.deploy.get_function(project, target).await {
Ok(Some(f)) => f,
Ok(None) => return Err(InvokeError::NotFound),
Err(err) => return Err(InvokeError::Failed(err.to_string())),
};
let Some(component) = function.resolve(&function.active).map(str::to_owned) else {
return Err(InvokeError::NotFound);
};
let bytes_in = request.body.len() as u64;
let axum_request = match build_internal_request(request) {
Ok(req) => req,
Err(err) => return Err(InvokeError::Failed(err)),
};
if let Err(response) = admit_by_quota(&inner, &self.deploy, project, &function).await {
return Ok(buffer_invoke_response(response).await);
}
let (response, duration_ms) = execute_function(
&inner,
&self.deploy,
project,
&function,
&component,
axum_request,
depth,
)
.await;
let invoke_response = buffer_invoke_response(response).await;
let sample = boatramp_core::function::MeteringSample {
success: invoke_response.status < 500,
duration_ms,
bytes_in,
bytes_out: invoke_response.body.len() as u64,
};
record_metering(&inner, &self.deploy, project, &function.name, &sample).await;
Ok(invoke_response)
}
async fn invoke_streaming(
&self,
target: &str,
request: boatramp_handlers::InvokeRequest,
depth: u32,
) -> Result<boatramp_handlers::InvokeStreamResponse, boatramp_handlers::InvokeError> {
use boatramp_handlers::InvokeError;
let Some(inner) = self.runtime.upgrade() else {
return Err(InvokeError::Failed(
"handler runtime is shutting down".into(),
));
};
let project = ProjectRef::new(&self.project);
let function = match self.deploy.get_function(project, target).await {
Ok(Some(f)) => f,
Ok(None) => return Err(InvokeError::NotFound),
Err(err) => return Err(InvokeError::Failed(err.to_string())),
};
let Some(component) = function.resolve(&function.active).map(str::to_owned) else {
return Err(InvokeError::NotFound);
};
let bytes_in = request.body.len() as u64;
let axum_request = match build_internal_request(request) {
Ok(req) => req,
Err(err) => return Err(InvokeError::Failed(err)),
};
if let Err(response) = admit_by_quota(&inner, &self.deploy, project, &function).await {
return Ok(stream_invoke_response(response));
}
let (response, duration_ms) = execute_function(
&inner,
&self.deploy,
project,
&function,
&component,
axum_request,
depth,
)
.await;
let stream_response = stream_invoke_response(response);
let bytes_out = stream_response
.headers
.iter()
.find(|(n, _)| n.eq_ignore_ascii_case("content-length"))
.and_then(|(_, v)| std::str::from_utf8(v).ok())
.and_then(|s| s.trim().parse::<u64>().ok())
.unwrap_or(0);
let sample = boatramp_core::function::MeteringSample {
success: stream_response.status < 500,
duration_ms,
bytes_in,
bytes_out,
};
record_metering(&inner, &self.deploy, project, &function.name, &sample).await;
Ok(stream_response)
}
}
#[cfg(feature = "handlers")]
fn build_internal_request(request: boatramp_handlers::InvokeRequest) -> Result<Request, String> {
let path = if request.path.starts_with('/') {
request.path.clone()
} else {
format!("/{}", request.path)
};
let mut builder = axum::http::Request::builder()
.method(request.method.as_str())
.uri(format!("http://{INVOKE_AUTHORITY}{path}"));
for (name, value) in &request.headers {
if name.eq_ignore_ascii_case("host") {
continue;
}
builder = builder.header(name.as_str(), value.as_slice());
}
builder
.body(axum::body::Body::from(request.body))
.map_err(|err| err.to_string())
}
#[cfg(feature = "handlers")]
async fn buffer_invoke_response(response: Response) -> boatramp_handlers::InvokeResponse {
let status = response.status().as_u16();
let headers: Vec<(String, Vec<u8>)> = response
.headers()
.iter()
.map(|(name, value)| (name.as_str().to_string(), value.as_bytes().to_vec()))
.collect();
let body = axum::body::to_bytes(response.into_body(), MAX_ASYNC_BODY_BYTES)
.await
.map(|b| b.to_vec())
.unwrap_or_default();
boatramp_handlers::InvokeResponse {
status,
headers,
body,
}
}
#[cfg(feature = "handlers")]
#[derive(Debug)]
pub(super) enum SubgraphSdlError {
Unavailable,
InvokeFailed(String),
NotASubgraph,
}
#[cfg(feature = "handlers")]
pub(super) async fn introspect_service_sdl(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
component: &str,
) -> Result<String, SubgraphSdlError> {
let body = serde_json::json!({ "query": "{ _service { sdl } }" })
.to_string()
.into_bytes();
let invoke = boatramp_handlers::InvokeRequest {
method: "POST".to_string(),
path: "/".to_string(),
headers: vec![("content-type".to_string(), b"application/json".to_vec())],
body,
};
let request = build_internal_request(invoke).map_err(SubgraphSdlError::InvokeFailed)?;
let run = tokio::time::timeout(
std::time::Duration::from_secs(10),
execute_function(inner, deploy, project, function, component, request, 0),
)
.await;
let (response, _ms) = match run {
Ok(pair) => pair,
Err(_elapsed) => {
return Err(SubgraphSdlError::InvokeFailed(
"timed out answering `_service { sdl }`".to_string(),
))
}
};
let buffered = buffer_invoke_response(response).await;
if buffered.status >= 500 {
return Err(SubgraphSdlError::InvokeFailed(format!(
"status {}",
buffered.status
)));
}
let parsed: serde_json::Value =
serde_json::from_slice(&buffered.body).unwrap_or(serde_json::Value::Null);
match parsed
.pointer("/data/_service/sdl")
.and_then(|v| v.as_str())
{
Some(sdl) if !sdl.trim().is_empty() => Ok(sdl.to_string()),
_ => Err(SubgraphSdlError::NotASubgraph),
}
}
#[cfg(feature = "handlers")]
fn stream_invoke_response(response: Response) -> boatramp_handlers::InvokeStreamResponse {
use futures::StreamExt as _;
let status = response.status().as_u16();
let headers: Vec<(String, Vec<u8>)> = response
.headers()
.iter()
.map(|(name, value)| (name.as_str().to_string(), value.as_bytes().to_vec()))
.collect();
let body = response
.into_body()
.into_data_stream()
.map(|chunk| chunk.map_err(|e| e.to_string()))
.boxed();
boatramp_handlers::InvokeStreamResponse {
status,
headers,
body,
}
}
#[cfg(feature = "handlers")]
pub(super) async fn capture_response(response: Response) -> (StatusCode, Option<String>, Vec<u8>) {
let status = response.status();
let content_type = response
.headers()
.get(header::CONTENT_TYPE)
.and_then(|v| v.to_str().ok())
.map(str::to_owned);
let body = axum::body::to_bytes(response.into_body(), usize::MAX)
.await
.map(|b| b.to_vec())
.unwrap_or_default();
(status, content_type, body)
}
#[cfg(feature = "handlers")]
fn rebuild_response(status: StatusCode, content_type: Option<&str>, body: Vec<u8>) -> Response {
let mut builder = axum::http::Response::builder().status(status);
if let Some(ct) = content_type {
if let Ok(value) = HeaderValue::from_str(ct) {
builder = builder.header(header::CONTENT_TYPE, value);
}
}
builder
.body(axum::body::Body::from(body))
.unwrap_or_else(|_| handler_unavailable())
}
#[cfg(feature = "handlers")]
pub(super) fn b64_encode(bytes: &[u8]) -> String {
use base64::Engine;
base64::engine::general_purpose::STANDARD.encode(bytes)
}
#[cfg(feature = "handlers")]
pub(super) fn b64_decode(s: &str) -> Vec<u8> {
use base64::Engine;
base64::engine::general_purpose::STANDARD
.decode(s)
.unwrap_or_default()
}
#[cfg(feature = "handlers")]
pub(super) async fn drain_function_invocations(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
) {
let queued = match deploy.list_invocations(project, &function.name).await {
Ok(list) => list,
Err(err) => {
tracing::warn!(function = %function.name, %err, "listing invocations failed");
return;
}
};
for inv in queued {
if !matches!(
inv.status,
boatramp_core::function::InvocationStatus::Queued
) {
continue;
}
run_queued_invocation(inner, deploy, project, function, inv).await;
}
}
#[cfg(feature = "handlers")]
async fn run_queued_invocation(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
mut inv: boatramp_core::function::Invocation,
) {
use boatramp_core::function::InvocationStatus;
let Some(component) = function.resolve(&inv.version).map(str::to_owned) else {
inv.status = InvocationStatus::Failed;
inv.updated = now_unix();
let _ = deploy.put_invocation(project, &inv).await;
return;
};
inv.status = InvocationStatus::Running;
inv.attempts = inv.attempts.saturating_add(1);
inv.updated = now_unix();
if let Err(err) = deploy.put_invocation(project, &inv).await {
tracing::warn!(function = %function.name, %err, "marking invocation running failed");
return;
}
let bytes_in = inv
.request_b64
.as_deref()
.map(|b| b64_decode(b).len() as u64)
.unwrap_or(0);
let request = build_stored_request(&inv);
let (response, duration_ms) =
execute_function(inner, deploy, project, function, &component, request, 0).await;
let (status, content_type, body) = capture_response(response).await;
let delivered = status != StatusCode::INTERNAL_SERVER_ERROR
&& status != StatusCode::GATEWAY_TIMEOUT
&& status != StatusCode::SERVICE_UNAVAILABLE;
if delivered {
inv.status = InvocationStatus::Succeeded;
inv.result = Some(boatramp_core::function::InvocationResult {
status: status.as_u16(),
content_type,
body_b64: b64_encode(&body),
});
} else if inv.attempts >= MAX_INVOKE_ATTEMPTS {
inv.status = InvocationStatus::Failed;
} else {
inv.status = InvocationStatus::Queued;
}
inv.updated = now_unix();
let _ = deploy.put_invocation(project, &inv).await;
if matches!(
inv.status,
InvocationStatus::Succeeded | InvocationStatus::Failed
) {
let sample = boatramp_core::function::MeteringSample {
success: matches!(inv.status, InvocationStatus::Succeeded),
duration_ms,
bytes_in,
bytes_out: body.len() as u64,
};
record_metering(inner, deploy, project, &function.name, &sample).await;
}
}
#[cfg(feature = "handlers")]
fn build_stored_request(inv: &boatramp_core::function::Invocation) -> Request {
let body = inv
.request_b64
.as_deref()
.map(b64_decode)
.unwrap_or_default();
let mut builder = axum::http::Request::builder()
.method(axum::http::Method::POST)
.uri(format!("http://{INVOKE_AUTHORITY}/"))
.header(header::HOST, INVOKE_AUTHORITY);
if let Some(ct) = &inv.request_content_type {
if let Ok(value) = HeaderValue::from_str(ct) {
builder = builder.header(header::CONTENT_TYPE, value);
}
}
builder
.body(axum::body::Body::from(body))
.unwrap_or_else(|_| Request::new(axum::body::Body::empty()))
}
#[cfg(feature = "handlers")]
fn function_meter_lock(inner: &HandlerRuntimeInner, name: &str) -> Arc<tokio::sync::Mutex<()>> {
inner
.function_meter_locks
.lock()
.unwrap()
.entry(name.to_string())
.or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
.clone()
}
#[cfg(feature = "handlers")]
fn acquire_function_permit(
inner: &HandlerRuntimeInner,
name: &str,
quota: &boatramp_core::function::FunctionQuota,
) -> Result<Option<tokio::sync::OwnedSemaphorePermit>, ()> {
let Some(max) = quota.max_concurrent else {
return Ok(None);
};
let semaphore = {
let mut map = inner.function_semaphores.lock().unwrap();
map.entry(name.to_string())
.or_insert_with(|| Arc::new(tokio::sync::Semaphore::new(max as usize)))
.clone()
};
semaphore.try_acquire_owned().map(Some).map_err(|_| ())
}
#[cfg(feature = "handlers")]
async fn admit_by_quota(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
) -> Result<(), Response> {
let quota = &function.config.quota;
if quota.max_invocations.is_none() {
return Ok(());
}
let lock = function_meter_lock(inner, &project.qualified(&function.name));
let _guard = lock.lock().await;
let now = now_unix();
let mut metering = match deploy.get_metering(project, &function.name).await {
Ok(Some(m)) => m,
Ok(None) => boatramp_core::function::Metering::new(&function.name),
Err(err) => return Err(deploy_error_response(err)),
};
if !metering.admit(quota, now) {
return Err((
StatusCode::TOO_MANY_REQUESTS,
"function invocation quota exceeded\n",
)
.into_response());
}
if let Err(err) = deploy.put_metering(project, &metering).await {
return Err(deploy_error_response(err));
}
Ok(())
}
#[cfg(feature = "handlers")]
async fn record_metering(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &str,
sample: &boatramp_core::function::MeteringSample,
) {
let lock = function_meter_lock(inner, &project.qualified(function));
let _guard = lock.lock().await;
let now = now_unix();
let mut metering = match deploy.get_metering(project, function).await {
Ok(Some(m)) => m,
Ok(None) => boatramp_core::function::Metering::new(function),
Err(err) => {
tracing::warn!(function, %err, "reading metering failed");
return;
}
};
metering.record(sample, now);
if let Err(err) = deploy.put_metering(project, &metering).await {
tracing::warn!(function, %err, "writing metering failed");
}
}
#[cfg(feature = "handlers")]
pub(super) async fn get_function_usage(
State(deploy): State<DeployStore>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path(name): Path<String>,
) -> Response {
match deploy.get_metering(project.as_ref(), &name).await {
Ok(Some(m)) => Json(m).into_response(),
Ok(None) => Json(boatramp_core::function::Metering::new(name)).into_response(),
Err(err) => deploy_error_response(err),
}
}
#[cfg(feature = "handlers")]
pub(super) async fn put_trigger_handler(
State(deploy): State<DeployStore>,
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path((name, id)): Path<(String, String)>,
Json(kind): Json<boatramp_core::function::TriggerKind>,
) -> Response {
match deploy.get_function(project.as_ref(), &name).await {
Ok(Some(_)) => {}
Ok(None) => {
return (StatusCode::NOT_FOUND, format!("no function {name:?}\n")).into_response()
}
Err(err) => return deploy_error_response(err),
}
if let boatramp_core::function::TriggerKind::Blob { prefix } = &kind {
let Some(inner) = handlers.inner.as_ref() else {
return (
StatusCode::BAD_REQUEST,
"this storage backend does not support blob-change triggers\n",
)
.into_response();
};
if !inner.storage.supports_watch() {
return (
StatusCode::BAD_REQUEST,
"this storage backend does not support blob-change triggers\n",
)
.into_response();
}
if let Some(provider) = inner.watch_provider.get() {
let storage_prefix = blob_storage_prefix(project.as_ref(), &name, prefix);
let tier = inner.provision_tier.get().copied().unwrap_or_default();
match boatramp_core::blob_provision::ensure_watch(
provider.as_ref(),
tier,
&name,
&storage_prefix,
&deploy,
now_unix(),
)
.await
{
Ok(boatramp_core::blob_provision::ProvisionOutcome::Ready) => {}
Ok(boatramp_core::blob_provision::ProvisionOutcome::Recipe(recipe)) => {
return (StatusCode::BAD_REQUEST, format!("{recipe}\n")).into_response();
}
Ok(boatramp_core::blob_provision::ProvisionOutcome::Refused(msg)) => {
return (StatusCode::BAD_REQUEST, format!("{msg}\n")).into_response();
}
Err(err) => {
return (StatusCode::BAD_GATEWAY, format!("{err}\n")).into_response();
}
}
}
}
let trigger = boatramp_core::function::FunctionTrigger {
id: id.clone(),
kind,
last_fired_minute: None,
};
if let Err(err) = deploy.put_trigger(project.as_ref(), &name, &trigger).await {
return deploy_error_response(err);
}
Json(trigger).into_response()
}
#[cfg(feature = "handlers")]
pub(super) async fn list_triggers_handler(
State(deploy): State<DeployStore>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path(name): Path<String>,
) -> Response {
match deploy.list_triggers(project.as_ref(), &name).await {
Ok(mut list) => {
list.sort_by(|a, b| a.id.cmp(&b.id));
Json(list).into_response()
}
Err(err) => deploy_error_response(err),
}
}
#[cfg(feature = "handlers")]
pub(super) async fn delete_trigger_handler(
State(deploy): State<DeployStore>,
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path((name, id)): Path<(String, String)>,
) -> Response {
if let (Some(inner), Ok(Some(trigger))) = (
handlers.inner.as_ref(),
deploy.get_trigger(project.as_ref(), &name, &id).await,
) {
if let (boatramp_core::function::TriggerKind::Blob { prefix }, Some(provider)) =
(&trigger.kind, inner.watch_provider.get())
{
let storage_prefix = blob_storage_prefix(project.as_ref(), &name, prefix);
if let Ok(Some(record)) = deploy
.get_managed_notification(project.as_ref(), &name, &storage_prefix)
.await
{
if let Err(err) = boatramp_core::blob_provision::retract_watch(
provider.as_ref(),
&record,
&deploy,
)
.await
{
tracing::warn!(function = %name, %err, "retracting blob notification failed");
}
}
}
}
match deploy.delete_trigger(project.as_ref(), &name, &id).await {
Ok(_) => StatusCode::NO_CONTENT.into_response(),
Err(err) => deploy_error_response(err),
}
}
#[cfg(feature = "handlers")]
pub(super) fn blob_storage_prefix(
project: ProjectRef<'_>,
function: &str,
trigger_prefix: &str,
) -> String {
format!(
"hblob/{}/{trigger_prefix}",
project.qualified(&format!("fn/{function}"))
)
}
#[cfg(feature = "handlers")]
pub(super) async fn dispatch_function_triggers(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
now: &CronNow,
) {
use boatramp_core::function::TriggerKind;
let triggers = match deploy.list_triggers(project, &function.name).await {
Ok(t) => t,
Err(err) => {
tracing::warn!(function = %function.name, %err, "listing triggers failed");
return;
}
};
for mut trigger in triggers {
match &trigger.kind {
TriggerKind::Cron { schedule, .. } => {
let Ok(parsed) = boatramp_core::cron::CronSchedule::parse(schedule) else {
continue;
};
if !parsed.fires_at(now.minute, now.hour, now.dom, now.month, now.dow) {
continue;
}
if trigger.last_fired_minute == Some(now.minute_stamp) {
continue; }
enqueue_scheduled_invocation(deploy, project, function).await;
trigger.last_fired_minute = Some(now.minute_stamp);
let _ = deploy.put_trigger(project, &function.name, &trigger).await;
}
TriggerKind::Queue { topic } => {
dispatch_function_queue(inner, deploy, project, function, topic).await;
}
_ => {}
}
}
}
#[cfg(feature = "handlers")]
async fn enqueue_scheduled_invocation(
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
) {
let now = now_unix();
let inv = boatramp_core::function::Invocation {
id: new_invocation_id(),
function: function.name.clone(),
version: function.active.clone(),
mode: boatramp_core::function::InvokeMode::Async,
status: boatramp_core::function::InvocationStatus::Queued,
idempotency_key: None,
attempts: 0,
request_b64: None,
request_content_type: None,
result: None,
created: now,
updated: now,
};
if let Err(err) = deploy.put_invocation(project, &inv).await {
tracing::warn!(function = %function.name, %err, "enqueuing scheduled invocation failed");
}
}
#[cfg(feature = "handlers")]
async fn dispatch_function_queue(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
function: &boatramp_core::function::Function,
topic: &str,
) {
let Some(messaging) = inner.messaging.clone() else {
return;
};
let namespaced = project.qualified(&format!("fn/{}/{topic}", function.name));
let batch = match messaging
.claim(
&namespaced,
CONSUMER_LEASE,
CONSUMER_BATCH,
CONSUMER_MAX_ATTEMPTS,
)
.await
{
Ok(batch) => batch,
Err(err) => {
tracing::warn!(function = %function.name, topic, %err, "claiming queue messages failed");
return;
}
};
if batch.is_empty() {
return;
}
let Some(component) = function.resolve(&function.active).map(str::to_owned) else {
return;
};
for msg in batch {
let bytes_in = msg.payload.len() as u64;
let request = build_webhook_request(None, msg.payload.clone());
let (response, duration_ms) =
execute_function(inner, deploy, project, function, &component, request, 0).await;
let (status, _content_type, body) = capture_response(response).await;
let delivered = status != StatusCode::INTERNAL_SERVER_ERROR
&& status != StatusCode::GATEWAY_TIMEOUT
&& status != StatusCode::SERVICE_UNAVAILABLE;
let sample = boatramp_core::function::MeteringSample {
success: delivered,
duration_ms,
bytes_in,
bytes_out: body.len() as u64,
};
record_metering(inner, deploy, project, &function.name, &sample).await;
if delivered {
let _ = messaging.ack(&msg).await;
} else {
let _ = messaging.nack(&msg).await;
}
}
}
#[cfg(feature = "handlers")]
pub(super) async fn webhook_ingress(
State(deploy): State<DeployStore>,
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path(name): Path<String>,
request: Request,
) -> Response {
let Some(inner) = handlers.inner.as_ref() else {
return handler_unavailable();
};
let function = match deploy.get_function(project.as_ref(), &name).await {
Ok(Some(f)) => f,
Ok(None) => {
return (StatusCode::NOT_FOUND, format!("no function {name:?}\n")).into_response()
}
Err(err) => return deploy_error_response(err),
};
let Some(webhook) = function.config.webhook.clone() else {
return (
StatusCode::NOT_FOUND,
format!("function {name:?} has no webhook\n"),
)
.into_response();
};
let Ok(secret) = std::env::var(&webhook.secret_env) else {
tracing::warn!(
function = %name,
env = %webhook.secret_env,
"webhook secret env var is not set; refusing",
);
return (StatusCode::SERVICE_UNAVAILABLE, "webhook not configured\n").into_response();
};
let provided = request
.headers()
.get(webhook.header())
.and_then(|v| v.to_str().ok())
.map(str::to_owned);
let content_type = request
.headers()
.get(header::CONTENT_TYPE)
.and_then(|v| v.to_str().ok())
.map(str::to_owned);
let body = match axum::body::to_bytes(request.into_body(), webhook.body_cap() as usize).await {
Ok(bytes) => bytes,
Err(_) => {
return (
StatusCode::PAYLOAD_TOO_LARGE,
"webhook body exceeds the cap\n",
)
.into_response()
}
};
let Some(provided) = provided else {
return (StatusCode::UNAUTHORIZED, "missing webhook signature\n").into_response();
};
if !verify_webhook_signature(webhook.algorithm, secret.as_bytes(), &body, &provided) {
return (StatusCode::UNAUTHORIZED, "invalid webhook signature\n").into_response();
}
if let Err(response) = admit_by_quota(inner, &deploy, project.as_ref(), &function).await {
return response;
}
let Some(component) = function.resolve(&function.active).map(str::to_owned) else {
return handler_unavailable();
};
let bytes_in = body.len() as u64;
let request = build_webhook_request(content_type, body.to_vec());
let (response, duration_ms) = execute_function(
inner,
&deploy,
project.as_ref(),
&function,
&component,
request,
0,
)
.await;
let sample = boatramp_core::function::MeteringSample {
success: response.status().as_u16() < 500,
duration_ms,
bytes_in,
bytes_out: 0,
};
record_metering(inner, &deploy, project.as_ref(), &function.name, &sample).await;
response
}
#[cfg(feature = "handlers")]
fn verify_webhook_signature(
algorithm: boatramp_core::function::WebhookAlgorithm,
secret: &[u8],
body: &[u8],
provided: &str,
) -> bool {
use boatramp_core::function::WebhookAlgorithm;
use hmac::{Hmac, Mac};
use subtle::ConstantTimeEq;
match algorithm {
WebhookAlgorithm::HmacSha256 => {
let provided = provided.strip_prefix("sha256=").unwrap_or(provided);
let Ok(provided_bytes) = hex::decode(provided) else {
return false;
};
let Ok(mut mac) = <Hmac<sha2::Sha256> as Mac>::new_from_slice(secret) else {
return false;
};
mac.update(body);
let expected = mac.finalize().into_bytes();
provided_bytes.ct_eq(&expected).into()
}
}
}
#[cfg(feature = "handlers")]
fn build_webhook_request(content_type: Option<String>, body: Vec<u8>) -> Request {
let mut builder = axum::http::Request::builder()
.method(axum::http::Method::POST)
.uri("/");
if let Some(ct) = &content_type {
if let Ok(value) = HeaderValue::from_str(ct) {
builder = builder.header(header::CONTENT_TYPE, value);
}
}
builder
.body(axum::body::Body::from(body))
.unwrap_or_else(|_| Request::new(axum::body::Body::empty()))
}