use std::{ io, sync::{Arc, Mutex}, }; use crank_observability::{ ObservabilityConfig, RedactionLimits, ServiceIdentity, build_subscriber, safe_json, }; use serde_json::Value; use time::{OffsetDateTime, format_description::well_known::Rfc3339}; use tracing_subscriber::fmt::MakeWriter; #[derive(Clone, Default)] struct SharedWriter { buffer: Arc>>, } impl SharedWriter { fn output(&self) -> String { String::from_utf8(self.buffer.lock().expect("test writer lock").clone()) .expect("log output must be UTF-8") } } impl<'a> MakeWriter<'a> for SharedWriter { type Writer = SharedWriterGuard; fn make_writer(&'a self) -> Self::Writer { SharedWriterGuard { buffer: Arc::clone(&self.buffer), } } } struct SharedWriterGuard { buffer: Arc>>, } impl io::Write for SharedWriterGuard { fn write(&mut self, bytes: &[u8]) -> io::Result { self.buffer .lock() .map_err(|_| io::Error::other("test writer lock poisoned"))? .extend_from_slice(bytes); Ok(bytes.len()) } fn flush(&mut self) -> io::Result<()> { Ok(()) } } fn capture(service: &'static str, emit: impl FnOnce()) -> Vec { capture_with_limits(service, RedactionLimits::default(), emit) } fn capture_with_limits( service: &'static str, limits: RedactionLimits, emit: impl FnOnce(), ) -> Vec { let writer = SharedWriter::default(); let config = ObservabilityConfig::new( ServiceIdentity::try_new(service, "0.3.1", "test").expect("valid test identity"), "info", limits, ); let subscriber = build_subscriber(config, writer.clone()).expect("test subscriber must be built"); tracing::subscriber::with_default(subscriber, emit); writer .output() .lines() .map(|line| { assert!(!line.contains('\u{1b}'), "ANSI is forbidden: {line}"); serde_json::from_str(line).expect("every line must be one JSON object") }) .collect() } #[test] fn schema_contract_is_identical_for_both_services() { let outputs = ["admin-api", "mcp-server"].map(|service| { capture(service, || { tracing::info!( name: "service.started", target: "crank::startup", port = 3101_u64, "service started" ); }) }); for (service, events) in ["admin-api", "mcp-server"].into_iter().zip(outputs.iter()) { assert_eq!(events.len(), 1); let event = &events[0]; assert_eq!(event["service"], service); assert_eq!(event["version"], "0.3.1"); assert_eq!(event["environment"], "test"); assert_eq!(event["level"], "INFO"); assert_eq!(event["target"], "crank::startup"); assert_eq!(event["event"], "service.started"); assert!(event["fields"].is_object()); assert_eq!(event["fields"]["port"], 3101); assert_eq!(event["fields"]["message"], "service started"); let timestamp = event["timestamp"].as_str().expect("timestamp string"); let parsed = OffsetDateTime::parse(timestamp, &Rfc3339).expect("timestamp must be RFC 3339"); assert_eq!(parsed.offset(), time::UtcOffset::UTC); } let first_keys: Vec<_> = outputs[0][0] .as_object() .expect("object") .keys() .cloned() .collect(); let second_keys: Vec<_> = outputs[1][0] .as_object() .expect("object") .keys() .cloned() .collect(); assert_eq!(first_keys, second_keys); } #[test] fn correlation_fields_are_distinct_and_only_present_when_recorded() { let present = capture("admin-api", || { tracing::info!( name: "admin.request.completed", request_id = "req-123", trace_id = "0af7651916cd43dd8448eb211c80319c" ); }); assert_eq!(present[0]["request_id"], "req-123"); assert_eq!(present[0]["trace_id"], "0af7651916cd43dd8448eb211c80319c"); assert!(present[0]["fields"].get("request_id").is_none()); assert!(present[0]["fields"].get("trace_id").is_none()); let absent = capture("mcp-server", || { tracing::info!(name: "mcp.request.completed", status = 200_u64); }); assert!(absent[0].get("request_id").is_none()); assert!(absent[0].get("trace_id").is_none()); } #[test] fn correlation_fields_accept_valid_strings_and_reject_non_string_values() { let limits = RedactionLimits { max_object_fields: 1, ..RedactionLimits::default() }; let request_id = "123"; let events = capture_with_limits("admin-api", limits, || { tracing::info!( name: "admin.request.completed", alpha = "field that consumes the object budget", request_id = %request_id, trace_id = true, ); }); assert_eq!(events[0]["request_id"], "123"); assert!(events[0].get("trace_id").is_none()); } #[test] fn empty_correlation_fields_are_omitted() { let events = capture("admin-api", || { tracing::info!( name: "admin.request.completed", request_id = "", trace_id = "" ); }); assert!(events[0].get("request_id").is_none()); assert!(events[0].get("trace_id").is_none()); } #[test] fn formatter_redacts_fields_before_serialization() { let context = safe_json( &serde_json::json!({ "nested": { "access_token": "nested-canary-secret", "endpoint": "https://example.test/private?key=nested-canary-secret" } }), RedactionLimits::default(), ) .expect("safe nested context"); let events = capture("admin-api", || { tracing::warn!( name: "admin.request.rejected", password = "canary-secret", url = "https://example.test/path?token=canary-secret", safe_fields = %context, unsafe_context = ?serde_json::json!({"password": "debug-canary-secret"}), error_code = "invalid_request" ); }); let serialized = serde_json::to_string(&events[0]).expect("event JSON"); assert_eq!(events[0]["fields"]["password"], "[REDACTED]"); assert_eq!(events[0]["fields"]["url"], "https://example.test/path"); assert_eq!( events[0]["fields"]["safe_fields"]["nested"]["access_token"], "[REDACTED]" ); assert_eq!( events[0]["fields"]["safe_fields"]["nested"]["endpoint"], "https://example.test/private" ); assert_eq!(events[0]["fields"]["unsafe_context"], "[REDACTED]"); assert_eq!(events[0]["fields"]["error_code"], "invalid_request"); assert!(!serialized.contains("canary-secret")); } #[test] fn arbitrary_debug_text_is_never_written_verbatim() { struct Credentials; impl std::fmt::Debug for Credentials { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { formatter.write_str("Credentials { password: \"debug-canary-secret\" }") } } let events = capture("admin-api", || { tracing::warn!( name: "admin.debug.inspected", details = ?Credentials ); }); let serialized = serde_json::to_string(&events[0]).expect("event JSON"); assert_eq!(events[0]["fields"]["details"], "[REDACTED]"); assert!(!serialized.contains("debug-canary-secret")); } #[test] fn compound_sensitive_event_fields_are_redacted() { let events = capture("admin-api", || { tracing::warn!( name: "admin.request.rejected", client_api_key = "client-canary-secret", authorization_header = "Bearer auth-canary-secret", response_body = "response-canary-secret", tool_arguments = "argument-canary-secret", query_params = "query-canary-secret", ); }); let serialized = serde_json::to_string(&events[0]).expect("event JSON"); for key in [ "client_api_key", "authorization_header", "response_body", "tool_arguments", "query_params", ] { assert_eq!(events[0]["fields"][key], "[REDACTED]"); } assert!(!serialized.contains("canary-secret")); } #[test] fn oversized_event_falls_back_to_valid_bounded_json() { let limits = RedactionLimits { max_event_bytes: 512, ..RedactionLimits::default() }; let events = capture_with_limits("admin-api", limits, || { tracing::info!( name: "admin.payload.inspected", description = %"x".repeat(1024) ); }); let serialized = serde_json::to_vec(&events[0]).expect("bounded event JSON"); assert!(serialized.len() < limits.max_event_bytes); assert_eq!(events[0]["fields"]["truncated"], true); } #[test] fn safe_json_honours_the_total_event_budget() { let limits = RedactionLimits { max_event_bytes: 512, ..RedactionLimits::default() }; let serialized = safe_json( &serde_json::json!({"description": "x".repeat(4096)}), limits, ) .expect("safe JSON must remain serializable"); assert!(serialized.len() <= limits.max_event_bytes); assert_eq!( serde_json::from_str::(&serialized).expect("valid JSON")["truncated"], true ); } #[test] fn subscriber_rejects_limits_that_cannot_hold_an_event() { let config = ObservabilityConfig::new( ServiceIdentity::try_new("admin-api", "0.3.1", "test").expect("valid identity"), "info", RedactionLimits { max_event_bytes: 16, ..RedactionLimits::default() }, ); assert!(build_subscriber(config, SharedWriter::default()).is_err()); } #[test] fn event_budget_rejects_maximum_identity_labels_when_correlation_cannot_fit() { let writer = SharedWriter::default(); let config = ObservabilityConfig::new( ServiceIdentity::try_new("s".repeat(64), "v".repeat(64), "e".repeat(64)) .expect("maximum identity labels are valid"), "info", RedactionLimits { max_event_bytes: 512, ..RedactionLimits::default() }, ); assert!(build_subscriber(config, writer).is_err()); } #[test] fn env_filter_is_applied_and_invalid_filter_is_safe() { let writer = SharedWriter::default(); let config = ObservabilityConfig::new( ServiceIdentity::try_new("admin-api", "0.3.1", "test").expect("valid identity"), "warn", RedactionLimits::default(), ); let subscriber = build_subscriber(config, writer.clone()).expect("test subscriber must be built"); tracing::subscriber::with_default(subscriber, || { tracing::info!(name: "filtered.info", "filtered"); tracing::warn!(name: "visible.warning", "visible"); }); let output = writer.output(); assert!(!output.contains("filtered.info")); assert!(output.contains("visible.warning")); let invalid = ObservabilityConfig::new( ServiceIdentity::try_new("admin-api", "0.3.1", "test").expect("valid identity"), "[not a valid filter", RedactionLimits::default(), ); let error = build_subscriber(invalid, SharedWriter::default()) .err() .expect("invalid filter must fail"); assert_eq!(error.to_string(), "invalid log filter"); assert!(!error.to_string().contains("not a valid filter")); } #[test] fn service_identity_rejects_empty_or_unsafe_labels() { for (service, version, environment) in [ ("", "0.3.1", "test"), ("admin api", "0.3.1", "test"), ("admin-api", "", "test"), ("admin-api", "0.3.1", "prod\nsecret"), ] { assert!(ServiceIdentity::try_new(service, version, environment).is_err()); } }