#![allow(clippy::unwrap_used)]
use axum::http::{HeaderMap, HeaderValue, StatusCode};
use serde_json::json;
use super::{
cache_control::{CacheContext, apply_cache_headers},
sse::{
accepts_sse, event_kind_to_sse_type, extract_last_event_id, extract_stream_resource,
format_heartbeat, format_sse_event, is_stream_path, observers_not_available,
},
};
#[test]
fn get_public_default_ttl() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: true,
authenticated: false,
query_ttl: None,
default_ttl: 60,
cdn_max_age: None,
},
);
assert_eq!(headers.get("cache-control").unwrap().to_str().unwrap(), "public, max-age=60");
assert_eq!(headers.get("vary").unwrap().to_str().unwrap(), "Authorization, Accept, Prefer");
}
#[test]
fn get_private_with_auth() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: true,
authenticated: true,
query_ttl: None,
default_ttl: 60,
cdn_max_age: None,
},
);
assert_eq!(headers.get("cache-control").unwrap().to_str().unwrap(), "private, max-age=60");
}
#[test]
fn get_custom_ttl_from_query() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: true,
authenticated: false,
query_ttl: Some(120),
default_ttl: 60,
cdn_max_age: None,
},
);
assert_eq!(headers.get("cache-control").unwrap().to_str().unwrap(), "public, max-age=120");
}
#[test]
fn mutation_no_store() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: false,
authenticated: false,
query_ttl: None,
default_ttl: 60,
cdn_max_age: None,
},
);
assert_eq!(headers.get("cache-control").unwrap().to_str().unwrap(), "no-store");
assert!(headers.get("vary").is_none());
}
#[test]
fn mutation_no_store_with_auth() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: false,
authenticated: true,
query_ttl: None,
default_ttl: 60,
cdn_max_age: None,
},
);
assert_eq!(headers.get("cache-control").unwrap().to_str().unwrap(), "no-store");
}
#[test]
fn zero_ttl_disables_caching() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: true,
authenticated: false,
query_ttl: Some(0),
default_ttl: 60,
cdn_max_age: None,
},
);
assert_eq!(headers.get("cache-control").unwrap().to_str().unwrap(), "public, max-age=0");
}
#[test]
fn s_maxage_on_public_get() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: true,
authenticated: false,
query_ttl: None,
default_ttl: 60,
cdn_max_age: Some(300),
},
);
assert_eq!(
headers.get("cache-control").unwrap().to_str().unwrap(),
"public, max-age=60, s-maxage=300"
);
}
#[test]
fn no_s_maxage_on_private_get() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: true,
authenticated: true,
query_ttl: None,
default_ttl: 60,
cdn_max_age: Some(300),
},
);
assert_eq!(headers.get("cache-control").unwrap().to_str().unwrap(), "private, max-age=60");
}
#[test]
fn no_s_maxage_when_none() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: true,
authenticated: false,
query_ttl: None,
default_ttl: 60,
cdn_max_age: None,
},
);
assert_eq!(headers.get("cache-control").unwrap().to_str().unwrap(), "public, max-age=60");
}
#[test]
fn no_s_maxage_on_mutations() {
let mut headers = HeaderMap::new();
apply_cache_headers(
&mut headers,
&CacheContext {
is_get: false,
authenticated: false,
query_ttl: None,
default_ttl: 60,
cdn_max_age: Some(300),
},
);
assert_eq!(headers.get("cache-control").unwrap().to_str().unwrap(), "no-store");
}
#[test]
fn accepts_sse_true_for_exact_match() {
let mut headers = HeaderMap::new();
headers.insert("accept", HeaderValue::from_static("text/event-stream"));
assert!(accepts_sse(&headers));
}
#[test]
fn accepts_sse_true_in_list() {
let mut headers = HeaderMap::new();
headers.insert("accept", HeaderValue::from_static("application/json, text/event-stream"));
assert!(accepts_sse(&headers));
}
#[test]
fn accepts_sse_false_for_json() {
let mut headers = HeaderMap::new();
headers.insert("accept", HeaderValue::from_static("application/json"));
assert!(!accepts_sse(&headers));
}
#[test]
fn accepts_sse_false_when_missing() {
let headers = HeaderMap::new();
assert!(!accepts_sse(&headers));
}
#[test]
fn is_stream_path_true() {
assert!(is_stream_path("/users/stream"));
}
#[test]
fn is_stream_path_false_collection() {
assert!(!is_stream_path("/users"));
}
#[test]
fn is_stream_path_false_single() {
assert!(!is_stream_path("/users/123"));
}
#[test]
fn is_stream_path_false_nested() {
assert!(!is_stream_path("/users/123/stream/extra"));
}
#[test]
fn extract_stream_resource_users() {
assert_eq!(extract_stream_resource("/users/stream"), Some("users"));
}
#[test]
fn extract_stream_resource_orders() {
assert_eq!(extract_stream_resource("/orders/stream"), Some("orders"));
}
#[test]
fn extract_stream_resource_none_for_collection() {
assert_eq!(extract_stream_resource("/users"), None);
}
#[test]
fn extract_stream_resource_none_for_single() {
assert_eq!(extract_stream_resource("/users/123"), None);
}
#[test]
fn extract_last_event_id_present() {
let mut headers = HeaderMap::new();
headers.insert("last-event-id", HeaderValue::from_static("evt-42"));
assert_eq!(extract_last_event_id(&headers), Some("evt-42".to_string()));
}
#[test]
fn extract_last_event_id_missing() {
let headers = HeaderMap::new();
assert_eq!(extract_last_event_id(&headers), None);
}
#[test]
fn format_sse_insert_event() {
let data = json!({"id": 1, "name": "Alice"});
let output = format_sse_event("insert", "evt-1", &data);
assert!(output.starts_with("event: insert\n"));
assert!(output.contains("id: evt-1\n"));
assert!(output.contains("data: "));
assert!(output.ends_with("\n\n"));
let data_line = output.lines().find(|l| l.starts_with("data: ")).unwrap();
let json_str = data_line.strip_prefix("data: ").unwrap();
let parsed: serde_json::Value = serde_json::from_str(json_str).unwrap();
assert_eq!(parsed["name"], "Alice");
}
#[test]
fn format_sse_update_event() {
let data = json!({"id": 1, "name": "Alice Updated"});
let output = format_sse_event("update", "evt-2", &data);
assert!(output.starts_with("event: update\n"));
}
#[test]
fn format_sse_delete_event() {
let data = json!({"entity_id": "abc-123"});
let output = format_sse_event("delete", "evt-3", &data);
assert!(output.starts_with("event: delete\n"));
assert!(output.contains("\"entity_id\""));
}
#[test]
fn format_heartbeat_event() {
let output = format_heartbeat();
assert!(output.starts_with("event: ping\n"));
assert!(output.contains("data: \n"));
assert!(output.ends_with("\n\n"));
}
#[test]
fn event_kind_insert() {
assert_eq!(event_kind_to_sse_type("INSERT"), "insert");
}
#[test]
fn event_kind_update() {
assert_eq!(event_kind_to_sse_type("UPDATE"), "update");
}
#[test]
fn event_kind_delete() {
assert_eq!(event_kind_to_sse_type("DELETE"), "delete");
}
#[test]
fn event_kind_custom() {
assert_eq!(event_kind_to_sse_type("CUSTOM"), "custom");
}
#[test]
fn event_kind_unknown() {
assert_eq!(event_kind_to_sse_type("SOMETHING"), "unknown");
}
#[test]
fn observers_not_available_returns_501() {
let err = observers_not_available();
assert_eq!(err.status, StatusCode::NOT_IMPLEMENTED);
assert_eq!(err.code, "NOT_IMPLEMENTED");
}
mod export_config {
use super::super::export_config::{ExportConfig, ExportFormat};
#[test]
fn default_values_match_spec() {
let cfg = ExportConfig::default();
assert_eq!(cfg.csv_delimiter, ',');
assert!(cfg.csv_include_bom);
assert_eq!(cfg.xlsx_max_rows, 100_000);
assert_eq!(cfg.parquet_max_rows, 1_000_000);
assert!(cfg.xlsx_temp_dir.is_none());
assert_eq!(cfg.max_concurrent_xlsx, 10);
assert_eq!(
cfg.export_formats,
vec![ExportFormat::Csv, ExportFormat::Xlsx, ExportFormat::Parquet],
"an unconfigured server must serve every export format"
);
}
#[test]
fn an_explicitly_empty_format_list_disables_every_export() {
let cfg: ExportConfig = toml::from_str("export_formats = []").unwrap();
assert!(!cfg.serves(ExportFormat::Csv));
assert!(!cfg.serves(ExportFormat::Xlsx));
assert!(!cfg.serves(ExportFormat::Parquet));
}
#[test]
fn a_partial_format_list_serves_only_what_it_names() {
let cfg: ExportConfig = toml::from_str(r#"export_formats = ["csv"]"#).unwrap();
assert!(cfg.serves(ExportFormat::Csv));
assert!(!cfg.serves(ExportFormat::Xlsx));
assert!(!cfg.serves(ExportFormat::Parquet));
}
#[test]
fn an_absent_export_table_serves_every_format() {
let cfg: ExportConfig = toml::from_str("").unwrap();
assert!(cfg.serves(ExportFormat::Csv));
assert!(cfg.serves(ExportFormat::Xlsx));
assert!(cfg.serves(ExportFormat::Parquet));
}
#[test]
fn deserializes_empty_toml_to_defaults() {
let cfg: ExportConfig = toml::from_str("").unwrap();
let default_cfg = ExportConfig::default();
assert_eq!(cfg.csv_delimiter, default_cfg.csv_delimiter);
assert_eq!(cfg.csv_include_bom, default_cfg.csv_include_bom);
assert_eq!(cfg.xlsx_max_rows, default_cfg.xlsx_max_rows);
assert_eq!(cfg.parquet_max_rows, default_cfg.parquet_max_rows);
assert_eq!(cfg.xlsx_temp_dir, default_cfg.xlsx_temp_dir);
assert_eq!(cfg.max_concurrent_xlsx, default_cfg.max_concurrent_xlsx);
assert_eq!(cfg.export_formats, default_cfg.export_formats);
}
#[test]
fn deserializes_full_toml_overrides_defaults() {
let toml_src = r#"
csv_delimiter = ";"
csv_include_bom = false
xlsx_max_rows = 50000
parquet_max_rows = 250000
xlsx_temp_dir = "/var/tmp/xlsx"
max_concurrent_xlsx = 4
export_formats = ["csv", "xlsx", "parquet"]
"#;
let cfg: ExportConfig = toml::from_str(toml_src).unwrap();
assert_eq!(cfg.csv_delimiter, ';');
assert!(!cfg.csv_include_bom);
assert_eq!(cfg.xlsx_max_rows, 50_000);
assert_eq!(cfg.parquet_max_rows, 250_000);
assert_eq!(cfg.xlsx_temp_dir.as_deref(), Some(std::path::Path::new("/var/tmp/xlsx")));
assert_eq!(cfg.max_concurrent_xlsx, 4);
assert_eq!(
cfg.export_formats,
vec![ExportFormat::Csv, ExportFormat::Xlsx, ExportFormat::Parquet],
);
}
#[test]
fn export_format_deserializes_lowercase_kebab_strings() {
let cfg: ExportConfig =
toml::from_str(r#"export_formats = ["csv", "xlsx", "parquet"]"#).unwrap();
assert_eq!(
cfg.export_formats,
vec![ExportFormat::Csv, ExportFormat::Xlsx, ExportFormat::Parquet],
);
}
#[test]
fn export_format_rejects_unknown_variant() {
let result: Result<ExportConfig, _> = toml::from_str(r#"export_formats = ["yaml"]"#);
assert!(result.is_err(), "unknown export format should fail to deserialize");
}
}
#[cfg(feature = "observers")]
mod stream_decisions {
use axum::http::{HeaderMap, HeaderValue, StatusCode};
use fraiseql_core::{
runtime::subscription::{ChangeSpineEnvelope, SubscriptionOperation},
security::SecurityContext,
};
use fraiseql_observers::transport::TenantScope;
use crate::{
routes::rest::sse::{
ResumeRequest, StreamEvent, replay_scope, resume_point_unknown, resume_too_far_behind,
resumption_unsupported, stream_event_matches, stream_resume_request,
stream_tenant_scope,
},
subscriptions::EntityEvent as BridgeEvent,
};
fn principal(tenant: Option<&str>) -> SecurityContext {
let user = fraiseql_core::security::auth_middleware::AuthenticatedUser {
user_id: fraiseql_core::types::UserId::new("user-1"),
scopes: vec!["user".to_string()],
expires_at: chrono::Utc::now() + chrono::Duration::hours(1),
email: None,
display_name: None,
extra_claims: std::collections::HashMap::new(),
};
let ctx = SecurityContext::from_user(&user, "req-1113".to_string());
match tenant {
Some(t) => ctx.with_tenant(t.to_string()),
None => ctx,
}
}
#[test]
fn a_multi_tenant_principal_scopes_the_subscription_to_its_own_tenant() {
let ctx = principal(Some("tenant-a"));
assert_eq!(
stream_tenant_scope(Some(&ctx), true).expect("a tenanted principal is servable"),
TenantScope::Tenant("tenant-a".to_string())
);
}
#[test]
fn a_multi_tenant_deployment_refuses_a_principal_with_no_tenant() {
let ctx = principal(None);
let refusal =
stream_tenant_scope(Some(&ctx), true).expect_err("an untenanted principal is refused");
assert_eq!(refusal.status, StatusCode::FORBIDDEN);
assert_eq!(refusal.code, "TENANT_SCOPE_REQUIRED");
}
#[test]
fn a_multi_tenant_deployment_refuses_a_request_with_no_principal() {
let refusal = stream_tenant_scope(None, true).expect_err("an anonymous request is refused");
assert_eq!(refusal.status, StatusCode::FORBIDDEN);
assert_eq!(refusal.code, "TENANT_SCOPE_REQUIRED");
}
#[test]
fn a_single_tenant_deployment_subscribes_unscoped() {
assert_eq!(
stream_tenant_scope(Some(&principal(None)), false).expect("servable"),
TenantScope::AllTenants
);
assert_eq!(stream_tenant_scope(None, false).expect("servable"), TenantScope::AllTenants);
}
#[test]
fn the_deployments_declaration_decides_not_the_credential() {
assert_eq!(
stream_tenant_scope(Some(&principal(Some("tenant-a"))), false).expect("servable"),
TenantScope::AllTenants
);
}
fn headers_with(last_event_id: &str) -> HeaderMap {
let mut headers = HeaderMap::new();
headers.insert("last-event-id", HeaderValue::from_str(last_event_id).unwrap());
headers
}
#[test]
fn a_resume_request_names_the_sequence_to_resume_after() {
let token = crate::routes::rest::stream_token::seal(41, "Order").unwrap();
assert_eq!(
stream_resume_request(&headers_with(&token), "Order")
.expect("a token this stream issued"),
ResumeRequest::From { seq: 41, token }
);
}
#[test]
fn a_fresh_delivery_is_not_a_resume_request() {
assert_eq!(
stream_resume_request(&HeaderMap::new(), "Order").expect("no header is no resume"),
ResumeRequest::Fresh
);
}
#[test]
fn an_empty_last_event_id_is_not_a_resume_request() {
assert_eq!(
stream_resume_request(&headers_with(""), "Order").unwrap(),
ResumeRequest::Fresh
);
assert_eq!(
stream_resume_request(&headers_with(" "), "Order").unwrap(),
ResumeRequest::Fresh
);
}
#[test]
fn an_id_this_stream_never_issues_is_refused() {
let other_stream = crate::routes::rest::stream_token::seal(41, "Invoice").unwrap();
for value in [
"41",
"not-a-number",
"3f2a1b",
"12.5",
other_stream.as_str(),
] {
let refusal = stream_resume_request(&headers_with(value), "Order")
.expect_err("not a token this stream issued");
assert_eq!(refusal.status, StatusCode::GONE);
assert_eq!(refusal.code, "RESUME_POINT_UNKNOWN");
assert!(
refusal.message.contains(value),
"the refusal must name the id it could not read: {}",
refusal.message
);
}
}
#[test]
fn each_refusal_names_its_own_reason_and_the_id() {
let unsupported = resumption_unsupported("41");
assert_eq!(unsupported.status, StatusCode::NOT_IMPLEMENTED);
assert_eq!(unsupported.code, "RESUMPTION_UNSUPPORTED");
let unknown = resume_point_unknown("41");
assert_eq!(unknown.status, StatusCode::GONE);
assert_eq!(unknown.code, "RESUME_POINT_UNKNOWN");
let too_far = resume_too_far_behind("41", 10_000);
assert_eq!(too_far.status, StatusCode::PAYLOAD_TOO_LARGE);
assert_eq!(too_far.code, "RESUME_TOO_FAR_BEHIND");
assert!(
too_far.message.contains("10000"),
"the bound a client exceeded must be named so it can be raised: {}",
too_far.message
);
for refusal in [&unsupported, &unknown, &too_far] {
assert!(
refusal.message.contains("41"),
"every refusal names the id it could not honour: {}",
refusal.message
);
assert!(
refusal.message.contains("without the header"),
"every refusal states the way forward: {}",
refusal.message
);
}
}
#[test]
fn a_replay_is_scoped_by_the_live_gates() {
assert_eq!(
replay_scope("Order", &TenantScope::Tenant("tenant-a".to_string())),
fraiseql_observers::listener::ReplayScope {
object_type: "Order".to_string(),
tenant: Some("tenant-a".to_string()),
}
);
assert_eq!(
replay_scope("Order", &TenantScope::AllTenants),
fraiseql_observers::listener::ReplayScope {
object_type: "Order".to_string(),
tenant: None,
}
);
}
fn event(operation: SubscriptionOperation, seq: Option<i64>) -> BridgeEvent {
let e = BridgeEvent::new(
"Order",
uuid::Uuid::new_v4().to_string(),
operation,
serde_json::json!({"total": 100}),
);
match seq {
Some(seq) => e.with_change_spine(ChangeSpineEnvelope {
seq: Some(seq),
..ChangeSpineEnvelope::default()
}),
None => e,
}
}
#[test]
fn the_wire_id_is_the_change_spine_sequence_not_the_event_uuid() {
let e = event(SubscriptionOperation::Create, Some(4_120));
let wire = StreamEvent::from_bridge_event(&e);
assert_eq!(wire.id.as_deref(), Some("4120"));
assert_ne!(
wire.id.as_deref(),
Some(e.entity_id.as_str()),
"the event UUID must not be the resume id"
);
}
#[test]
fn an_event_with_no_sequence_carries_no_wire_id() {
assert_eq!(
StreamEvent::from_bridge_event(&event(SubscriptionOperation::Create, None)).id,
None
);
let with_empty_seq =
BridgeEvent::new("Order", "1", SubscriptionOperation::Create, serde_json::json!({}))
.with_change_spine(ChangeSpineEnvelope {
actor_type: Some("human_user".to_string()),
..ChangeSpineEnvelope::default()
});
assert_eq!(StreamEvent::from_bridge_event(&with_empty_seq).id, None);
}
#[test]
fn the_wire_event_type_follows_the_operation() {
for (operation, expected) in [
(SubscriptionOperation::Create, "insert"),
(SubscriptionOperation::Update, "update"),
(SubscriptionOperation::Delete, "delete"),
] {
let e = event(operation, Some(1));
assert_eq!(StreamEvent::from_bridge_event(&e).event_type, expected, "{operation:?}");
}
}
#[test]
fn the_wire_payload_is_the_events_data() {
let e = event(SubscriptionOperation::Update, Some(7));
assert_eq!(StreamEvent::from_bridge_event(&e).data, &serde_json::json!({"total": 100}));
}
fn tenanted(entity_type: &str, tenant: Option<&str>) -> BridgeEvent {
let e = BridgeEvent::new(
entity_type,
"1",
SubscriptionOperation::Create,
serde_json::json!({}),
);
match tenant {
Some(t) => e.with_tenant_id(t),
None => e,
}
}
#[test]
fn the_entity_gate_matches_the_graphql_type_name() {
let scope = TenantScope::AllTenants;
assert!(stream_event_matches(&tenanted("Order", None), "Order", &scope));
assert!(
!stream_event_matches(&tenanted("Order", None), "orders", &scope),
"the route name must not match — nothing stamps it"
);
assert!(!stream_event_matches(&tenanted("Invoice", None), "Order", &scope));
}
#[test]
fn a_tenant_scoped_stream_receives_only_its_own_tenants_events() {
let scope = TenantScope::Tenant("tenant-a".to_string());
assert!(stream_event_matches(&tenanted("Order", Some("tenant-a")), "Order", &scope));
assert!(!stream_event_matches(&tenanted("Order", Some("tenant-b")), "Order", &scope));
}
#[test]
fn an_untagged_event_does_not_reach_a_tenant_scoped_stream() {
let scope = TenantScope::Tenant("tenant-a".to_string());
assert!(!stream_event_matches(&tenanted("Order", None), "Order", &scope));
}
#[test]
fn the_lag_frame_names_the_code_and_how_much_was_lost() {
let payload = crate::routes::rest::sse::stream_lagged_payload(17);
assert_eq!(payload["code"], "STREAM_LAGGED");
assert_eq!(payload["skipped"], 17);
assert!(
payload["message"].as_str().unwrap().contains("17"),
"the human-readable half must carry the count too: {payload}"
);
}
#[test]
fn the_lag_frame_is_typed_error_not_an_entity_event() {
use crate::routes::rest::sse::STREAM_LAGGED_EVENT;
assert_eq!(STREAM_LAGGED_EVENT, "error");
for entity_event in ["insert", "update", "delete", "unknown"] {
assert_ne!(STREAM_LAGGED_EVENT, entity_event);
}
}
#[test]
fn an_unscoped_stream_receives_every_tenants_events() {
let scope = TenantScope::AllTenants;
assert!(stream_event_matches(&tenanted("Order", None), "Order", &scope));
assert!(stream_event_matches(&tenanted("Order", Some("tenant-a")), "Order", &scope));
assert!(stream_event_matches(&tenanted("Order", Some("tenant-b")), "Order", &scope));
}
}