From e94cbee12bcc6b36ec45a6d9533d3b5fb35db3ef Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Tue, 29 Sep 2026 13:07:01 -0600 Subject: [PATCH 01/15] feat(datadog_logs): optionally truncate oversized log events Adds opt-in encoded-log and retained-message limits to the Datadog logs sink. When enabled, this preserves standard fields, mark shortened messages, and shrinks encoded messages to fit configured payload limits. --- ...dog_logs_truncate_oversized.enhancement.md | 7 + .../src/internal_event/metric_name.rs | 2 + src/internal_events/datadog_logs.rs | 20 + src/sinks/datadog/logs/config.rs | 121 +++- src/sinks/datadog/logs/mod.rs | 3 + src/sinks/datadog/logs/sink.rs | 604 ++++++++++++++++-- src/sinks/datadog/logs/tests.rs | 66 +- .../components/sinks/datadog_logs.cue | 11 +- .../sinks/generated/datadog_logs.cue | 28 +- 9 files changed, 788 insertions(+), 74 deletions(-) create mode 100644 changelog.d/datadog_logs_truncate_oversized.enhancement.md diff --git a/changelog.d/datadog_logs_truncate_oversized.enhancement.md b/changelog.d/datadog_logs_truncate_oversized.enhancement.md new file mode 100644 index 0000000000000..f8394eb4054f5 --- /dev/null +++ b/changelog.d/datadog_logs_truncate_oversized.enhancement.md @@ -0,0 +1,7 @@ +The `datadog_logs` sink can now optionally truncate logs that exceed Datadog's +per-log size limit. Configure `truncate_oversized_logs` to set the encoded log +and retained message limits, mark shortened messages, tag reduced logs, and +preserve standard Datadog fields. Logs that cannot be reduced below the +configured limit are dropped. + +authors: bruceg diff --git a/lib/vector-common/src/internal_event/metric_name.rs b/lib/vector-common/src/internal_event/metric_name.rs index b05ebb651ea02..abfd82c49a400 100644 --- a/lib/vector-common/src/internal_event/metric_name.rs +++ b/lib/vector-common/src/internal_event/metric_name.rs @@ -109,6 +109,7 @@ pub enum CounterName { MemoryEnrichmentTableTtlExpirations, MemoryEnrichmentTableTtlExpirationsTotal, ComponentCpuUsageNsTotal, + DatadogLogsEventsTruncatedTotal, DatadogLogsReservedAttributeConflictsTotal, // Data-plane counter names emitted by the `host_metrics` source. CpuSecondsTotal, @@ -503,6 +504,7 @@ impl CounterName { "memory_enrichment_table_ttl_expirations_total" } Self::ComponentCpuUsageNsTotal => "component_cpu_usage_ns_total", + Self::DatadogLogsEventsTruncatedTotal => "datadog_logs_events_truncated_total", Self::DatadogLogsReservedAttributeConflictsTotal => { "datadog_logs_reserved_attribute_conflicts_total" } diff --git a/src/internal_events/datadog_logs.rs b/src/internal_events/datadog_logs.rs index 49c5ff92421f3..242be9b2aa216 100644 --- a/src/internal_events/datadog_logs.rs +++ b/src/internal_events/datadog_logs.rs @@ -4,6 +4,26 @@ use vector_lib::{ }; use vrl::path::OwnedTargetPath; +#[derive(Debug, NamedInternalEvent)] +pub struct DatadogLogsEventTruncated { + pub max_log_bytes: usize, + pub max_message_bytes: usize, + pub original_encoded_size: usize, +} + +impl InternalEvent for DatadogLogsEventTruncated { + fn emit(self) { + warn!( + message = "Truncated a Datadog log event that exceeded the per-log size limit.", + max_log_bytes = self.max_log_bytes, + max_message_bytes = self.max_message_bytes, + original_encoded_size = self.original_encoded_size, + internal_log_rate_limit = true, + ); + counter!(CounterName::DatadogLogsEventsTruncatedTotal).increment(1); + } +} + #[derive(Debug, NamedInternalEvent)] pub struct DatadogLogsReservedAttributeConflict<'a> { pub meaning: &'static str, diff --git a/src/sinks/datadog/logs/config.rs b/src/sinks/datadog/logs/config.rs index 0df5afc82954c..070251b6c6e59 100644 --- a/src/sinks/datadog/logs/config.rs +++ b/src/sinks/datadog/logs/config.rs @@ -34,6 +34,8 @@ use crate::{ // of escaped double-quotes -- but we believe this should be very rare in // practice. pub const MAX_PAYLOAD_BYTES: usize = 5_000_000; +pub(super) const DEFAULT_MAX_LOG_BYTES: usize = 1_000_000; +pub(super) const DEFAULT_MAX_MESSAGE_BYTES: usize = 900_000; pub const BATCH_HEADROOM_BYTES: usize = 750_000; pub const BATCH_MAX_EVENTS: usize = 1_000; pub const BATCH_DEFAULT_TIMEOUT_SECS: f64 = 5.0; @@ -48,6 +50,23 @@ impl SinkBatchSettings for DatadogLogsDefaultBatchSettings { const TIMEOUT_SECS: f64 = BATCH_DEFAULT_TIMEOUT_SECS; } +/// Options for truncating logs that exceed Datadog's per-log size limit. +#[configurable_component] +#[derive(Clone, Copy, Debug, Derivative)] +#[derivative(Default)] +#[serde(deny_unknown_fields)] +pub struct DatadogLogsTruncationConfig { + /// Maximum encoded size, in bytes, of a log before truncation is applied. + #[derivative(Default(value = "default_max_log_bytes()"))] + #[serde(default = "default_max_log_bytes")] + pub max_log_bytes: usize, + + /// Maximum number of message bytes to retain before appending the truncation marker. + #[derivative(Default(value = "default_max_message_bytes()"))] + #[serde(default = "default_max_message_bytes")] + pub max_message_bytes: usize, +} + /// Configuration for the `datadog_logs` sink. #[configurable_component(sink("datadog_logs", "Publish log events to Datadog."))] #[derive(Clone, Debug, Derivative)] @@ -81,17 +100,34 @@ pub struct DatadogLogsConfig { /// to not set it above 5,000,000 (5 MB, the standard Datadog API limit). Increase /// this when targeting a compatible endpoint that accepts larger payloads. The batch /// goal is derived as `max_payload_bytes - 750,000` bytes; events larger than the - /// batch goal are sent alone in their batch. Events exceeding `max_payload_bytes` are - /// dropped. + /// batch goal are sent alone in their batch. Single events that still exceed + /// `max_payload_bytes` after optional truncation are dropped. #[derivative(Default(value = "default_max_payload_bytes()"))] #[serde(default = "default_max_payload_bytes")] pub max_payload_bytes: Option, + + /// Truncate logs whose encoded JSON exceeds `max_log_bytes`. + /// + /// When a message is shortened, at most `max_message_bytes` raw bytes are retained and + /// `...TRUNCATED...` is appended. Every reduced log is tagged with `truncated:single_line`. + /// If the log remains oversized after the initial message cap, non-standard fields are removed + /// before the message is shortened further to account for JSON encoding. Logs that still + /// exceed the limit, or have no string message to truncate, are dropped. + pub truncate_oversized_logs: Option, } const fn default_max_payload_bytes() -> Option { Some(MAX_PAYLOAD_BYTES) } +const fn default_max_log_bytes() -> usize { + DEFAULT_MAX_LOG_BYTES +} + +const fn default_max_message_bytes() -> usize { + DEFAULT_MAX_MESSAGE_BYTES +} + const fn default_compression() -> Option { Some(Compression::zstd_default()) } @@ -175,6 +211,7 @@ impl DatadogLogsConfig { self.max_payload_bytes.unwrap_or(MAX_PAYLOAD_BYTES), ) .compression(self.compression.or_else(default_compression).unwrap()) + .truncation(self.truncate_oversized_logs) .build(); Ok(VectorSink::from_event_streamsink(sink)) @@ -253,6 +290,27 @@ impl ValidatedSink for DatadogLogsConfig { .into()); } + if let Some(truncation) = self.truncate_oversized_logs { + let max_payload_bytes = self.max_payload_bytes.unwrap_or(MAX_PAYLOAD_BYTES); + // A single log is wrapped in `[` and `]`, so leave two bytes for the JSON array. + let maximum_log_bytes = max_payload_bytes - 2; + if !(1..=maximum_log_bytes).contains(&truncation.max_log_bytes) { + return Err(format!( + "truncate_oversized_logs.max_log_bytes ({}) must be between 1 and {}", + truncation.max_log_bytes, maximum_log_bytes + ) + .into()); + } + let maximum_message_bytes = truncation.max_log_bytes - 1; + if !(1..=maximum_message_bytes).contains(&truncation.max_message_bytes) { + return Err(format!( + "truncate_oversized_logs.max_message_bytes ({}) must be between 1 and {}", + truncation.max_message_bytes, maximum_message_bytes + ) + .into()); + } + } + let batch_goal_bytes = self.max_payload_bytes.unwrap_or(MAX_PAYLOAD_BYTES) - BATCH_HEADROOM_BYTES; @@ -309,6 +367,65 @@ mod test { crate::test_util::test_generate_config::(); } + #[test] + fn truncate_oversized_logs_accepts_empty_configuration() { + let config = serde_yaml::from_str::(indoc::indoc! {r#" + default_api_key: "test_key" + truncate_oversized_logs: {} + "#}) + .expect("truncation configuration should deserialize"); + + let truncation = config + .truncate_oversized_logs + .expect("truncation should be enabled"); + assert_eq!(truncation.max_log_bytes, 1_000_000); + assert_eq!(truncation.max_message_bytes, 900_000); + } + + #[test] + fn validate_rejects_invalid_truncation_limits() { + for max_message_bytes in [0, DEFAULT_MAX_LOG_BYTES] { + let config: DatadogLogsConfig = serde_yaml::from_str(&format!( + r#" + default_api_key: "test_key" + truncate_oversized_logs: + max_message_bytes: {max_message_bytes} + "# + )) + .unwrap(); + + assert!(config.validate().is_err()); + } + } + + #[test] + fn validate_rejects_truncation_limit_without_array_wrapper_room() { + let config: DatadogLogsConfig = serde_yaml::from_str(indoc::indoc! {r#" + default_api_key: "test_key" + max_payload_bytes: 800000 + truncate_oversized_logs: + max_log_bytes: 799999 + max_message_bytes: 700000 + "#}) + .unwrap(); + + assert!(config.validate().is_err()); + } + + #[test] + fn validate_accepts_custom_log_limit_above_datadog_default() { + let config: DatadogLogsConfig = serde_yaml::from_str(indoc::indoc! {r#" + default_api_key: "test_key" + max_payload_bytes: 3000000 + truncate_oversized_logs: + max_log_bytes: 2000000 + max_message_bytes: 1800000 + "#}) + .unwrap(); + + assert!(config.validate().is_ok()); + } + #[test] fn validate_produces_usable_batch_settings() { let config = DatadogLogsConfig::default(); diff --git a/src/sinks/datadog/logs/mod.rs b/src/sinks/datadog/logs/mod.rs index 8e3b69c66f85b..a95d3e4d5de53 100644 --- a/src/sinks/datadog/logs/mod.rs +++ b/src/sinks/datadog/logs/mod.rs @@ -5,10 +5,13 @@ //! Datadog Log API. The log API is relatively generous in terms of its //! constraints, except that: //! +//! * an individual log over 1 MB is accepted but truncated //! * a 'payload' is comprised of no more than 1,000 array members //! * a 'payload' may not be more than 5Mb in size, uncompressed and //! * a 'payload' may not mix API keys //! +//! The sink can optionally perform per-log truncation before sending. +//! //! Otherwise per [the //! docs](https://docs.datadoghq.com/api/latest/logs/#send-logs) there aren't //! other major constraints we have to follow in this implementation. The sink diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index ffe4ee405c486..6928c1ad96e4d 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -1,5 +1,6 @@ use std::{collections::VecDeque, fmt::Debug, io, sync::Arc}; +use bytes::Bytes; use itertools::Itertools; use snafu::Snafu; use tracing::Instrument; @@ -10,10 +11,10 @@ use vector_lib::{ }; use vrl::path::{OwnedSegment, OwnedTargetPath, PathPrefix}; -use super::service::LogApiRequest; +use super::{config::DatadogLogsTruncationConfig, service::LogApiRequest}; use crate::{ common::datadog::{DD_RESERVED_SEMANTIC_ATTRS, DDTAGS, MESSAGE, is_reserved_attribute}, - internal_events::DatadogLogsReservedAttributeConflict, + internal_events::{DatadogLogsEventTruncated, DatadogLogsReservedAttributeConflict}, sinks::{ prelude::*, util::{Compressor, http::HttpJsonBatchSizer}, @@ -41,6 +42,7 @@ pub struct LogSinkBuilder { protocol: String, conforms_as_agent: bool, max_payload_bytes: usize, + truncation: Option, } impl LogSinkBuilder { @@ -62,6 +64,7 @@ impl LogSinkBuilder { protocol, conforms_as_agent, max_payload_bytes, + truncation: None, } } @@ -70,6 +73,11 @@ impl LogSinkBuilder { self } + pub const fn truncation(mut self, truncation: Option) -> Self { + self.truncation = truncation; + self + } + pub fn build(self) -> LogSink { LogSink { default_api_key: self.default_api_key, @@ -80,6 +88,7 @@ impl LogSinkBuilder { protocol: self.protocol, conforms_as_agent: self.conforms_as_agent, max_payload_bytes: self.max_payload_bytes, + truncation: self.truncation, } } } @@ -106,6 +115,8 @@ pub struct LogSink { conforms_as_agent: bool, /// Maximum uncompressed payload size in bytes max_payload_bytes: usize, + /// Limits used when oversized log truncation is enabled. + truncation: Option, } // The Datadog logs intake does not require the fields that are set in this @@ -254,6 +265,7 @@ struct LogRequestBuilder { pub compression: Compression, pub conforms_as_agent: bool, pub max_payload_bytes: usize, + pub truncation: Option, } impl LogRequestBuilder { @@ -280,24 +292,67 @@ impl LogRequestBuilder { let mut requests: Vec = Vec::new(); while !events_with_estimated_size.is_empty() { let (events_serialized, body, byte_size) = - serialize_with_capacity(&mut events_with_estimated_size, self.max_payload_bytes)?; + self.serialize_with_capacity(&mut events_with_estimated_size)?; if events_serialized.is_empty() { - // first event was too large for whole request - let _too_big = events_with_estimated_size.pop_front(); - emit!(ComponentEventsDropped:: { - count: 1, - reason: "Event too large to encode." - }); - } else { - let request = - self.finish_request(body, events_serialized, byte_size, Arc::clone(&api_key))?; - requests.push(request); + if events_with_estimated_size.pop_front().is_some() { + emit!(ComponentEventsDropped:: { + count: 1, + reason: "Event too large to encode." + }); + } + continue; } + let request = + self.finish_request(body, events_serialized, byte_size, Arc::clone(&api_key))?; + requests.push(request); } Ok(requests) } + fn serialize_with_capacity( + &self, + events: &mut VecDeque<(Event, JsonSize)>, + ) -> Result<(Vec, Vec, GroupedCountByteSize), io::Error> { + let total_estimated = + events.iter().map(|(_, size)| size.get()).sum::() + events.len() * 2; + let mut buf = Vec::with_capacity(total_estimated); + let mut byte_size = telemetry().create_request_count_byte_size(); + let mut events_serialized = Vec::with_capacity(events.len()); + + buf.push(b'['); + while let Some((mut event, mut estimated_json_size)) = events.pop_front() { + let existing_len = buf.len(); + match encode_log( + &mut buf, + &mut event, + !events_serialized.is_empty(), + self.truncation, + self.conforms_as_agent, + )? { + LogEncoding::Unchanged => {} + LogEncoding::Truncated => { + estimated_json_size = event.estimated_json_encoded_size_of(); + } + LogEncoding::Dropped { reason } => { + emit!(ComponentEventsDropped:: { count: 1, reason }); + continue; + } + } + + if buf.len() >= self.max_payload_bytes { + events.push_front((event, estimated_json_size)); + buf.truncate(existing_len); + break; + } + byte_size.add_event(&event, estimated_json_size); + events_serialized.push(event); + } + buf.push(b']'); + + Ok((events_serialized, buf, byte_size)) + } + fn finish_request( &self, buf: Vec, @@ -333,49 +388,142 @@ impl LogRequestBuilder { } } -/// Serialize events into a buffer as a JSON array that has a maximum size of -/// `max_payload_bytes`. -/// -/// Returns the serialized events, the buffer, and the byte size of the events. -/// Events that are not serialized remain in the `events` parameter. -pub fn serialize_with_capacity( - events: &mut VecDeque<(Event, JsonSize)>, - max_payload_bytes: usize, -) -> Result<(Vec, Vec, GroupedCountByteSize), io::Error> { - // Compute estimated size, accounting for the size of the brackets and commas. - let total_estimated = - events.iter().map(|(_, size)| size.get()).sum::() + events.len() * 2; - - // Initialize state. - let mut buf = Vec::with_capacity(total_estimated); - let mut byte_size = telemetry().create_request_count_byte_size(); - let mut events_serialized = Vec::with_capacity(events.len()); - - // Write entries until the buffer is full. - buf.push(b'['); - let mut first = true; - while let Some((event, estimated_json_size)) = events.pop_front() { - // Track the existing length of the buffer so we can truncate it if we need to. - let existing_len = buf.len(); - if first { - first = false; - } else { - buf.push(b','); - } - serde_json::to_writer(&mut buf, event.as_log())?; - // If the buffer is too big, truncate it and break out of the loop. - if buf.len() >= max_payload_bytes { - events.push_front((event, estimated_json_size)); - buf.truncate(existing_len); - break; +enum LogEncoding { + Unchanged, + Truncated, + Dropped { reason: &'static str }, +} + +fn encode_log( + buf: &mut Vec, + event: &mut Event, + include_comma: bool, + truncation: Option, + conforms_as_agent: bool, +) -> Result { + let existing_len = buf.len(); + let original_encoded_size = write_log(buf, event, include_comma)?; + let Some(truncation) = + truncation.filter(|truncation| original_encoded_size > truncation.max_log_bytes) + else { + return Ok(LogEncoding::Unchanged); + }; + + let Some(message) = message_bytes_mut(event.as_mut_log()).map(|message| message.clone()) else { + buf.truncate(existing_len); + return Ok(LogEncoding::Dropped { + reason: "Oversized event has no string message to truncate.", + }); + }; + let mut body_len = + floor_char_boundary(&message, message.len().min(truncation.max_message_bytes)); + if body_len < message.len() { + set_truncated_message(event.as_mut_log(), &message, body_len); + } + ensure_truncated_tag(event.as_mut_log()); + + let mut encoded_size = rewrite_log(buf, existing_len, event, include_comma)?; + if encoded_size > truncation.max_log_bytes { + strip_non_standard_fields(event.as_mut_log(), conforms_as_agent); + encoded_size = rewrite_log(buf, existing_len, event, include_comma)?; + } + + if encoded_size > truncation.max_log_bytes { + let marker_bytes_to_add = TRUNCATION_MARKER.len() * usize::from(body_len == message.len()); + let bytes_to_remove = (encoded_size - truncation.max_log_bytes) + marker_bytes_to_add; + body_len = floor_char_boundary(&message, body_len.saturating_sub(bytes_to_remove)); + set_truncated_message(event.as_mut_log(), &message, body_len); + encoded_size = rewrite_log(buf, existing_len, event, include_comma)?; + } + + if encoded_size > truncation.max_log_bytes { + buf.truncate(existing_len); + return Ok(LogEncoding::Dropped { + reason: "Event remains too large after truncation.", + }); + } + emit!(DatadogLogsEventTruncated { + max_log_bytes: truncation.max_log_bytes, + max_message_bytes: truncation.max_message_bytes, + original_encoded_size, + }); + Ok(LogEncoding::Truncated) +} + +fn write_log(buf: &mut Vec, event: &Event, include_comma: bool) -> Result { + if include_comma { + buf.push(b','); + } + let object_start = buf.len(); + serde_json::to_writer(&mut *buf, event.as_log())?; + Ok(buf.len() - object_start) +} + +fn rewrite_log( + buf: &mut Vec, + existing_len: usize, + event: &Event, + include_comma: bool, +) -> Result { + buf.truncate(existing_len); + write_log(buf, event, include_comma) +} + +const TRUNCATION_MARKER: &str = "...TRUNCATED..."; +const TRUNCATED_TAG: &str = "truncated:single_line"; + +fn set_truncated_message(log: &mut LogEvent, message: &Bytes, body_len: usize) { + let marker = TRUNCATION_MARKER.as_bytes(); + let mut truncated = Vec::with_capacity(body_len + marker.len()); + truncated.extend_from_slice(&message[..body_len]); + truncated.extend_from_slice(marker); + *message_bytes_mut(log).expect("the message was previously found") = Bytes::from(truncated); +} + +fn ensure_truncated_tag(log: &mut LogEvent) { + let tags_path = event_path!(DDTAGS); + if let Some(tags) = log.get(tags_path).and_then(Value::as_bytes) + && !tags.is_empty() + { + let tags = String::from_utf8_lossy(tags); + if tags.split(',').any(|tag| tag.trim() == TRUNCATED_TAG) { + return; } - // Otherwise, track the size of the event and continue. - byte_size.add_event(&event, estimated_json_size); - events_serialized.push(event); + log.insert(tags_path, format!("{tags},{TRUNCATED_TAG}")); + } else { + log.insert(tags_path, TRUNCATED_TAG); } - buf.push(b']'); +} - Ok((events_serialized, buf, byte_size)) +fn message_bytes_mut(log: &mut LogEvent) -> Option<&mut Bytes> { + match log.as_map_mut()?.get_mut(MESSAGE)? { + Value::Bytes(message) => Some(message), + Value::Object(fields) => match fields.get_mut(MESSAGE)? { + Value::Bytes(message) => Some(message), + _ => None, + }, + _ => None, + } +} + +fn strip_non_standard_fields(log: &mut LogEvent, conforms_as_agent: bool) { + let Some(fields) = log.as_map_mut() else { + return; + }; + fields.retain(|field, _| is_reserved_attribute(field.as_str()) || field.as_str() == MESSAGE); + if conforms_as_agent && let Some(Value::Object(nested)) = fields.get_mut(MESSAGE) { + nested.retain(|field, _| field.as_str() == MESSAGE); + } +} + +fn floor_char_boundary(bytes: &[u8], mut index: usize) -> usize { + if index >= bytes.len() { + return bytes.len(); + } + while index > 0 && (bytes[index] & 0b1100_0000) == 0b1000_0000 { + index -= 1; + } + index } impl LogSink @@ -396,6 +544,7 @@ where compression: self.compression, conforms_as_agent: self.conforms_as_agent, max_payload_bytes: self.max_payload_bytes, + truncation: self.truncation, }); let input = input.batched_partitioned(partitioner, batch_settings.timeout, |_| { @@ -466,8 +615,353 @@ mod tests { value::{Kind, kind::Collection}, }; - use super::{normalize_as_agent_event, normalize_event}; - use crate::common::datadog::DD_RESERVED_SEMANTIC_ATTRS; + use super::{LogRequestBuilder, normalize_as_agent_event, normalize_event}; + use crate::{ + common::datadog::DD_RESERVED_SEMANTIC_ATTRS, + sinks::{ + datadog::logs::config::{ + DEFAULT_MAX_LOG_BYTES as MAX_LOG_BYTES, DatadogLogsTruncationConfig, + }, + util::Compression, + }, + }; + + const TRUNCATION_MARKER: &str = "...TRUNCATED..."; + const TRUNCATED_TAG: &str = "truncated:single_line"; + + fn encode_logs( + events: Vec, + max_message_bytes: Option, + conforms_as_agent: bool, + ) -> Vec { + encode_logs_with_limits( + events, + max_message_bytes.map(|max_message_bytes| DatadogLogsTruncationConfig { + max_log_bytes: MAX_LOG_BYTES, + max_message_bytes, + }), + conforms_as_agent, + ) + } + + fn encode_logs_with_limits( + events: Vec, + truncation: Option, + conforms_as_agent: bool, + ) -> Vec { + let requests = LogRequestBuilder { + default_api_key: Arc::from("unused"), + transformer: Default::default(), + compression: Compression::None, + conforms_as_agent, + max_payload_bytes: 5_000_000, + truncation, + } + .build_request(events, Arc::from("api-key")) + .expect("request should build"); + + requests + .into_iter() + .flat_map(|request| { + serde_json::from_slice::>(&request.body) + .expect("payload should be a JSON array") + }) + .collect() + } + + #[test] + fn truncates_at_configured_log_limit() { + let log = LogEvent::from("a".repeat(600)); + + let logs = encode_logs_with_limits( + vec![Event::Log(log)], + Some(DatadogLogsTruncationConfig { + max_log_bytes: 500, + max_message_bytes: 400, + }), + false, + ); + + let message = logs[0]["message"] + .as_str() + .expect("message should be a string"); + assert_eq!(message.len(), 400 + TRUNCATION_MARKER.len()); + assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= 500); + } + + #[test] + fn truncates_oversized_message_and_preserves_other_fields() { + let original = "a".repeat(MAX_LOG_BYTES + 1); + let mut log = LogEvent::from(original); + log.insert(event_path!("service"), "payments"); + log.insert(event_path!("custom"), "keep-me"); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert_eq!(logs.len(), 1); + assert_eq!( + logs[0]["message"] + .as_str() + .expect("message should be a string"), + format!("{}{TRUNCATION_MARKER}", "a".repeat(900_000)) + ); + assert_eq!(logs[0]["service"], "payments"); + assert_eq!(logs[0]["custom"], "keep-me"); + assert_eq!(logs[0]["ddtags"], TRUNCATED_TAG); + assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= MAX_LOG_BYTES); + } + + #[test] + fn truncates_agent_normalized_message_and_preserves_nested_fields() { + let original = "b".repeat(MAX_LOG_BYTES + 1); + let mut log = LogEvent::from(original); + log.insert(event_path!("service"), "payments"); + log.insert(event_path!("custom"), "keep-me"); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), true); + + let nested = logs[0]["message"] + .as_object() + .expect("agent message should be an object"); + let message = nested["message"] + .as_str() + .expect("nested message should be a string"); + assert_eq!(message.len(), 900_000 + TRUNCATION_MARKER.len()); + assert!(message.ends_with(TRUNCATION_MARKER)); + assert_eq!(nested["custom"], "keep-me"); + assert_eq!(logs[0]["service"], "payments"); + assert_eq!(logs[0]["ddtags"], TRUNCATED_TAG); + assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= MAX_LOG_BYTES); + } + + #[test] + fn removes_non_standard_fields_if_truncated_log_is_still_oversized() { + let mut log = LogEvent::from("c".repeat(MAX_LOG_BYTES + 1)); + log.insert(event_path!("service"), "payments"); + log.insert(event_path!("custom"), "x".repeat(200_000)); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert_eq!(logs.len(), 1); + assert!(logs[0].get("custom").is_none()); + assert_eq!(logs[0]["service"], "payments"); + assert_eq!(logs[0]["ddtags"], TRUNCATED_TAG); + assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= MAX_LOG_BYTES); + } + + #[test] + fn leaves_short_message_unmarked_when_only_custom_fields_are_removed() { + let mut log = LogEvent::from("hello"); + log.insert(event_path!("service"), "payments"); + log.insert(event_path!("custom"), "x".repeat(MAX_LOG_BYTES + 1)); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert_eq!(logs.len(), 1); + assert_eq!(logs[0]["message"], "hello"); + assert!(logs[0].get("custom").is_none()); + assert_eq!(logs[0]["ddtags"], TRUNCATED_TAG); + } + + #[test] + fn removes_nested_non_standard_fields_from_agent_normalized_log() { + let mut log = LogEvent::from("d".repeat(MAX_LOG_BYTES + 1)); + log.insert(event_path!("service"), "payments"); + log.insert(event_path!("custom"), "x".repeat(200_000)); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), true); + + assert_eq!(logs.len(), 1); + let nested = logs[0]["message"] + .as_object() + .expect("agent message should be an object"); + assert!(nested.get("custom").is_none()); + assert!( + nested["message"] + .as_str() + .expect("message should be a string") + .ends_with(TRUNCATION_MARKER) + ); + assert_eq!(logs[0]["service"], "payments"); + assert_eq!(logs[0]["ddtags"], TRUNCATED_TAG); + assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= MAX_LOG_BYTES); + } + + #[test] + fn drops_log_that_remains_oversized_after_reduction() { + let mut log = LogEvent::from("e".repeat(MAX_LOG_BYTES + 1)); + log.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert!(logs.is_empty()); + } + + #[test] + fn counts_irreducible_log_drop_once() { + vector_lib::metrics::init_test(); + let controller = vector_lib::metrics::Controller::get().unwrap(); + controller.reset(); + + let mut log = LogEvent::from("e".repeat(MAX_LOG_BYTES + 1)); + log.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert!(logs.is_empty()); + let discarded_events = controller + .capture_metrics() + .iter() + .filter(|metric| metric.name() == "component_discarded_events_total") + .map(|metric| match metric.value() { + vector_lib::event::MetricValue::Counter { value } => *value, + _ => panic!("discarded events metric must be a counter"), + }) + .sum::(); + assert_eq!(discarded_events, 1.0); + } + + #[test] + fn drops_oversized_log_with_non_string_message() { + let mut log = LogEvent::default(); + log.insert( + event_path!("message"), + value!({ "body": ("f".repeat(MAX_LOG_BYTES + 1)) }), + ); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert!(logs.is_empty()); + } + + #[test] + fn drops_oversized_log_without_string_message_before_removing_fields() { + for conforms_as_agent in [false, true] { + let mut log = LogEvent::default(); + log.insert(event_path!("custom"), "f".repeat(MAX_LOG_BYTES + 1)); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), conforms_as_agent); + + assert!(logs.is_empty()); + } + } + + #[test] + fn shrinks_escaped_message_again_to_fit_default_limits() { + for conforms_as_agent in [false, true] { + let log = LogEvent::from("\"".repeat(600_000)); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), conforms_as_agent); + + assert_eq!(logs.len(), 1); + let message = if conforms_as_agent { + &logs[0]["message"]["message"] + } else { + &logs[0]["message"] + } + .as_str() + .expect("message should be a string"); + assert!(message.len() < 600_000 + TRUNCATION_MARKER.len()); + assert!(message.ends_with(TRUNCATION_MARKER)); + assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= MAX_LOG_BYTES); + } + } + + #[test] + fn second_message_shrink_reserves_space_for_new_marker() { + let message = format!("{}{}", "\"".repeat(450_000), "a".repeat(110_000)); + let log = LogEvent::from(message); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert_eq!(logs.len(), 1); + assert!( + logs[0]["message"] + .as_str() + .expect("message should be a string") + .ends_with(TRUNCATION_MARKER) + ); + assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= MAX_LOG_BYTES); + } + + #[test] + fn does_not_duplicate_truncation_tag() { + let mut log = LogEvent::from("g".repeat(MAX_LOG_BYTES + 1)); + log.insert(event_path!("ddtags"), format!("env:test,{TRUNCATED_TAG}")); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert_eq!(logs[0]["ddtags"], format!("env:test,{TRUNCATED_TAG}")); + } + + #[test] + fn appends_truncation_tag_to_existing_tags() { + let mut log = LogEvent::from("g".repeat(MAX_LOG_BYTES + 1)); + log.insert(event_path!("ddtags"), "env:test"); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert_eq!(logs[0]["ddtags"], format!("env:test,{TRUNCATED_TAG}")); + } + + #[test] + fn leaves_oversized_message_unchanged_when_truncation_is_disabled() { + let original = "h".repeat(MAX_LOG_BYTES + 1); + let log = LogEvent::from(original.clone()); + + let logs = encode_logs(vec![Event::Log(log)], None, false); + + assert_eq!(logs.len(), 1); + assert_eq!(logs[0]["message"].as_str(), Some(original.as_str())); + assert!(logs[0].get("ddtags").is_none()); + } + + #[test] + fn leaves_log_at_exact_encoded_limit_unchanged() { + // `{"message":""}` contributes 14 bytes around the message body. + let original = "h".repeat(MAX_LOG_BYTES - 14); + let mut log = LogEvent::default(); + log.insert(event_path!("message"), original.clone()); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + + assert_eq!(serde_json::to_vec(&logs[0]).unwrap().len(), MAX_LOG_BYTES); + assert_eq!(logs[0]["message"].as_str(), Some(original.as_str())); + assert!(logs[0].get("ddtags").is_none()); + } + + #[test] + fn truncation_preserves_utf8_boundaries() { + let original = "😀".repeat(300_000); + let log = LogEvent::from(original.clone()); + + let logs = encode_logs(vec![Event::Log(log)], Some(900_001), false); + + let message = logs[0]["message"] + .as_str() + .expect("message should be valid UTF-8"); + let body = message + .strip_suffix(TRUNCATION_MARKER) + .expect("message should have truncation marker"); + assert!(original.starts_with(body)); + assert_eq!(body.len(), 900_000); + } + + #[test] + fn delivers_following_log_after_dropping_irreducible_log() { + let mut oversized = LogEvent::from("i".repeat(MAX_LOG_BYTES + 1)); + oversized.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); + let small = LogEvent::from("ok"); + + let logs = encode_logs( + vec![Event::Log(oversized), Event::Log(small)], + Some(900_000), + false, + ); + + assert_eq!(logs.len(), 1); + assert_eq!(logs[0]["message"], "ok"); + } fn assert_normalized_log_has_expected_attrs(log: &LogEvent) { assert!( diff --git a/src/sinks/datadog/logs/tests.rs b/src/sinks/datadog/logs/tests.rs index d394a71e0b3b6..4289cb82c6b2a 100644 --- a/src/sinks/datadog/logs/tests.rs +++ b/src/sinks/datadog/logs/tests.rs @@ -19,6 +19,7 @@ use crate::{ extra_context::ExtraContext, http::HttpError, sinks::{ + VectorSink, datadog::test_utils::{ApiStatus, test_server}, util::{ retries::RetryLogic, @@ -26,13 +27,13 @@ use crate::{ }, }, test_util::{ - addr::next_addr, + addr::{PortGuard, next_addr}, components::{ COMPONENT_ERROR_TAGS, DATA_VOLUME_SINK_TAGS, SINK_TAGS, run_and_assert_data_volume_sink_compliance, run_and_assert_sink_compliance, run_and_assert_sink_error, }, - random_lines_with_stream, + generate_lines_with_stream, random_lines_with_stream, }, tls::TlsError, }; @@ -52,6 +53,25 @@ enum TestType { Error, } +async fn start_test_sink( + config: &str, + api_status: ApiStatus, +) -> ( + PortGuard, + stream_cancel::Trigger, + VectorSink, + Receiver<(Parts, Bytes)>, +) { + let (mut config, cx) = load_sink::(config).unwrap(); + let (guard, addr) = next_addr(); + config.local_dd_common.endpoint = Some(format!("http://{addr}")); + let (sink, _) = config.build(cx).await.unwrap(); + + let (rx, trigger, server) = test_server(addr, api_status); + tokio::spawn(server); + (guard, trigger, sink, rx) +} + /// Starts a test sink with random lines running into it /// /// This function starts a Datadog Logs sink with a simplistic configuration and @@ -82,18 +102,7 @@ async fn start_test_detail( default_api_key = "atoken" compression = "none" "#}; - let (mut config, cx) = load_sink::(config).unwrap(); - - let (_guard, addr) = next_addr(); - // Swap out the endpoint so we can force send it - // to our local server - let endpoint = format!("http://{addr}"); - config.local_dd_common.endpoint = Some(endpoint.clone()); - - let (sink, _) = config.build(cx).await.unwrap(); - - let (rx, _trigger, server) = test_server(addr, api_status); - tokio::spawn(server); + let (_guard, _trigger, sink, rx) = start_test_sink(config, api_status).await; let (batch, receiver) = BatchNotifier::new_with_receiver(); let (expected, events) = random_lines_with_stream(100, 10, Some(batch)); @@ -176,6 +185,35 @@ async fn smoke() { } } +#[tokio::test] +async fn truncates_oversized_log_over_http() { + let config = indoc! {r#" + default_api_key = "atoken" + compression = "none" + + [truncate_oversized_logs] + max_log_bytes = 1000 + max_message_bytes = 900 + "#}; + let (_guard, _trigger, sink, rx) = start_test_sink(config, ApiStatus::OKv2).await; + + let (batch, receiver) = BatchNotifier::new_with_receiver(); + let (_, events) = generate_lines_with_stream(|_| "\"".repeat(600), 1, Some(batch)); + run_and_assert_sink_compliance(sink, events, &SINK_TAGS).await; + assert_eq!(receiver.await, BatchStatus::Delivered); + + let output = rx.take(1).collect::>().await; + let logs: Vec = + serde_json::from_slice(&output[0].1).expect("request body should contain a JSON array"); + assert_eq!(logs.len(), 1); + let message = logs[0]["message"] + .as_str() + .expect("message should be a string"); + assert!(message.ends_with("...TRUNCATED...")); + assert_eq!(logs[0]["ddtags"], "truncated:single_line"); + assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= 1000); +} + /// Assert the sink emits source and service tags when run with telemetry configured. #[tokio::test] async fn telemetry() { diff --git a/website/cue/reference/components/sinks/datadog_logs.cue b/website/cue/reference/components/sinks/datadog_logs.cue index 3da6d41ad18ed..77f6def6b0842 100644 --- a/website/cue/reference/components/sinks/datadog_logs.cue +++ b/website/cue/reference/components/sinks/datadog_logs.cue @@ -82,7 +82,16 @@ components: sinks: datadog_logs: { limit. Increase it when targeting a compatible endpoint that accepts larger payloads. A batch that exceeds `max_payload_bytes` is split across multiple requests. A single event - that exceeds `max_payload_bytes` is dropped. + that still exceeds `max_payload_bytes` after optional truncation is dropped. + + Set `truncate_oversized_logs` to opt in to reducing individual logs whose encoded JSON + exceeds `truncate_oversized_logs.max_log_bytes` (default 1,000,000). When the message + is shortened, at most `truncate_oversized_logs.max_message_bytes` raw bytes (default + 900,000) are retained and `...TRUNCATED...` is appended. Every reduced log is tagged + `truncated:single_line`. If the log remains oversized after the initial message cap, + non-standard fields are removed before the message is shortened further to account for + JSON encoding. Logs that still exceed the limit, or have no string message to truncate, + are dropped. """ } } diff --git a/website/cue/reference/components/sinks/generated/datadog_logs.cue b/website/cue/reference/components/sinks/generated/datadog_logs.cue index 27aaa9f5e1d94..84c15d78952dd 100644 --- a/website/cue/reference/components/sinks/generated/datadog_logs.cue +++ b/website/cue/reference/components/sinks/generated/datadog_logs.cue @@ -128,8 +128,8 @@ generated: components: sinks: datadog_logs: configuration: { to not set it above 5,000,000 (5 MB, the standard Datadog API limit). Increase this when targeting a compatible endpoint that accepts larger payloads. The batch goal is derived as `max_payload_bytes - 750,000` bytes; events larger than the - batch goal are sent alone in their batch. Events exceeding `max_payload_bytes` are - dropped. + batch goal are sent alone in their batch. Single events that still exceed + `max_payload_bytes` after optional truncation are dropped. """ required: false type: uint: default: 5000000 @@ -159,4 +159,28 @@ generated: components: sinks: datadog_logs: configuration: { required: false type: _schemaDefinitions["core::option::Option"] } + truncate_oversized_logs: { + description: """ + Truncate logs whose encoded JSON exceeds `max_log_bytes`. + + When a message is shortened, at most `max_message_bytes` raw bytes are retained and + `...TRUNCATED...` is appended. Every reduced log is tagged with `truncated:single_line`. + If the log remains oversized after the initial message cap, non-standard fields are removed + before the message is shortened further to account for JSON encoding. Logs that still + exceed the limit, or have no string message to truncate, are dropped. + """ + required: false + type: object: options: { + max_log_bytes: { + description: "Maximum encoded size, in bytes, of a log before truncation is applied." + required: false + type: uint: default: 1000000 + } + max_message_bytes: { + description: "Maximum number of message bytes to retain before appending the truncation marker." + required: false + type: uint: default: 900000 + } + } + } } From 1ce3a7bfda904f268b78c89c31a3ece42f88961b Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Wed, 30 Sep 2026 07:15:48 -0600 Subject: [PATCH 02/15] Mark events that cannot be encoded as rejected for E2E acks --- src/sinks/datadog/logs/sink.rs | 73 +++++++++++++++++++++++++++++----- 1 file changed, 62 insertions(+), 11 deletions(-) diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index 6928c1ad96e4d..7f951fce7c99a 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -294,7 +294,8 @@ impl LogRequestBuilder { let (events_serialized, body, byte_size) = self.serialize_with_capacity(&mut events_with_estimated_size)?; if events_serialized.is_empty() { - if events_with_estimated_size.pop_front().is_some() { + if let Some((event, _)) = events_with_estimated_size.pop_front() { + event.metadata().update_status(EventStatus::Rejected); emit!(ComponentEventsDropped:: { count: 1, reason: "Event too large to encode." @@ -310,6 +311,12 @@ impl LogRequestBuilder { Ok(requests) } + /// Serialize events into a buffer as a JSON array that has a maximum size of + /// `max_payload_bytes`. + /// + /// Returns the serialized events, the buffer, and the byte size of the events. Events that do + /// not fit remain in the `events` parameter. Events rejected during encoding are removed. + #[doc(hidden)] fn serialize_with_capacity( &self, events: &mut VecDeque<(Event, JsonSize)>, @@ -335,6 +342,7 @@ impl LogRequestBuilder { estimated_json_size = event.estimated_json_encoded_size_of(); } LogEncoding::Dropped { reason } => { + event.metadata().update_status(EventStatus::Rejected); emit!(ComponentEventsDropped:: { count: 1, reason }); continue; } @@ -607,6 +615,7 @@ mod tests { use vector_lib::{ config::{LegacyKey, LogNamespace}, event::{Event, EventMetadata, LogEvent}, + finalization::{BatchNotifier, BatchStatus}, schema::{Definition, meaning}, }; use vrl::{ @@ -649,16 +658,7 @@ mod tests { truncation: Option, conforms_as_agent: bool, ) -> Vec { - let requests = LogRequestBuilder { - default_api_key: Arc::from("unused"), - transformer: Default::default(), - compression: Compression::None, - conforms_as_agent, - max_payload_bytes: 5_000_000, - truncation, - } - .build_request(events, Arc::from("api-key")) - .expect("request should build"); + let requests = build_requests(events, truncation, conforms_as_agent, 5_000_000); requests .into_iter() @@ -669,6 +669,24 @@ mod tests { .collect() } + fn build_requests( + events: Vec, + truncation: Option, + conforms_as_agent: bool, + max_payload_bytes: usize, + ) -> Vec { + LogRequestBuilder { + default_api_key: Arc::from("unused"), + transformer: Default::default(), + compression: Compression::None, + conforms_as_agent, + max_payload_bytes, + truncation, + } + .build_request(events, Arc::from("api-key")) + .expect("request should build") + } + #[test] fn truncates_at_configured_log_limit() { let log = LogEvent::from("a".repeat(600)); @@ -797,6 +815,39 @@ mod tests { assert!(logs.is_empty()); } + #[test] + fn rejects_log_that_remains_oversized_after_reduction() { + let (batch, mut receiver) = BatchNotifier::new_with_receiver(); + let mut log = LogEvent::from("e".repeat(MAX_LOG_BYTES + 1)).with_batch_notifier(&batch); + log.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); + drop(batch); + + let requests = build_requests( + vec![Event::Log(log)], + Some(DatadogLogsTruncationConfig { + max_log_bytes: MAX_LOG_BYTES, + max_message_bytes: 900_000, + }), + false, + 5_000_000, + ); + + assert!(requests.is_empty()); + assert_eq!(receiver.try_recv(), Ok(BatchStatus::Rejected)); + } + + #[test] + fn rejects_log_that_exceeds_payload_limit() { + let (batch, mut receiver) = BatchNotifier::new_with_receiver(); + let log = LogEvent::from("oversized").with_batch_notifier(&batch); + drop(batch); + + let requests = build_requests(vec![Event::Log(log)], None, false, 2); + + assert!(requests.is_empty()); + assert_eq!(receiver.try_recv(), Ok(BatchStatus::Rejected)); + } + #[test] fn counts_irreducible_log_drop_once() { vector_lib::metrics::init_test(); From 6e356db22abe0bd7db41580866318667965bc2ca Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Wed, 30 Sep 2026 11:13:55 -0600 Subject: [PATCH 03/15] Optimize encoding by calculating the required truncated message length --- Cargo.lock | 1 + Cargo.toml | 4 +- src/sinks/datadog/logs/config.rs | 7 +- src/sinks/datadog/logs/sink.rs | 226 ++++++++++++++---- .../sinks/generated/datadog_logs.cue | 7 +- 5 files changed, 192 insertions(+), 53 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index dc737f69e588f..3449f3d96e45d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -13532,6 +13532,7 @@ dependencies = [ "serde_with", "serde_yaml", "serial_test", + "simdutf8", "similar-asserts", "smallvec", "smpl_jwt", diff --git a/Cargo.toml b/Cargo.toml index c27081250a296..ef735e1917e93 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -223,6 +223,7 @@ serde_json = { version = "1.0.150", default-features = false, features = ["prese serde_path_to_error = "0.1.14" serde_with = { version = "3.21.0", default-features = false, features = ["std", "macros", "chrono_0_4"] } serde_yaml = { version = "0.9.34", default-features = false } +simdutf8 = { version = "0.1.5", default-features = false } snafu = { version = "0.9.0", default-features = false, features = ["futures", "std"] } socket2 = { version = "0.6.3", default-features = false } tempfile = "3.27.0" @@ -361,6 +362,7 @@ serde_json.workspace = true serde_path_to_error.workspace = true serde_with = { version = "3.14.0", default-features = false, features = ["macros", "std"] } serde_yaml.workspace = true +simdutf8 = { workspace = true, optional = true } # Messagepack rmp-serde = { version = "1.3.0", default-features = false, optional = true } @@ -977,7 +979,7 @@ sinks-console = [] sinks-databend = ["dep:databend-client"] sinks-databricks_zerobus = ["dep:databricks-zerobus-ingest-sdk", "codecs-arrow", "arrow/ipc_compression"] sinks-datadog_events = [] -sinks-datadog_logs = [] +sinks-datadog_logs = ["dep:simdutf8"] sinks-datadog_metrics = ["dep:datadog-proto", "dep:prost", "dep:prost-reflect", "dep:datadog-agent-metrics-v3", "dep:protobuf"] sinks-datadog_traces = ["protobuf-build", "dep:datadog-proto", "dep:prost", "dep:rmpv", "dep:rmp-serde", "dep:serde_bytes"] sinks-doris = ["sqlx/mysql"] diff --git a/src/sinks/datadog/logs/config.rs b/src/sinks/datadog/logs/config.rs index 070251b6c6e59..7a0576405959a 100644 --- a/src/sinks/datadog/logs/config.rs +++ b/src/sinks/datadog/logs/config.rs @@ -110,9 +110,10 @@ pub struct DatadogLogsConfig { /// /// When a message is shortened, at most `max_message_bytes` raw bytes are retained and /// `...TRUNCATED...` is appended. Every reduced log is tagged with `truncated:single_line`. - /// If the log remains oversized after the initial message cap, non-standard fields are removed - /// before the message is shortened further to account for JSON encoding. Logs that still - /// exceed the limit, or have no string message to truncate, are dropped. + /// The message is sized to account for JSON encoding while preserving non-standard fields when + /// possible. If the non-message fields leave no room for a truncated message, non-standard + /// fields are removed and the message is sized again. Logs that still exceed the limit, or have + /// no string message to truncate, are dropped. pub truncate_oversized_logs: Option, } diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index 7f951fce7c99a..69cfffdfda0f1 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -1,4 +1,4 @@ -use std::{collections::VecDeque, fmt::Debug, io, sync::Arc}; +use std::{borrow::Cow, collections::VecDeque, fmt::Debug, io, sync::Arc}; use bytes::Bytes; use itertools::Itertools; @@ -423,27 +423,46 @@ fn encode_log( reason: "Oversized event has no string message to truncate.", }); }; - let mut body_len = - floor_char_boundary(&message, message.len().min(truncation.max_message_bytes)); - if body_len < message.len() { - set_truncated_message(event.as_mut_log(), &message, body_len); - } - ensure_truncated_tag(event.as_mut_log()); + let message = simdutf8_lossy(&message); + let message_encoded_size = json_string_encoded_size(&message); + let tagged_encoded_size = ensure_truncated_tag(event.as_mut_log(), original_encoded_size)?; + // The encoded event consists of a fixed non-message portion and the message value. Size them + // separately so the final message can be selected without repeatedly encoding the whole event. + let mut non_message_encoded_size = tagged_encoded_size - message_encoded_size; + let select_rewrite = |non_message_encoded_size| { + truncation + .max_log_bytes + .checked_sub(non_message_encoded_size) + .and_then(|budget| { + select_message_rewrite( + &message, + message_encoded_size, + truncation.max_message_bytes, + budget, + ) + }) + }; + let mut rewrite = select_rewrite(non_message_encoded_size); - let mut encoded_size = rewrite_log(buf, existing_len, event, include_comma)?; - if encoded_size > truncation.max_log_bytes { - strip_non_standard_fields(event.as_mut_log(), conforms_as_agent); - encoded_size = rewrite_log(buf, existing_len, event, include_comma)?; + if rewrite.is_none() { + let stripped_size = strip_non_standard_fields(event.as_mut_log(), conforms_as_agent)?; + non_message_encoded_size -= stripped_size; + rewrite = select_rewrite(non_message_encoded_size); } - if encoded_size > truncation.max_log_bytes { - let marker_bytes_to_add = TRUNCATION_MARKER.len() * usize::from(body_len == message.len()); - let bytes_to_remove = (encoded_size - truncation.max_log_bytes) + marker_bytes_to_add; - body_len = floor_char_boundary(&message, body_len.saturating_sub(bytes_to_remove)); + let Some(body_len) = rewrite else { + buf.truncate(existing_len); + return Ok(LogEncoding::Dropped { + reason: "Event remains too large after truncation.", + }); + }; + if body_len < message.len() { set_truncated_message(event.as_mut_log(), &message, body_len); - encoded_size = rewrite_log(buf, existing_len, event, include_comma)?; } + buf.truncate(existing_len); + let encoded_size = write_log(buf, event, include_comma)?; + if encoded_size > truncation.max_log_bytes { buf.truncate(existing_len); return Ok(LogEncoding::Dropped { @@ -458,6 +477,40 @@ fn encode_log( Ok(LogEncoding::Truncated) } +fn simdutf8_lossy(bytes: &[u8]) -> Cow<'_, str> { + match simdutf8::basic::from_utf8(bytes) { + Ok(value) => Cow::Borrowed(value), + Err(_) => String::from_utf8_lossy(bytes), + } +} + +fn select_message_rewrite( + message: &str, + message_encoded_size: usize, + max_message_bytes: usize, + encoded_budget: usize, +) -> Option { + if message.len() <= max_message_bytes && message_encoded_size <= encoded_budget { + Some(message.len()) + } else { + let content_budget = + encoded_budget.checked_sub(2 + json_string_content_size(TRUNCATION_MARKER))?; + let mut body_len = 0; + let mut encoded_size = 0; + for (index, character) in message.char_indices() { + let next_body_len = index + character.len_utf8(); + let next_encoded_size = encoded_size + json_character_encoded_size(character); + if next_body_len > max_message_bytes || next_encoded_size > content_budget { + break; + } + body_len = next_body_len; + encoded_size = next_encoded_size; + } + + (body_len < message.len()).then_some(body_len) + } +} + fn write_log(buf: &mut Vec, event: &Event, include_comma: bool) -> Result { if include_comma { buf.push(b','); @@ -467,40 +520,43 @@ fn write_log(buf: &mut Vec, event: &Event, include_comma: bool) -> Result, - existing_len: usize, - event: &Event, - include_comma: bool, -) -> Result { - buf.truncate(existing_len); - write_log(buf, event, include_comma) -} - const TRUNCATION_MARKER: &str = "...TRUNCATED..."; const TRUNCATED_TAG: &str = "truncated:single_line"; -fn set_truncated_message(log: &mut LogEvent, message: &Bytes, body_len: usize) { +fn set_truncated_message(log: &mut LogEvent, message: &str, body_len: usize) { let marker = TRUNCATION_MARKER.as_bytes(); let mut truncated = Vec::with_capacity(body_len + marker.len()); - truncated.extend_from_slice(&message[..body_len]); + truncated.extend_from_slice(&message.as_bytes()[..body_len]); truncated.extend_from_slice(marker); *message_bytes_mut(log).expect("the message was previously found") = Bytes::from(truncated); } -fn ensure_truncated_tag(log: &mut LogEvent) { +fn ensure_truncated_tag(log: &mut LogEvent, encoded_size: usize) -> Result { let tags_path = event_path!(DDTAGS); + let previous_size = log + .get(tags_path) + .map(json_value_encoded_size) + .transpose()?; if let Some(tags) = log.get(tags_path).and_then(Value::as_bytes) && !tags.is_empty() { let tags = String::from_utf8_lossy(tags); if tags.split(',').any(|tag| tag.trim() == TRUNCATED_TAG) { - return; + return Ok(encoded_size); } log.insert(tags_path, format!("{tags},{TRUNCATED_TAG}")); } else { log.insert(tags_path, TRUNCATED_TAG); } + let new_size = + json_value_encoded_size(log.get(tags_path).expect("the truncation tag was inserted"))?; + + Ok(if let Some(previous_size) = previous_size { + encoded_size - previous_size + new_size + } else { + // The log already contains the message field, so adding ddtags also adds one comma. + encoded_size + json_string_encoded_size(DDTAGS) + 1 + new_size + 1 + }) } fn message_bytes_mut(log: &mut LogEvent) -> Option<&mut Bytes> { @@ -514,24 +570,66 @@ fn message_bytes_mut(log: &mut LogEvent) -> Option<&mut Bytes> { } } -fn strip_non_standard_fields(log: &mut LogEvent, conforms_as_agent: bool) { +fn strip_non_standard_fields( + log: &mut LogEvent, + conforms_as_agent: bool, +) -> Result { let Some(fields) = log.as_map_mut() else { - return; + return Ok(0); }; - fields.retain(|field, _| is_reserved_attribute(field.as_str()) || field.as_str() == MESSAGE); + let mut stripped_size = strip_fields(fields, |field| { + is_reserved_attribute(field) || field == MESSAGE + })?; if conforms_as_agent && let Some(Value::Object(nested)) = fields.get_mut(MESSAGE) { - nested.retain(|field, _| field.as_str() == MESSAGE); + stripped_size += strip_fields(nested, |field| field == MESSAGE)?; } + Ok(stripped_size) } -fn floor_char_boundary(bytes: &[u8], mut index: usize) -> usize { - if index >= bytes.len() { - return bytes.len(); +// Calculate the exact encoded size removed rather than using `retain`, so the message budget can +// be recomputed without encoding the entire log again. +fn strip_fields(fields: &mut ObjectMap, retain: impl Fn(&str) -> bool) -> Result { + let fields_to_strip = fields + .iter() + .filter(|(field, _)| !retain(field.as_str())) + .map(|(field, value)| { + Ok(( + field.clone(), + json_string_encoded_size(field.as_str()) + 1 + json_value_encoded_size(value)?, + )) + }) + .collect::, io::Error>>()?; + debug_assert!(fields.len() > fields_to_strip.len()); + // At least one field remains, so removing each member also removes one separating comma. + let mut stripped_size = fields_to_strip.len(); + for (field, member_size) in fields_to_strip { + let _ = fields.remove(field.as_str()); + stripped_size += member_size; } - while index > 0 && (bytes[index] & 0b1100_0000) == 0b1000_0000 { - index -= 1; + Ok(stripped_size) +} + +fn json_string_encoded_size(value: &str) -> usize { + 2 + json_string_content_size(value) +} + +fn json_string_content_size(value: &str) -> usize { + value.chars().map(json_character_encoded_size).sum() +} + +const fn json_character_encoded_size(character: char) -> usize { + match character { + '"' | '\\' | '\u{0008}' | '\t' | '\n' | '\u{000c}' | '\r' => 2, + '\u{0000}'..='\u{001f}' => 6, + _ => character.len_utf8(), } - index +} + +fn json_value_encoded_size(value: &Value) -> Result { + let mut sink = io::sink(); + encoding::as_tracked_write(&mut sink, value, |writer, value| { + serde_json::to_writer(writer, value) + }) } impl LogSink @@ -624,7 +722,10 @@ mod tests { value::{Kind, kind::Collection}, }; - use super::{LogRequestBuilder, normalize_as_agent_event, normalize_event}; + use super::{ + LogRequestBuilder, json_string_encoded_size, normalize_as_agent_event, normalize_event, + simdutf8_lossy, + }; use crate::{ common::datadog::DD_RESERVED_SEMANTIC_ATTRS, sinks::{ @@ -638,6 +739,27 @@ mod tests { const TRUNCATION_MARKER: &str = "...TRUNCATED..."; const TRUNCATED_TAG: &str = "truncated:single_line"; + #[test] + fn calculates_exact_json_string_size() { + for value in [ + "plain text", + "\"quoted\\text\"", + "\0\u{0001}\u{0008}\t\n\u{000c}\r", + "multibyte 😀 text", + ] { + assert_eq!( + json_string_encoded_size(value), + serde_json::to_vec(value).unwrap().len() + ); + } + + let invalid = bytes::Bytes::from_static(&[b'a', 0xff, b'b']); + assert_eq!( + json_string_encoded_size(&simdutf8_lossy(&invalid)), + serde_json::to_vec(&Value::Bytes(invalid)).unwrap().len() + ); + } + fn encode_logs( events: Vec, max_message_bytes: Option, @@ -753,7 +875,7 @@ mod tests { } #[test] - fn removes_non_standard_fields_if_truncated_log_is_still_oversized() { + fn preserves_non_standard_fields_if_message_can_be_shortened_further() { let mut log = LogEvent::from("c".repeat(MAX_LOG_BYTES + 1)); log.insert(event_path!("service"), "payments"); log.insert(event_path!("custom"), "x".repeat(200_000)); @@ -761,9 +883,14 @@ mod tests { let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); assert_eq!(logs.len(), 1); - assert!(logs[0].get("custom").is_none()); + assert_eq!(logs[0]["custom"], "x".repeat(200_000)); assert_eq!(logs[0]["service"], "payments"); assert_eq!(logs[0]["ddtags"], TRUNCATED_TAG); + let message = logs[0]["message"] + .as_str() + .expect("message should be a string"); + assert!(message.ends_with(TRUNCATION_MARKER)); + assert!(message.len() < 900_000 + TRUNCATION_MARKER.len()); assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= MAX_LOG_BYTES); } @@ -771,7 +898,7 @@ mod tests { fn leaves_short_message_unmarked_when_only_custom_fields_are_removed() { let mut log = LogEvent::from("hello"); log.insert(event_path!("service"), "payments"); - log.insert(event_path!("custom"), "x".repeat(MAX_LOG_BYTES + 1)); + log.insert(event_path!("custom"), "\u{0001}".repeat(200_000)); let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); @@ -782,7 +909,7 @@ mod tests { } #[test] - fn removes_nested_non_standard_fields_from_agent_normalized_log() { + fn preserves_nested_non_standard_fields_if_message_can_be_shortened_further() { let mut log = LogEvent::from("d".repeat(MAX_LOG_BYTES + 1)); log.insert(event_path!("service"), "payments"); log.insert(event_path!("custom"), "x".repeat(200_000)); @@ -793,13 +920,20 @@ mod tests { let nested = logs[0]["message"] .as_object() .expect("agent message should be an object"); - assert!(nested.get("custom").is_none()); + assert_eq!(nested["custom"], "x".repeat(200_000)); assert!( nested["message"] .as_str() .expect("message should be a string") .ends_with(TRUNCATION_MARKER) ); + assert!( + nested["message"] + .as_str() + .expect("message should be a string") + .len() + < 900_000 + TRUNCATION_MARKER.len() + ); assert_eq!(logs[0]["service"], "payments"); assert_eq!(logs[0]["ddtags"], TRUNCATED_TAG); assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= MAX_LOG_BYTES); diff --git a/website/cue/reference/components/sinks/generated/datadog_logs.cue b/website/cue/reference/components/sinks/generated/datadog_logs.cue index 84c15d78952dd..fae334a2da816 100644 --- a/website/cue/reference/components/sinks/generated/datadog_logs.cue +++ b/website/cue/reference/components/sinks/generated/datadog_logs.cue @@ -165,9 +165,10 @@ generated: components: sinks: datadog_logs: configuration: { When a message is shortened, at most `max_message_bytes` raw bytes are retained and `...TRUNCATED...` is appended. Every reduced log is tagged with `truncated:single_line`. - If the log remains oversized after the initial message cap, non-standard fields are removed - before the message is shortened further to account for JSON encoding. Logs that still - exceed the limit, or have no string message to truncate, are dropped. + The message is sized to account for JSON encoding while preserving non-standard fields when + possible. If the non-message fields leave no room for a truncated message, non-standard + fields are removed and the message is sized again. Logs that still exceed the limit, or have + no string message to truncate, are dropped. """ required: false type: object: options: { From 2d453ebe3dc3590b40188281b5ae6315f9b1eee9 Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Wed, 30 Sep 2026 15:38:18 -0600 Subject: [PATCH 04/15] Drop no longer needed `max_message_bytes` config --- src/internal_events/datadog_logs.rs | 2 - src/sinks/datadog/logs/config.rs | 47 +--------- src/sinks/datadog/logs/sink.rs | 87 ++++++++----------- src/sinks/datadog/logs/tests.rs | 1 - .../sinks/generated/datadog_logs.cue | 25 ++---- 5 files changed, 46 insertions(+), 116 deletions(-) diff --git a/src/internal_events/datadog_logs.rs b/src/internal_events/datadog_logs.rs index 242be9b2aa216..9f859aa6c3d10 100644 --- a/src/internal_events/datadog_logs.rs +++ b/src/internal_events/datadog_logs.rs @@ -7,7 +7,6 @@ use vrl::path::OwnedTargetPath; #[derive(Debug, NamedInternalEvent)] pub struct DatadogLogsEventTruncated { pub max_log_bytes: usize, - pub max_message_bytes: usize, pub original_encoded_size: usize, } @@ -16,7 +15,6 @@ impl InternalEvent for DatadogLogsEventTruncated { warn!( message = "Truncated a Datadog log event that exceeded the per-log size limit.", max_log_bytes = self.max_log_bytes, - max_message_bytes = self.max_message_bytes, original_encoded_size = self.original_encoded_size, internal_log_rate_limit = true, ); diff --git a/src/sinks/datadog/logs/config.rs b/src/sinks/datadog/logs/config.rs index 7a0576405959a..169e12fb8e123 100644 --- a/src/sinks/datadog/logs/config.rs +++ b/src/sinks/datadog/logs/config.rs @@ -35,7 +35,6 @@ use crate::{ // practice. pub const MAX_PAYLOAD_BYTES: usize = 5_000_000; pub(super) const DEFAULT_MAX_LOG_BYTES: usize = 1_000_000; -pub(super) const DEFAULT_MAX_MESSAGE_BYTES: usize = 900_000; pub const BATCH_HEADROOM_BYTES: usize = 750_000; pub const BATCH_MAX_EVENTS: usize = 1_000; pub const BATCH_DEFAULT_TIMEOUT_SECS: f64 = 5.0; @@ -60,11 +59,6 @@ pub struct DatadogLogsTruncationConfig { #[derivative(Default(value = "default_max_log_bytes()"))] #[serde(default = "default_max_log_bytes")] pub max_log_bytes: usize, - - /// Maximum number of message bytes to retain before appending the truncation marker. - #[derivative(Default(value = "default_max_message_bytes()"))] - #[serde(default = "default_max_message_bytes")] - pub max_message_bytes: usize, } /// Configuration for the `datadog_logs` sink. @@ -108,12 +102,10 @@ pub struct DatadogLogsConfig { /// Truncate logs whose encoded JSON exceeds `max_log_bytes`. /// - /// When a message is shortened, at most `max_message_bytes` raw bytes are retained and - /// `...TRUNCATED...` is appended. Every reduced log is tagged with `truncated:single_line`. - /// The message is sized to account for JSON encoding while preserving non-standard fields when - /// possible. If the non-message fields leave no room for a truncated message, non-standard - /// fields are removed and the message is sized again. Logs that still exceed the limit, or have - /// no string message to truncate, are dropped. + /// The message is shortened to the largest size that fits and `...TRUNCATED...` is appended. + /// Every reduced log is tagged with `truncated:single_line`. Non-standard fields are preserved + /// when possible, but removed when they leave no room for a truncated message. Logs that still + /// exceed the limit, or have no string message to truncate, are dropped. pub truncate_oversized_logs: Option, } @@ -125,10 +117,6 @@ const fn default_max_log_bytes() -> usize { DEFAULT_MAX_LOG_BYTES } -const fn default_max_message_bytes() -> usize { - DEFAULT_MAX_MESSAGE_BYTES -} - const fn default_compression() -> Option { Some(Compression::zstd_default()) } @@ -302,14 +290,6 @@ impl ValidatedSink for DatadogLogsConfig { ) .into()); } - let maximum_message_bytes = truncation.max_log_bytes - 1; - if !(1..=maximum_message_bytes).contains(&truncation.max_message_bytes) { - return Err(format!( - "truncate_oversized_logs.max_message_bytes ({}) must be between 1 and {}", - truncation.max_message_bytes, maximum_message_bytes - ) - .into()); - } } let batch_goal_bytes = @@ -380,23 +360,6 @@ mod test { .truncate_oversized_logs .expect("truncation should be enabled"); assert_eq!(truncation.max_log_bytes, 1_000_000); - assert_eq!(truncation.max_message_bytes, 900_000); - } - - #[test] - fn validate_rejects_invalid_truncation_limits() { - for max_message_bytes in [0, DEFAULT_MAX_LOG_BYTES] { - let config: DatadogLogsConfig = serde_yaml::from_str(&format!( - r#" - default_api_key: "test_key" - truncate_oversized_logs: - max_message_bytes: {max_message_bytes} - "# - )) - .unwrap(); - - assert!(config.validate().is_err()); - } } #[test] @@ -406,7 +369,6 @@ mod test { max_payload_bytes: 800000 truncate_oversized_logs: max_log_bytes: 799999 - max_message_bytes: 700000 "#}) .unwrap(); @@ -420,7 +382,6 @@ mod test { max_payload_bytes: 3000000 truncate_oversized_logs: max_log_bytes: 2000000 - max_message_bytes: 1800000 "#}) .unwrap(); diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index 69cfffdfda0f1..03e5b6ba6abc6 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -429,28 +429,21 @@ fn encode_log( // The encoded event consists of a fixed non-message portion and the message value. Size them // separately so the final message can be selected without repeatedly encoding the whole event. let mut non_message_encoded_size = tagged_encoded_size - message_encoded_size; - let select_rewrite = |non_message_encoded_size| { + let select_message_body_len = |non_message_encoded_size| { truncation .max_log_bytes .checked_sub(non_message_encoded_size) - .and_then(|budget| { - select_message_rewrite( - &message, - message_encoded_size, - truncation.max_message_bytes, - budget, - ) - }) + .and_then(|budget| select_message_body_len(&message, message_encoded_size, budget)) }; - let mut rewrite = select_rewrite(non_message_encoded_size); + let mut body_len = select_message_body_len(non_message_encoded_size); - if rewrite.is_none() { + if body_len.is_none() { let stripped_size = strip_non_standard_fields(event.as_mut_log(), conforms_as_agent)?; non_message_encoded_size -= stripped_size; - rewrite = select_rewrite(non_message_encoded_size); + body_len = select_message_body_len(non_message_encoded_size); } - let Some(body_len) = rewrite else { + let Some(body_len) = body_len else { buf.truncate(existing_len); return Ok(LogEncoding::Dropped { reason: "Event remains too large after truncation.", @@ -471,7 +464,6 @@ fn encode_log( } emit!(DatadogLogsEventTruncated { max_log_bytes: truncation.max_log_bytes, - max_message_bytes: truncation.max_message_bytes, original_encoded_size, }); Ok(LogEncoding::Truncated) @@ -484,13 +476,12 @@ fn simdutf8_lossy(bytes: &[u8]) -> Cow<'_, str> { } } -fn select_message_rewrite( +fn select_message_body_len( message: &str, message_encoded_size: usize, - max_message_bytes: usize, encoded_budget: usize, ) -> Option { - if message.len() <= max_message_bytes && message_encoded_size <= encoded_budget { + if message_encoded_size <= encoded_budget { Some(message.len()) } else { let content_budget = @@ -500,7 +491,7 @@ fn select_message_rewrite( for (index, character) in message.char_indices() { let next_body_len = index + character.len_utf8(); let next_encoded_size = encoded_size + json_character_encoded_size(character); - if next_body_len > max_message_bytes || next_encoded_size > content_budget { + if next_encoded_size > content_budget { break; } body_len = next_body_len; @@ -762,14 +753,13 @@ mod tests { fn encode_logs( events: Vec, - max_message_bytes: Option, + truncate_oversized_logs: bool, conforms_as_agent: bool, ) -> Vec { encode_logs_with_limits( events, - max_message_bytes.map(|max_message_bytes| DatadogLogsTruncationConfig { + truncate_oversized_logs.then_some(DatadogLogsTruncationConfig { max_log_bytes: MAX_LOG_BYTES, - max_message_bytes, }), conforms_as_agent, ) @@ -815,17 +805,14 @@ mod tests { let logs = encode_logs_with_limits( vec![Event::Log(log)], - Some(DatadogLogsTruncationConfig { - max_log_bytes: 500, - max_message_bytes: 400, - }), + Some(DatadogLogsTruncationConfig { max_log_bytes: 500 }), false, ); let message = logs[0]["message"] .as_str() .expect("message should be a string"); - assert_eq!(message.len(), 400 + TRUNCATION_MARKER.len()); + assert!(message.ends_with(TRUNCATION_MARKER)); assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= 500); } @@ -836,14 +823,14 @@ mod tests { log.insert(event_path!("service"), "payments"); log.insert(event_path!("custom"), "keep-me"); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert_eq!(logs.len(), 1); - assert_eq!( + assert!( logs[0]["message"] .as_str() - .expect("message should be a string"), - format!("{}{TRUNCATION_MARKER}", "a".repeat(900_000)) + .expect("message should be a string") + .ends_with(TRUNCATION_MARKER) ); assert_eq!(logs[0]["service"], "payments"); assert_eq!(logs[0]["custom"], "keep-me"); @@ -858,7 +845,7 @@ mod tests { log.insert(event_path!("service"), "payments"); log.insert(event_path!("custom"), "keep-me"); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), true); + let logs = encode_logs(vec![Event::Log(log)], true, true); let nested = logs[0]["message"] .as_object() @@ -866,7 +853,6 @@ mod tests { let message = nested["message"] .as_str() .expect("nested message should be a string"); - assert_eq!(message.len(), 900_000 + TRUNCATION_MARKER.len()); assert!(message.ends_with(TRUNCATION_MARKER)); assert_eq!(nested["custom"], "keep-me"); assert_eq!(logs[0]["service"], "payments"); @@ -880,7 +866,7 @@ mod tests { log.insert(event_path!("service"), "payments"); log.insert(event_path!("custom"), "x".repeat(200_000)); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert_eq!(logs.len(), 1); assert_eq!(logs[0]["custom"], "x".repeat(200_000)); @@ -900,7 +886,7 @@ mod tests { log.insert(event_path!("service"), "payments"); log.insert(event_path!("custom"), "\u{0001}".repeat(200_000)); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert_eq!(logs.len(), 1); assert_eq!(logs[0]["message"], "hello"); @@ -914,7 +900,7 @@ mod tests { log.insert(event_path!("service"), "payments"); log.insert(event_path!("custom"), "x".repeat(200_000)); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), true); + let logs = encode_logs(vec![Event::Log(log)], true, true); assert_eq!(logs.len(), 1); let nested = logs[0]["message"] @@ -944,7 +930,7 @@ mod tests { let mut log = LogEvent::from("e".repeat(MAX_LOG_BYTES + 1)); log.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert!(logs.is_empty()); } @@ -960,7 +946,6 @@ mod tests { vec![Event::Log(log)], Some(DatadogLogsTruncationConfig { max_log_bytes: MAX_LOG_BYTES, - max_message_bytes: 900_000, }), false, 5_000_000, @@ -991,7 +976,7 @@ mod tests { let mut log = LogEvent::from("e".repeat(MAX_LOG_BYTES + 1)); log.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert!(logs.is_empty()); let discarded_events = controller @@ -1014,7 +999,7 @@ mod tests { value!({ "body": ("f".repeat(MAX_LOG_BYTES + 1)) }), ); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert!(logs.is_empty()); } @@ -1025,7 +1010,7 @@ mod tests { let mut log = LogEvent::default(); log.insert(event_path!("custom"), "f".repeat(MAX_LOG_BYTES + 1)); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), conforms_as_agent); + let logs = encode_logs(vec![Event::Log(log)], true, conforms_as_agent); assert!(logs.is_empty()); } @@ -1036,7 +1021,7 @@ mod tests { for conforms_as_agent in [false, true] { let log = LogEvent::from("\"".repeat(600_000)); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), conforms_as_agent); + let logs = encode_logs(vec![Event::Log(log)], true, conforms_as_agent); assert_eq!(logs.len(), 1); let message = if conforms_as_agent { @@ -1057,7 +1042,7 @@ mod tests { let message = format!("{}{}", "\"".repeat(450_000), "a".repeat(110_000)); let log = LogEvent::from(message); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert_eq!(logs.len(), 1); assert!( @@ -1074,7 +1059,7 @@ mod tests { let mut log = LogEvent::from("g".repeat(MAX_LOG_BYTES + 1)); log.insert(event_path!("ddtags"), format!("env:test,{TRUNCATED_TAG}")); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert_eq!(logs[0]["ddtags"], format!("env:test,{TRUNCATED_TAG}")); } @@ -1084,7 +1069,7 @@ mod tests { let mut log = LogEvent::from("g".repeat(MAX_LOG_BYTES + 1)); log.insert(event_path!("ddtags"), "env:test"); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert_eq!(logs[0]["ddtags"], format!("env:test,{TRUNCATED_TAG}")); } @@ -1094,7 +1079,7 @@ mod tests { let original = "h".repeat(MAX_LOG_BYTES + 1); let log = LogEvent::from(original.clone()); - let logs = encode_logs(vec![Event::Log(log)], None, false); + let logs = encode_logs(vec![Event::Log(log)], false, false); assert_eq!(logs.len(), 1); assert_eq!(logs[0]["message"].as_str(), Some(original.as_str())); @@ -1108,7 +1093,7 @@ mod tests { let mut log = LogEvent::default(); log.insert(event_path!("message"), original.clone()); - let logs = encode_logs(vec![Event::Log(log)], Some(900_000), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); assert_eq!(serde_json::to_vec(&logs[0]).unwrap().len(), MAX_LOG_BYTES); assert_eq!(logs[0]["message"].as_str(), Some(original.as_str())); @@ -1120,7 +1105,7 @@ mod tests { let original = "😀".repeat(300_000); let log = LogEvent::from(original.clone()); - let logs = encode_logs(vec![Event::Log(log)], Some(900_001), false); + let logs = encode_logs(vec![Event::Log(log)], true, false); let message = logs[0]["message"] .as_str() @@ -1129,7 +1114,7 @@ mod tests { .strip_suffix(TRUNCATION_MARKER) .expect("message should have truncation marker"); assert!(original.starts_with(body)); - assert_eq!(body.len(), 900_000); + assert!(body.len() < MAX_LOG_BYTES); } #[test] @@ -1138,11 +1123,7 @@ mod tests { oversized.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); let small = LogEvent::from("ok"); - let logs = encode_logs( - vec![Event::Log(oversized), Event::Log(small)], - Some(900_000), - false, - ); + let logs = encode_logs(vec![Event::Log(oversized), Event::Log(small)], true, false); assert_eq!(logs.len(), 1); assert_eq!(logs[0]["message"], "ok"); diff --git a/src/sinks/datadog/logs/tests.rs b/src/sinks/datadog/logs/tests.rs index 4289cb82c6b2a..f08da8b9bedeb 100644 --- a/src/sinks/datadog/logs/tests.rs +++ b/src/sinks/datadog/logs/tests.rs @@ -193,7 +193,6 @@ async fn truncates_oversized_log_over_http() { [truncate_oversized_logs] max_log_bytes = 1000 - max_message_bytes = 900 "#}; let (_guard, _trigger, sink, rx) = start_test_sink(config, ApiStatus::OKv2).await; diff --git a/website/cue/reference/components/sinks/generated/datadog_logs.cue b/website/cue/reference/components/sinks/generated/datadog_logs.cue index fae334a2da816..12c708f5bb774 100644 --- a/website/cue/reference/components/sinks/generated/datadog_logs.cue +++ b/website/cue/reference/components/sinks/generated/datadog_logs.cue @@ -163,25 +163,16 @@ generated: components: sinks: datadog_logs: configuration: { description: """ Truncate logs whose encoded JSON exceeds `max_log_bytes`. - When a message is shortened, at most `max_message_bytes` raw bytes are retained and - `...TRUNCATED...` is appended. Every reduced log is tagged with `truncated:single_line`. - The message is sized to account for JSON encoding while preserving non-standard fields when - possible. If the non-message fields leave no room for a truncated message, non-standard - fields are removed and the message is sized again. Logs that still exceed the limit, or have - no string message to truncate, are dropped. + The message is shortened to the largest size that fits and `...TRUNCATED...` is appended. + Every reduced log is tagged with `truncated:single_line`. Non-standard fields are preserved + when possible, but removed when they leave no room for a truncated message. Logs that still + exceed the limit, or have no string message to truncate, are dropped. """ required: false - type: object: options: { - max_log_bytes: { - description: "Maximum encoded size, in bytes, of a log before truncation is applied." - required: false - type: uint: default: 1000000 - } - max_message_bytes: { - description: "Maximum number of message bytes to retain before appending the truncation marker." - required: false - type: uint: default: 900000 - } + type: object: options: max_log_bytes: { + description: "Maximum encoded size, in bytes, of a log before truncation is applied." + required: false + type: uint: default: 1000000 } } } From f9a620af5d681f70e6a70bc116715e064741caec Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Wed, 30 Sep 2026 17:57:30 -0600 Subject: [PATCH 05/15] Drop the custom field removal, just focus on truncating `message` --- ...dog_logs_truncate_oversized.enhancement.md | 6 +- src/sinks/datadog/logs/config.rs | 5 +- src/sinks/datadog/logs/sink.rs | 72 ++----------------- .../components/sinks/datadog_logs.cue | 11 ++- .../sinks/generated/datadog_logs.cue | 5 +- 5 files changed, 16 insertions(+), 83 deletions(-) diff --git a/changelog.d/datadog_logs_truncate_oversized.enhancement.md b/changelog.d/datadog_logs_truncate_oversized.enhancement.md index f8394eb4054f5..428f6e7ee914c 100644 --- a/changelog.d/datadog_logs_truncate_oversized.enhancement.md +++ b/changelog.d/datadog_logs_truncate_oversized.enhancement.md @@ -1,7 +1,7 @@ The `datadog_logs` sink can now optionally truncate logs that exceed Datadog's per-log size limit. Configure `truncate_oversized_logs` to set the encoded log -and retained message limits, mark shortened messages, tag reduced logs, and -preserve standard Datadog fields. Logs that cannot be reduced below the -configured limit are dropped. +limit, mark shortened messages, and tag reduced logs. Logs whose non-message +fields leave no room for a truncated message, or have no string message to +truncate, are dropped. authors: bruceg diff --git a/src/sinks/datadog/logs/config.rs b/src/sinks/datadog/logs/config.rs index 169e12fb8e123..34a5f4e24ac84 100644 --- a/src/sinks/datadog/logs/config.rs +++ b/src/sinks/datadog/logs/config.rs @@ -103,9 +103,8 @@ pub struct DatadogLogsConfig { /// Truncate logs whose encoded JSON exceeds `max_log_bytes`. /// /// The message is shortened to the largest size that fits and `...TRUNCATED...` is appended. - /// Every reduced log is tagged with `truncated:single_line`. Non-standard fields are preserved - /// when possible, but removed when they leave no room for a truncated message. Logs that still - /// exceed the limit, or have no string message to truncate, are dropped. + /// Every reduced log is tagged with `truncated:single_line`. Logs whose non-message fields + /// leave no room for a truncated message, or have no string message to truncate, are dropped. pub truncate_oversized_logs: Option, } diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index 03e5b6ba6abc6..881b32e851120 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -335,7 +335,6 @@ impl LogRequestBuilder { &mut event, !events_serialized.is_empty(), self.truncation, - self.conforms_as_agent, )? { LogEncoding::Unchanged => {} LogEncoding::Truncated => { @@ -407,7 +406,6 @@ fn encode_log( event: &mut Event, include_comma: bool, truncation: Option, - conforms_as_agent: bool, ) -> Result { let existing_len = buf.len(); let original_encoded_size = write_log(buf, event, include_comma)?; @@ -428,22 +426,14 @@ fn encode_log( let tagged_encoded_size = ensure_truncated_tag(event.as_mut_log(), original_encoded_size)?; // The encoded event consists of a fixed non-message portion and the message value. Size them // separately so the final message can be selected without repeatedly encoding the whole event. - let mut non_message_encoded_size = tagged_encoded_size - message_encoded_size; + let non_message_encoded_size = tagged_encoded_size - message_encoded_size; let select_message_body_len = |non_message_encoded_size| { truncation .max_log_bytes .checked_sub(non_message_encoded_size) .and_then(|budget| select_message_body_len(&message, message_encoded_size, budget)) }; - let mut body_len = select_message_body_len(non_message_encoded_size); - - if body_len.is_none() { - let stripped_size = strip_non_standard_fields(event.as_mut_log(), conforms_as_agent)?; - non_message_encoded_size -= stripped_size; - body_len = select_message_body_len(non_message_encoded_size); - } - - let Some(body_len) = body_len else { + let Some(body_len) = select_message_body_len(non_message_encoded_size) else { buf.truncate(existing_len); return Ok(LogEncoding::Dropped { reason: "Event remains too large after truncation.", @@ -561,45 +551,6 @@ fn message_bytes_mut(log: &mut LogEvent) -> Option<&mut Bytes> { } } -fn strip_non_standard_fields( - log: &mut LogEvent, - conforms_as_agent: bool, -) -> Result { - let Some(fields) = log.as_map_mut() else { - return Ok(0); - }; - let mut stripped_size = strip_fields(fields, |field| { - is_reserved_attribute(field) || field == MESSAGE - })?; - if conforms_as_agent && let Some(Value::Object(nested)) = fields.get_mut(MESSAGE) { - stripped_size += strip_fields(nested, |field| field == MESSAGE)?; - } - Ok(stripped_size) -} - -// Calculate the exact encoded size removed rather than using `retain`, so the message budget can -// be recomputed without encoding the entire log again. -fn strip_fields(fields: &mut ObjectMap, retain: impl Fn(&str) -> bool) -> Result { - let fields_to_strip = fields - .iter() - .filter(|(field, _)| !retain(field.as_str())) - .map(|(field, value)| { - Ok(( - field.clone(), - json_string_encoded_size(field.as_str()) + 1 + json_value_encoded_size(value)?, - )) - }) - .collect::, io::Error>>()?; - debug_assert!(fields.len() > fields_to_strip.len()); - // At least one field remains, so removing each member also removes one separating comma. - let mut stripped_size = fields_to_strip.len(); - for (field, member_size) in fields_to_strip { - let _ = fields.remove(field.as_str()); - stripped_size += member_size; - } - Ok(stripped_size) -} - fn json_string_encoded_size(value: &str) -> usize { 2 + json_string_content_size(value) } @@ -881,17 +832,14 @@ mod tests { } #[test] - fn leaves_short_message_unmarked_when_only_custom_fields_are_removed() { + fn drops_log_when_custom_fields_leave_no_room_for_message() { let mut log = LogEvent::from("hello"); log.insert(event_path!("service"), "payments"); log.insert(event_path!("custom"), "\u{0001}".repeat(200_000)); let logs = encode_logs(vec![Event::Log(log)], true, false); - assert_eq!(logs.len(), 1); - assert_eq!(logs[0]["message"], "hello"); - assert!(logs[0].get("custom").is_none()); - assert_eq!(logs[0]["ddtags"], TRUNCATED_TAG); + assert!(logs.is_empty()); } #[test] @@ -925,16 +873,6 @@ mod tests { assert!(serde_json::to_vec(&logs[0]).unwrap().len() <= MAX_LOG_BYTES); } - #[test] - fn drops_log_that_remains_oversized_after_reduction() { - let mut log = LogEvent::from("e".repeat(MAX_LOG_BYTES + 1)); - log.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); - - let logs = encode_logs(vec![Event::Log(log)], true, false); - - assert!(logs.is_empty()); - } - #[test] fn rejects_log_that_remains_oversized_after_reduction() { let (batch, mut receiver) = BatchNotifier::new_with_receiver(); @@ -1005,7 +943,7 @@ mod tests { } #[test] - fn drops_oversized_log_without_string_message_before_removing_fields() { + fn drops_oversized_log_without_string_message() { for conforms_as_agent in [false, true] { let mut log = LogEvent::default(); log.insert(event_path!("custom"), "f".repeat(MAX_LOG_BYTES + 1)); diff --git a/website/cue/reference/components/sinks/datadog_logs.cue b/website/cue/reference/components/sinks/datadog_logs.cue index 77f6def6b0842..315da1b63729e 100644 --- a/website/cue/reference/components/sinks/datadog_logs.cue +++ b/website/cue/reference/components/sinks/datadog_logs.cue @@ -85,13 +85,10 @@ components: sinks: datadog_logs: { that still exceeds `max_payload_bytes` after optional truncation is dropped. Set `truncate_oversized_logs` to opt in to reducing individual logs whose encoded JSON - exceeds `truncate_oversized_logs.max_log_bytes` (default 1,000,000). When the message - is shortened, at most `truncate_oversized_logs.max_message_bytes` raw bytes (default - 900,000) are retained and `...TRUNCATED...` is appended. Every reduced log is tagged - `truncated:single_line`. If the log remains oversized after the initial message cap, - non-standard fields are removed before the message is shortened further to account for - JSON encoding. Logs that still exceed the limit, or have no string message to truncate, - are dropped. + exceeds `truncate_oversized_logs.max_log_bytes` (default 1,000,000). The message is + shortened to the largest size that fits and `...TRUNCATED...` is appended. Every reduced + log is tagged `truncated:single_line`. Logs whose non-message fields leave no room for + a truncated message, or have no string message to truncate, are dropped. """ } } diff --git a/website/cue/reference/components/sinks/generated/datadog_logs.cue b/website/cue/reference/components/sinks/generated/datadog_logs.cue index 12c708f5bb774..1f428bd04d0f1 100644 --- a/website/cue/reference/components/sinks/generated/datadog_logs.cue +++ b/website/cue/reference/components/sinks/generated/datadog_logs.cue @@ -164,9 +164,8 @@ generated: components: sinks: datadog_logs: configuration: { Truncate logs whose encoded JSON exceeds `max_log_bytes`. The message is shortened to the largest size that fits and `...TRUNCATED...` is appended. - Every reduced log is tagged with `truncated:single_line`. Non-standard fields are preserved - when possible, but removed when they leave no room for a truncated message. Logs that still - exceed the limit, or have no string message to truncate, are dropped. + Every reduced log is tagged with `truncated:single_line`. Logs whose non-message fields + leave no room for a truncated message, or have no string message to truncate, are dropped. """ required: false type: object: options: max_log_bytes: { From 151413fb1a83834fdb3e48e8f63ab657945af881 Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Wed, 30 Sep 2026 16:56:40 -0600 Subject: [PATCH 06/15] Messages that cannot be truncated are dropped, not rejected --- src/sinks/datadog/logs/sink.rs | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index 881b32e851120..edf192b01f869 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -294,8 +294,7 @@ impl LogRequestBuilder { let (events_serialized, body, byte_size) = self.serialize_with_capacity(&mut events_with_estimated_size)?; if events_serialized.is_empty() { - if let Some((event, _)) = events_with_estimated_size.pop_front() { - event.metadata().update_status(EventStatus::Rejected); + if events_with_estimated_size.pop_front().is_some() { emit!(ComponentEventsDropped:: { count: 1, reason: "Event too large to encode." @@ -341,7 +340,6 @@ impl LogRequestBuilder { estimated_json_size = event.estimated_json_encoded_size_of(); } LogEncoding::Dropped { reason } => { - event.metadata().update_status(EventStatus::Rejected); emit!(ComponentEventsDropped:: { count: 1, reason }); continue; } @@ -874,7 +872,7 @@ mod tests { } #[test] - fn rejects_log_that_remains_oversized_after_reduction() { + fn drops_log_that_remains_oversized_after_reduction() { let (batch, mut receiver) = BatchNotifier::new_with_receiver(); let mut log = LogEvent::from("e".repeat(MAX_LOG_BYTES + 1)).with_batch_notifier(&batch); log.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); @@ -890,11 +888,11 @@ mod tests { ); assert!(requests.is_empty()); - assert_eq!(receiver.try_recv(), Ok(BatchStatus::Rejected)); + assert_eq!(receiver.try_recv(), Ok(BatchStatus::Delivered)); } #[test] - fn rejects_log_that_exceeds_payload_limit() { + fn drops_log_that_exceeds_payload_limit() { let (batch, mut receiver) = BatchNotifier::new_with_receiver(); let log = LogEvent::from("oversized").with_batch_notifier(&batch); drop(batch); @@ -902,7 +900,7 @@ mod tests { let requests = build_requests(vec![Event::Log(log)], None, false, 2); assert!(requests.is_empty()); - assert_eq!(receiver.try_recv(), Ok(BatchStatus::Rejected)); + assert_eq!(receiver.try_recv(), Ok(BatchStatus::Delivered)); } #[test] From 22a148fe3ced2651a396fa49ef538fa2fb81df7b Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Wed, 30 Sep 2026 17:07:27 -0600 Subject: [PATCH 07/15] Miscellaneous minor review feedback --- Cargo.lock | 1 - Cargo.toml | 3 +-- src/internal_events/datadog_logs.rs | 2 +- src/sinks/datadog/logs/sink.rs | 34 ++++++++++++----------------- 4 files changed, 16 insertions(+), 24 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 3449f3d96e45d..dc737f69e588f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -13532,7 +13532,6 @@ dependencies = [ "serde_with", "serde_yaml", "serial_test", - "simdutf8", "similar-asserts", "smallvec", "smpl_jwt", diff --git a/Cargo.toml b/Cargo.toml index ef735e1917e93..e41145a3cab31 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -362,7 +362,6 @@ serde_json.workspace = true serde_path_to_error.workspace = true serde_with = { version = "3.14.0", default-features = false, features = ["macros", "std"] } serde_yaml.workspace = true -simdutf8 = { workspace = true, optional = true } # Messagepack rmp-serde = { version = "1.3.0", default-features = false, optional = true } @@ -979,7 +978,7 @@ sinks-console = [] sinks-databend = ["dep:databend-client"] sinks-databricks_zerobus = ["dep:databricks-zerobus-ingest-sdk", "codecs-arrow", "arrow/ipc_compression"] sinks-datadog_events = [] -sinks-datadog_logs = ["dep:simdutf8"] +sinks-datadog_logs = [] sinks-datadog_metrics = ["dep:datadog-proto", "dep:prost", "dep:prost-reflect", "dep:datadog-agent-metrics-v3", "dep:protobuf"] sinks-datadog_traces = ["protobuf-build", "dep:datadog-proto", "dep:prost", "dep:rmpv", "dep:rmp-serde", "dep:serde_bytes"] sinks-doris = ["sqlx/mysql"] diff --git a/src/internal_events/datadog_logs.rs b/src/internal_events/datadog_logs.rs index 9f859aa6c3d10..c5d06f2d73dab 100644 --- a/src/internal_events/datadog_logs.rs +++ b/src/internal_events/datadog_logs.rs @@ -12,7 +12,7 @@ pub struct DatadogLogsEventTruncated { impl InternalEvent for DatadogLogsEventTruncated { fn emit(self) { - warn!( + debug!( message = "Truncated a Datadog log event that exceeded the per-log size limit.", max_log_bytes = self.max_log_bytes, original_encoded_size = self.original_encoded_size, diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index edf192b01f869..bc446df9e936f 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -1,4 +1,4 @@ -use std::{borrow::Cow, collections::VecDeque, fmt::Debug, io, sync::Arc}; +use std::{collections::VecDeque, fmt::Debug, io, sync::Arc}; use bytes::Bytes; use itertools::Itertools; @@ -9,7 +9,10 @@ use vector_lib::{ internal_event::{ComponentEventsDropped, UNINTENTIONAL}, lookup::event_path, }; -use vrl::path::{OwnedSegment, OwnedTargetPath, PathPrefix}; +use vrl::{ + path::{OwnedSegment, OwnedTargetPath, PathPrefix}, + value::value::simdutf_bytes_utf8_lossy, +}; use super::{config::DatadogLogsTruncationConfig, service::LogApiRequest}; use crate::{ @@ -419,19 +422,17 @@ fn encode_log( reason: "Oversized event has no string message to truncate.", }); }; - let message = simdutf8_lossy(&message); + let message = simdutf_bytes_utf8_lossy(&message); let message_encoded_size = json_string_encoded_size(&message); let tagged_encoded_size = ensure_truncated_tag(event.as_mut_log(), original_encoded_size)?; // The encoded event consists of a fixed non-message portion and the message value. Size them // separately so the final message can be selected without repeatedly encoding the whole event. let non_message_encoded_size = tagged_encoded_size - message_encoded_size; - let select_message_body_len = |non_message_encoded_size| { - truncation - .max_log_bytes - .checked_sub(non_message_encoded_size) - .and_then(|budget| select_message_body_len(&message, message_encoded_size, budget)) - }; - let Some(body_len) = select_message_body_len(non_message_encoded_size) else { + let Some(body_len) = truncation + .max_log_bytes + .checked_sub(non_message_encoded_size) + .and_then(|budget| select_message_body_len(&message, message_encoded_size, budget)) + else { buf.truncate(existing_len); return Ok(LogEncoding::Dropped { reason: "Event remains too large after truncation.", @@ -457,13 +458,6 @@ fn encode_log( Ok(LogEncoding::Truncated) } -fn simdutf8_lossy(bytes: &[u8]) -> Cow<'_, str> { - match simdutf8::basic::from_utf8(bytes) { - Ok(value) => Cow::Borrowed(value), - Err(_) => String::from_utf8_lossy(bytes), - } -} - fn select_message_body_len( message: &str, message_encoded_size: usize, @@ -664,7 +658,7 @@ mod tests { use super::{ LogRequestBuilder, json_string_encoded_size, normalize_as_agent_event, normalize_event, - simdutf8_lossy, + simdutf_bytes_utf8_lossy, }; use crate::{ common::datadog::DD_RESERVED_SEMANTIC_ATTRS, @@ -695,7 +689,7 @@ mod tests { let invalid = bytes::Bytes::from_static(&[b'a', 0xff, b'b']); assert_eq!( - json_string_encoded_size(&simdutf8_lossy(&invalid)), + json_string_encoded_size(&simdutf_bytes_utf8_lossy(&invalid)), serde_json::to_vec(&Value::Bytes(invalid)).unwrap().len() ); } @@ -974,7 +968,7 @@ mod tests { } #[test] - fn second_message_shrink_reserves_space_for_new_marker() { + fn truncates_escaped_message_with_plain_suffix_to_fit() { let message = format!("{}{}", "\"".repeat(450_000), "a".repeat(110_000)); let log = LogEvent::from(message); From b9b4d1d73db06dd4e91b86591f6d388bb36f5829 Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Thu, 1 Oct 2026 12:27:10 -0600 Subject: [PATCH 08/15] Mark dropped events as `INTENTIONAL` drops since it is a config policy choice --- src/sinks/datadog/logs/sink.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index bc446df9e936f..6e28604ad1ade 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -6,7 +6,7 @@ use snafu::Snafu; use tracing::Instrument; use vector_lib::{ event::{ObjectMap, Value}, - internal_event::{ComponentEventsDropped, UNINTENTIONAL}, + internal_event::{ComponentEventsDropped, INTENTIONAL}, lookup::event_path, }; use vrl::{ @@ -298,7 +298,7 @@ impl LogRequestBuilder { self.serialize_with_capacity(&mut events_with_estimated_size)?; if events_serialized.is_empty() { if events_with_estimated_size.pop_front().is_some() { - emit!(ComponentEventsDropped:: { + emit!(ComponentEventsDropped:: { count: 1, reason: "Event too large to encode." }); @@ -343,7 +343,7 @@ impl LogRequestBuilder { estimated_json_size = event.estimated_json_encoded_size_of(); } LogEncoding::Dropped { reason } => { - emit!(ComponentEventsDropped:: { count: 1, reason }); + emit!(ComponentEventsDropped:: { count: 1, reason }); continue; } } From 7f19aad3e438c167a51a142356ad4c1db6dd4a59 Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Thu, 1 Oct 2026 14:46:25 -0600 Subject: [PATCH 09/15] Gate nested-message truncation on agent normalization --- src/sinks/datadog/logs/sink.rs | 35 ++++++++++++++++++++++++++++------ 1 file changed, 29 insertions(+), 6 deletions(-) diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index 6e28604ad1ade..681a2f7b10286 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -337,6 +337,7 @@ impl LogRequestBuilder { &mut event, !events_serialized.is_empty(), self.truncation, + self.conforms_as_agent, )? { LogEncoding::Unchanged => {} LogEncoding::Truncated => { @@ -407,6 +408,7 @@ fn encode_log( event: &mut Event, include_comma: bool, truncation: Option, + conforms_as_agent: bool, ) -> Result { let existing_len = buf.len(); let original_encoded_size = write_log(buf, event, include_comma)?; @@ -416,7 +418,9 @@ fn encode_log( return Ok(LogEncoding::Unchanged); }; - let Some(message) = message_bytes_mut(event.as_mut_log()).map(|message| message.clone()) else { + let Some(message) = + message_bytes_mut(event.as_mut_log(), conforms_as_agent).map(|message| message.clone()) + else { buf.truncate(existing_len); return Ok(LogEncoding::Dropped { reason: "Oversized event has no string message to truncate.", @@ -439,7 +443,7 @@ fn encode_log( }); }; if body_len < message.len() { - set_truncated_message(event.as_mut_log(), &message, body_len); + set_truncated_message(event.as_mut_log(), &message, body_len, conforms_as_agent); } buf.truncate(existing_len); @@ -496,12 +500,18 @@ fn write_log(buf: &mut Vec, event: &Event, include_comma: bool) -> Result Result { @@ -532,10 +542,10 @@ fn ensure_truncated_tag(log: &mut LogEvent, encoded_size: usize) -> Result Option<&mut Bytes> { +fn message_bytes_mut(log: &mut LogEvent, conforms_as_agent: bool) -> Option<&mut Bytes> { match log.as_map_mut()?.get_mut(MESSAGE)? { Value::Bytes(message) => Some(message), - Value::Object(fields) => match fields.get_mut(MESSAGE)? { + Value::Object(fields) if conforms_as_agent => match fields.get_mut(MESSAGE)? { Value::Bytes(message) => Some(message), _ => None, }, @@ -934,6 +944,19 @@ mod tests { assert!(logs.is_empty()); } + #[test] + fn drops_oversized_structured_message_without_agent_normalization() { + let mut log = LogEvent::default(); + log.insert( + event_path!("message"), + value!({ "message": ("f".repeat(MAX_LOG_BYTES + 1)), "other": "preserve-me" }), + ); + + let logs = encode_logs(vec![Event::Log(log)], true, false); + + assert!(logs.is_empty()); + } + #[test] fn drops_oversized_log_without_string_message() { for conforms_as_agent in [false, true] { From 7e7154b6835ed274649af27f1ccaf1ea34b28d89 Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Thu, 1 Oct 2026 16:15:12 -0600 Subject: [PATCH 10/15] Replace unreachable drop conditional --- src/sinks/datadog/logs/sink.rs | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index 681a2f7b10286..b64fcb710a1c0 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -449,12 +449,10 @@ fn encode_log( buf.truncate(existing_len); let encoded_size = write_log(buf, event, include_comma)?; - if encoded_size > truncation.max_log_bytes { - buf.truncate(existing_len); - return Ok(LogEncoding::Dropped { - reason: "Event remains too large after truncation.", - }); - } + // The message budget uses the same UTF-8 bytes and JSON sizing as serialization, so a + // mismatch here indicates that the size calculation has diverged. Request serialization + // separately enforces `max_payload_bytes` in release builds. + debug_assert!(encoded_size <= truncation.max_log_bytes); emit!(DatadogLogsEventTruncated { max_log_bytes: truncation.max_log_bytes, original_encoded_size, From 8305c63d79733ce4824b405eb7829ae8d906e9df Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Fri, 2 Oct 2026 08:29:42 -0600 Subject: [PATCH 11/15] Don't drop events that cannot be truncated to fit in max_log_bytes --- src/sinks/datadog/logs/config.rs | 6 ++- src/sinks/datadog/logs/sink.rs | 39 ++++++++++++------- .../sinks/generated/datadog_logs.cue | 6 ++- 3 files changed, 33 insertions(+), 18 deletions(-) diff --git a/src/sinks/datadog/logs/config.rs b/src/sinks/datadog/logs/config.rs index 34a5f4e24ac84..302caf6e79b78 100644 --- a/src/sinks/datadog/logs/config.rs +++ b/src/sinks/datadog/logs/config.rs @@ -100,11 +100,13 @@ pub struct DatadogLogsConfig { #[serde(default = "default_max_payload_bytes")] pub max_payload_bytes: Option, - /// Truncate logs whose encoded JSON exceeds `max_log_bytes`. + /// Attempt to truncate logs whose encoded JSON exceeds `max_log_bytes`. /// /// The message is shortened to the largest size that fits and `...TRUNCATED...` is appended. /// Every reduced log is tagged with `truncated:single_line`. Logs whose non-message fields - /// leave no room for a truncated message, or have no string message to truncate, are dropped. + /// leave no room for a truncated message are sent unchanged if they fit + /// `max_payload_bytes`; the Datadog intake can further truncate them. Logs without a string + /// message to truncate are dropped. pub truncate_oversized_logs: Option, } diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index b64fcb710a1c0..55538a5e3e658 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -428,6 +428,7 @@ fn encode_log( }; let message = simdutf_bytes_utf8_lossy(&message); let message_encoded_size = json_string_encoded_size(&message); + let previous_tags = event.as_log().get(event_path!(DDTAGS)).cloned(); let tagged_encoded_size = ensure_truncated_tag(event.as_mut_log(), original_encoded_size)?; // The encoded event consists of a fixed non-message portion and the message value. Size them // separately so the final message can be selected without repeatedly encoding the whole event. @@ -437,10 +438,13 @@ fn encode_log( .checked_sub(non_message_encoded_size) .and_then(|budget| select_message_body_len(&message, message_encoded_size, budget)) else { - buf.truncate(existing_len); - return Ok(LogEncoding::Dropped { - reason: "Event remains too large after truncation.", - }); + if let Some(tags) = previous_tags { + event.as_mut_log().insert(event_path!(DDTAGS), tags); + } else { + // Remove the temporary tag that `ensure_truncated_tag` inserted for size calculation. + event.as_mut_log().remove(event_path!(DDTAGS)); + } + return Ok(LogEncoding::Unchanged); }; if body_len < message.len() { set_truncated_message(event.as_mut_log(), &message, body_len, conforms_as_agent); @@ -673,6 +677,7 @@ mod tests { sinks::{ datadog::logs::config::{ DEFAULT_MAX_LOG_BYTES as MAX_LOG_BYTES, DatadogLogsTruncationConfig, + MAX_PAYLOAD_BYTES, }, util::Compression, }, @@ -832,14 +837,18 @@ mod tests { } #[test] - fn drops_log_when_custom_fields_leave_no_room_for_message() { + fn forwards_log_when_custom_fields_leave_no_room_for_message() { let mut log = LogEvent::from("hello"); log.insert(event_path!("service"), "payments"); log.insert(event_path!("custom"), "\u{0001}".repeat(200_000)); let logs = encode_logs(vec![Event::Log(log)], true, false); - assert!(logs.is_empty()); + assert_eq!(logs.len(), 1); + assert_eq!(logs[0]["message"], "hello"); + assert_eq!(logs[0]["service"], "payments"); + assert_eq!(logs[0]["custom"], "\u{0001}".repeat(200_000)); + assert!(logs[0].get("ddtags").is_none()); } #[test] @@ -874,10 +883,10 @@ mod tests { } #[test] - fn drops_log_that_remains_oversized_after_reduction() { + fn drops_log_that_exceeds_payload_limit_after_best_effort_truncation() { let (batch, mut receiver) = BatchNotifier::new_with_receiver(); let mut log = LogEvent::from("e".repeat(MAX_LOG_BYTES + 1)).with_batch_notifier(&batch); - log.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); + log.insert(event_path!("service"), "x".repeat(MAX_PAYLOAD_BYTES + 1)); drop(batch); let requests = build_requests( @@ -906,7 +915,7 @@ mod tests { } #[test] - fn counts_irreducible_log_drop_once() { + fn does_not_count_best_effort_log_as_dropped() { vector_lib::metrics::init_test(); let controller = vector_lib::metrics::Controller::get().unwrap(); controller.reset(); @@ -916,7 +925,7 @@ mod tests { let logs = encode_logs(vec![Event::Log(log)], true, false); - assert!(logs.is_empty()); + assert_eq!(logs.len(), 1); let discarded_events = controller .capture_metrics() .iter() @@ -926,7 +935,7 @@ mod tests { _ => panic!("discarded events metric must be a counter"), }) .sum::(); - assert_eq!(discarded_events, 1.0); + assert_eq!(discarded_events, 0.0); } #[test] @@ -1069,15 +1078,17 @@ mod tests { } #[test] - fn delivers_following_log_after_dropping_irreducible_log() { + fn delivers_following_log_after_forwarding_best_effort_log() { let mut oversized = LogEvent::from("i".repeat(MAX_LOG_BYTES + 1)); oversized.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); let small = LogEvent::from("ok"); let logs = encode_logs(vec![Event::Log(oversized), Event::Log(small)], true, false); - assert_eq!(logs.len(), 1); - assert_eq!(logs[0]["message"], "ok"); + assert_eq!(logs.len(), 2); + assert_eq!(logs[0]["message"], "i".repeat(MAX_LOG_BYTES + 1)); + assert!(logs[0].get("ddtags").is_none()); + assert_eq!(logs[1]["message"], "ok"); } fn assert_normalized_log_has_expected_attrs(log: &LogEvent) { diff --git a/website/cue/reference/components/sinks/generated/datadog_logs.cue b/website/cue/reference/components/sinks/generated/datadog_logs.cue index 1f428bd04d0f1..52ab116748e40 100644 --- a/website/cue/reference/components/sinks/generated/datadog_logs.cue +++ b/website/cue/reference/components/sinks/generated/datadog_logs.cue @@ -161,11 +161,13 @@ generated: components: sinks: datadog_logs: configuration: { } truncate_oversized_logs: { description: """ - Truncate logs whose encoded JSON exceeds `max_log_bytes`. + Attempt to truncate logs whose encoded JSON exceeds `max_log_bytes`. The message is shortened to the largest size that fits and `...TRUNCATED...` is appended. Every reduced log is tagged with `truncated:single_line`. Logs whose non-message fields - leave no room for a truncated message, or have no string message to truncate, are dropped. + leave no room for a truncated message are sent unchanged if they fit + `max_payload_bytes`; the Datadog intake can further truncate them. Logs without a string + message to truncate are dropped. """ required: false type: object: options: max_log_bytes: { From ee5750ca9269a31be4c3ccd5170111ed0bd91b53 Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Fri, 2 Oct 2026 09:11:09 -0600 Subject: [PATCH 12/15] Fix changelog --- changelog.d/datadog_logs_truncate_oversized.enhancement.md | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/changelog.d/datadog_logs_truncate_oversized.enhancement.md b/changelog.d/datadog_logs_truncate_oversized.enhancement.md index 428f6e7ee914c..df663cdfee879 100644 --- a/changelog.d/datadog_logs_truncate_oversized.enhancement.md +++ b/changelog.d/datadog_logs_truncate_oversized.enhancement.md @@ -1,7 +1,8 @@ The `datadog_logs` sink can now optionally truncate logs that exceed Datadog's per-log size limit. Configure `truncate_oversized_logs` to set the encoded log limit, mark shortened messages, and tag reduced logs. Logs whose non-message -fields leave no room for a truncated message, or have no string message to -truncate, are dropped. +fields leave no room for a truncated message are forwarded unchanged when they +fit the payload limit, allowing the Datadog intake to truncate them. Logs +without a string message to truncate are dropped. authors: bruceg From 025a7e60fed3186f6af746cd2088e231c986ad96 Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Fri, 2 Oct 2026 10:47:59 -0600 Subject: [PATCH 13/15] Apply same best-effort for logs without `message` --- ...dog_logs_truncate_oversized.enhancement.md | 6 +- src/sinks/datadog/logs/config.rs | 7 +- src/sinks/datadog/logs/sink.rs | 102 +++++++++--------- .../sinks/generated/datadog_logs.cue | 7 +- 4 files changed, 57 insertions(+), 65 deletions(-) diff --git a/changelog.d/datadog_logs_truncate_oversized.enhancement.md b/changelog.d/datadog_logs_truncate_oversized.enhancement.md index df663cdfee879..ad24bd86ec3fa 100644 --- a/changelog.d/datadog_logs_truncate_oversized.enhancement.md +++ b/changelog.d/datadog_logs_truncate_oversized.enhancement.md @@ -1,8 +1,8 @@ The `datadog_logs` sink can now optionally truncate logs that exceed Datadog's per-log size limit. Configure `truncate_oversized_logs` to set the encoded log limit, mark shortened messages, and tag reduced logs. Logs whose non-message -fields leave no room for a truncated message are forwarded unchanged when they -fit the payload limit, allowing the Datadog intake to truncate them. Logs -without a string message to truncate are dropped. +fields leave no room for a truncated message, or have no string message, are +forwarded unchanged when they fit the payload limit, allowing the Datadog +intake to truncate them. authors: bruceg diff --git a/src/sinks/datadog/logs/config.rs b/src/sinks/datadog/logs/config.rs index 302caf6e79b78..fbff6e141ad75 100644 --- a/src/sinks/datadog/logs/config.rs +++ b/src/sinks/datadog/logs/config.rs @@ -103,10 +103,9 @@ pub struct DatadogLogsConfig { /// Attempt to truncate logs whose encoded JSON exceeds `max_log_bytes`. /// /// The message is shortened to the largest size that fits and `...TRUNCATED...` is appended. - /// Every reduced log is tagged with `truncated:single_line`. Logs whose non-message fields - /// leave no room for a truncated message are sent unchanged if they fit - /// `max_payload_bytes`; the Datadog intake can further truncate them. Logs without a string - /// message to truncate are dropped. + /// Every reduced log is tagged with `truncated:single_line`. Logs with no string message, or + /// whose non-message fields leave no room for a truncated message, are sent unchanged if they + /// fit `max_payload_bytes`; the Datadog intake can further truncate them. pub truncate_oversized_logs: Option, } diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index 55538a5e3e658..b2b9c60162852 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -332,21 +332,14 @@ impl LogRequestBuilder { buf.push(b'['); while let Some((mut event, mut estimated_json_size)) = events.pop_front() { let existing_len = buf.len(); - match encode_log( + if encode_log( &mut buf, &mut event, !events_serialized.is_empty(), self.truncation, self.conforms_as_agent, )? { - LogEncoding::Unchanged => {} - LogEncoding::Truncated => { - estimated_json_size = event.estimated_json_encoded_size_of(); - } - LogEncoding::Dropped { reason } => { - emit!(ComponentEventsDropped:: { count: 1, reason }); - continue; - } + estimated_json_size = event.estimated_json_encoded_size_of(); } if buf.len() >= self.max_payload_bytes { @@ -397,34 +390,29 @@ impl LogRequestBuilder { } } -enum LogEncoding { - Unchanged, - Truncated, - Dropped { reason: &'static str }, -} - +/// Encodes a log event and returns whether it was locally truncated. +/// +/// Returns `true` if the message was shortened and tagged, so the caller must recompute its +/// estimated encoded size. fn encode_log( buf: &mut Vec, event: &mut Event, include_comma: bool, truncation: Option, conforms_as_agent: bool, -) -> Result { +) -> Result { let existing_len = buf.len(); let original_encoded_size = write_log(buf, event, include_comma)?; let Some(truncation) = truncation.filter(|truncation| original_encoded_size > truncation.max_log_bytes) else { - return Ok(LogEncoding::Unchanged); + return Ok(false); }; let Some(message) = message_bytes_mut(event.as_mut_log(), conforms_as_agent).map(|message| message.clone()) else { - buf.truncate(existing_len); - return Ok(LogEncoding::Dropped { - reason: "Oversized event has no string message to truncate.", - }); + return Ok(false); }; let message = simdutf_bytes_utf8_lossy(&message); let message_encoded_size = json_string_encoded_size(&message); @@ -444,7 +432,7 @@ fn encode_log( // Remove the temporary tag that `ensure_truncated_tag` inserted for size calculation. event.as_mut_log().remove(event_path!(DDTAGS)); } - return Ok(LogEncoding::Unchanged); + return Ok(false); }; if body_len < message.len() { set_truncated_message(event.as_mut_log(), &message, body_len, conforms_as_agent); @@ -461,7 +449,7 @@ fn encode_log( max_log_bytes: truncation.max_log_bytes, original_encoded_size, }); - Ok(LogEncoding::Truncated) + Ok(true) } fn select_message_body_len( @@ -939,20 +927,7 @@ mod tests { } #[test] - fn drops_oversized_log_with_non_string_message() { - let mut log = LogEvent::default(); - log.insert( - event_path!("message"), - value!({ "body": ("f".repeat(MAX_LOG_BYTES + 1)) }), - ); - - let logs = encode_logs(vec![Event::Log(log)], true, false); - - assert!(logs.is_empty()); - } - - #[test] - fn drops_oversized_structured_message_without_agent_normalization() { + fn forwards_oversized_structured_message_without_agent_normalization() { let mut log = LogEvent::default(); log.insert( event_path!("message"), @@ -961,21 +936,54 @@ mod tests { let logs = encode_logs(vec![Event::Log(log)], true, false); - assert!(logs.is_empty()); + assert_eq!(logs.len(), 1); + assert_eq!(logs[0]["message"]["message"], "f".repeat(MAX_LOG_BYTES + 1)); + assert_eq!(logs[0]["message"]["other"], "preserve-me"); + assert!(logs[0].get("ddtags").is_none()); } #[test] - fn drops_oversized_log_without_string_message() { + fn forwards_oversized_log_without_string_message() { for conforms_as_agent in [false, true] { let mut log = LogEvent::default(); log.insert(event_path!("custom"), "f".repeat(MAX_LOG_BYTES + 1)); let logs = encode_logs(vec![Event::Log(log)], true, conforms_as_agent); - assert!(logs.is_empty()); + assert_eq!(logs.len(), 1); + let custom = if conforms_as_agent { + &logs[0]["message"]["custom"] + } else { + &logs[0]["custom"] + }; + assert_eq!(custom, &"f".repeat(MAX_LOG_BYTES + 1)); + assert!(logs[0].get("ddtags").is_none()); } } + #[test] + fn drops_non_string_message_that_exceeds_payload_limit() { + let (batch, mut receiver) = BatchNotifier::new_with_receiver(); + let mut log = LogEvent::default().with_batch_notifier(&batch); + log.insert( + event_path!("message"), + value!({ "body": ("f".repeat(MAX_PAYLOAD_BYTES + 1)) }), + ); + drop(batch); + + let requests = build_requests( + vec![Event::Log(log)], + Some(DatadogLogsTruncationConfig { + max_log_bytes: MAX_LOG_BYTES, + }), + false, + MAX_PAYLOAD_BYTES, + ); + + assert!(requests.is_empty()); + assert_eq!(receiver.try_recv(), Ok(BatchStatus::Delivered)); + } + #[test] fn shrinks_escaped_message_again_to_fit_default_limits() { for conforms_as_agent in [false, true] { @@ -1077,20 +1085,6 @@ mod tests { assert!(body.len() < MAX_LOG_BYTES); } - #[test] - fn delivers_following_log_after_forwarding_best_effort_log() { - let mut oversized = LogEvent::from("i".repeat(MAX_LOG_BYTES + 1)); - oversized.insert(event_path!("service"), "x".repeat(MAX_LOG_BYTES + 1)); - let small = LogEvent::from("ok"); - - let logs = encode_logs(vec![Event::Log(oversized), Event::Log(small)], true, false); - - assert_eq!(logs.len(), 2); - assert_eq!(logs[0]["message"], "i".repeat(MAX_LOG_BYTES + 1)); - assert!(logs[0].get("ddtags").is_none()); - assert_eq!(logs[1]["message"], "ok"); - } - fn assert_normalized_log_has_expected_attrs(log: &LogEvent) { assert!( log.get(event_path!("timestamp")) diff --git a/website/cue/reference/components/sinks/generated/datadog_logs.cue b/website/cue/reference/components/sinks/generated/datadog_logs.cue index 52ab116748e40..67605e273bd9d 100644 --- a/website/cue/reference/components/sinks/generated/datadog_logs.cue +++ b/website/cue/reference/components/sinks/generated/datadog_logs.cue @@ -164,10 +164,9 @@ generated: components: sinks: datadog_logs: configuration: { Attempt to truncate logs whose encoded JSON exceeds `max_log_bytes`. The message is shortened to the largest size that fits and `...TRUNCATED...` is appended. - Every reduced log is tagged with `truncated:single_line`. Logs whose non-message fields - leave no room for a truncated message are sent unchanged if they fit - `max_payload_bytes`; the Datadog intake can further truncate them. Logs without a string - message to truncate are dropped. + Every reduced log is tagged with `truncated:single_line`. Logs with no string message, or + whose non-message fields leave no room for a truncated message, are sent unchanged if they + fit `max_payload_bytes`; the Datadog intake can further truncate them. """ required: false type: object: options: max_log_bytes: { From 1c2be343295ea201a64b37f7c918212571072f68 Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Fri, 2 Oct 2026 15:14:52 -0600 Subject: [PATCH 14/15] Address review feedback --- src/sinks/datadog/logs/sink.rs | 35 +++++++++++++++++----------------- 1 file changed, 18 insertions(+), 17 deletions(-) diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index b2b9c60162852..b08d524e2ec6d 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -318,11 +318,10 @@ impl LogRequestBuilder { /// /// Returns the serialized events, the buffer, and the byte size of the events. Events that do /// not fit remain in the `events` parameter. Events rejected during encoding are removed. - #[doc(hidden)] fn serialize_with_capacity( &self, events: &mut VecDeque<(Event, JsonSize)>, - ) -> Result<(Vec, Vec, GroupedCountByteSize), io::Error> { + ) -> io::Result<(Vec, Vec, GroupedCountByteSize)> { let total_estimated = events.iter().map(|(_, size)| size.get()).sum::() + events.len() * 2; let mut buf = Vec::with_capacity(total_estimated); @@ -400,7 +399,7 @@ fn encode_log( include_comma: bool, truncation: Option, conforms_as_agent: bool, -) -> Result { +) -> io::Result { let existing_len = buf.len(); let original_encoded_size = write_log(buf, event, include_comma)?; let Some(truncation) = @@ -478,7 +477,7 @@ fn select_message_body_len( } } -fn write_log(buf: &mut Vec, event: &Event, include_comma: bool) -> Result { +fn write_log(buf: &mut Vec, event: &Event, include_comma: bool) -> io::Result { if include_comma { buf.push(b','); } @@ -504,7 +503,7 @@ fn set_truncated_message( Bytes::from(truncated); } -fn ensure_truncated_tag(log: &mut LogEvent, encoded_size: usize) -> Result { +fn ensure_truncated_tag(log: &mut LogEvent, encoded_size: usize) -> io::Result { let tags_path = event_path!(DDTAGS); let previous_size = log .get(tags_path) @@ -559,7 +558,7 @@ const fn json_character_encoded_size(character: char) -> usize { } } -fn json_value_encoded_size(value: &Value) -> Result { +fn json_value_encoded_size(value: &Value) -> io::Result { let mut sink = io::sink(); encoding::as_tracked_write(&mut sink, value, |writer, value| { serde_json::to_writer(writer, value) @@ -1070,19 +1069,21 @@ mod tests { #[test] fn truncation_preserves_utf8_boundaries() { - let original = "😀".repeat(300_000); - let log = LogEvent::from(original.clone()); + for prefix_len in 0..4 { + let original = format!("{}{}", "x".repeat(prefix_len), "😀".repeat(300_000)); + let log = LogEvent::from(original.clone()); - let logs = encode_logs(vec![Event::Log(log)], true, false); + let logs = encode_logs(vec![Event::Log(log)], true, false); - let message = logs[0]["message"] - .as_str() - .expect("message should be valid UTF-8"); - let body = message - .strip_suffix(TRUNCATION_MARKER) - .expect("message should have truncation marker"); - assert!(original.starts_with(body)); - assert!(body.len() < MAX_LOG_BYTES); + let message = logs[0]["message"] + .as_str() + .expect("message should be valid UTF-8"); + let body = message + .strip_suffix(TRUNCATION_MARKER) + .expect("message should have truncation marker"); + assert!(original.starts_with(body)); + assert!(body.len() < MAX_LOG_BYTES); + } } fn assert_normalized_log_has_expected_attrs(log: &LogEvent) { From 0fdf6d44e8311098efd494f8c47bcce0cd0afe87 Mon Sep 17 00:00:00 2001 From: Bruce Guenter Date: Fri, 2 Oct 2026 15:40:48 -0600 Subject: [PATCH 15/15] Keep payload-limit drops unintentional --- src/sinks/datadog/logs/sink.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/sinks/datadog/logs/sink.rs b/src/sinks/datadog/logs/sink.rs index b08d524e2ec6d..6f77622544117 100644 --- a/src/sinks/datadog/logs/sink.rs +++ b/src/sinks/datadog/logs/sink.rs @@ -6,7 +6,7 @@ use snafu::Snafu; use tracing::Instrument; use vector_lib::{ event::{ObjectMap, Value}, - internal_event::{ComponentEventsDropped, INTENTIONAL}, + internal_event::{ComponentEventsDropped, UNINTENTIONAL}, lookup::event_path, }; use vrl::{ @@ -298,7 +298,7 @@ impl LogRequestBuilder { self.serialize_with_capacity(&mut events_with_estimated_size)?; if events_serialized.is_empty() { if events_with_estimated_size.pop_front().is_some() { - emit!(ComponentEventsDropped:: { + emit!(ComponentEventsDropped:: { count: 1, reason: "Event too large to encode." });