1use std::time::SystemTime;
5
6use crate::{OsInfo, SigInfo};
7
8use super::{build_crash_ping_message, CrashInfo, Experimental, Metadata, StackTrace};
9use anyhow::Context;
10use chrono::{DateTime, Utc};
11use http::{uri::PathAndQuery, Uri};
12use libdd_common::{config::parse_env, parse_uri, Endpoint};
13use serde::{Deserialize, Serialize};
14use std::{borrow::Cow, time::Duration};
15
16pub const DEFAULT_DD_SITE: &str = "datadoghq.com";
17pub const PROD_ERRORS_INTAKE_SUBDOMAIN: &str = "error-tracking-intake";
18
19const DIRECT_ERRORS_INTAKE_URL_PATH: &str = "/api/v2/errorsintake";
20const AGENT_ERRORS_INTAKE_URL_PATH: &str = "/evp_proxy/v4/api/v2/errorsintake";
21
22const DEFAULT_AGENT_HOST: &str = "localhost";
23const DEFAULT_AGENT_PORT: u16 = 8126;
24
25#[derive(Clone, Debug, Default, Serialize, Deserialize)]
26pub struct ErrorsIntakeConfig {
27 pub(crate) endpoint: Option<Endpoint>,
28 pub direct_submission_enabled: bool,
29 pub debug_enabled: bool,
30 pub errors_intake_enabled: bool,
31}
32
33fn endpoint_with_errors_intake_path(
34 mut endpoint: Endpoint,
35 direct_submission_enabled: bool,
36) -> anyhow::Result<Endpoint> {
37 let mut uri_parts = endpoint.url.into_parts();
38 if uri_parts
39 .scheme
40 .as_ref()
41 .is_some_and(|scheme| scheme.as_str() != "file")
42 {
43 uri_parts.path_and_query = Some(PathAndQuery::from_static(
44 if endpoint.api_key.is_some() && direct_submission_enabled {
45 DIRECT_ERRORS_INTAKE_URL_PATH
46 } else {
47 AGENT_ERRORS_INTAKE_URL_PATH
48 },
49 ));
50 }
51
52 endpoint.url = Uri::from_parts(uri_parts)?;
53 Ok(endpoint)
54}
55
56#[derive(Debug, Default)]
58pub struct ErrorsIntakeSettings {
59 pub agent_host: Option<String>,
61 pub trace_agent_port: Option<u16>,
62 pub trace_agent_url: Option<String>,
63 pub trace_pipe_name: Option<String>,
64 pub direct_submission_enabled: bool,
65 pub api_key: Option<String>,
66 pub site: Option<String>,
67 pub errors_intake_dd_url: Option<String>,
68 pub shared_lib_debug: bool,
69 pub errors_intake_enabled: bool,
70
71 pub agent_uds_socket_found: bool,
73}
74
75impl ErrorsIntakeSettings {
76 const DD_TRACE_AGENT_URL: &'static str = "DD_TRACE_AGENT_URL";
78 const DD_AGENT_HOST: &'static str = "DD_AGENT_HOST";
79 const DD_TRACE_AGENT_PORT: &'static str = "DD_TRACE_AGENT_PORT";
80 const DD_TRACE_PIPE_NAME: &'static str = "DD_TRACE_PIPE_NAME";
81
82 const _DD_DIRECT_SUBMISSION_ENABLED: &'static str = "_DD_DIRECT_SUBMISSION_ENABLED";
84 const DD_API_KEY: &'static str = "DD_API_KEY";
85 const DD_SITE: &'static str = "DD_SITE";
86 const DD_ERRORS_INTAKE_DD_URL: &'static str = "DD_ERRORS_INTAKE_DD_URL";
87
88 const _DD_SHARED_LIB_DEBUG: &'static str = "_DD_SHARED_LIB_DEBUG";
90
91 const DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED: &'static str =
93 "DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED";
94
95 pub fn from_env() -> Self {
96 let default = Self::default();
97 Self {
98 agent_host: parse_env::str_not_empty(Self::DD_AGENT_HOST),
99 trace_agent_port: parse_env::int(Self::DD_TRACE_AGENT_PORT),
100 trace_agent_url: parse_env::str_not_empty(Self::DD_TRACE_AGENT_URL)
101 .or(default.trace_agent_url),
102 trace_pipe_name: parse_env::str_not_empty(Self::DD_TRACE_PIPE_NAME)
103 .or(default.trace_pipe_name),
104 direct_submission_enabled: parse_env::bool(Self::_DD_DIRECT_SUBMISSION_ENABLED)
105 .unwrap_or(default.direct_submission_enabled),
106 api_key: parse_env::str_not_empty(Self::DD_API_KEY),
107 site: parse_env::str_not_empty(Self::DD_SITE),
108 errors_intake_dd_url: parse_env::str_not_empty(Self::DD_ERRORS_INTAKE_DD_URL),
109 shared_lib_debug: parse_env::bool(Self::_DD_SHARED_LIB_DEBUG).unwrap_or(false),
110 errors_intake_enabled: parse_env::bool(Self::DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED)
111 .unwrap_or(false),
112
113 agent_uds_socket_found: (|| {
114 #[cfg(unix)]
115 return std::fs::metadata("/var/run/datadog/apm.socket").is_ok();
116 #[cfg(not(unix))]
117 return false;
118 })(),
119 }
120 }
121}
122
123impl ErrorsIntakeConfig {
124 fn trace_agent_url_from_setting(settings: &ErrorsIntakeSettings) -> String {
125 None.or_else(|| {
126 settings
127 .trace_agent_url
128 .as_deref()
129 .filter(|u| {
130 u.starts_with("unix://")
131 || u.starts_with("http://")
132 || u.starts_with("https://")
133 })
134 .map(ToString::to_string)
135 })
136 .or_else(|| {
137 #[cfg(windows)]
138 return settings
139 .trace_pipe_name
140 .as_ref()
141 .map(|pipe_name| format!("windows:{pipe_name}"));
142 #[cfg(not(windows))]
143 return None;
144 })
145 .or_else(|| {
146 #[cfg(unix)]
147 return settings
148 .agent_uds_socket_found
149 .then(|| "unix:///var/run/datadog/apm.socket".to_string());
150 #[cfg(not(unix))]
151 return None;
152 })
153 .or_else(|| match (&settings.agent_host, settings.trace_agent_port) {
154 (None, None) => None,
155 _ => Some(format!(
156 "http://{}:{}",
157 settings.agent_host.as_deref().unwrap_or(DEFAULT_AGENT_HOST),
158 settings.trace_agent_port.unwrap_or(DEFAULT_AGENT_PORT),
159 )),
160 })
161 .unwrap_or_else(|| format!("http://{DEFAULT_AGENT_HOST}:{DEFAULT_AGENT_PORT}"))
162 }
163
164 fn api_key_from_settings(settings: &ErrorsIntakeSettings) -> Option<Cow<'static, str>> {
165 if !settings.direct_submission_enabled {
166 return None;
167 }
168 settings.api_key.clone().map(Cow::Owned)
169 }
170
171 pub fn endpoint(&self) -> Option<&Endpoint> {
172 self.endpoint.as_ref()
173 }
174
175 pub fn is_errors_intake_enabled(&self) -> bool {
176 self.errors_intake_enabled
177 }
178
179 pub fn set_endpoint(&mut self, endpoint: Endpoint) -> anyhow::Result<()> {
180 self.endpoint = Some(endpoint_with_errors_intake_path(
181 endpoint,
182 self.direct_submission_enabled,
183 )?);
184 Ok(())
185 }
186
187 pub fn from_settings(settings: &ErrorsIntakeSettings) -> Self {
188 let api_key = Self::api_key_from_settings(settings);
189
190 let mut this = Self {
191 endpoint: None,
192 direct_submission_enabled: settings.direct_submission_enabled,
193 debug_enabled: settings.shared_lib_debug,
194 errors_intake_enabled: settings.errors_intake_enabled,
195 };
196
197 let url = if settings.direct_submission_enabled && settings.api_key.is_some() {
199 if let Some(ref errors_intake_url) = settings.errors_intake_dd_url {
201 errors_intake_url.clone()
202 } else {
203 let site = settings.site.as_deref().unwrap_or(DEFAULT_DD_SITE);
205 format!("https://{}.{}", PROD_ERRORS_INTAKE_SUBDOMAIN, site)
206 }
207 } else {
208 Self::trace_agent_url_from_setting(settings)
209 };
210
211 if let Ok(parsed_url) = parse_uri(&url) {
212 let _res = this.set_endpoint(Endpoint {
213 url: parsed_url,
214 api_key,
215 ..Default::default()
216 });
217 }
218
219 this
220 }
221
222 pub fn from_env() -> Self {
223 let settings = ErrorsIntakeSettings::from_env();
224 Self::from_settings(&settings)
225 }
226
227 pub fn set_host_from_url(&mut self, host_url: &str) -> anyhow::Result<()> {
228 let endpoint = self.endpoint.take().unwrap_or_default();
229
230 self.set_endpoint(Endpoint {
231 url: parse_uri(host_url)?,
232 ..endpoint
233 })
234 }
235}
236
237#[derive(serde::Serialize, Debug)]
238pub struct ErrorObject {
239 #[serde(rename = "type", skip_serializing_if = "Option::is_none")]
240 pub error_type: Option<String>,
241 #[serde(skip_serializing_if = "Option::is_none")]
242 pub message: Option<String>,
243 #[serde(skip_serializing_if = "Option::is_none")]
244 pub stack: Option<StackTrace>,
245 #[serde(skip_serializing_if = "Option::is_none")]
246 pub is_crash: Option<bool>,
247 #[serde(skip_serializing_if = "Option::is_none")]
248 pub fingerprint: Option<String>,
249 #[serde(skip_serializing_if = "Option::is_none")]
250 pub source_type: Option<String>,
251 #[serde(skip_serializing_if = "Option::is_none")]
252 pub experimental: Option<Experimental>,
253}
254
255#[derive(serde::Serialize, Debug)]
256pub struct ErrorsIntakePayload {
257 pub timestamp: u64,
258 pub ddsource: String,
259 pub ddtags: String,
260 pub error: ErrorObject,
261 #[serde(skip_serializing_if = "Option::is_none")]
262 pub trace_id: Option<String>,
263 pub os_info: OsInfo,
264 #[serde(skip_serializing_if = "Option::is_none")]
265 pub sig_info: Option<SigInfo>,
266}
267
268#[derive(Debug, Default)]
269struct ExtractedMetadata {
270 service_name: String,
271 env: Option<String>,
272 service_version: Option<String>,
273 language_name: Option<String>,
274 language_version: Option<String>,
275 tracer_version: Option<String>,
276}
277
278impl ExtractedMetadata {
279 fn from_metadata(metadata: &Metadata) -> Self {
280 let mut result = Self {
281 service_name: "unknown".to_string(),
282 ..Default::default()
283 };
284
285 for tag in &metadata.tags {
286 if let Some((key, value)) = tag.split_once(':') {
287 match key {
288 "service" => result.service_name = value.to_string(),
289 "env" => result.env = Some(value.to_string()),
290 "version" | "service_version" => {
291 result.service_version = Some(value.to_string())
292 }
293 "language" => result.language_name = Some(value.to_string()),
294 "language_version" | "runtime_version" => {
295 result.language_version = Some(value.to_string())
296 }
297 "library_version" | "profiler_version" => {
298 result.tracer_version = Some(value.to_string())
299 }
300 _ => {}
301 }
302 }
303 }
304
305 result
306 }
307
308 fn append_base_tags(&self, tags: &mut String) {
309 tags.push_str(&format!("service:{}", self.service_name));
310 if let Some(env) = &self.env {
311 tags.push_str(&format!(",env:{env}"));
312 }
313 if let Some(version) = &self.service_version {
314 tags.push_str(&format!(",version:{version}"));
315 }
316 }
317
318 fn append_runtime_tags(&self, tags: &mut String) {
319 if let Some(language_name) = &self.language_name {
320 tags.push_str(&format!(",language_name:{language_name}"));
321 }
322 if let Some(language_version) = &self.language_version {
323 tags.push_str(&format!(",language_version:{language_version}"));
324 }
325 if let Some(tracer_version) = &self.tracer_version {
326 tags.push_str(&format!(",tracer_version:{tracer_version}"));
327 }
328 }
329}
330
331fn append_signal_tags(tags: &mut String, sig_info: &SigInfo) {
332 tags.push_str(&format!(
333 ",si_code_human_readable:{:?}",
334 sig_info.si_code_human_readable
335 ));
336 tags.push_str(&format!(",si_signo:{}", sig_info.si_signo));
337 tags.push_str(&format!(
338 ",si_signo_human_readable:{:?}",
339 sig_info.si_signo_human_readable
340 ));
341}
342
343fn build_crash_info_tags(crash_info: &CrashInfo) -> String {
344 let mut tags = format!("data_schema_version:{}", crash_info.data_schema_version);
345
346 if let Some(fingerprint) = &crash_info.fingerprint {
347 tags.push_str(&format!(",fingerprint:{fingerprint}"));
348 }
349
350 tags.push_str(&format!(",incomplete:{}", crash_info.incomplete));
351 tags.push_str(&format!(",is_crash:{}", crash_info.error.is_crash));
352 tags.push_str(&format!(",uuid:{}", crash_info.uuid));
353
354 for (counter, value) in &crash_info.counters {
355 tags.push_str(&format!(",{counter}:{value}"));
356 }
357
358 if let Some(siginfo) = &crash_info.sig_info {
359 if let Some(si_addr) = &siginfo.si_addr {
360 tags.push_str(&format!(",si_addr:{si_addr}"));
361 }
362 tags.push_str(&format!(",si_code:{}", siginfo.si_code));
363 append_signal_tags(&mut tags, siginfo);
364 }
365
366 tags
367}
368
369impl ErrorsIntakePayload {
370 pub fn from_crash_info(crash_info: &CrashInfo) -> anyhow::Result<Self> {
371 let timestamp = crash_info.timestamp.parse::<DateTime<Utc>>().map_or_else(
372 |_| {
373 SystemTime::now()
374 .duration_since(SystemTime::UNIX_EPOCH)
375 .map(|d| d.as_millis() as u64)
376 .unwrap_or(0)
377 },
378 |ts| ts.timestamp_millis() as u64,
379 );
380
381 let metadata = ExtractedMetadata::from_metadata(&crash_info.metadata);
382 let mut ddtags = String::new();
383 metadata.append_base_tags(&mut ddtags);
384 metadata.append_runtime_tags(&mut ddtags);
385
386 let crash_tags = build_crash_info_tags(crash_info);
387 ddtags.push_str(&format!(",{crash_tags}"));
388
389 let (error_type, error_message) = if let Some(sig_info) = &crash_info.sig_info {
390 (
391 Some(format!("{:?}", sig_info.si_signo_human_readable)),
392 Some(format!(
393 "Process terminated by signal {:?}",
394 sig_info.si_signo_human_readable
395 )),
396 )
397 } else {
398 (
399 Some("Unknown".to_string()),
400 crash_info.error.message.clone(),
401 )
402 };
403
404 let error_stack = if !crash_info.error.stack.frames.is_empty() {
406 Some(crash_info.error.stack.clone())
407 } else {
408 None
409 };
410
411 let sig_info = crash_info.sig_info.clone();
412
413 Ok(Self {
414 timestamp,
415 ddsource: "crashtracker".to_string(),
416 ddtags,
417 error: ErrorObject {
418 error_type,
419 message: error_message,
420 stack: error_stack,
421 is_crash: Some(true),
422 fingerprint: crash_info.fingerprint.clone(),
423 source_type: Some("Crashtracking".to_string()),
424 experimental: crash_info.experimental.clone(),
425 },
426 trace_id: None,
427 os_info: ::os_info::get().into(),
428 sig_info,
429 })
430 }
431
432 pub fn from_crash_ping(
433 crash_uuid: &str,
434 sig_info: Option<&SigInfo>,
435 metadata: &Metadata,
436 ) -> anyhow::Result<Self> {
437 let timestamp = SystemTime::now()
438 .duration_since(SystemTime::UNIX_EPOCH)
439 .map(|d| d.as_millis() as u64)
440 .unwrap_or(0);
441
442 let extracted_metadata = ExtractedMetadata::from_metadata(metadata);
443 let mut ddtags = format!(
444 "uuid:{},is_crash_ping:true,service:{}",
445 crash_uuid, extracted_metadata.service_name
446 );
447 extracted_metadata.append_runtime_tags(&mut ddtags);
448 if let Some(env) = &extracted_metadata.env {
449 ddtags.push_str(&format!(",env:{env}"));
450 }
451 if let Some(version) = &extracted_metadata.service_version {
452 ddtags.push_str(&format!(",version:{version}"));
453 }
454
455 if let Some(sig_info) = sig_info {
456 append_signal_tags(&mut ddtags, sig_info);
457 }
458
459 let (error_type, message) = if let Some(sig_info) = sig_info {
460 (
461 Some(format!("{:?}", sig_info.si_signo_human_readable)),
462 Some(build_crash_ping_message(sig_info)),
463 )
464 } else {
465 (
466 Some("Unknown".to_string()),
467 Some(
468 "Crashtracker crash ping: crash processing started - Process terminated"
469 .to_string(),
470 ),
471 )
472 };
473
474 Ok(Self {
475 timestamp,
476 ddsource: "crashtracker".to_string(),
477 ddtags,
478 error: ErrorObject {
479 error_type,
480 message,
481 stack: None,
482 is_crash: Some(false),
483 fingerprint: None,
484 source_type: Some("Crashtracking".to_string()),
485 experimental: None,
486 },
487 sig_info: sig_info.cloned(),
488 trace_id: None,
489 os_info: ::os_info::get().into(),
492 })
493 }
494}
495
496pub struct ErrorsIntakeUploader {
497 cfg: ErrorsIntakeConfig,
498}
499
500impl ErrorsIntakeUploader {
501 pub fn new(endpoint: &Option<Endpoint>) -> anyhow::Result<Self> {
502 let mut cfg = ErrorsIntakeConfig::from_env();
503
504 if let Some(endpoint) = endpoint {
505 cfg.set_endpoint(endpoint.clone())?;
506 }
507 Ok(Self { cfg })
508 }
509
510 pub fn is_enabled(&self) -> bool {
511 self.cfg.is_errors_intake_enabled()
512 }
513
514 pub async fn upload_crash_ping(
515 &self,
516 crash_uuid: &str,
517 sig_info: Option<&SigInfo>,
518 metadata: &Metadata,
519 ) -> anyhow::Result<()> {
520 let payload = ErrorsIntakePayload::from_crash_ping(crash_uuid, sig_info, metadata)?;
521 self.send_payload(&payload).await
522 }
523
524 pub async fn upload_to_errors_intake(&self, crash_info: &CrashInfo) -> anyhow::Result<()> {
525 let payload = ErrorsIntakePayload::from_crash_info(crash_info)?;
526 self.send_payload(&payload).await
527 }
528
529 async fn send_payload(&self, payload: &ErrorsIntakePayload) -> anyhow::Result<()> {
530 let Some(endpoint) = self.cfg.endpoint() else {
531 return Ok(());
532 };
533
534 if endpoint.url.scheme_str() == Some("file") {
536 let path = libdd_common::decode_uri_path_in_authority(&endpoint.url)
537 .context("errors intake file path is not valid")?;
538
539 let file_path = path.with_extension("errors");
540 let file = std::fs::File::create(&file_path).with_context(|| {
541 format!(
542 "Failed to create errors intake file {}",
543 file_path.display()
544 )
545 })?;
546
547 serde_json::to_writer_pretty(file, payload).with_context(|| {
548 format!(
549 "Failed to write errors intake JSON to {}",
550 file_path.display()
551 )
552 })?;
553
554 return Ok(());
555 }
556
557 let mut req_builder =
559 endpoint.to_request_builder(concat!("crashtracker/", env!("CARGO_PKG_VERSION")))?;
560
561 if endpoint.api_key.is_some() {
563 } else {
565 req_builder =
567 req_builder.header("X-Datadog-EVP-Subdomain", PROD_ERRORS_INTAKE_SUBDOMAIN);
568 }
569
570 let req = req_builder
571 .method(http::Method::POST)
572 .header(
573 http::header::CONTENT_TYPE,
574 libdd_common::header::APPLICATION_JSON,
575 )
576 .body(serde_json::to_string(payload)?.into())?;
577
578 let client = libdd_common::hyper_migration::new_client_periodic();
580
581 tokio::time::timeout(
582 Duration::from_millis(endpoint.timeout_ms),
583 client.request(req),
584 )
585 .await??;
586
587 Ok(())
588 }
589}
590
591#[cfg(test)]
592mod tests {
593 use super::*;
594 use crate::crash_info::test_utils::TestInstance;
595 use std::sync::Mutex;
596
597 static ENV_TEST_LOCK: Mutex<()> = Mutex::new(());
599
600 fn clear_errors_intake_env() {
601 std::env::remove_var("DD_TRACE_AGENT_URL");
602 std::env::remove_var("DD_AGENT_HOST");
603 std::env::remove_var("DD_TRACE_AGENT_PORT");
604 std::env::remove_var("DD_TRACE_PIPE_NAME");
605 std::env::remove_var("_DD_DIRECT_SUBMISSION_ENABLED");
606 std::env::remove_var("DD_API_KEY");
607 std::env::remove_var("DD_SITE");
608 std::env::remove_var("DD_ERRORS_INTAKE_DD_URL");
609 std::env::remove_var("_DD_SHARED_LIB_DEBUG");
610 std::env::remove_var("DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED");
611 }
612
613 #[cfg_attr(miri, ignore)]
614 #[test]
615 fn test_errors_payload_from_crash_info() {
616 let crash_info = CrashInfo::test_instance(1);
617 let payload = ErrorsIntakePayload::from_crash_info(&crash_info).unwrap();
618
619 assert_eq!(payload.ddsource, "crashtracker");
620 assert_eq!(payload.error.source_type, Some("Crashtracking".to_string()));
621 assert_eq!(payload.error.is_crash, Some(true));
622
623 let ddtags = &payload.ddtags;
624
625 assert!(ddtags.contains("service:foo"));
626 assert!(ddtags.contains("version:bar"));
627 assert!(ddtags.contains("language_name:native"));
628
629 assert!(ddtags.contains("data_schema_version:1.4"));
630 assert!(ddtags.contains("incomplete:true"));
631 assert!(ddtags.contains("is_crash:true"));
632 assert!(ddtags.contains("uuid:1d6b97cb-968c-40c9-af6e-e4b4d71e8781"));
633
634 assert!(ddtags.contains("collecting_sample:1"));
635 assert!(ddtags.contains("not_profiling:0"));
636
637 assert!(ddtags.contains("si_addr:0x0000000000001234"));
638 assert!(ddtags.contains("si_code:1"));
639 assert!(ddtags.contains("si_code_human_readable:SEGV_BNDERR"));
640 assert!(ddtags.contains("si_signo:11"));
641 assert!(ddtags.contains("si_signo_human_readable:SIGSEGV"));
642 }
643
644 #[cfg_attr(miri, ignore)]
645 #[test]
646 fn test_errors_payload_from_crash_ping() {
647 let metadata = Metadata::test_instance(1);
648 let sig_info = crate::SigInfo::test_instance(42);
649 let crash_uuid = "test-uuid-123";
650
651 let payload =
652 ErrorsIntakePayload::from_crash_ping(crash_uuid, Some(&sig_info), &metadata).unwrap();
653
654 assert_eq!(payload.ddsource, "crashtracker");
655 assert_eq!(payload.error.source_type, Some("Crashtracking".to_string()));
656 assert_eq!(payload.error.is_crash, Some(false));
657 assert!(payload.error.stack.is_none());
658
659 let ddtags = &payload.ddtags;
660
661 assert!(ddtags.contains("uuid:test-uuid-123"));
662 assert!(ddtags.contains("is_crash_ping:true"));
663 assert!(ddtags.contains("service:foo"));
664
665 assert!(ddtags.contains("language_name:native"));
666
667 assert!(ddtags.contains("version:bar"));
668
669 assert!(ddtags.contains("si_code_human_readable:SEGV_BNDERR"));
670 assert!(ddtags.contains("si_signo:11"));
671 assert!(ddtags.contains("si_signo_human_readable:SIGSEGV"));
672 }
673
674 #[cfg_attr(miri, ignore)]
675 #[test]
676 fn test_errors_intake_has_all_telemetry_tags() {
677 let crash_info = CrashInfo::test_instance(1);
678 let payload = ErrorsIntakePayload::from_crash_info(&crash_info).unwrap();
679
680 let expected_crash_tags = [
681 "data_schema_version:1.4",
682 "incomplete:true",
683 "is_crash:true",
684 "uuid:1d6b97cb-968c-40c9-af6e-e4b4d71e8781",
685 "collecting_sample:1",
686 "not_profiling:0",
687 "si_addr:0x0000000000001234",
688 "si_code:1",
689 "si_code_human_readable:SEGV_BNDERR",
690 "si_signo:11",
691 "si_signo_human_readable:SIGSEGV",
692 ];
693
694 let expected_metadata_tags = ["service:foo", "version:bar", "language_name:native"];
695
696 for tag in expected_crash_tags
697 .iter()
698 .chain(expected_metadata_tags.iter())
699 {
700 assert!(
701 payload.ddtags.contains(tag),
702 "Missing expected tag: {} in ddtags: {}",
703 tag,
704 payload.ddtags
705 );
706 }
707 }
708
709 #[cfg_attr(miri, ignore)]
710 #[test]
711 fn test_crash_ping_has_all_telemetry_tags() {
712 let metadata = Metadata::test_instance(1);
713 let sig_info = crate::SigInfo::test_instance(42);
714 let crash_uuid = "test-crash-ping-uuid";
715
716 let payload =
717 ErrorsIntakePayload::from_crash_ping(crash_uuid, Some(&sig_info), &metadata).unwrap();
718
719 let expected_tags = [
721 "uuid:test-crash-ping-uuid",
722 "is_crash_ping:true",
723 "service:foo",
724 "language_name:native",
725 "version:bar",
726 "si_code_human_readable:SEGV_BNDERR",
727 "si_signo:11",
728 "si_signo_human_readable:SIGSEGV",
729 ];
730
731 for tag in expected_tags {
732 assert!(
733 payload.ddtags.contains(tag),
734 "Missing expected tag: {} in ddtags: {}",
735 tag,
736 payload.ddtags
737 );
738 }
739 }
740
741 #[test]
742 fn test_errors_intake_config_from_env() {
743 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
744
745 clear_errors_intake_env();
746
747 std::env::set_var("DD_API_KEY", "test-key");
749 std::env::set_var("_DD_DIRECT_SUBMISSION_ENABLED", "true");
750
751 let cfg = ErrorsIntakeConfig::from_env();
752 let endpoint = cfg.endpoint().unwrap();
753
754 assert_eq!(
756 endpoint.url.host(),
757 Some("error-tracking-intake.datadoghq.com")
758 );
759 assert_eq!(endpoint.url.scheme_str(), Some("https"));
760 assert!(endpoint.api_key.is_some());
761
762 assert_eq!(endpoint.url.path(), DIRECT_ERRORS_INTAKE_URL_PATH);
764 }
765
766 #[test]
767 fn test_errors_intake_config_custom_site() {
768 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
769
770 clear_errors_intake_env();
771
772 std::env::set_var("DD_API_KEY", "test-key");
774 std::env::set_var("_DD_DIRECT_SUBMISSION_ENABLED", "true");
775 std::env::set_var("DD_SITE", "us3.datadoghq.com");
776
777 let cfg = ErrorsIntakeConfig::from_env();
778 let endpoint = cfg.endpoint().unwrap();
779
780 assert_eq!(
782 endpoint.url.host(),
783 Some("error-tracking-intake.us3.datadoghq.com")
784 );
785 assert_eq!(endpoint.url.scheme_str(), Some("https"));
786 assert!(endpoint.api_key.is_some());
787 assert_eq!(endpoint.url.path(), DIRECT_ERRORS_INTAKE_URL_PATH);
788 }
789
790 #[test]
791 fn test_errors_intake_config_agent_proxy() {
792 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
793
794 clear_errors_intake_env();
795
796 std::env::set_var("DD_TRACE_AGENT_URL", "http://localhost:9126");
797
798 let cfg = ErrorsIntakeConfig::from_env();
799 let endpoint = cfg.endpoint().unwrap();
800
801 assert_eq!(endpoint.url.host(), Some("localhost"));
802 assert_eq!(endpoint.url.port_u16(), Some(9126));
803
804 assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
806 }
807
808 #[test]
809 fn test_errors_intake_config_agent_with_api_key_but_no_direct() {
810 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
811
812 clear_errors_intake_env();
813
814 std::env::set_var("DD_TRACE_AGENT_URL", "http://localhost:9126");
817 std::env::set_var("DD_API_KEY", "test-key");
818
819 let cfg = ErrorsIntakeConfig::from_env();
820 let endpoint = cfg.endpoint().unwrap();
821
822 assert_eq!(endpoint.url.host(), Some("localhost"));
824 assert_eq!(endpoint.url.port_u16(), Some(9126));
825
826 assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
828 assert!(endpoint.api_key.is_none());
829 }
830
831 #[test]
832 fn test_errors_intake_enabled_flag() {
833 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
834
835 clear_errors_intake_env();
837 let cfg = ErrorsIntakeConfig::from_env();
838 assert!(!cfg.is_errors_intake_enabled());
839
840 std::env::set_var("DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED", "true");
842 let cfg = ErrorsIntakeConfig::from_env();
843 assert!(cfg.is_errors_intake_enabled());
844
845 std::env::set_var("DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED", "false");
847 let cfg = ErrorsIntakeConfig::from_env();
848 assert!(!cfg.is_errors_intake_enabled());
849
850 let uploader = ErrorsIntakeUploader::new(&None).unwrap();
852 assert!(!uploader.is_enabled());
853
854 std::env::set_var("DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED", "true");
855 let uploader = ErrorsIntakeUploader::new(&None).unwrap();
856 assert!(uploader.is_enabled());
857 }
858
859 #[test]
860 #[cfg(unix)]
861 fn test_errors_intake_config_uds_socket() {
862 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
863
864 clear_errors_intake_env();
865
866 let settings = ErrorsIntakeSettings {
868 agent_uds_socket_found: true,
869 ..Default::default()
870 };
871
872 let cfg = ErrorsIntakeConfig::from_settings(&settings);
873 let endpoint = cfg.endpoint().unwrap();
874
875 assert_eq!(endpoint.url.scheme_str(), Some("unix"));
876 let decoded_path = libdd_common::decode_uri_path_in_authority(&endpoint.url).unwrap();
878 assert_eq!(
879 decoded_path.to_string_lossy(),
880 "/var/run/datadog/apm.socket"
881 );
882
883 assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
885 assert!(endpoint.api_key.is_none());
886 }
887
888 #[test]
889 #[cfg(windows)]
890 fn test_errors_intake_config_named_pipe() {
891 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
892
893 clear_errors_intake_env();
894
895 std::env::set_var("DD_TRACE_PIPE_NAME", "my_custom_pipe");
897
898 let cfg = ErrorsIntakeConfig::from_env();
899 let endpoint = cfg.endpoint().unwrap();
900
901 assert_eq!(endpoint.url.scheme_str(), Some("windows"));
902 let decoded_path = libdd_common::decode_uri_path_in_authority(&endpoint.url).unwrap();
904 assert_eq!(decoded_path.to_string_lossy(), "my_custom_pipe");
905
906 assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
908 assert!(endpoint.api_key.is_none());
909 }
910
911 #[test]
912 fn test_errors_intake_config_unix_url() {
913 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
914
915 clear_errors_intake_env();
916
917 std::env::set_var("DD_TRACE_AGENT_URL", "unix:///tmp/custom.socket");
919
920 let cfg = ErrorsIntakeConfig::from_env();
921 let endpoint = cfg.endpoint().unwrap();
922
923 assert_eq!(endpoint.url.scheme_str(), Some("unix"));
924 let decoded_path = libdd_common::decode_uri_path_in_authority(&endpoint.url).unwrap();
925 assert_eq!(decoded_path.to_string_lossy(), "/tmp/custom.socket");
926
927 assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
929 assert!(endpoint.api_key.is_none());
930 }
931
932 #[test]
933 fn test_errors_intake_endpoint_priority_order() {
934 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
935
936 clear_errors_intake_env();
937
938 std::env::set_var("DD_TRACE_AGENT_URL", "http://priority-url:9999");
940 std::env::set_var("DD_AGENT_HOST", "ignored-host");
941 std::env::set_var("DD_TRACE_AGENT_PORT", "1111");
942
943 let cfg = ErrorsIntakeConfig::from_env();
944 let endpoint = cfg.endpoint().unwrap();
945
946 assert_eq!(endpoint.url.host(), Some("priority-url"));
947 assert_eq!(endpoint.url.port_u16(), Some(9999));
948
949 clear_errors_intake_env();
950
951 std::env::set_var("DD_AGENT_HOST", "custom-host");
953 std::env::set_var("DD_TRACE_AGENT_PORT", "7777");
954
955 let cfg = ErrorsIntakeConfig::from_env();
956 let endpoint = cfg.endpoint().unwrap();
957
958 assert_eq!(endpoint.url.host(), Some("custom-host"));
959 assert_eq!(endpoint.url.port_u16(), Some(7777));
960
961 clear_errors_intake_env();
962
963 let cfg = ErrorsIntakeConfig::from_env();
965 let endpoint = cfg.endpoint().unwrap();
966
967 assert_eq!(endpoint.url.host(), Some(DEFAULT_AGENT_HOST));
968 assert_eq!(endpoint.url.port_u16(), Some(DEFAULT_AGENT_PORT));
969 }
970
971 #[test]
972 fn test_errors_intake_direct_submission_vs_agent_priority() {
973 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
974
975 clear_errors_intake_env();
976
977 std::env::set_var("DD_TRACE_AGENT_URL", "http://agent-host:8888");
979 std::env::set_var("DD_API_KEY", "test-key");
980 std::env::set_var("_DD_DIRECT_SUBMISSION_ENABLED", "true");
981
982 let cfg = ErrorsIntakeConfig::from_env();
983 let endpoint = cfg.endpoint().unwrap();
984
985 assert_eq!(
987 endpoint.url.host(),
988 Some("error-tracking-intake.datadoghq.com")
989 );
990 assert_eq!(endpoint.url.scheme_str(), Some("https"));
991 assert!(endpoint.api_key.is_some());
992 assert_eq!(endpoint.url.path(), DIRECT_ERRORS_INTAKE_URL_PATH);
993
994 clear_errors_intake_env();
995
996 std::env::set_var("DD_TRACE_AGENT_URL", "http://agent-host:8888");
998 std::env::set_var("DD_API_KEY", "test-key");
999 let cfg = ErrorsIntakeConfig::from_env();
1002 let endpoint = cfg.endpoint().unwrap();
1003
1004 assert_eq!(endpoint.url.host(), Some("agent-host"));
1006 assert_eq!(endpoint.url.port_u16(), Some(8888));
1007 assert!(endpoint.api_key.is_none());
1008 assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
1009 }
1010
1011 #[test]
1012 #[cfg(unix)]
1013 fn test_errors_intake_uds_priority_over_host_port() {
1014 let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
1015
1016 clear_errors_intake_env();
1017
1018 std::env::set_var("DD_AGENT_HOST", "ignored-host");
1020 std::env::set_var("DD_TRACE_AGENT_PORT", "9999");
1021
1022 let settings = ErrorsIntakeSettings {
1023 agent_host: Some("ignored-host".to_string()),
1024 trace_agent_port: Some(9999),
1025 agent_uds_socket_found: true,
1026 ..Default::default()
1027 };
1028
1029 let cfg = ErrorsIntakeConfig::from_settings(&settings);
1030 let endpoint = cfg.endpoint().unwrap();
1031
1032 assert_eq!(endpoint.url.scheme_str(), Some("unix"));
1034 let decoded_path = libdd_common::decode_uri_path_in_authority(&endpoint.url).unwrap();
1035 assert_eq!(
1036 decoded_path.to_string_lossy(),
1037 "/var/run/datadog/apm.socket"
1038 );
1039 }
1040}