use axum::http::{HeaderMap, StatusCode};
use super::handler::RestError;
pub const SSE_CONTENT_TYPE: &str = "text/event-stream";
pub const DEFAULT_SSE_HEARTBEAT_SECONDS: u64 = 30;
pub const DEFAULT_SSE_MAX_REPLAY_EVENTS: u64 = 10_000;
#[must_use]
pub fn accepts_sse(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(SSE_CONTENT_TYPE))
})
}
#[must_use]
pub fn is_stream_path(relative_path: &str) -> bool {
let segments: Vec<&str> = relative_path
.trim_start_matches('/')
.split('/')
.filter(|s| !s.is_empty())
.collect();
segments.last().is_some_and(|s| *s == "stream")
}
#[must_use]
pub fn extract_stream_resource(relative_path: &str) -> Option<&str> {
let segments: Vec<&str> = relative_path
.trim_start_matches('/')
.split('/')
.filter(|s| !s.is_empty())
.collect();
if segments.len() == 2 && segments[1] == "stream" {
Some(segments[0])
} else {
None
}
}
#[must_use]
pub fn extract_last_event_id(headers: &HeaderMap) -> Option<String> {
headers.get("last-event-id").and_then(|v| v.to_str().ok()).map(String::from)
}
#[must_use]
pub fn format_sse_event(event_type: &str, event_id: &str, data: &serde_json::Value) -> String {
let data_str = serde_json::to_string(data).unwrap_or_default();
format!("event: {event_type}\nid: {event_id}\ndata: {data_str}\n\n")
}
#[must_use]
pub fn format_heartbeat() -> String {
"event: ping\ndata: \n\n".to_string()
}
#[must_use]
pub fn observers_not_available() -> RestError {
RestError {
status: StatusCode::NOT_IMPLEMENTED,
code: "NOT_IMPLEMENTED",
message: "SSE streaming requires the observers feature".to_string(),
details: None,
}
}
#[must_use]
pub fn event_kind_to_sse_type(kind: &str) -> &str {
match kind {
"INSERT" => "insert",
"UPDATE" => "update",
"DELETE" => "delete",
"CUSTOM" => "custom",
_ => "unknown",
}
}
#[cfg(feature = "observers")]
pub fn stream_tenant_scope(
security_ctx: Option<&fraiseql_core::security::SecurityContext>,
multi_tenant: bool,
) -> Result<fraiseql_observers::transport::TenantScope, RestError> {
use fraiseql_observers::transport::TenantScope;
if !multi_tenant {
return Ok(TenantScope::AllTenants);
}
security_ctx.and_then(|ctx| ctx.tenant_id.as_ref()).map_or_else(
|| {
Err(RestError {
status: StatusCode::FORBIDDEN,
code: "TENANT_SCOPE_REQUIRED",
message: "This deployment is multi-tenant and the request carries no tenant, \
so an event stream cannot be scoped to one. Present a credential \
carrying a tenant."
.to_string(),
details: None,
})
},
|tenant| Ok(TenantScope::Tenant(tenant.as_str().to_string())),
)
}
#[cfg(feature = "observers")]
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ResumeRequest {
Fresh,
From {
seq: i64,
token: String,
},
}
#[cfg(feature = "observers")]
pub fn stream_resume_request(
headers: &HeaderMap,
stream: &str,
) -> Result<ResumeRequest, RestError> {
let Some(raw) = extract_last_event_id(headers) else {
return Ok(ResumeRequest::Fresh);
};
let trimmed = raw.trim();
if trimmed.is_empty() {
return Ok(ResumeRequest::Fresh);
}
super::stream_token::open(trimmed, stream).map_or_else(
|| Err(resume_point_unknown(trimmed)),
|seq| {
Ok(ResumeRequest::From {
seq,
token: trimmed.to_string(),
})
},
)
}
#[cfg(feature = "observers")]
#[must_use]
pub fn resumption_unsupported(id: &str) -> RestError {
RestError {
status: StatusCode::NOT_IMPLEMENTED,
code: "RESUMPTION_UNSUPPORTED",
message: format!(
"Last-Event-ID {id} cannot be honoured: this deployment keeps no record of \
what this stream delivered, so the events since {id} cannot be \
established. Reconnect without the header to receive events from now on."
),
details: None,
}
}
#[cfg(feature = "observers")]
#[must_use]
pub fn resume_point_unknown(id: &str) -> RestError {
RestError {
status: StatusCode::GONE,
code: "RESUME_POINT_UNKNOWN",
message: format!(
"Last-Event-ID {id} names no event on this stream: it has aged out of the \
change log, or it was issued by a different stream. What followed it \
cannot be established, so it is refused rather than answered with a replay \
that might skip. Reconnect without the header to receive events from now on."
),
details: None,
}
}
#[cfg(feature = "observers")]
#[must_use]
pub fn resume_too_far_behind(id: &str, cap: u64) -> RestError {
RestError {
status: StatusCode::PAYLOAD_TOO_LARGE,
code: "RESUME_TOO_FAR_BEHIND",
message: format!(
"Last-Event-ID {id} is more than {cap} delivered events behind, which is \
this deployment's replay bound (`[rest].sse_max_replay_events`). Refusing \
rather than replaying part of the gap. Reconnect without the header to \
receive events from now on, or raise the bound."
),
details: None,
}
}
#[cfg(feature = "observers")]
#[must_use]
pub fn replay_scope(
entity_type: &str,
scope: &fraiseql_observers::transport::TenantScope,
) -> fraiseql_observers::listener::ReplayScope {
use fraiseql_observers::transport::TenantScope;
fraiseql_observers::listener::ReplayScope {
object_type: entity_type.to_string(),
tenant: match scope {
TenantScope::AllTenants => None,
TenantScope::Tenant(tenant) => Some(tenant.clone()),
},
}
}
#[cfg(feature = "observers")]
#[derive(Debug, PartialEq, Eq)]
pub struct StreamEvent<'a> {
pub event_type: &'static str,
pub id: Option<String>,
pub data: &'a serde_json::Value,
}
#[cfg(feature = "observers")]
impl<'a> StreamEvent<'a> {
#[must_use]
pub fn from_bridge_event(event: &'a crate::subscriptions::EntityEvent) -> Self {
Self {
event_type: event_kind_to_sse_type(operation_kind(event.operation)),
id: event.change_spine.as_ref().and_then(|env| env.seq).map(|s| s.to_string()),
data: &event.data,
}
}
}
#[cfg(feature = "observers")]
#[must_use]
pub fn operation_kind(
operation: fraiseql_core::runtime::subscription::SubscriptionOperation,
) -> &'static str {
use fraiseql_core::runtime::subscription::SubscriptionOperation as Op;
match operation {
Op::Create => "INSERT",
Op::Update => "UPDATE",
Op::Delete => "DELETE",
other => {
tracing::warn!(
operation = %other,
"REST /{{resource}}/stream received a subscription operation it has no SSE \
event name for; delivering it as `event: unknown`. Add it to \
`routes::rest::sse::operation_kind`."
);
"UNRECOGNISED"
},
}
}
#[cfg(feature = "observers")]
pub const STREAM_LAGGED_EVENT: &str = "error";
#[cfg(feature = "observers")]
#[must_use]
pub fn stream_lagged_payload(skipped: u64) -> serde_json::Value {
serde_json::json!({
"code": "STREAM_LAGGED",
"skipped": skipped,
"message": format!(
"This stream fell behind and {skipped} event(s) were dropped before they \
could be delivered. It is ending rather than resuming past the gap, which \
you could not have seen. Reconnect to receive events from now on."
),
})
}
#[cfg(feature = "observers")]
#[must_use]
pub fn stream_event_matches(
event: &crate::subscriptions::EntityEvent,
entity_type: &str,
scope: &fraiseql_observers::transport::TenantScope,
) -> bool {
use fraiseql_observers::transport::TenantScope;
if event.entity_type != entity_type {
return false;
}
match scope {
TenantScope::AllTenants => true,
TenantScope::Tenant(tenant) => event.tenant_id.as_deref() == Some(tenant.as_str()),
}
}