From 31f72a28312547e0d662f060ed3239a4132b54fc Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Thu, 27 Aug 2026 10:37:39 -0400 Subject: [PATCH 01/11] Send oversized trace payload by itself instead of requeuing it --- crates/datadog-trace-agent/src/aggregator.rs | 39 +++++++++++++++++++- 1 file changed, 37 insertions(+), 2 deletions(-) diff --git a/crates/datadog-trace-agent/src/aggregator.rs b/crates/datadog-trace-agent/src/aggregator.rs index b9a5b0d..ce8dac6 100644 --- a/crates/datadog-trace-agent/src/aggregator.rs +++ b/crates/datadog-trace-agent/src/aggregator.rs @@ -50,9 +50,21 @@ impl TraceAggregator { if let Some(payload) = self.queue.pop_front() { let payload_size = payload.len(); - // Put stats back in the queue + // Put payload back in the queue if it doesn't fit in this batch if batch_size + payload_size > self.max_content_size_bytes { - self.queue.push_front(payload); + if self.buffer.is_empty() { + // This payload alone exceeds max_content_size_bytes. Send it on + // its own rather than requeuing it, otherwise it would permanently + // block every payload queued behind it from ever being flushed. + tracing::warn!( + payload_size, + max_content_size_bytes = self.max_content_size_bytes, + "Trace payload exceeds max batch size; sending it standalone" + ); + self.buffer.push(payload); + } else { + self.queue.push_front(payload); + } break; } batch_size += payload_size; @@ -141,4 +153,27 @@ mod tests { assert_eq!(second_batch.len(), 1); assert_eq!(aggregator.queue.len(), 0); } + + #[test] + fn test_get_batch_oversized_payload_is_sent_standalone() { + let mut aggregator = TraceAggregator::new(10); + + // A payload larger than max_content_size_bytes on its own, followed by a + // normal-sized payload behind it. + aggregator.add(create_test_send_data(20)); + aggregator.add(create_test_send_data(5)); + + // The oversized payload must be flushed on its own rather than blocking + // the queue forever. + let first_batch = aggregator.get_batch(); + assert_eq!(first_batch.len(), 1); + assert_eq!(first_batch[0].len(), 20); + assert_eq!(aggregator.queue.len(), 1); + + // The payload queued behind it must still be reachable. + let second_batch = aggregator.get_batch(); + assert_eq!(second_batch.len(), 1); + assert_eq!(second_batch[0].len(), 5); + assert_eq!(aggregator.queue.len(), 0); + } } From 801d4b4a5161c83e4032e351ce55e9decdd339e1 Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Fri, 28 Aug 2026 16:32:36 -0400 Subject: [PATCH 02/11] Split oversized trace payload into trace chunks --- Cargo.lock | 1 + crates/datadog-trace-agent/Cargo.toml | 1 + crates/datadog-trace-agent/src/aggregator.rs | 37 ++- .../src/trace_processor.rs | 247 +++++++++++++++++- 4 files changed, 254 insertions(+), 32 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index fef94d1..c4e7919 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -676,6 +676,7 @@ dependencies = [ "libdd-trace-protobuf", "libdd-trace-stats", "libdd-trace-utils", + "prost 0.14.3", "reqwest", "rmp-serde", "serde", diff --git a/crates/datadog-trace-agent/Cargo.toml b/crates/datadog-trace-agent/Cargo.toml index c4f9577..291ff87 100644 --- a/crates/datadog-trace-agent/Cargo.toml +++ b/crates/datadog-trace-agent/Cargo.toml @@ -39,6 +39,7 @@ reqwest = { version = "0.12.23", features = [ "http2", ], default-features = false } bytes = "1.10.1" +prost = "0.14.1" [dev-dependencies] flate2 = "1" diff --git a/crates/datadog-trace-agent/src/aggregator.rs b/crates/datadog-trace-agent/src/aggregator.rs index ce8dac6..eb8eed7 100644 --- a/crates/datadog-trace-agent/src/aggregator.rs +++ b/crates/datadog-trace-agent/src/aggregator.rs @@ -47,31 +47,24 @@ impl TraceAggregator { // Fill the batch while batch_size < self.max_content_size_bytes { - if let Some(payload) = self.queue.pop_front() { - let payload_size = payload.len(); - - // Put payload back in the queue if it doesn't fit in this batch - if batch_size + payload_size > self.max_content_size_bytes { - if self.buffer.is_empty() { - // This payload alone exceeds max_content_size_bytes. Send it on - // its own rather than requeuing it, otherwise it would permanently - // block every payload queued behind it from ever being flushed. - tracing::warn!( - payload_size, - max_content_size_bytes = self.max_content_size_bytes, - "Trace payload exceeds max batch size; sending it standalone" - ); - self.buffer.push(payload); - } else { - self.queue.push_front(payload); - } - break; + let Some(payload) = self.queue.pop_front() else { + break; + }; + let payload_size = payload.len(); + + // Put payload back in the queue if it doesn't fit in this batch + if batch_size + payload_size > self.max_content_size_bytes { + if self.buffer.is_empty() { + // A single payload larger than max_content_size_bytes still needs to be + // flushed — form a batch of just this one + self.buffer.push(payload); + } else { + self.queue.push_front(payload); } - batch_size += payload_size; - self.buffer.push(payload); - } else { break; } + batch_size += payload_size; + self.buffer.push(payload); } std::mem::take(&mut self.buffer) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index 361c66a..ecbf5b7 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -8,15 +8,17 @@ use hyper::{StatusCode, http}; use libdd_common::http_common; use libdd_library_config::tracer_metadata::TracerMetadata; use tokio::sync::mpsc::Sender; -use tracing::{debug, error}; +use tracing::{debug, error, warn}; use libdd_trace_obfuscation::obfuscate::obfuscate_span; use libdd_trace_protobuf::pb; use libdd_trace_utils::trace_utils::{self}; use libdd_trace_utils::trace_utils::{EnvironmentType, SendData}; use libdd_trace_utils::tracer_payload::{TraceChunkProcessor, TracerPayloadCollection}; +use prost::Message; use crate::{ + aggregator::MAX_CONTENT_SIZE_BYTES, config::Config, http_utils::{self, log_and_create_http_response, log_and_create_traces_success_http_response}, stats_concentrator_service::StatsConcentratorHandle, @@ -24,6 +26,45 @@ use crate::{ const TRACER_PAYLOAD_FUNCTION_TAGS_TAG_KEY: &str = "_dd.tags.function"; +/// Splits `payloads` so that each returned `TracerPayload`'s encoded size fits within +/// `max_size` where possible. Recursively bisects by trace-chunk boundary. A single chunk +/// that's still oversized is returned as-is and gets sent standalone. +fn split_oversized_payloads( + payloads: Vec, + max_size: usize, +) -> Vec { + payloads + .into_iter() + .flat_map(|tp| split_tracer_payload(tp, max_size)) + .collect() +} + +fn split_tracer_payload(tp: pb::TracerPayload, max_size: usize) -> Vec { + if tp.encoded_len() <= max_size { + return vec![tp]; + } + + if tp.chunks.len() > 1 { + let mid = tp.chunks.len() / 2; + let mut first = tp.clone(); + let second_chunks = first.chunks.split_off(mid); + let mut second = tp; + second.chunks = second_chunks; + + let mut result = split_tracer_payload(first, max_size); + result.extend(split_tracer_payload(second, max_size)); + return result; + } + + vec![tp] +} + +/// Computes the total encoded size of `payloads`, matching the metric `MAX_CONTENT_SIZE_BYTES` +/// is measured in. +fn encoded_size(payloads: &[pb::TracerPayload]) -> usize { + payloads.iter().map(Message::encoded_len).sum() +} + #[async_trait] pub trait TraceProcessor { /// Deserializes traces from a hyper request body and sends them through the provided tokio mpsc @@ -196,23 +237,54 @@ impl TraceProcessor for ServerlessTraceProcessor { Self::send_to_concentrator(concentrator, &payload); } - let send_data = SendData::new(body_size, payload, tracer_header_tags, &config.trace_intake); + let pieces: Vec<(TracerPayloadCollection, usize)> = match payload { + TracerPayloadCollection::V07(payloads) => { + split_oversized_payloads(payloads, MAX_CONTENT_SIZE_BYTES) + .into_iter() + .map(|tp| { + let size = encoded_size(std::slice::from_ref(&tp)); + (TracerPayloadCollection::V07(vec![tp]), size) + }) + .collect() + } + other => vec![(other, body_size)], + }; - // send trace payload to our trace flusher - match tx.send(send_data).await { - Ok(_) => { - return log_and_create_traces_success_http_response( - "Successfully buffered traces to be flushed.", - StatusCode::OK, + if pieces.len() > 1 { + debug!( + piece_count = pieces.len(), + "Oversized trace payload split into multiple pieces" + ); + } + + for (piece, size) in pieces { + let send_data = SendData::new( + size, + piece, + tracer_header_tags.clone(), + &config.trace_intake, + ); + + if size > MAX_CONTENT_SIZE_BYTES { + warn!( + payload_size = size, + max_content_size_bytes = MAX_CONTENT_SIZE_BYTES, + "Trace chunk is over max batch size; sending standalone" ); } - Err(err) => { + + if let Err(err) = tx.send(send_data).await { return log_and_create_http_response( &format!("Error sending traces to the trace flusher: {err}"), StatusCode::INTERNAL_SERVER_ERROR, ); } } + + log_and_create_traces_success_http_response( + "Successfully buffered traces to be flushed.", + StatusCode::OK, + ) } } @@ -224,9 +296,13 @@ mod tests { use tokio::sync::mpsc::{self, Receiver, Sender}; use crate::{ + aggregator::MAX_CONTENT_SIZE_BYTES, config::{Config, Tags}, peer_tags::peer_tag_keys, - trace_processor::{self, TRACER_PAYLOAD_FUNCTION_TAGS_TAG_KEY, TraceProcessor}, + trace_processor::{ + self, TRACER_PAYLOAD_FUNCTION_TAGS_TAG_KEY, TraceProcessor, encoded_size, + split_oversized_payloads, + }, }; use libdd_common::{Endpoint, http_common}; use libdd_trace_protobuf::pb; @@ -303,6 +379,89 @@ mod tests { } } + fn make_span(meta: HashMap) -> pb::Span { + pb::Span { + meta, + ..Default::default() + } + } + + fn make_chunk(spans: Vec) -> pb::TraceChunk { + pb::TraceChunk { + spans, + ..Default::default() + } + } + + fn make_payload(chunks: Vec) -> pb::TracerPayload { + pb::TracerPayload { + chunks, + ..Default::default() + } + } + + fn big_span() -> pb::Span { + make_span(HashMap::from([("blob".to_string(), "x".repeat(50))])) + } + + #[test] + fn test_no_split_needed_when_under_max() { + let payload = make_payload(vec![make_chunk(vec![big_span()])]); + let size = encoded_size(std::slice::from_ref(&payload)); + + let result = split_oversized_payloads(vec![payload], size); + + assert_eq!(result.len(), 1); + } + + #[test] + fn test_splits_multiple_chunks_when_collectively_oversized() { + // Two chunks, each individually small, but together over max_size. + let payload = make_payload(vec![ + make_chunk(vec![big_span()]), + make_chunk(vec![big_span()]), + ]); + let one_chunk_size = + encoded_size(std::slice::from_ref(&make_payload(vec![make_chunk(vec![ + big_span(), + ])]))); + let max_size = one_chunk_size + 10; // fits one chunk, not both + + let result = split_oversized_payloads(vec![payload], max_size); + + assert_eq!(result.len(), 2); + for piece in &result { + assert_eq!(piece.chunks.len(), 1); + assert!(encoded_size(std::slice::from_ref(piece)) <= max_size); + } + } + + #[test] + fn test_single_oversized_span_returned_as_is() { + // One chunk, one span, that span alone already exceeds max_size. + let payload = make_payload(vec![make_chunk(vec![big_span()])]); + let actual_size = encoded_size(std::slice::from_ref(&payload)); + let max_size = actual_size - 1; // impossible to fit, even alone + + let result = split_oversized_payloads(vec![payload], max_size); + + assert_eq!(result.len(), 1); + assert_eq!(result[0].chunks.len(), 1); + assert_eq!(result[0].chunks[0].spans.len(), 1); + // Still oversized - this is the signal the caller logs a warning for. + assert!(encoded_size(&result) > max_size); + } + + #[test] + fn test_encoded_size_sums_multiple_payloads() { + let a = make_payload(vec![make_chunk(vec![big_span()])]); + let b = make_payload(vec![make_chunk(vec![big_span()])]); + let a_size = encoded_size(std::slice::from_ref(&a)); + let b_size = encoded_size(std::slice::from_ref(&b)); + + assert_eq!(encoded_size(&[a, b]), a_size + b_size); + } + #[tokio::test] async fn test_process_trace() { let (tx, mut rx): ( @@ -453,4 +612,72 @@ mod tests { assert_eq!(expected_tracer_payload, received_payload.unwrap()); } + + #[tokio::test] + async fn test_process_trace_sends_oversized_single_chunk_standalone() { + let (tx, mut rx): ( + Sender, + Receiver, + ) = mpsc::channel(10); + + let start = get_current_timestamp_nanos(); + + // One trace (one chunk) with enough spans to exceed MAX_CONTENT_SIZE_BYTES after + // protobuf encoding. Each span gets a large meta tag to ensure the payload is + // sufficiently large when encoded. + let mut spans = Vec::new(); + for i in 0..6000 { + let mut span = create_test_json_span(11, 222, 333 + i as u64, start, false); + if let Some(obj) = span.as_object_mut() { + obj.insert( + "meta".to_string(), + serde_json::json!({ + "large_field": "x".repeat(1024) // 1KB per span x 6000 = ~6MB before encoding + }), + ); + } + spans.push(span); + } + + let bytes = rmp_serde::to_vec(&vec![spans]).unwrap(); + let request = Request::builder() + .header("datadog-meta-tracer-version", "4.0.0") + .header("datadog-meta-lang", "nodejs") + .header("datadog-meta-lang-version", "v19.7.0") + .header("datadog-meta-lang-interpreter", "v8") + .header("datadog-container-id", "33") + .header("content-length", "100") + .body(http_common::Body::from(bytes)) + .unwrap(); + + let trace_processor = trace_processor::ServerlessTraceProcessor { + stats_concentrator: None, + }; + let res = trace_processor + .process_traces( + Arc::new(create_test_config()), + request, + tx, + Arc::new(create_test_metadata()), + ) + .await; + assert!(res.is_ok()); + + let mut received = Vec::new(); + while let Ok(send_data) = rx.try_recv() { + received.push(send_data); + } + + assert_eq!( + received.len(), + 1, + "expected the oversized single-chunk trace to be sent standalone as one piece, got {}", + received.len() + ); + assert!( + received[0].len() > MAX_CONTENT_SIZE_BYTES, + "expected the standalone piece to still be reported as oversized (size {})", + received[0].len() + ); + } } From 7dd47adfbc494cd4eb386cab3afd5cf9010be09a Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Fri, 28 Aug 2026 16:55:07 -0400 Subject: [PATCH 03/11] Clone the metadata-only base payload instead of the chunk, improve test --- .../src/trace_processor.rs | 36 ++++++++++--------- 1 file changed, 19 insertions(+), 17 deletions(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index ecbf5b7..7c26238 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -46,9 +46,14 @@ fn split_tracer_payload(tp: pb::TracerPayload, max_size: usize) -> Vec 1 { let mid = tp.chunks.len() / 2; - let mut first = tp.clone(); - let second_chunks = first.chunks.split_off(mid); - let mut second = tp; + + // Avoid cloning large chunk/span data on each bisection: clone only the metadata. + let mut base = tp; + let mut first_chunks = std::mem::take(&mut base.chunks); + let second_chunks = first_chunks.split_off(mid); + let mut first = base.clone(); + first.chunks = first_chunks; + let mut second = base; second.chunks = second_chunks; let mut result = split_tracer_payload(first, max_size); @@ -622,22 +627,19 @@ mod tests { let start = get_current_timestamp_nanos(); - // One trace (one chunk) with enough spans to exceed MAX_CONTENT_SIZE_BYTES after - // protobuf encoding. Each span gets a large meta tag to ensure the payload is - // sufficiently large when encoded. + // One trace (one chunk) with a single span whose meta field alone exceeds + // MAX_CONTENT_SIZE_BYTES once encoded, but stays under max_request_content_length. let mut spans = Vec::new(); - for i in 0..6000 { - let mut span = create_test_json_span(11, 222, 333 + i as u64, start, false); - if let Some(obj) = span.as_object_mut() { - obj.insert( - "meta".to_string(), - serde_json::json!({ - "large_field": "x".repeat(1024) // 1KB per span x 6000 = ~6MB before encoding - }), - ); - } - spans.push(span); + let mut span = create_test_json_span(11, 222, 333, start, false); + if let Some(obj) = span.as_object_mut() { + obj.insert( + "meta".to_string(), + serde_json::json!({ + "large_field": "x".repeat(MAX_CONTENT_SIZE_BYTES) + }), + ); } + spans.push(span); let bytes = rmp_serde::to_vec(&vec![spans]).unwrap(); let request = Request::builder() From 5d8739c56e8b9eec3cafbfa1c86ff5519595017d Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Mon, 31 Aug 2026 11:12:43 -0400 Subject: [PATCH 04/11] Add safety buffer for tracer payload size --- crates/datadog-trace-agent/src/trace_processor.rs | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index 7c26238..f9ebe78 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -26,6 +26,10 @@ use crate::{ const TRACER_PAYLOAD_FUNCTION_TAGS_TAG_KEY: &str = "_dd.tags.function"; +/// Rough upper bound on the protobuf framing overhead added when a V07 `TracerPayload` is +/// wrapped in the outer `AgentPayload` envelope before being sent +const V07_ENVELOPE_OVERHEAD_BYTES: usize = 64; + /// Splits `payloads` so that each returned `TracerPayload`'s encoded size fits within /// `max_size` where possible. Recursively bisects by trace-chunk boundary. A single chunk /// that's still oversized is returned as-is and gets sent standalone. @@ -244,7 +248,9 @@ impl TraceProcessor for ServerlessTraceProcessor { let pieces: Vec<(TracerPayloadCollection, usize)> = match payload { TracerPayloadCollection::V07(payloads) => { - split_oversized_payloads(payloads, MAX_CONTENT_SIZE_BYTES) + let split_budget = + MAX_CONTENT_SIZE_BYTES.saturating_sub(V07_ENVELOPE_OVERHEAD_BYTES); + split_oversized_payloads(payloads, split_budget) .into_iter() .map(|tp| { let size = encoded_size(std::slice::from_ref(&tp)); From 6fe37c2808a6979b6fff15cab718485f9f50da76 Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Mon, 31 Aug 2026 12:49:30 -0400 Subject: [PATCH 05/11] Update comment --- crates/datadog-trace-agent/src/trace_processor.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index f9ebe78..8804704 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -68,8 +68,7 @@ fn split_tracer_payload(tp: pb::TracerPayload, max_size: usize) -> Vec usize { payloads.iter().map(Message::encoded_len).sum() } @@ -277,6 +276,7 @@ impl TraceProcessor for ServerlessTraceProcessor { ); if size > MAX_CONTENT_SIZE_BYTES { + // `size` is the inner-TracerPayload-only encoded_size() - this is an approximate check warn!( payload_size = size, max_content_size_bytes = MAX_CONTENT_SIZE_BYTES, From db6b375051d70930d0de2a3120cb4373e45f82df Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Mon, 31 Aug 2026 13:05:29 -0400 Subject: [PATCH 06/11] Include size buffer in SendData input --- crates/datadog-trace-agent/src/trace_processor.rs | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index 8804704..6126d31 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -252,7 +252,8 @@ impl TraceProcessor for ServerlessTraceProcessor { split_oversized_payloads(payloads, split_budget) .into_iter() .map(|tp| { - let size = encoded_size(std::slice::from_ref(&tp)); + let size = + encoded_size(std::slice::from_ref(&tp)) + V07_ENVELOPE_OVERHEAD_BYTES; (TracerPayloadCollection::V07(vec![tp]), size) }) .collect() @@ -276,7 +277,8 @@ impl TraceProcessor for ServerlessTraceProcessor { ); if size > MAX_CONTENT_SIZE_BYTES { - // `size` is the inner-TracerPayload-only encoded_size() - this is an approximate check + // For V07, `size` includes V07_ENVELOPE_OVERHEAD_BYTES; for other + // formats it's the raw body_size - both approximate checks warn!( payload_size = size, max_content_size_bytes = MAX_CONTENT_SIZE_BYTES, From 1cda70d2fdf2ce2659db4aeff733f8a526edf4d8 Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Mon, 31 Aug 2026 13:11:27 -0400 Subject: [PATCH 07/11] Update warning log --- crates/datadog-trace-agent/src/trace_processor.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index 6126d31..b3c6c20 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -282,7 +282,7 @@ impl TraceProcessor for ServerlessTraceProcessor { warn!( payload_size = size, max_content_size_bytes = MAX_CONTENT_SIZE_BYTES, - "Trace chunk is over max batch size; sending standalone" + "Trace payload is over max batch size; sending standalone" ); } From bfafff3809a1f07c2badbc840c60ab0b37614362 Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Wed, 2 Sep 2026 15:00:59 -0400 Subject: [PATCH 08/11] Decouple tx.send() from the HTTP response to the tracer --- .../src/trace_processor.rs | 71 ++++++++++--------- 1 file changed, 37 insertions(+), 34 deletions(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index b3c6c20..8d9bddf 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -268,31 +268,36 @@ impl TraceProcessor for ServerlessTraceProcessor { ); } - for (piece, size) in pieces { - let send_data = SendData::new( - size, - piece, - tracer_header_tags.clone(), - &config.trace_intake, - ); - - if size > MAX_CONTENT_SIZE_BYTES { - // For V07, `size` includes V07_ENVELOPE_OVERHEAD_BYTES; for other - // formats it's the raw body_size - both approximate checks - warn!( - payload_size = size, - max_content_size_bytes = MAX_CONTENT_SIZE_BYTES, - "Trace payload is over max batch size; sending standalone" - ); - } + let send_datas: Vec = pieces + .into_iter() + .map(|(piece, size)| { + if size > MAX_CONTENT_SIZE_BYTES { + // For V07, `size` includes V07_ENVELOPE_OVERHEAD_BYTES; for other + // formats it's the raw body_size - both approximate checks + warn!( + payload_size = size, + max_content_size_bytes = MAX_CONTENT_SIZE_BYTES, + "Trace payload is over max batch size; sending standalone" + ); + } - if let Err(err) = tx.send(send_data).await { - return log_and_create_http_response( - &format!("Error sending traces to the trace flusher: {err}"), - StatusCode::INTERNAL_SERVER_ERROR, - ); + SendData::new( + size, + piece, + tracer_header_tags.clone(), + &config.trace_intake, + ) + }) + .collect(); + + tokio::spawn(async move { + for send_data in send_datas { + if let Err(err) = tx.send(send_data).await { + error!("Error sending traces to the trace flusher: {err}"); + return; + } } - } + }); log_and_create_traces_success_http_response( "Successfully buffered traces to be flushed.", @@ -673,21 +678,19 @@ mod tests { .await; assert!(res.is_ok()); - let mut received = Vec::new(); - while let Ok(send_data) = rx.try_recv() { - received.push(send_data); - } + let send_data = rx + .recv() + .await + .expect("expected the oversized single-chunk trace to be sent standalone"); - assert_eq!( - received.len(), - 1, - "expected the oversized single-chunk trace to be sent standalone as one piece, got {}", - received.len() + assert!( + rx.try_recv().is_err(), + "expected exactly one piece to be sent" ); assert!( - received[0].len() > MAX_CONTENT_SIZE_BYTES, + send_data.len() > MAX_CONTENT_SIZE_BYTES, "expected the standalone piece to still be reported as oversized (size {})", - received[0].len() + send_data.len() ); } } From 358ec20a2be5cbebd9e889c44d7c642b008df51d Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Tue, 8 Sep 2026 16:20:29 -0400 Subject: [PATCH 09/11] Add bounded semaphore --- crates/datadog-serverless-compat/src/main.rs | 6 +-- crates/datadog-trace-agent/src/config.rs | 4 ++ .../src/stats_processor.rs | 1 + .../src/trace_processor.rs | 49 +++++++++++++++---- .../tests/integration_test.rs | 18 +++---- 5 files changed, 54 insertions(+), 24 deletions(-) diff --git a/crates/datadog-serverless-compat/src/main.rs b/crates/datadog-serverless-compat/src/main.rs index bda4aee..cdcdcab 100644 --- a/crates/datadog-serverless-compat/src/main.rs +++ b/crates/datadog-serverless-compat/src/main.rs @@ -184,9 +184,9 @@ pub async fn main() { None }; - let trace_processor = Arc::new(trace_processor::ServerlessTraceProcessor { - stats_concentrator: stats_concentrator.as_ref().map(|c| c.handle.clone()), - }); + let trace_processor = Arc::new(trace_processor::ServerlessTraceProcessor::new( + stats_concentrator.as_ref().map(|c| c.handle.clone()), + )); let stats_flusher = Arc::new(stats_flusher::ServerlessStatsFlusher { stats_concentrator: stats_concentrator.as_ref().map(|c| c.handle.clone()), diff --git a/crates/datadog-trace-agent/src/config.rs b/crates/datadog-trace-agent/src/config.rs index d097208..186ce0a 100644 --- a/crates/datadog-trace-agent/src/config.rs +++ b/crates/datadog-trace-agent/src/config.rs @@ -112,6 +112,8 @@ pub struct Config { pub proxy_request_retry_backoff_base_ms: u64, /// timeout for environment verification, in milliseconds pub verify_env_timeout_ms: u64, + /// How long to wait for a trace-enqueue permit before shedding load, in seconds + pub enqueue_permit_timeout_secs: u64, pub proxy_url: Option, pub env: String, pub peer_tags: Vec, @@ -241,6 +243,7 @@ impl Config { proxy_request_max_retries: 3, proxy_request_retry_backoff_base_ms: 100, verify_env_timeout_ms: 100, + enqueue_permit_timeout_secs: 2, dd_apm_receiver_port, #[cfg(any(all(windows, feature = "windows-pipes"), test))] dd_apm_windows_pipe_name, @@ -919,6 +922,7 @@ pub mod test_helpers { proxy_request_max_retries: 3, proxy_request_retry_backoff_base_ms: 100, verify_env_timeout_ms: 1000, + enqueue_permit_timeout_secs: 2, proxy_url: None, env: "none".to_string(), peer_tags: peer_tag_keys().unwrap(), diff --git a/crates/datadog-trace-agent/src/stats_processor.rs b/crates/datadog-trace-agent/src/stats_processor.rs index 2bf6fa3..990d706 100644 --- a/crates/datadog-trace-agent/src/stats_processor.rs +++ b/crates/datadog-trace-agent/src/stats_processor.rs @@ -123,6 +123,7 @@ mod tests { proxy_request_max_retries: 3, proxy_request_retry_backoff_base_ms: 100, verify_env_timeout_ms: 100, + enqueue_permit_timeout_secs: 2, trace_intake: Endpoint { url: hyper::Uri::from_static("https://trace.agent.notdog.com/traces"), api_key: Some("dummy_api_key".into()), diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index 8d9bddf..79d0154 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -115,12 +115,25 @@ impl TraceChunkProcessor for ChunkProcessor { } } } +/// Maximum number of trace-enqueue tasks that may be in flight (spawned but not yet finished +/// handing their pieces to the flusher) at once +const MAX_IN_FLIGHT_ENQUEUES: usize = 10; + #[derive(Clone)] pub struct ServerlessTraceProcessor { pub stats_concentrator: Option, + enqueue_permits: Arc, } impl ServerlessTraceProcessor { + #[allow(clippy::must_use_candidate)] + pub fn new(stats_concentrator: Option) -> Self { + ServerlessTraceProcessor { + stats_concentrator, + enqueue_permits: Arc::new(tokio::sync::Semaphore::new(MAX_IN_FLIGHT_ENQUEUES)), + } + } + fn send_to_concentrator( concentrator: &StatsConcentratorHandle, payload: &TracerPayloadCollection, @@ -290,7 +303,30 @@ impl TraceProcessor for ServerlessTraceProcessor { }) .collect(); + // Bound how many enqueue tasks can be in flight at once, so a slow/retrying flush + // (which holds the aggregator lock, see aggregator.rs) can't let unbounded spawned + // tasks pile up in memory. If no permit frees up in time, drop the traces rather + // than hold this request open indefinitely. + let permit = match tokio::time::timeout( + std::time::Duration::from_secs(config.enqueue_permit_timeout_secs), + self.enqueue_permits.clone().acquire_owned(), + ) + .await + { + Ok(Ok(permit)) => permit, + Ok(Err(_)) | Err(_) => { + warn!( + "Could not acquire an enqueue permit in time; shedding load, dropping traces" + ); + return log_and_create_traces_success_http_response( + "Successfully buffered traces to be flushed.", + StatusCode::OK, + ); + } + }; + tokio::spawn(async move { + let _permit = permit; // released when this task ends for send_data in send_datas { if let Err(err) = tx.send(send_data).await { error!("Error sending traces to the trace flusher: {err}"); @@ -344,6 +380,7 @@ mod tests { proxy_request_max_retries: 3, proxy_request_retry_backoff_base_ms: 100, verify_env_timeout_ms: 100, + enqueue_permit_timeout_secs: 2, trace_intake: Endpoint { url: hyper::Uri::from_static("https://trace.agent.notdog.com/traces"), api_key: Some("dummy_api_key".into()), @@ -502,9 +539,7 @@ mod tests { .body(http_common::Body::from(bytes)) .unwrap(); - let trace_processor = trace_processor::ServerlessTraceProcessor { - stats_concentrator: None, - }; + let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); let res = trace_processor .process_traces( Arc::new(create_test_config()), @@ -577,9 +612,7 @@ mod tests { .body(http_common::Body::from(bytes)) .unwrap(); - let trace_processor = trace_processor::ServerlessTraceProcessor { - stats_concentrator: None, - }; + let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); let res = trace_processor .process_traces( Arc::new(create_test_config()), @@ -665,9 +698,7 @@ mod tests { .body(http_common::Body::from(bytes)) .unwrap(); - let trace_processor = trace_processor::ServerlessTraceProcessor { - stats_concentrator: None, - }; + let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); let res = trace_processor .process_traces( Arc::new(create_test_config()), diff --git a/crates/datadog-trace-agent/tests/integration_test.rs b/crates/datadog-trace-agent/tests/integration_test.rs index 4e528ac..cb0117b 100644 --- a/crates/datadog-trace-agent/tests/integration_test.rs +++ b/crates/datadog-trace-agent/tests/integration_test.rs @@ -97,9 +97,9 @@ pub fn create_mini_agent_with_real_flushers( let aggregator = Arc::new(tokio::sync::Mutex::new(TraceAggregator::default())); let mini_agent = MiniAgent { config: config.clone(), - trace_processor: Arc::new(ServerlessTraceProcessor { - stats_concentrator: Some(stats_concentrator_handle.clone()), - }), + trace_processor: Arc::new(ServerlessTraceProcessor::new(Some( + stats_concentrator_handle.clone(), + ))), trace_flusher: Arc::new(ServerlessTraceFlusher::new( aggregator.clone(), config.clone(), @@ -249,9 +249,7 @@ async fn test_mini_agent_tcp_handles_requests() { let test_port = config.dd_apm_receiver_port; let mini_agent = MiniAgent { config: config.clone(), - trace_processor: Arc::new(ServerlessTraceProcessor { - stats_concentrator: None, - }), + trace_processor: Arc::new(ServerlessTraceProcessor::new(None)), trace_flusher: Arc::new(MockTraceFlusher), stats_processor: Arc::new(MockStatsProcessor), stats_flusher: Arc::new(MockStatsFlusher), @@ -365,9 +363,7 @@ async fn test_mini_agent_named_pipe_handles_requests() { let mini_agent = MiniAgent { config: config.clone(), - trace_processor: Arc::new(ServerlessTraceProcessor { - stats_concentrator: None, - }), + trace_processor: Arc::new(ServerlessTraceProcessor::new(None)), trace_flusher: Arc::new(MockTraceFlusher), stats_processor: Arc::new(MockStatsProcessor), stats_flusher: Arc::new(MockStatsFlusher), @@ -539,9 +535,7 @@ async fn test_mini_agent_tcp_proxies_dsm_requests() { let mini_agent = MiniAgent { config: config.clone(), - trace_processor: Arc::new(ServerlessTraceProcessor { - stats_concentrator: None, - }), + trace_processor: Arc::new(ServerlessTraceProcessor::new(None)), trace_flusher: Arc::new(MockTraceFlusher), stats_processor: Arc::new(MockStatsProcessor), stats_flusher: Arc::new(MockStatsFlusher), From da4892546553ef8502fa11f6f932ca0ce4a84528 Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Tue, 8 Sep 2026 17:36:19 -0400 Subject: [PATCH 10/11] Acquire the enqueue permit before decoding --- .../src/trace_processor.rs | 188 +++++++++++++++--- 1 file changed, 163 insertions(+), 25 deletions(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index 79d0154..ca6a740 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -198,6 +198,24 @@ impl TraceProcessor for ServerlessTraceProcessor { return response; } + // Bound how many requests can be in the decode/enrich/split/enqueue pipeline at once + // and carried through to the spawned enqueue task at the end + let permit = match tokio::time::timeout( + std::time::Duration::from_secs(config.enqueue_permit_timeout_secs), + self.enqueue_permits.clone().acquire_owned(), + ) + .await + { + Ok(Ok(permit)) => permit, + Ok(Err(_)) | Err(_) => { + warn!("Could not acquire an enqueue permit in time; dropping traces"); + return log_and_create_traces_success_http_response( + "Dropped traces due to enqueue capacity", + StatusCode::OK, + ); + } + }; + let tracer_header_tags = (&parts.headers).into(); // deserialize traces from the request body, convert to protobuf structs (see trace-protobuf @@ -303,28 +321,6 @@ impl TraceProcessor for ServerlessTraceProcessor { }) .collect(); - // Bound how many enqueue tasks can be in flight at once, so a slow/retrying flush - // (which holds the aggregator lock, see aggregator.rs) can't let unbounded spawned - // tasks pile up in memory. If no permit frees up in time, drop the traces rather - // than hold this request open indefinitely. - let permit = match tokio::time::timeout( - std::time::Duration::from_secs(config.enqueue_permit_timeout_secs), - self.enqueue_permits.clone().acquire_owned(), - ) - .await - { - Ok(Ok(permit)) => permit, - Ok(Err(_)) | Err(_) => { - warn!( - "Could not acquire an enqueue permit in time; shedding load, dropping traces" - ); - return log_and_create_traces_success_http_response( - "Successfully buffered traces to be flushed.", - StatusCode::OK, - ); - } - }; - tokio::spawn(async move { let _permit = permit; // released when this task ends for send_data in send_datas { @@ -354,14 +350,15 @@ mod tests { config::{Config, Tags}, peer_tags::peer_tag_keys, trace_processor::{ - self, TRACER_PAYLOAD_FUNCTION_TAGS_TAG_KEY, TraceProcessor, encoded_size, - split_oversized_payloads, + self, MAX_IN_FLIGHT_ENQUEUES, TRACER_PAYLOAD_FUNCTION_TAGS_TAG_KEY, TraceProcessor, + encoded_size, split_oversized_payloads, }, }; use libdd_common::{Endpoint, http_common}; use libdd_trace_protobuf::pb; use libdd_trace_utils::test_utils::{create_test_gcp_json_span, create_test_gcp_span}; - use libdd_trace_utils::trace_utils::MiniAgentMetadata; + use libdd_trace_utils::trace_utils::{MiniAgentMetadata, SendData}; + use libdd_trace_utils::tracer_header_tags::TracerHeaderTags; use libdd_trace_utils::{ test_utils::create_test_json_span, trace_utils, tracer_payload::TracerPayloadCollection, }; @@ -724,4 +721,145 @@ mod tests { send_data.len() ); } + + fn small_trace_request() -> http_common::HttpRequest { + let start = get_current_timestamp_nanos(); + let json_span = create_test_json_span(11, 222, 333, start, false); + let bytes = rmp_serde::to_vec(&vec![vec![json_span]]).unwrap(); + Request::builder() + .header("datadog-meta-tracer-version", "4.0.0") + .header("datadog-meta-lang", "nodejs") + .header("datadog-meta-lang-version", "v19.7.0") + .header("datadog-meta-lang-interpreter", "v8") + .header("datadog-container-id", "33") + .header("content-length", "100") + .body(http_common::Body::from(bytes)) + .unwrap() + } + + fn dummy_send_data(config: &Config) -> trace_utils::SendData { + SendData::new( + 1, + TracerPayloadCollection::V07(Vec::new()), + TracerHeaderTags::default(), + &config.trace_intake, + ) + } + + #[tokio::test] + async fn test_enqueue_permits_bound_in_flight_tasks_and_shed_load() { + let (tx, mut rx): ( + Sender, + Receiver, + ) = mpsc::channel(1); + + let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); + let config = Arc::new(Config { + enqueue_permit_timeout_secs: 0, + ..create_test_config() + }); + + // Pre-fill the channel's one slot so every enqueue attempt below blocks on tx.send() + // (and holds its permit) instead of succeeding immediately. + tx.try_send(dummy_send_data(&config)).unwrap(); + + // Saturate all MAX_IN_FLIGHT_ENQUEUES permits: each call acquires a permit and spawns + // a task that then blocks forever on tx.send(), since nothing drains rx yet. Permit + // acquisition happens synchronously inside process_traces before it returns, so by the + // time this loop finishes, all permits are deterministically held - no race. + for _ in 0..MAX_IN_FLIGHT_ENQUEUES { + let res = trace_processor + .process_traces( + config.clone(), + small_trace_request(), + tx.clone(), + Arc::new(create_test_metadata()), + ) + .await; + assert!(res.is_ok()); + } + assert_eq!( + trace_processor.enqueue_permits.available_permits(), + 0, + "expected all permits to be held after saturating requests" + ); + + // All permits are held, so this request can't get one in time (timeout is 0) and + // should shed load rather than hang - still a 200, but its trace is dropped before + // ever reaching tx.send() (permit acquisition happens before decode). + let res = trace_processor + .process_traces( + config.clone(), + small_trace_request(), + tx.clone(), + Arc::new(create_test_metadata()), + ) + .await; + assert!(res.is_ok()); + + // Drop the original sender - the shed 11th request never created a clone of its own, + // since it sheds before decoding/building anything. The channel will only close once + // every clone held by the 10 blocked tasks is also dropped, which happens as each one + // is unblocked in turn by draining the item ahead of it. + drop(tx); + + let mut received = 0; + while rx.recv().await.is_some() { + received += 1; + } + assert_eq!( + received, + MAX_IN_FLIGHT_ENQUEUES + 1, + "expected exactly the dummy plus the 10 saturating payloads, nothing from the shed 11th" + ); + } + + #[tokio::test] + async fn test_enqueue_permit_releases_after_send_completes() { + let (tx, mut rx): ( + Sender, + Receiver, + ) = mpsc::channel(MAX_IN_FLIGHT_ENQUEUES + 1); + + let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); + // Uses the default (non-zero) enqueue_permit_timeout_secs, so the 11th request's permit + // acquire has real time to succeed once an earlier request's send completes and + // releases its permit, rather than racing a 0-second timeout against task scheduling. + let config = Arc::new(create_test_config()); + + // With enough channel capacity for every send to complete immediately, all + // MAX_IN_FLIGHT_ENQUEUES requests should succeed and their permits should be released + // right away - leaving room for one more request to succeed too, not shed. + for _ in 0..=MAX_IN_FLIGHT_ENQUEUES { + let res = trace_processor + .process_traces( + config.clone(), + small_trace_request(), + tx.clone(), + Arc::new(create_test_metadata()), + ) + .await; + assert!(res.is_ok()); + // Give the previous request's spawned enqueue task a chance to run its (instant, + // since the channel has room) send and release its permit before the next request + // tries to acquire one. + tokio::task::yield_now().await; + } + + // Drop the original sender so the channel closes once every clone held by the 11 + // completed send tasks is also dropped, then drain to closure for a deterministic + // count instead of a try_recv() snapshot that could race a still-completing task. + drop(tx); + + let mut received = 0; + while rx.recv().await.is_some() { + received += 1; + } + assert_eq!( + received, + MAX_IN_FLIGHT_ENQUEUES + 1, + "expected every request's permit to be released after its send completed, \ + allowing all of them to succeed rather than shedding load" + ); + } } From 11a6fad85bf7fa761f22c4e7d56a9b24b9cbcee4 Mon Sep 17 00:00:00 2001 From: Kathie Huang Date: Tue, 8 Sep 2026 18:26:08 -0400 Subject: [PATCH 11/11] Collect body if trace gets dropped --- crates/datadog-trace-agent/src/trace_processor.rs | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index ca6a740..e26465c 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -4,6 +4,7 @@ use std::sync::Arc; use async_trait::async_trait; +use http_body_util::BodyExt; use hyper::{StatusCode, http}; use libdd_common::http_common; use libdd_library_config::tracer_metadata::TracerMetadata; @@ -208,7 +209,12 @@ impl TraceProcessor for ServerlessTraceProcessor { { Ok(Ok(permit)) => permit, Ok(Err(_)) | Err(_) => { + // The traces are dropped. The body is still drained so the connection can be + // kept alive for the next request instead of being closed. warn!("Could not acquire an enqueue permit in time; dropping traces"); + if let Err(err) = body.collect().await { + debug!("Error draining /v0.4/traces request body while dropping traces: {err}"); + } return log_and_create_traces_success_http_response( "Dropped traces due to enqueue capacity", StatusCode::OK,