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-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/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 b9a5b0d..eb8eed7 100644 --- a/crates/datadog-trace-agent/src/aggregator.rs +++ b/crates/datadog-trace-agent/src/aggregator.rs @@ -47,19 +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 stats back in the queue - if batch_size + payload_size > self.max_content_size_bytes { + 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); - break; } - batch_size += payload_size; - self.buffer.push(payload); - } else { break; } + batch_size += payload_size; + self.buffer.push(payload); } std::mem::take(&mut self.buffer) @@ -141,4 +146,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); + } } 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 361c66a..e26465c 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -4,19 +4,22 @@ 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; 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 +27,53 @@ 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. +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; + + // 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); + result.extend(split_tracer_payload(second, max_size)); + return result; + } + + vec![tp] +} + +/// Computes the total encoded size of the inner TracerPayloads +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 @@ -66,12 +116,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, @@ -136,6 +199,29 @@ 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(_) => { + // 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, + ); + } + }; + let tracer_header_tags = (&parts.headers).into(); // deserialize traces from the request body, convert to protobuf structs (see trace-protobuf @@ -196,23 +282,65 @@ impl TraceProcessor for ServerlessTraceProcessor { Self::send_to_concentrator(concentrator, &payload); } - let send_data = SendData::new(body_size, payload, tracer_header_tags, &config.trace_intake); - - // 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, - ); - } - Err(err) => { - return log_and_create_http_response( - &format!("Error sending traces to the trace flusher: {err}"), - StatusCode::INTERNAL_SERVER_ERROR, - ); + let pieces: Vec<(TracerPayloadCollection, usize)> = match payload { + TracerPayloadCollection::V07(payloads) => { + 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)) + V07_ENVELOPE_OVERHEAD_BYTES; + (TracerPayloadCollection::V07(vec![tp]), size) + }) + .collect() } + other => vec![(other, body_size)], + }; + + if pieces.len() > 1 { + debug!( + piece_count = pieces.len(), + "Oversized trace payload split into multiple pieces" + ); } + + 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" + ); + } + + SendData::new( + size, + piece, + tracer_header_tags.clone(), + &config.trace_intake, + ) + }) + .collect(); + + 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}"); + return; + } + } + }); + + log_and_create_traces_success_http_response( + "Successfully buffered traces to be flushed.", + StatusCode::OK, + ) } } @@ -224,14 +352,19 @@ 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, 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, }; @@ -250,6 +383,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()), @@ -303,6 +437,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): ( @@ -325,9 +542,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()), @@ -400,9 +615,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()), @@ -453,4 +666,206 @@ 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 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(); + 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() + .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::new(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 send_data = rx + .recv() + .await + .expect("expected the oversized single-chunk trace to be sent standalone"); + + assert!( + rx.try_recv().is_err(), + "expected exactly one piece to be sent" + ); + assert!( + send_data.len() > MAX_CONTENT_SIZE_BYTES, + "expected the standalone piece to still be reported as oversized (size {})", + 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" + ); + } } 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),