pub mod helpers;
#[cfg(test)]
mod tests;
use std::{future::Future, sync::Arc};
use axum::{
Router,
body::Body,
extract::{FromRequestParts, Request, State},
http::{StatusCode, request::Parts},
response::{IntoResponse, Response},
routing::{delete, get, patch, post, put},
};
use fraiseql_core::{runtime::Executor, security::SecurityContext};
use helpers::{
error_response, parse_query_pairs, rest_result_to_response, strip_base_path, to_axum_path,
};
use serde_json::json;
use tower_http::compression::{CompressionLayer, predicate::SizeAbove};
use tracing::info;
use super::{
handler::{RestError, RestHandler, RestResponse},
resource::{HttpMethod, MountedRoutes, RestRouteTable},
};
use crate::{extractors::OptionalSecurityContext, routes::graphql::AppState};
#[derive(Debug, Clone, Default)]
pub struct RestMountConfig {
pub compression_enabled: bool,
pub auth_layer_attached: bool,
pub export: Arc<super::export_config::ExportConfig>,
}
fn derive_rest_context(
state: &AppState,
#[cfg_attr(
not(any(feature = "export-csv", feature = "export-xlsx")),
allow(unused_variables)
)]
mount: &RestMountConfig,
) -> Option<(String, Arc<RestRouteTable>, RestState)> {
let executor = state.executor();
let schema = executor.schema();
let config = match &schema.rest_config {
Some(cfg) if cfg.enabled => cfg.clone(),
Some(_) => {
info!("REST transport disabled (rest.enabled = false)");
return None;
},
None => {
return None;
},
};
warn_if_embeds_are_refused(&executor);
let route_table = match RestRouteTable::from_compiled_schema(schema) {
Ok(rt) => Arc::new(rt),
Err(e) => {
tracing::warn!(error = %e, "REST route derivation failed — REST transport disabled");
return None;
},
};
for diag in &route_table.diagnostics {
match diag.level {
super::resource::DiagnosticLevel::Info => {
tracing::debug!(message = %diag.message, "REST derivation");
},
super::resource::DiagnosticLevel::Warning => {
tracing::warn!(message = %diag.message, "REST derivation");
},
super::resource::DiagnosticLevel::Error => {
tracing::error!(message = %diag.message, "REST derivation");
},
}
}
let base_path = config.path;
let idempotency_store = super::idempotency::create_store(config.idempotency_ttl_seconds);
let rest_state = RestState {
executor: state.executor.load_full(),
route_table: route_table.clone(),
idempotency_store,
error_sanitizer: Arc::clone(&state.error_sanitizer),
#[cfg(feature = "auth")]
identity_resolver: state.identity_resolver.clone(),
#[cfg(feature = "observers")]
event_fanout: state.entity_event_fanout.clone(),
#[cfg(feature = "observers")]
stream_replay: state.stream_replay.clone(),
#[cfg(feature = "export-xlsx")]
xlsx_semaphore: Arc::new(tokio::sync::Semaphore::new(mount.export.max_concurrent_xlsx)),
#[cfg(any(feature = "export-csv", feature = "export-xlsx"))]
export: Arc::clone(&mount.export),
};
Some((base_path, route_table, rest_state))
}
pub fn warn_if_embeds_are_refused(executor: &fraiseql_core::runtime::Executor) {
if !executor.supports_composed_reads() {
tracing::warn!(
"REST is mounted over a database adapter that cannot compose reads: `?select=` \
embeds and `rel.count` will be refused with 501 Not Implemented. Plain reads are \
unaffected. Embedding needs the PostgreSQL backend."
);
}
}
pub fn rest_query_router(state: &AppState, mount: &RestMountConfig) -> Option<Router> {
let (base_path, route_table, rest_state) = derive_rest_context(state, mount)?;
let executor = state.executor();
let schema = executor.schema();
let mounted = MountedRoutes::read_surface(&route_table);
let mut router = Router::new();
for (path, method) in mounted.iter() {
let axum_path = to_axum_path(&base_path, path);
router = if path.ends_with("/stream") {
router.route(&axum_path, get(rest_sse_handler))
} else {
debug_assert_eq!(method, HttpMethod::Get, "read surface must be GET-only");
router.route(&axum_path, get(rest_get_handler))
};
}
let openapi_path = format!("{}/openapi.json", base_path.trim_end_matches('/'));
let openapi_spec =
build_openapi_spec(schema, &route_table, mount.auth_layer_attached, &mounted);
let router = router.route(&openapi_path, get(serve_openapi(openapi_spec)));
let mut router = router.with_state(rest_state);
if mount.compression_enabled {
router = router.layer(CompressionLayer::new().compress_when(SizeAbove::new(1024)));
}
let resource_count = route_table.resources.len();
let get_route_count = mounted.len();
let paths: Vec<String> = route_table
.resources
.iter()
.map(|r| format!("{}/{}", base_path, r.name))
.collect();
info!(
resources = resource_count,
routes = get_route_count,
base_path = %base_path,
paths = ?paths,
"REST query transport enabled (read-only)"
);
Some(router)
}
pub fn rest_router(state: &AppState, mount: &RestMountConfig) -> Option<Router> {
let (base_path, route_table, rest_state) = derive_rest_context(state, mount)?;
let executor = state.executor();
let schema = executor.schema();
let mounted = MountedRoutes::write_surface(schema, &route_table);
let mut router = Router::new();
for (path, method) in mounted.iter() {
let axum_path = to_axum_path(&base_path, path);
router = match method {
HttpMethod::Get if path.ends_with("/stream") => {
router.route(&axum_path, get(rest_sse_handler))
},
HttpMethod::Get => router.route(&axum_path, get(rest_get_handler)),
HttpMethod::Post => router.route(&axum_path, post(rest_post_handler)),
HttpMethod::Put => router.route(&axum_path, put(rest_put_handler)),
HttpMethod::Patch => router.route(&axum_path, patch(rest_patch_handler)),
HttpMethod::Delete => router.route(&axum_path, delete(rest_delete_handler)),
};
}
let openapi_path = format!("{}/openapi.json", base_path.trim_end_matches('/'));
let openapi_spec =
build_openapi_spec(schema, &route_table, mount.auth_layer_attached, &mounted);
let router = router.route(&openapi_path, get(serve_openapi(openapi_spec)));
let mut router = router.with_state(rest_state);
if mount.compression_enabled {
router = router.layer(CompressionLayer::new().compress_when(SizeAbove::new(1024)));
}
let resource_count = route_table.resources.len();
let route_count = mounted.len();
let paths: Vec<String> = route_table
.resources
.iter()
.map(|r| format!("{}/{}", base_path, r.name))
.collect();
info!(
resources = resource_count,
routes = route_count,
base_path = %base_path,
paths = ?paths,
"REST transport enabled"
);
Some(router)
}
#[cfg(any(feature = "export-csv", feature = "export-xlsx"))]
fn refuse_disabled_export(
rest: &RestState,
format: super::export_config::ExportFormat,
) -> Option<Response> {
if rest.export.serves(format) {
return None;
}
Some(
Response::builder()
.status(StatusCode::NOT_ACCEPTABLE)
.header("content-type", "application/json")
.body(Body::from(format!(
r#"{{"error":{{"code":"EXPORT_FORMAT_DISABLED","message":"export format {format:?} is not enabled; see [export] export_formats"}}}}"#
)))
.unwrap_or_else(|_| {
Response::builder()
.status(StatusCode::NOT_ACCEPTABLE)
.body(Body::empty())
.expect("fallback response with empty body is infallible")
}),
)
}
fn build_openapi_spec(
schema: &fraiseql_core::schema::CompiledSchema,
route_table: &RestRouteTable,
auth_layer_attached: bool,
mounted: &MountedRoutes,
) -> Arc<serde_json::Value> {
match super::openapi::generate_openapi(schema, route_table, auth_layer_attached, mounted) {
Ok(spec) => Arc::new(spec),
Err(e) => {
tracing::warn!(error = %e, "OpenAPI spec generation failed");
Arc::new(json!({"error": "OpenAPI generation failed"}))
},
}
}
fn serve_openapi(
spec: Arc<serde_json::Value>,
) -> impl Fn(RestSecurityContext) -> std::future::Ready<axum::Json<serde_json::Value>> + Clone {
move |RestSecurityContext(_)| std::future::ready(axum::Json((*spec).clone()))
}
#[derive(Clone)]
struct RestState {
executor: Arc<Executor>,
route_table: Arc<RestRouteTable>,
idempotency_store: Arc<dyn super::idempotency::IdempotencyStore>,
error_sanitizer: Arc<crate::config::error_sanitization::ErrorSanitizer>,
#[cfg(feature = "auth")]
identity_resolver: Option<Arc<crate::identity::IdentityResolver>>,
#[cfg(any(feature = "export-csv", feature = "export-xlsx"))]
export: Arc<super::export_config::ExportConfig>,
#[cfg(feature = "observers")]
event_fanout: Option<crate::subscriptions::EntityEventFanout>,
#[cfg(feature = "observers")]
stream_replay: Option<std::sync::Arc<fraiseql_observers::listener::ChangeLogReplayReader>>,
#[cfg(feature = "export-xlsx")]
xlsx_semaphore: Arc<tokio::sync::Semaphore>,
}
struct RestSecurityContext(Option<SecurityContext>);
impl FromRequestParts<RestState> for RestSecurityContext {
type Rejection = Response;
#[allow(clippy::manual_async_fn)] fn from_request_parts(
parts: &mut Parts,
state: &RestState,
) -> impl Future<Output = Result<Self, Self::Rejection>> + Send {
let require_auth =
state.executor.schema().rest_config.as_ref().is_some_and(|c| c.require_auth);
let sanitizer = Arc::clone(&state.error_sanitizer);
#[cfg(feature = "auth")]
let resolver = state.identity_resolver.clone();
async move {
let OptionalSecurityContext(security_ctx) =
OptionalSecurityContext::from_request_parts(parts, &())
.await
.map_err(IntoResponse::into_response)?;
if require_auth && security_ctx.is_none() {
return Err(rest_result_to_response(Err(RestError::unauthenticated()), &sanitizer));
}
#[cfg(feature = "auth")]
let security_ctx = {
let mut security_ctx = security_ctx;
match crate::identity::resolve_request_identity(
resolver.as_deref(),
security_ctx.as_mut(),
)
.await
{
crate::identity::EnrichmentOutcome::Proceed => security_ctx,
crate::identity::EnrichmentOutcome::Denied => {
return Err(rest_result_to_response(
Err(RestError::forbidden()),
&sanitizer,
));
},
crate::identity::EnrichmentOutcome::Unavailable => {
return Err(rest_result_to_response(
Err(RestError::service_unavailable(
crate::identity::EnrichmentOutcome::UNAVAILABLE_MESSAGE,
)),
&sanitizer,
));
},
}
};
Ok(Self(security_ctx))
}
}
}
async fn rest_get_handler(
State(rest): State<RestState>,
RestSecurityContext(security_ctx): RestSecurityContext,
request: Request<Body>,
) -> Response {
let (parts, _body) = request.into_parts();
let relative_path = strip_base_path(&rest.route_table.base_path, parts.uri.path());
let query_string = parts.uri.query().unwrap_or("");
let query_pairs = parse_query_pairs(query_string);
let query_refs: Vec<(&str, &str)> =
query_pairs.iter().map(|(k, v)| (k.as_str(), v.as_str())).collect();
if super::streaming::accepts_ndjson(&parts.headers) {
let schema = rest.executor.schema();
let config = schema.rest_config.as_ref().expect("REST config must exist: handler is only reached via a matched REST route, which requires rest_config to be present in the schema");
let handler = RestHandler::new(&rest.executor, schema, config, &rest.route_table);
let result = super::streaming::handle_ndjson_get(
&handler,
&relative_path,
&query_refs,
&parts.headers,
security_ctx.as_ref(),
)
.await;
return match result {
Ok(ndjson) => {
let mut builder = Response::builder().status(StatusCode::OK);
for (key, value) in &ndjson.headers {
builder = builder.header(key, value);
}
builder.body(ndjson.body.into_body()).unwrap_or_else(|_| {
Response::builder()
.status(StatusCode::INTERNAL_SERVER_ERROR)
.body(Body::empty())
.expect("fallback response: Response::builder() with INTERNAL_SERVER_ERROR status and empty body is infallible")
})
},
Err(rest_err) => rest_result_to_response(Err(rest_err), &rest.error_sanitizer),
};
}
#[cfg(feature = "export-xlsx")]
if super::streaming::xlsx::accepts_xlsx(&parts.headers) {
if let Some(refusal) =
refuse_disabled_export(&rest, super::export_config::ExportFormat::Xlsx)
{
return refusal;
}
let Ok(_permit) = Arc::clone(&rest.xlsx_semaphore).try_acquire_owned() else {
return Response::builder()
.status(StatusCode::SERVICE_UNAVAILABLE)
.header("retry-after", "1")
.header("content-type", "application/json")
.body(Body::from(
r#"{"error":{"code":"XLSX_BUSY","message":"max concurrent XLSX exports reached; try again shortly"}}"#,
))
.unwrap_or_else(|_| {
Response::builder()
.status(StatusCode::SERVICE_UNAVAILABLE)
.body(Body::empty())
.expect("fallback response: Response::builder() with SERVICE_UNAVAILABLE status and empty body is infallible")
});
};
let schema = rest.executor.schema();
let config = schema.rest_config.as_ref().expect("REST config must exist: handler is only reached via a matched REST route, which requires rest_config to be present in the schema");
let handler = RestHandler::new(&rest.executor, schema, config, &rest.route_table);
let export_config = rest.export.as_ref();
let result = super::streaming::xlsx::handle_xlsx_get(
&handler,
export_config,
&relative_path,
&query_refs,
&parts.headers,
security_ctx.as_ref(),
)
.await;
return match result {
Ok(xlsx) => {
let mut builder = Response::builder().status(StatusCode::OK);
for (key, value) in &xlsx.headers {
builder = builder.header(key, value);
}
builder.body(xlsx.body.into_body()).unwrap_or_else(|_| {
Response::builder()
.status(StatusCode::INTERNAL_SERVER_ERROR)
.body(Body::empty())
.expect("fallback response: Response::builder() with INTERNAL_SERVER_ERROR status and empty body is infallible")
})
},
Err(rest_err) => rest_result_to_response(Err(rest_err), &rest.error_sanitizer),
};
}
#[cfg(feature = "export-csv")]
if super::streaming::csv::accepts_csv(&parts.headers) {
if let Some(refusal) =
refuse_disabled_export(&rest, super::export_config::ExportFormat::Csv)
{
return refusal;
}
let schema = rest.executor.schema();
let config = schema.rest_config.as_ref().expect("REST config must exist: handler is only reached via a matched REST route, which requires rest_config to be present in the schema");
let handler = RestHandler::new(&rest.executor, schema, config, &rest.route_table);
let export_config = rest.export.as_ref();
let result = super::streaming::csv::handle_csv_get(
&handler,
export_config,
&relative_path,
&query_refs,
&parts.headers,
security_ctx.as_ref(),
)
.await;
return match result {
Ok(csv) => {
let mut builder = Response::builder().status(StatusCode::OK);
for (key, value) in &csv.headers {
builder = builder.header(key, value);
}
builder.body(csv.body.into_body()).unwrap_or_else(|_| {
Response::builder()
.status(StatusCode::INTERNAL_SERVER_ERROR)
.body(Body::empty())
.expect("fallback response: Response::builder() with INTERNAL_SERVER_ERROR status and empty body is infallible")
})
},
Err(rest_err) => rest_result_to_response(Err(rest_err), &rest.error_sanitizer),
};
}
let schema = rest.executor.schema();
let config = schema.rest_config.as_ref().expect("REST config must exist: handler is only reached via a matched REST route, which requires rest_config to be present in the schema");
let handler = RestHandler::new(&rest.executor, schema, config, &rest.route_table);
let result = handler
.handle_get(&relative_path, &query_refs, &parts.headers, security_ctx.as_ref())
.await;
rest_result_to_response(result, &rest.error_sanitizer)
}
async fn rest_post_handler(
State(rest): State<RestState>,
RestSecurityContext(security_ctx): RestSecurityContext,
request: Request<Body>,
) -> Response {
let (parts, body) = request.into_parts();
let relative_path = strip_base_path(&rest.route_table.base_path, parts.uri.path());
let body_value = match read_json_body(body).await {
Ok(v) => v,
Err(resp) => return resp,
};
let schema = rest.executor.schema();
let config = schema.rest_config.as_ref().expect("REST config must exist: handler is only reached via a matched REST route, which requires rest_config to be present in the schema");
let handler = RestHandler::new(&rest.executor, schema, config, &rest.route_table)
.with_idempotency_store(&rest.idempotency_store);
let result = handler
.handle_post(&relative_path, &body_value, &parts.headers, security_ctx.as_ref())
.await;
rest_result_to_response(result, &rest.error_sanitizer)
}
async fn rest_put_handler(
State(rest): State<RestState>,
RestSecurityContext(security_ctx): RestSecurityContext,
request: Request<Body>,
) -> Response {
let (parts, body) = request.into_parts();
let relative_path = strip_base_path(&rest.route_table.base_path, parts.uri.path());
let body_value = match read_json_body(body).await {
Ok(v) => v,
Err(resp) => return resp,
};
let schema = rest.executor.schema();
let config = schema.rest_config.as_ref().expect("REST config must exist: handler is only reached via a matched REST route, which requires rest_config to be present in the schema");
let handler = RestHandler::new(&rest.executor, schema, config, &rest.route_table);
let result = handler
.handle_put(&relative_path, &body_value, &parts.headers, security_ctx.as_ref())
.await;
rest_result_to_response(result, &rest.error_sanitizer)
}
async fn rest_patch_handler(
State(rest): State<RestState>,
RestSecurityContext(security_ctx): RestSecurityContext,
request: Request<Body>,
) -> Response {
let (parts, body) = request.into_parts();
let relative_path = strip_base_path(&rest.route_table.base_path, parts.uri.path());
let query_string = parts.uri.query().unwrap_or("");
let query_pairs = parse_query_pairs(query_string);
let query_refs: Vec<(&str, &str)> =
query_pairs.iter().map(|(k, v)| (k.as_str(), v.as_str())).collect();
let body_value = match read_json_body(body).await {
Ok(v) => v,
Err(resp) => return resp,
};
let schema = rest.executor.schema();
let config = schema.rest_config.as_ref().expect("REST config must exist: handler is only reached via a matched REST route, which requires rest_config to be present in the schema");
let handler = RestHandler::new(&rest.executor, schema, config, &rest.route_table);
let result = handler
.handle_patch(
&relative_path,
&body_value,
&query_refs,
&parts.headers,
security_ctx.as_ref(),
)
.await;
rest_result_to_response(result, &rest.error_sanitizer)
}
async fn rest_delete_handler(
State(rest): State<RestState>,
RestSecurityContext(security_ctx): RestSecurityContext,
request: Request<Body>,
) -> Response {
let (parts, _body) = request.into_parts();
let relative_path = strip_base_path(&rest.route_table.base_path, parts.uri.path());
let query_string = parts.uri.query().unwrap_or("");
let query_pairs = parse_query_pairs(query_string);
let query_refs: Vec<(&str, &str)> =
query_pairs.iter().map(|(k, v)| (k.as_str(), v.as_str())).collect();
let schema = rest.executor.schema();
let config = schema.rest_config.as_ref().expect("REST config must exist: handler is only reached via a matched REST route, which requires rest_config to be present in the schema");
let handler = RestHandler::new(&rest.executor, schema, config, &rest.route_table);
let result = handler
.handle_delete(&relative_path, &query_refs, &parts.headers, security_ctx.as_ref())
.await;
rest_result_to_response(result, &rest.error_sanitizer)
}
#[cfg(feature = "observers")]
async fn resume_state(
rest: &RestState,
entity_type: &str,
tenant: &fraiseql_observers::transport::TenantScope,
resume: super::sse::ResumeRequest,
) -> Result<Option<super::resumable_stream::ResumeState>, super::handler::RestError> {
use fraiseql_observers::listener::ResumeAnchor;
let super::sse::ResumeRequest::From { seq, token } = resume else {
return Ok(None);
};
let Some(reader) = rest.stream_replay.clone() else {
return Err(super::sse::resumption_unsupported(&token));
};
let scope = super::sse::replay_scope(entity_type, tenant);
let anchor = reader.anchor(&scope, seq).await.map_err(|error| {
tracing::error!(%error, seq, "failed to resolve a stream resume point");
super::handler::RestError::internal("Could not resolve the resume point")
})?;
let ResumeAnchor::Found(origin) = anchor else {
return Err(super::sse::resume_point_unknown(&token));
};
let cap = rest
.executor
.schema()
.rest_config
.as_ref()
.map_or(super::sse::DEFAULT_SSE_MAX_REPLAY_EVENTS, |c| c.sse_max_replay_events);
if cap > 0 {
let backlog = reader.count_since(&origin, cap + 1).await.map_err(|error| {
tracing::error!(%error, seq, "failed to measure a stream resume backlog");
super::handler::RestError::internal("Could not measure the resume backlog")
})?;
let in_flight =
reader.count_in_flight(&scope, &origin, cap + 1).await.map_err(|error| {
tracing::error!(%error, seq, "failed to measure a stream in-flight tail");
super::handler::RestError::internal("Could not measure the resume backlog")
})?;
if backlog > cap || in_flight > cap {
return Err(super::sse::resume_too_far_behind(&token, cap));
}
}
Ok(Some(super::resumable_stream::ResumeState {
reader,
scope,
origin,
}))
}
async fn rest_sse_handler(
State(rest): State<RestState>,
#[cfg_attr(not(feature = "observers"), allow(unused_variables))]
RestSecurityContext(security_ctx): RestSecurityContext,
request: Request<Body>,
) -> Response {
let (parts, _body) = request.into_parts();
let relative_path = strip_base_path(&rest.route_table.base_path, parts.uri.path());
let resource_name = match super::sse::extract_stream_resource(&relative_path) {
Some(name) => name.to_string(),
None => {
return rest_result_to_response(
Err(super::handler::RestError::not_found("Stream endpoint not found")),
&rest.error_sanitizer,
);
},
};
let schema = rest.executor.schema();
let has_resource = rest.route_table.resources.iter().any(|r| r.name == resource_name);
if !has_resource {
return rest_result_to_response(
Err(super::handler::RestError::not_found(format!(
"Resource not found: {resource_name}"
))),
&rest.error_sanitizer,
);
}
let heartbeat_secs = schema
.rest_config
.as_ref()
.map_or(super::sse::DEFAULT_SSE_HEARTBEAT_SECONDS, |c| c.sse_heartbeat_seconds);
#[cfg(not(feature = "observers"))]
{
let _ = heartbeat_secs; rest_result_to_response(Err(super::sse::observers_not_available()), &rest.error_sanitizer)
}
#[cfg(feature = "observers")]
{
let heartbeat_interval = std::time::Duration::from_secs(heartbeat_secs);
if let Some(ref fanout) = rest.event_fanout {
use futures::StreamExt;
let tenant = match super::sse::stream_tenant_scope(
security_ctx.as_ref(),
schema.is_multi_tenant(),
) {
Ok(scope) => scope,
Err(refusal) => {
return rest_result_to_response(Err(refusal), &rest.error_sanitizer);
},
};
let Some(entity_type) = rest
.route_table
.resources
.iter()
.find(|r| r.name == resource_name)
.map(|r| r.type_name.clone())
else {
return rest_result_to_response(
Err(super::handler::RestError::not_found(format!(
"Resource not found: {resource_name}"
))),
&rest.error_sanitizer,
);
};
let resume = match super::sse::stream_resume_request(&parts.headers, &entity_type) {
Ok(resume) => resume,
Err(refusal) => {
return rest_result_to_response(Err(refusal), &rest.error_sanitizer);
},
};
let plan = match rest.executor.plan_type_stream(&entity_type, security_ctx.as_ref()) {
Ok(plan) => std::sync::Arc::new(plan),
Err(refusal) => {
return rest_result_to_response(
Err(super::handler::RestError::from(refusal)),
&rest.error_sanitizer,
);
},
};
let receiver = fanout.subscribe();
let resume = match resume_state(&rest, &entity_type, &tenant, resume).await {
Ok(resume) => resume,
Err(refusal) => {
return rest_result_to_response(Err(refusal), &rest.error_sanitizer);
},
};
let entity_events = super::resumable_stream::resumable_event_stream(
receiver,
entity_type,
tenant,
resume,
plan,
);
let heartbeat = futures::stream::unfold((), move |()| async move {
tokio::time::sleep(heartbeat_interval).await;
let event = axum::response::sse::Event::default().event("ping").data("");
Some((event, ()))
});
let merged = futures::stream::select(entity_events, heartbeat)
.map(Ok::<_, std::convert::Infallible>);
let sse = axum::response::sse::Sse::new(merged).keep_alive(
axum::response::sse::KeepAlive::new().interval(heartbeat_interval).text(""),
);
return axum::response::IntoResponse::into_response(sse);
}
let _ = heartbeat_interval;
rest_result_to_response(Err(super::sse::observers_not_available()), &rest.error_sanitizer)
}
}
#[allow(clippy::result_large_err)]
async fn read_json_body(body: Body) -> Result<serde_json::Value, Response> {
let Ok(bytes) = axum::body::to_bytes(body, 1_048_576).await else {
return Err(error_response(
StatusCode::PAYLOAD_TOO_LARGE,
"PAYLOAD_TOO_LARGE",
"Request body too large",
));
};
if bytes.is_empty() {
return Ok(serde_json::Value::Object(serde_json::Map::new()));
}
serde_json::from_slice(&bytes).map_err(|e| {
error_response(StatusCode::BAD_REQUEST, "INVALID_JSON", &format!("Invalid JSON body: {e}"))
})
}