Skip to main content

libdd_crashtracker/crash_info/
errors_intake.rs

1// Copyright 2024-Present Datadog, Inc. https://www.datadoghq.com/
2// SPDX-License-Identifier: Apache-2.0
3
4use 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/// Settings gathers configuration options we receive from the environment
62#[derive(Debug, Default)]
63pub struct ErrorsIntakeSettings {
64    // Env parameter
65    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    // Filesystem check
77    pub agent_uds_socket_found: bool,
78}
79
80impl ErrorsIntakeSettings {
81    // Agent connection configuration
82    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    // Direct submission configuration
88    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    // Debug configuration
94    const _DD_SHARED_LIB_DEBUG: &'static str = "_DD_SHARED_LIB_DEBUG";
95
96    // Feature flags
97    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        // For direct submission, construct the proper intake URL
203        let url = if settings.direct_submission_enabled && settings.api_key.is_some() {
204            // Check for explicit errors intake URL first
205            if let Some(ref errors_intake_url) = settings.errors_intake_dd_url {
206                errors_intake_url.clone()
207            } else {
208                // Build direct submission URL using site configuration
209                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        // Handle file endpoint
542        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        // Build HTTP request using the same pattern as telemetry
565        let mut req_builder =
566            endpoint.to_request_builder(concat!("crashtracker/", env!("CARGO_PKG_VERSION")))?;
567
568        // Add errors intake specific headers
569        if endpoint.api_key.is_some() {
570            // Direct intake - DD-API-KEY is added by to_request_builder
571        } else {
572            // Agent proxy - add EvP subdomain header
573            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        // Create HTTP client and send request
586        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    // Mutex to ensure environment variable tests run sequentially
605    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        // Basic structure
627        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        // error.message inherits from crash_info.error.message
632        assert_eq!(payload.error.message, crash_info.error.message);
633
634        // error.thread_name inherits from crash_info.error.thread_name
635        assert_eq!(payload.error.thread_name, crash_info.error.thread_name);
636
637        // error.stack inherits from crash_info.error.stack
638        assert_eq!(payload.error.stack, Some(crash_info.error.stack.clone()));
639
640        // error_type: signal name when sig_info present, ErrorKind as fallback
641        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        // experimental inherits from crash_info
654        assert_eq!(payload.error.experimental, crash_info.experimental);
655
656        // os_info inherits from crash_info
657        assert_eq!(payload.os_info, crash_info.os_info);
658
659        // sig_info inherits from crash_info
660        assert_eq!(payload.sig_info, crash_info.sig_info);
661
662        // proc_info inherits from crash_info
663        assert_eq!(payload.proc_info, crash_info.proc_info);
664
665        // files inherits from crash_info
666        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        // This test ensures we have all the tags that telemetry crash ping produces
783        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        // Test direct submission configuration
816        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        // Should use error-tracking-intake.datadoghq.com for direct submission
823        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        // With direct submission enabled and API key, should use direct path
831        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        // Test direct submission with custom site
841        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        // Should use error-tracking-intake with custom site
849        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        // Should use agent proxy path
873        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        // API key is set but direct submission is NOT enabled
883        // Should still use agent proxy
884        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        // Should use agent URL, not direct submission
891        assert_eq!(endpoint.url.host(), Some("localhost"));
892        assert_eq!(endpoint.url.port_u16(), Some(9126));
893
894        // Should use agent proxy path, not direct path
895        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        // Default behavior is enabled
906        let cfg = ErrorsIntakeConfig::from_env();
907        assert!(cfg.is_errors_intake_enabled());
908
909        // Test with uploader
910        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        // Test UDS socket configuration
926        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        // The original socket path is preserved in authority for unix:// URLs (URL encoded)
936        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        // Should use agent proxy path
943        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        // Test named pipe configuration
955        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        // For windows: scheme, the pipe name is URL-encoded in the authority
962        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        // Should use agent proxy path
966        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        // Test unix:// URL in DD_TRACE_AGENT_URL
977        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        // Should use agent proxy path
987        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        // Test 1: DD_TRACE_AGENT_URL takes highest priority
998        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        // Test 2: DD_AGENT_HOST + DD_TRACE_AGENT_PORT used when no DD_TRACE_AGENT_URL
1011        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        // Test 3: Default fallback when nothing is set
1023        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        // Test that direct submission takes priority over agent configuration
1037        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        // Should use direct submission URL, not agent URL
1045        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        // Test that without direct submission enabled, agent URL is used
1056        std::env::set_var("DD_TRACE_AGENT_URL", "http://agent-host:8888");
1057        std::env::set_var("DD_API_KEY", "test-key");
1058        // _DD_DIRECT_SUBMISSION_ENABLED not set (defaults to false)
1059
1060        let cfg = ErrorsIntakeConfig::from_env();
1061        let endpoint = cfg.endpoint().unwrap();
1062
1063        // Should use agent URL
1064        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        // Test that UDS socket takes priority over DD_AGENT_HOST/DD_TRACE_AGENT_PORT
1078        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        // Should use UDS socket, not host/port
1092        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}