saluki_components/encoders/datadog/events/
mod.rs

1use async_trait::async_trait;
2use datadog_protos::events as proto;
3use http::{uri::PathAndQuery, HeaderValue, Method, Uri};
4use protobuf::{rt::WireType, CodedOutputStream};
5use saluki_common::iter::ReusableDeduplicator;
6use saluki_context::tags::Tag;
7use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder};
8use saluki_core::{
9    components::{encoders::*, ComponentContext},
10    data_model::{
11        event::{eventd::EventD, Event, EventType},
12        payload::{HttpPayload, Payload, PayloadMetadata, PayloadType},
13    },
14    observability::ComponentMetricsExt as _,
15    topology::PayloadsDispatcher,
16};
17use saluki_error::{ErrorContext as _, GenericError};
18use saluki_io::compression::CompressionScheme;
19use saluki_metrics::MetricsBuilder;
20use tracing::{debug, error, warn};
21
22use crate::common::datadog::{
23    clamp_payload_limits,
24    io::RB_BUFFER_CHUNK_SIZE,
25    request_builder::{EndpointEncoder, RequestBuilder},
26    telemetry::ComponentTelemetry,
27    DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT, DEFAULT_INTAKE_UNCOMPRESSED_SIZE_LIMIT,
28};
29
30const MAX_EVENTS_PER_PAYLOAD: usize = 100;
31const EVENTS_FIELD_NUMBER: u32 = 1;
32
33static CONTENT_TYPE_PROTOBUF: HeaderValue = HeaderValue::from_static("application/x-protobuf");
34
35/// Datadog Events incremental encoder.
36///
37/// Generates Datadog Events payloads for the Datadog platform.
38#[derive(Debug)]
39pub struct DatadogEventsConfiguration {
40    /// Maximum compressed size, in bytes, of an events payload.
41    ///
42    /// This uses the same generic event payload setting as the Datadog Agent. ADP sends events to
43    /// `/api/v1/events_batch`, so the effective value is clamped to that endpoint's global intake limit of 3,200,000
44    /// bytes. If set to `0`, every non-empty compressed payload exceeds the limit and is dropped during flush.
45    ///
46    /// Defaults to 2,621,440 bytes.
47    pub max_payload_size: usize,
48
49    /// Maximum uncompressed size, in bytes, of an events payload.
50    ///
51    /// This uses the same generic event payload setting as the Datadog Agent. ADP sends events to
52    /// `/api/v1/events_batch`, so the effective value is clamped to that endpoint's global intake limit of 62,914,560
53    /// bytes. Values smaller than the minimum endpoint framing size prevent the request builder from starting.
54    ///
55    /// Defaults to 4,194,304 bytes.
56    pub max_uncompressed_payload_size: usize,
57
58    /// Compression kind to use for the request payloads.
59    ///
60    /// Defaults to `zstd`.
61    pub compressor_kind: String,
62
63    /// The compression level to use when `zstd` is the algorithm.
64    ///
65    /// Ignored for algorithms other than `zstd`.
66    pub zstd_level: i32,
67
68    /// Whether to log event payload contents before encoding.
69    ///
70    /// This logs decoded event objects, not the encoded HTTP body.
71    ///
72    /// Defaults to `false`.
73    pub log_payloads: bool,
74}
75
76#[async_trait]
77impl IncrementalEncoderBuilder for DatadogEventsConfiguration {
78    type Output = DatadogEvents;
79
80    fn input_event_type(&self) -> EventType {
81        EventType::EventD
82    }
83
84    fn output_payload_type(&self) -> PayloadType {
85        PayloadType::Http
86    }
87
88    async fn build(&self, context: ComponentContext) -> Result<Self::Output, GenericError> {
89        let metrics_builder = MetricsBuilder::from_component_context(&context);
90        let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
91        let compression_scheme = CompressionScheme::new(&self.compressor_kind, self.zstd_level);
92
93        // Create our request builder.
94        let mut request_builder =
95            RequestBuilder::new(EventsEndpointEncoder::new(), compression_scheme, RB_BUFFER_CHUNK_SIZE).await?;
96        let (uncompressed_limit, compressed_limit) = clamp_payload_limits(
97            self.max_uncompressed_payload_size,
98            self.max_payload_size,
99            DEFAULT_INTAKE_UNCOMPRESSED_SIZE_LIMIT,
100            DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT,
101        );
102        request_builder.with_len_limits(uncompressed_limit, compressed_limit)?;
103        request_builder.with_max_inputs_per_payload(MAX_EVENTS_PER_PAYLOAD);
104
105        Ok(DatadogEvents {
106            request_builder,
107            telemetry,
108            log_payloads: self.log_payloads,
109        })
110    }
111}
112
113impl MemoryBounds for DatadogEventsConfiguration {
114    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
115        // TODO: How do we properly represent the requests we can generate that may be sitting around in-flight?
116        //
117        // Theoretically, we'll end up being limited by the size of the downstream forwarder's interconnect, and however
118        // many payloads it will buffer internally... so realistically the firm limit boils down to the forwarder itself
119        // but we'll have a hard time in the forwarder knowing the maximum size of any given payload being sent in, which
120        // then makes it hard to calculate a proper firm bound even though we know the rest of the values required to
121        // calculate the firm bound.
122        builder.minimum().with_single_value::<DatadogEvents>("component struct");
123
124        builder
125            .firm()
126            // Capture the size of the "split re-encode" buffer in the request builder, which is where we keep owned
127            // versions of events that we encode in case we need to actually re-encode them during a split operation.
128            .with_array::<EventD>("events split re-encode buffer", MAX_EVENTS_PER_PAYLOAD);
129    }
130}
131
132pub struct DatadogEvents {
133    request_builder: RequestBuilder<EventsEndpointEncoder>,
134    telemetry: ComponentTelemetry,
135    log_payloads: bool,
136}
137
138#[async_trait]
139impl IncrementalEncoder for DatadogEvents {
140    async fn process_event(&mut self, event: Event) -> Result<ProcessResult, GenericError> {
141        let eventd = match event.try_into_eventd() {
142            Some(eventd) => eventd,
143            None => return Ok(ProcessResult::Continue),
144        };
145
146        if self.log_payloads {
147            debug!(event = ?eventd, "Flushing event.");
148        }
149
150        match self.request_builder.encode(eventd).await {
151            Ok(None) => Ok(ProcessResult::Continue),
152            Ok(Some(eventd)) => Ok(ProcessResult::FlushRequired(Event::EventD(eventd))),
153            Err(e) => {
154                if e.is_recoverable() {
155                    warn!(error = %e, "Failed to encode Datadog event due to recoverable error. Continuing...");
156
157                    // TODO: Get the actual number of events dropped from the error itself.
158                    self.telemetry.events_dropped_encoder().increment(1);
159
160                    Ok(ProcessResult::Continue)
161                } else {
162                    Err(e).error_context("Failed to encode Datadog event due to unrecoverable error.")
163                }
164            }
165        }
166    }
167
168    async fn flush(&mut self, dispatcher: &PayloadsDispatcher) -> Result<(), GenericError> {
169        let maybe_requests = self.request_builder.flush().await;
170        for maybe_request in maybe_requests {
171            match maybe_request {
172                Ok((events, _data_points, request)) => {
173                    let payload_meta = PayloadMetadata::from_event_count(events);
174                    let http_payload = HttpPayload::new(payload_meta, request);
175                    let payload = Payload::Http(http_payload);
176
177                    dispatcher.dispatch(payload).await?;
178                }
179                Err(e) => error!(error = %e, "Failed to build Datadog events payload. Continuing..."),
180            }
181        }
182
183        Ok(())
184    }
185}
186
187#[derive(Debug)]
188struct EventsEndpointEncoder {
189    tags_deduplicator: ReusableDeduplicator<Tag>,
190}
191
192impl EventsEndpointEncoder {
193    fn new() -> Self {
194        Self {
195            tags_deduplicator: ReusableDeduplicator::new(),
196        }
197    }
198}
199
200impl EndpointEncoder for EventsEndpointEncoder {
201    type Input = EventD;
202    type EncodeError = protobuf::Error;
203
204    fn encoder_name() -> &'static str {
205        "events"
206    }
207
208    fn compressed_size_limit(&self) -> usize {
209        DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT
210    }
211
212    fn uncompressed_size_limit(&self) -> usize {
213        DEFAULT_INTAKE_UNCOMPRESSED_SIZE_LIMIT
214    }
215
216    fn encode(&mut self, input: &Self::Input, buffer: &mut Vec<u8>) -> Result<(), Self::EncodeError> {
217        encode_and_write_eventd(input, buffer, &mut self.tags_deduplicator)
218    }
219
220    fn endpoint_uri(&self) -> Uri {
221        PathAndQuery::from_static("/api/v1/events_batch").into()
222    }
223
224    fn endpoint_method(&self) -> Method {
225        Method::POST
226    }
227
228    fn content_type(&self) -> HeaderValue {
229        CONTENT_TYPE_PROTOBUF.clone()
230    }
231}
232
233fn encode_and_write_eventd(
234    eventd: &EventD, buf: &mut Vec<u8>, tags_deduplicator: &mut ReusableDeduplicator<Tag>,
235) -> Result<(), protobuf::Error> {
236    let mut output_stream = CodedOutputStream::vec(buf);
237
238    // Write the field tag.
239    output_stream.write_tag(EVENTS_FIELD_NUMBER, WireType::LengthDelimited)?;
240
241    // Write the message.
242    let encoded_eventd = encode_eventd(eventd, tags_deduplicator);
243    output_stream.write_message_no_tag(&encoded_eventd)
244}
245
246fn encode_eventd(eventd: &EventD, tags_deduplicator: &mut ReusableDeduplicator<Tag>) -> proto::Event {
247    let mut event = proto::Event::new();
248    event.set_title(eventd.title().into());
249    event.set_text(eventd.text().into());
250
251    if let Some(timestamp) = eventd.timestamp() {
252        event.set_ts(timestamp as i64);
253    }
254
255    if let Some(priority) = eventd.priority() {
256        event.set_priority(priority.as_str().into());
257    }
258
259    if let Some(alert_type) = eventd.alert_type() {
260        event.set_alert_type(alert_type.as_str().into());
261    }
262
263    if let Some(hostname) = eventd.hostname() {
264        event.set_host(hostname.into());
265    }
266
267    if let Some(aggregation_key) = eventd.aggregation_key() {
268        event.set_aggregation_key(aggregation_key.into());
269    }
270
271    if let Some(source_type_name) = eventd.source_type_name() {
272        event.set_source_type_name(source_type_name.into());
273    }
274
275    let chained_tags = eventd.tags().into_iter().chain(eventd.origin_tags());
276    let deduplicated_tags = tags_deduplicator.deduplicated(chained_tags);
277
278    event.set_tags(deduplicated_tags.map(|tag| tag.as_str().into()).collect());
279
280    event
281}
282
283#[cfg(test)]
284mod tests {
285    use std::collections::BTreeSet;
286
287    use saluki_common::iter::ReusableDeduplicator;
288    use saluki_context::tags::{Tag, TagSet};
289    use saluki_core::data_model::event::eventd::{AlertType, EventD, Priority};
290    use stringtheory::MetaString;
291
292    use super::encode_eventd;
293
294    fn tag_set<const N: usize>(tags: [&'static str; N]) -> TagSet {
295        tags.into_iter().map(Tag::from_static).collect()
296    }
297
298    #[test]
299    fn encode_eventd_maps_all_documented_fields() {
300        let eventd = EventD::new("deploy", "release rolled out")
301            .with_timestamp(1_700_000_000u64)
302            .with_priority(Priority::Low)
303            .with_alert_type(AlertType::Error)
304            .with_hostname(MetaString::from_static("host-a"))
305            .with_aggregation_key(MetaString::from_static("deploy-key"))
306            .with_source_type_name(MetaString::from_static("my-source"));
307
308        let mut tags_deduplicator = ReusableDeduplicator::new();
309        let encoded = encode_eventd(&eventd, &mut tags_deduplicator);
310
311        assert_eq!("deploy", encoded.title());
312        assert_eq!("release rolled out", encoded.text());
313        assert_eq!(1_700_000_000, encoded.ts());
314        assert_eq!("low", encoded.priority());
315        assert_eq!("error", encoded.alert_type());
316        assert_eq!("host-a", encoded.host());
317        assert_eq!("deploy-key", encoded.aggregation_key());
318        assert_eq!("my-source", encoded.source_type_name());
319    }
320
321    #[test]
322    fn encode_eventd_applies_defaults_and_skips_empty_string_fields() {
323        // `EventD::new` defaults the priority to `normal` and the alert type to `info`. An unset timestamp and the
324        // empty host/aggregation-key/source-type fields are treated as absent and left at their protobuf defaults.
325        let eventd = EventD::new("title-only", "body");
326        let mut tags_deduplicator = ReusableDeduplicator::new();
327        let encoded = encode_eventd(&eventd, &mut tags_deduplicator);
328
329        assert_eq!("title-only", encoded.title());
330        assert_eq!("body", encoded.text());
331        assert_eq!(0, encoded.ts());
332        assert_eq!("normal", encoded.priority());
333        assert_eq!("info", encoded.alert_type());
334        assert_eq!("", encoded.host());
335        assert_eq!("", encoded.aggregation_key());
336        assert_eq!("", encoded.source_type_name());
337        assert!(encoded.tags().is_empty());
338    }
339
340    #[test]
341    fn encode_eventd_deduplicates_tags_across_origin_tags() {
342        // Event tags and origin tags are chained then deduplicated, so an overlapping tag is written only once.
343        let eventd = EventD::new("dedup", "body")
344            .with_tags(tag_set(["env:prod", "team:core"]))
345            .with_origin_tags(tag_set(["env:prod", "region:us"]));
346        let mut tags_deduplicator = ReusableDeduplicator::new();
347        let encoded = encode_eventd(&eventd, &mut tags_deduplicator);
348
349        let tags = encoded.tags().iter().map(String::as_str).collect::<BTreeSet<_>>();
350        assert_eq!(BTreeSet::from(["env:prod", "team:core", "region:us"]), tags);
351        assert_eq!(
352            3,
353            encoded.tags().len(),
354            "the overlapping `env:prod` tag should not be duplicated"
355        );
356    }
357}