saluki_components/encoders/datadog/events/
mod.rs1use 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#[derive(Debug)]
39pub struct DatadogEventsConfiguration {
40 pub max_payload_size: usize,
48
49 pub max_uncompressed_payload_size: usize,
57
58 pub compressor_kind: String,
62
63 pub zstd_level: i32,
67
68 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 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 builder.minimum().with_single_value::<DatadogEvents>("component struct");
123
124 builder
125 .firm()
126 .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 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 output_stream.write_tag(EVENTS_FIELD_NUMBER, WireType::LengthDelimited)?;
240
241 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 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 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}