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