pub mod helpers;
#[cfg(feature = "export-csv")]
pub mod csv;
#[cfg(feature = "export-xlsx")]
pub mod xlsx;
#[cfg(test)]
mod tests;
#[cfg(all(test, feature = "export-csv", feature = "export-xlsx"))]
mod export_header_tests;
#[cfg(all(test, feature = "export-csv"))]
mod export_pagination_tests;
#[cfg(all(test, feature = "export-csv"))]
mod export_embedding_filter_tests;
use axum::http::{HeaderMap, HeaderValue};
use bytes::Bytes;
use fraiseql_core::security::SecurityContext;
use futures::{StreamExt as _, stream};
use super::handler::{ResolvedGetQuery, RestError, RestHandler, set_request_id};
pub const NDJSON_CONTENT_TYPE: &str = "application/x-ndjson";
#[must_use]
pub fn accepts_ndjson(headers: &HeaderMap) -> bool {
headers.get("accept").and_then(|v| v.to_str().ok()).is_some_and(|accept| {
accept
.split(',')
.any(|part| part.trim().eq_ignore_ascii_case(NDJSON_CONTENT_TYPE))
})
}
pub async fn handle_ndjson_get(
handler: &RestHandler<'_>,
relative_path: &str,
query_pairs: &[(&str, &str)],
headers: &HeaderMap,
security_context: Option<&SecurityContext>,
) -> Result<NdjsonResponse, RestError> {
let resolved = handler.resolve_streaming_get_query(
relative_path,
query_pairs,
headers,
security_context,
)?;
let ResolvedGetQuery {
query_match,
variables,
params,
..
} = resolved;
let batch_size = handler.config().ndjson_batch_size.max(1);
let rows = helpers::export_rows(
handler.executor(),
query_match,
variables,
security_context.cloned(),
params.requested_pagination.export_total(),
)
.await?;
let mut response_headers = HeaderMap::new();
set_request_id(headers, &mut response_headers);
response_headers.insert("content-type", HeaderValue::from_static(NDJSON_CONTENT_TYPE));
response_headers.insert(
"x-stream-batch-size",
HeaderValue::from_str(&batch_size.to_string())
.unwrap_or_else(|_| HeaderValue::from_static("500")),
);
let chunks = rows.ready_chunks(usize::try_from(batch_size).unwrap_or(usize::MAX));
let ndjson_stream = stream::unfold(Some(chunks), |state| async move {
let mut chunks = state?;
let chunk = chunks.next().await?;
let (bytes, failed) = helpers::ndjson_chunk(chunk);
Some((Ok(bytes), if failed { None } else { Some(chunks) }))
});
Ok(NdjsonResponse {
headers: response_headers,
body: NdjsonBody::Stream(Box::pin(ndjson_stream)),
})
}
pub struct NdjsonResponse {
pub headers: HeaderMap,
pub body: NdjsonBody,
}
#[non_exhaustive]
pub enum NdjsonBody {
Stream(
std::pin::Pin<
Box<dyn futures::Stream<Item = Result<Bytes, std::convert::Infallible>> + Send>,
>,
),
}
impl NdjsonBody {
pub fn into_body(self) -> axum::body::Body {
match self {
Self::Stream(stream) => axum::body::Body::from_stream(stream),
}
}
}
#[cfg(any(feature = "export-csv", feature = "export-xlsx"))]
const FORMULA_INJECTION_SENTINELS: [char; 6] = ['=', '+', '-', '@', '\t', '\r'];
#[cfg(any(feature = "export-csv", feature = "export-xlsx"))]
pub(crate) fn guard_formula_injection(value: &str) -> String {
match value.chars().next() {
Some(c) if FORMULA_INJECTION_SENTINELS.contains(&c) => {
let mut out = String::with_capacity(value.len() + 1);
out.push('\'');
out.push_str(value);
out
},
_ => value.to_owned(),
}
}