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 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/// Settings gathers configuration options we receive from the environment
57#[derive(Debug, Default)]
58pub struct ErrorsIntakeSettings {
59    // Env parameter
60    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    // Filesystem check
72    pub agent_uds_socket_found: bool,
73}
74
75impl ErrorsIntakeSettings {
76    // Agent connection configuration
77    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    // Direct submission configuration
83    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    // Debug configuration
89    const _DD_SHARED_LIB_DEBUG: &'static str = "_DD_SHARED_LIB_DEBUG";
90
91    // Feature flags
92    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        // For direct submission, construct the proper intake URL
198        let url = if settings.direct_submission_enabled && settings.api_key.is_some() {
199            // Check for explicit errors intake URL first
200            if let Some(ref errors_intake_url) = settings.errors_intake_dd_url {
201                errors_intake_url.clone()
202            } else {
203                // Build direct submission URL using site configuration
204                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        // Use crash stack if available
405        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            // Crash ping does not include os_info, but we can recalculate it here
490            // so that errors intake crash pings include this information
491            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        // Handle file endpoint
535        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        // Build HTTP request using the same pattern as telemetry
558        let mut req_builder =
559            endpoint.to_request_builder(concat!("crashtracker/", env!("CARGO_PKG_VERSION")))?;
560
561        // Add errors intake specific headers
562        if endpoint.api_key.is_some() {
563            // Direct intake - DD-API-KEY is added by to_request_builder
564        } else {
565            // Agent proxy - add EvP subdomain header
566            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        // Create HTTP client and send request
579        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    // Mutex to ensure environment variable tests run sequentially
598    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        // This test ensures we have all the tags that telemetry crash ping produces
720        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        // Test direct submission configuration
748        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        // Should use error-tracking-intake.datadoghq.com for direct submission
755        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        // With direct submission enabled and API key, should use direct path
763        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        // Test direct submission with custom site
773        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        // Should use error-tracking-intake with custom site
781        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        // Should use agent proxy path
805        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        // API key is set but direct submission is NOT enabled
815        // Should still use agent proxy
816        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        // Should use agent URL, not direct submission
823        assert_eq!(endpoint.url.host(), Some("localhost"));
824        assert_eq!(endpoint.url.port_u16(), Some(9126));
825
826        // Should use agent proxy path, not direct path
827        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        // Test default behavior (should be enabled)
836        clear_errors_intake_env();
837        let cfg = ErrorsIntakeConfig::from_env();
838        assert!(!cfg.is_errors_intake_enabled());
839
840        // Test explicitly enabled
841        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        // Test explicitly disabled
846        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        // Test with uploader
851        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        // Test UDS socket configuration
867        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        // The original socket path is preserved in authority for unix:// URLs (URL encoded)
877        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        // Should use agent proxy path
884        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        // Test named pipe configuration
896        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        // For windows: scheme, the pipe name is URL-encoded in the authority
903        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        // Should use agent proxy path
907        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        // Test unix:// URL in DD_TRACE_AGENT_URL
918        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        // Should use agent proxy path
928        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        // Test 1: DD_TRACE_AGENT_URL takes highest priority
939        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        // Test 2: DD_AGENT_HOST + DD_TRACE_AGENT_PORT used when no DD_TRACE_AGENT_URL
952        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        // Test 3: Default fallback when nothing is set
964        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        // Test that direct submission takes priority over agent configuration
978        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        // Should use direct submission URL, not agent URL
986        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        // Test that without direct submission enabled, agent URL is used
997        std::env::set_var("DD_TRACE_AGENT_URL", "http://agent-host:8888");
998        std::env::set_var("DD_API_KEY", "test-key");
999        // _DD_DIRECT_SUBMISSION_ENABLED not set (defaults to false)
1000
1001        let cfg = ErrorsIntakeConfig::from_env();
1002        let endpoint = cfg.endpoint().unwrap();
1003
1004        // Should use agent URL
1005        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        // Test that UDS socket takes priority over DD_AGENT_HOST/DD_TRACE_AGENT_PORT
1019        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        // Should use UDS socket, not host/port
1033        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}