use super::super::common::build_test_app_with_external_references; use super::*; use std::{ io, sync::{Arc, Mutex}, }; use metrics_util::debugging::DebuggingRecorder; use opentelemetry::trace::TracerProvider as _; use opentelemetry_sdk::{ error::OTelSdkResult, trace::{SdkTracerProvider, SpanData, SpanExporter}, }; use tracing_subscriber::{fmt::MakeWriter, layer::SubscriberExt}; #[tokio::test(flavor = "multi_thread")] #[serial] async fn parser_canary_never_reaches_diagnostics_logs_traces_or_metrics() { const CANARY: &str = "openapi-telemetry-secret-canary"; let recorder = DebuggingRecorder::new(); let snapshotter = recorder.snapshotter(); recorder .install() .expect("isolated integration test metrics recorder"); let writer = SharedLogWriter::default(); let exported = Arc::new(Mutex::new(Vec::new())); let provider = SdkTracerProvider::builder() .with_simple_exporter(CapturingExporter(Arc::clone(&exported))) .build(); let tracer = provider.tracer("admin-openapi-source-test"); let subscriber = tracing_subscriber::registry() .with( tracing_subscriber::fmt::layer() .with_ansi(false) .with_writer(writer.clone()), ) .with(tracing_opentelemetry::layer().with_tracer(tracer)); let dispatch = tracing::Dispatch::new(subscriber); tracing::dispatcher::set_global_default(dispatch) .expect("isolated integration test tracing subscriber"); let registry = test_registry().await; let app = build_test_app(registry, test_storage_root("openapi_telemetry_canary")); let server = spawn_admin_api(app).await; let client = authorized_client(&server).await; let document = format!( "openapi: 3.0.3\ninfo: {{ title: Canary }}\npaths:\n /broken:\n get:\n description: {CANARY}\n responses: [" ); let response = client .post(format!("{server}/imports/openapi/preview")) .multipart(Form::new().part( "file", file_part(document.as_bytes(), "openapi.yaml", "application/yaml"), )) .send() .await .unwrap(); assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST); let trace_id = response .headers() .get("x-trace-id") .unwrap() .to_str() .unwrap() .to_owned(); let body = response.text().await.unwrap(); assert!(!body.contains(CANARY)); let diagnostics: Value = serde_json::from_str(&body).unwrap(); assert_eq!(diagnostics["error"]["trace_id"], trace_id); let route = format!("/{CANARY}.yaml"); let external = axum::Router::new().route( &route, axum::routing::get(|| async { (axum::http::StatusCode::INTERNAL_SERVER_ERROR, CANARY) }), ); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); tokio::spawn(async move { axum::serve(listener, external).await.unwrap() }); let origin = format!("http://{address}"); let external_registry = test_registry().await; let external_app = build_test_app_with_external_references( external_registry, test_storage_root("openapi_materialization_telemetry_canary"), vec![format!("{origin}/")], ); let external_server = spawn_admin_api(external_app).await; let external_client = authorized_client(&external_server).await; let external_document = format!( r#" openapi: 3.1.0 info: {{ title: Safe materialization }} servers: [{{ url: https://api.example.test }}] paths: /items: get: responses: '200': description: ok content: application/json: schema: {{ $ref: '{origin}/{CANARY}.yaml#/Item' }} "# ); let response = external_client .post(format!("{external_server}/imports/openapi/preview")) .multipart(Form::new().part( "file", file_part( external_document.as_bytes(), "external.yaml", "application/yaml", ), )) .send() .await .unwrap(); assert_eq!(response.status(), reqwest::StatusCode::OK); provider.force_flush().unwrap(); let logs = writer.output(); assert!(!logs.contains(CANARY)); assert!(logs.contains(&trace_id)); let safe_materialization_failure = logs.lines().any(is_safe_materialization_failure); assert!( safe_materialization_failure, "expected one bounded, sanitized materialization failure" ); let spans = exported.lock().unwrap(); let rendered_spans = format!("{spans:?}"); assert!(!rendered_spans.contains(CANARY)); drop(spans); for (key, _, _, _) in snapshotter.snapshot().into_vec() { assert!(!key.key().name().contains(CANARY)); assert!( !key.key() .labels() .any(|label| { label.key().contains(CANARY) || label.value().contains(CANARY) }) ); } provider.shutdown().unwrap(); } fn is_safe_materialization_failure(line: &str) -> bool { if !line.contains("external OpenAPI materialization failed") || !line.contains("count=0") { return false; } let immediate_status = contains_log_field(line, "stage", "fetch") && contains_log_field(line, "error_code", "unexpected_status"); let bounded_timeout = contains_log_field(line, "stage", "chain") && contains_log_field(line, "error_code", "timeout"); immediate_status || bounded_timeout } fn contains_log_field(line: &str, name: &str, value: &str) -> bool { line.contains(&format!("{name}=\"{value}\"")) || line.contains(&format!("{name}={value}")) } #[test] fn materialization_log_accepts_only_coherent_safe_outcomes() { for line in [ r#"external OpenAPI materialization failed stage="fetch" error_code="unexpected_status" count=0"#, "external OpenAPI materialization failed stage=chain error_code=timeout count=0", ] { assert!(is_safe_materialization_failure(line)); } assert!(!is_safe_materialization_failure( r#"external OpenAPI materialization failed stage="fetch" error_code="timeout" count=0"# )); } #[derive(Clone, Default)] struct SharedLogWriter { buffer: Arc>>, } impl SharedLogWriter { fn output(&self) -> String { String::from_utf8(self.buffer.lock().unwrap().clone()).unwrap() } } impl<'a> MakeWriter<'a> for SharedLogWriter { type Writer = SharedLogGuard; fn make_writer(&'a self) -> Self::Writer { SharedLogGuard { buffer: Arc::clone(&self.buffer), pending: Vec::new(), } } } struct SharedLogGuard { buffer: Arc>>, pending: Vec, } impl io::Write for SharedLogGuard { fn write(&mut self, bytes: &[u8]) -> io::Result { self.pending.extend_from_slice(bytes); Ok(bytes.len()) } fn flush(&mut self) -> io::Result<()> { Ok(()) } } impl Drop for SharedLogGuard { fn drop(&mut self) { self.buffer.lock().unwrap().extend_from_slice(&self.pending); } } #[test] fn shared_log_writer_publishes_complete_records_only() { let writer = SharedLogWriter::default(); let mut guard = writer.make_writer(); io::Write::write_all(&mut guard, b"message").unwrap(); assert!(writer.output().is_empty()); io::Write::write_all(&mut guard, b" stage=\"fetch\"\n").unwrap(); drop(guard); assert_eq!(writer.output(), "message stage=\"fetch\"\n"); } #[derive(Clone, Debug)] struct CapturingExporter(Arc>>); impl SpanExporter for CapturingExporter { async fn export(&self, batch: Vec) -> OTelSdkResult { self.0.lock().unwrap().extend(batch); Ok(()) } }