Skip to content
Merged
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 3 additions & 3 deletions crates/datadog-serverless-compat/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()),
Expand Down
1 change: 1 addition & 0 deletions crates/datadog-trace-agent/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
46 changes: 37 additions & 9 deletions crates/datadog-trace-agent/src/aggregator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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);
}
}
4 changes: 4 additions & 0 deletions crates/datadog-trace-agent/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String>,
pub env: String,
pub peer_tags: Vec<String>,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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(),
Expand Down
1 change: 1 addition & 0 deletions crates/datadog-trace-agent/src/stats_processor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()),
Expand Down
Loading
Loading