diff --git a/apps/admin-api/src/request_context.rs b/apps/admin-api/src/request_context.rs index 61acf98..3e7da24 100644 --- a/apps/admin-api/src/request_context.rs +++ b/apps/admin-api/src/request_context.rs @@ -6,7 +6,7 @@ use axum::{ }; use crank_core::{CorrelationContext, RequestId, TraceContext}; use crank_metrics::ExemplarTraceId; -use crank_observability::{set_remote_trace_parent, with_request_correlation}; +use crank_observability::with_request_correlation; use tracing::{Instrument, info, info_span}; pub const REQUEST_ID_HEADER: HeaderName = HeaderName::from_static("x-request-id"); @@ -124,11 +124,7 @@ fn one_auxiliary_header_within_budget( } fn set_canonical_parent(span: &tracing::Span, context: &TraceContext) { - let mut headers = axum::http::HeaderMap::new(); - if let Ok(value) = HeaderValue::from_str(context.traceparent()) { - headers.insert("traceparent", value); - set_remote_trace_parent(span, &headers); - } + crank_trace::set_parent_from_trace_context(span, context); } #[cfg(test)] diff --git a/apps/admin-api/tests/integration/openapi_source/telemetry.rs b/apps/admin-api/tests/integration/openapi_source/telemetry.rs index 566aefb..dad5728 100644 --- a/apps/admin-api/tests/integration/openapi_source/telemetry.rs +++ b/apps/admin-api/tests/integration/openapi_source/telemetry.rs @@ -154,17 +154,19 @@ impl<'a> MakeWriter<'a> for SharedLogWriter { 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.buffer.lock().unwrap().extend_from_slice(bytes); + self.pending.extend_from_slice(bytes); Ok(bytes.len()) } @@ -173,6 +175,24 @@ impl io::Write for SharedLogGuard { } } +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>>); diff --git a/apps/admin-api/tests/request_context_parent.rs b/apps/admin-api/tests/request_context_parent.rs new file mode 100644 index 0000000..cab285f --- /dev/null +++ b/apps/admin-api/tests/request_context_parent.rs @@ -0,0 +1,39 @@ +use admin_api::request_context::{TRACE_ID_HEADER, apply_request_context}; +use axum::{Router, body::Body, http::Request, routing::get}; +use opentelemetry::trace::TracerProvider as _; +use opentelemetry_sdk::trace::SdkTracerProvider; +use tower::ServiceExt; +use tracing_subscriber::layer::SubscriberExt; + +#[tokio::test(flavor = "current_thread")] +async fn preserves_remote_trace_id_with_an_active_tracer_and_no_global_propagator() { + let provider = SdkTracerProvider::builder().build(); + let tracer = provider.tracer("admin-request-context-parent-test"); + let subscriber = + tracing_subscriber::registry().with(tracing_opentelemetry::layer().with_tracer(tracer)); + let dispatch = tracing::Dispatch::new(subscriber); + let _dispatch_guard = tracing::dispatcher::set_default(&dispatch); + let app = Router::new() + .route("/probe", get(|| async { "ok" })) + .layer(axum::middleware::from_fn(apply_request_context)); + + let response = app + .oneshot( + Request::builder() + .uri("/probe") + .header( + "traceparent", + "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01", + ) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!( + response.headers()[TRACE_ID_HEADER], + "0af7651916cd43dd8448eb211c80319c" + ); + provider.shutdown().unwrap(); +}