diff --git a/Cargo.toml b/Cargo.toml index c27081250a296..e41145a3cab31 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" 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..ad24bd86ec3fa --- /dev/null +++ b/changelog.d/datadog_logs_truncate_oversized.enhancement.md @@ -0,0 +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, are +forwarded unchanged when they fit the payload limit, allowing the Datadog +intake to truncate them. + +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..c5d06f2d73dab 100644 --- a/src/internal_events/datadog_logs.rs +++ b/src/internal_events/datadog_logs.rs @@ -4,6 +4,24 @@ use vector_lib::{ }; use vrl::path::OwnedTargetPath; +#[derive(Debug, NamedInternalEvent)] +pub struct DatadogLogsEventTruncated { + pub max_log_bytes: usize, + pub original_encoded_size: usize, +} + +impl InternalEvent for DatadogLogsEventTruncated { + fn emit(self) { + 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, + 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..fbff6e141ad75 100644 --- a/src/sinks/datadog/logs/config.rs +++ b/src/sinks/datadog/logs/config.rs @@ -34,6 +34,7 @@ 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 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 +49,18 @@ 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, +} + /// Configuration for the `datadog_logs` sink. #[configurable_component(sink("datadog_logs", "Publish log events to Datadog."))] #[derive(Clone, Debug, Derivative)] @@ -81,17 +94,29 @@ 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, + + /// 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 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, } 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_compression() -> Option { Some(Compression::zstd_default()) } @@ -175,6 +200,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 +279,19 @@ 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 batch_goal_bytes = self.max_payload_bytes.unwrap_or(MAX_PAYLOAD_BYTES) - BATCH_HEADROOM_BYTES; @@ -309,6 +348,46 @@ 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); + } + + #[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 + "#}) + .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 + "#}) + .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..6f77622544117 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; @@ -8,12 +9,15 @@ 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::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 +45,7 @@ pub struct LogSinkBuilder { protocol: String, conforms_as_agent: bool, max_payload_bytes: usize, + truncation: Option, } impl LogSinkBuilder { @@ -62,6 +67,7 @@ impl LogSinkBuilder { protocol, conforms_as_agent, max_payload_bytes, + truncation: None, } } @@ -70,6 +76,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 +91,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 +118,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 +268,7 @@ struct LogRequestBuilder { pub compression: Compression, pub conforms_as_agent: bool, pub max_payload_bytes: usize, + pub truncation: Option, } impl LogRequestBuilder { @@ -280,24 +295,65 @@ 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) } + /// 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. + fn serialize_with_capacity( + &self, + events: &mut VecDeque<(Event, JsonSize)>, + ) -> 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); + 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(); + if encode_log( + &mut buf, + &mut event, + !events_serialized.is_empty(), + self.truncation, + self.conforms_as_agent, + )? { + estimated_json_size = event.estimated_json_encoded_size_of(); + } + + 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 +389,180 @@ impl LogRequestBuilder { } } -/// Serialize events into a buffer as a JSON array that has a maximum size of -/// `max_payload_bytes`. +/// Encodes a log event and returns whether it was locally truncated. /// -/// 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; +/// 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, +) -> io::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(false); + }; + + let Some(message) = + message_bytes_mut(event.as_mut_log(), conforms_as_agent).map(|message| message.clone()) + else { + return Ok(false); + }; + 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. + let non_message_encoded_size = tagged_encoded_size - message_encoded_size; + 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 { + if let Some(tags) = previous_tags { + event.as_mut_log().insert(event_path!(DDTAGS), tags); } else { - buf.push(b','); + // Remove the temporary tag that `ensure_truncated_tag` inserted for size calculation. + event.as_mut_log().remove(event_path!(DDTAGS)); + } + return Ok(false); + }; + if body_len < message.len() { + set_truncated_message(event.as_mut_log(), &message, body_len, conforms_as_agent); + } + + buf.truncate(existing_len); + let encoded_size = write_log(buf, event, include_comma)?; + + // 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, + }); + Ok(true) +} + +fn select_message_body_len( + message: &str, + message_encoded_size: usize, + encoded_budget: usize, +) -> Option { + if 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_encoded_size > content_budget { + break; + } + body_len = next_body_len; + encoded_size = next_encoded_size; } - 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; + + (body_len < message.len()).then_some(body_len) + } +} + +fn write_log(buf: &mut Vec, event: &Event, include_comma: bool) -> io::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) +} + +const TRUNCATION_MARKER: &str = "...TRUNCATED..."; +const TRUNCATED_TAG: &str = "truncated:single_line"; + +fn set_truncated_message( + log: &mut LogEvent, + message: &str, + body_len: usize, + conforms_as_agent: bool, +) { + let marker = TRUNCATION_MARKER.as_bytes(); + let mut truncated = Vec::with_capacity(body_len + marker.len()); + truncated.extend_from_slice(&message.as_bytes()[..body_len]); + truncated.extend_from_slice(marker); + *message_bytes_mut(log, conforms_as_agent).expect("the message was previously found") = + Bytes::from(truncated); +} + +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) + .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 Ok(encoded_size); } - // 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); + } + 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, conforms_as_agent: bool) -> Option<&mut Bytes> { + match log.as_map_mut()?.get_mut(MESSAGE)? { + Value::Bytes(message) => Some(message), + Value::Object(fields) if conforms_as_agent => match fields.get_mut(MESSAGE)? { + Value::Bytes(message) => Some(message), + _ => None, + }, + _ => None, } - buf.push(b']'); +} - Ok((events_serialized, buf, byte_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(), + } +} + +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) + }) } impl LogSink @@ -396,6 +583,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, |_| { @@ -458,6 +646,7 @@ mod tests { use vector_lib::{ config::{LegacyKey, LogNamespace}, event::{Event, EventMetadata, LogEvent}, + finalization::{BatchNotifier, BatchStatus}, schema::{Definition, meaning}, }; use vrl::{ @@ -466,8 +655,436 @@ 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, json_string_encoded_size, normalize_as_agent_event, normalize_event, + simdutf_bytes_utf8_lossy, + }; + use crate::{ + common::datadog::DD_RESERVED_SEMANTIC_ATTRS, + sinks::{ + datadog::logs::config::{ + DEFAULT_MAX_LOG_BYTES as MAX_LOG_BYTES, DatadogLogsTruncationConfig, + MAX_PAYLOAD_BYTES, + }, + util::Compression, + }, + }; + + 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(&simdutf_bytes_utf8_lossy(&invalid)), + serde_json::to_vec(&Value::Bytes(invalid)).unwrap().len() + ); + } + + fn encode_logs( + events: Vec, + truncate_oversized_logs: bool, + conforms_as_agent: bool, + ) -> Vec { + encode_logs_with_limits( + events, + truncate_oversized_logs.then_some(DatadogLogsTruncationConfig { + max_log_bytes: MAX_LOG_BYTES, + }), + conforms_as_agent, + ) + } + + fn encode_logs_with_limits( + events: Vec, + truncation: Option, + conforms_as_agent: bool, + ) -> Vec { + let requests = build_requests(events, truncation, conforms_as_agent, 5_000_000); + + requests + .into_iter() + .flat_map(|request| { + serde_json::from_slice::>(&request.body) + .expect("payload should be a JSON array") + }) + .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)); + + let logs = encode_logs_with_limits( + vec![Event::Log(log)], + Some(DatadogLogsTruncationConfig { max_log_bytes: 500 }), + false, + ); + + let message = logs[0]["message"] + .as_str() + .expect("message should be a string"); + assert!(message.ends_with(TRUNCATION_MARKER)); + 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)], true, false); + + assert_eq!(logs.len(), 1); + assert!( + logs[0]["message"] + .as_str() + .expect("message should be a string") + .ends_with(TRUNCATION_MARKER) + ); + 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)], true, 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!(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 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)); + + let logs = encode_logs(vec![Event::Log(log)], true, false); + + assert_eq!(logs.len(), 1); + 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); + } + + #[test] + 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_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] + 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)); + + let logs = encode_logs(vec![Event::Log(log)], true, true); + + assert_eq!(logs.len(), 1); + let nested = logs[0]["message"] + .as_object() + .expect("agent message should be an object"); + 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); + } + + #[test] + 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_PAYLOAD_BYTES + 1)); + drop(batch); + + let requests = build_requests( + vec![Event::Log(log)], + Some(DatadogLogsTruncationConfig { + max_log_bytes: MAX_LOG_BYTES, + }), + false, + 5_000_000, + ); + + assert!(requests.is_empty()); + assert_eq!(receiver.try_recv(), Ok(BatchStatus::Delivered)); + } + + #[test] + 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); + + let requests = build_requests(vec![Event::Log(log)], None, false, 2); + + assert!(requests.is_empty()); + assert_eq!(receiver.try_recv(), Ok(BatchStatus::Delivered)); + } + + #[test] + fn does_not_count_best_effort_log_as_dropped() { + 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)], true, false); + + assert_eq!(logs.len(), 1); + 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, 0.0); + } + + #[test] + fn forwards_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_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 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_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] { + let log = LogEvent::from("\"".repeat(600_000)); + + let logs = encode_logs(vec![Event::Log(log)], true, 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 truncates_escaped_message_with_plain_suffix_to_fit() { + let message = format!("{}{}", "\"".repeat(450_000), "a".repeat(110_000)); + let log = LogEvent::from(message); + + let logs = encode_logs(vec![Event::Log(log)], true, 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)], true, 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)], true, 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)], false, 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)], 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())); + assert!(logs[0].get("ddtags").is_none()); + } + + #[test] + fn truncation_preserves_utf8_boundaries() { + 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 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) { assert!( diff --git a/src/sinks/datadog/logs/tests.rs b/src/sinks/datadog/logs/tests.rs index d394a71e0b3b6..f08da8b9bedeb 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,34 @@ 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 + "#}; + 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..315da1b63729e 100644 --- a/website/cue/reference/components/sinks/datadog_logs.cue +++ b/website/cue/reference/components/sinks/datadog_logs.cue @@ -82,7 +82,13 @@ 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). 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 27aaa9f5e1d94..67605e273bd9d 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,20 @@ generated: components: sinks: datadog_logs: configuration: { required: false type: _schemaDefinitions["core::option::Option"] } + truncate_oversized_logs: { + description: """ + 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 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: { + description: "Maximum encoded size, in bytes, of a log before truncation is applied." + required: false + type: uint: default: 1000000 + } + } }