use alloc::borrow::Cow;
use core::time::Duration;
use std::time::SystemTime;
use crate::{OsInfo, SigInfo, Ucontext};
use super::{
telemetry::CrashPing, CrashInfo, Experimental, Metadata, ProcInfo, StackTrace, ThreadData,
TARGET_TRIPLE,
};
use anyhow::Context;
use chrono::{DateTime, Utc};
use http::{uri::PathAndQuery, Uri};
use libdd_common::{config::parse_env, parse_uri, Endpoint};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
pub const DEFAULT_DD_SITE: &str = "datadoghq.com";
pub const PROD_ERRORS_INTAKE_SUBDOMAIN: &str = "error-tracking-intake";
const DIRECT_ERRORS_INTAKE_URL_PATH: &str = "/api/v2/errorsintake";
const AGENT_ERRORS_INTAKE_URL_PATH: &str = "/evp_proxy/v4/api/v2/errorsintake";
const DEFAULT_AGENT_HOST: &str = "localhost";
const DEFAULT_AGENT_PORT: u16 = 8126;
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct ErrorsIntakeConfig {
pub(crate) endpoint: Option<Endpoint>,
pub direct_submission_enabled: bool,
pub debug_enabled: bool,
pub errors_intake_enabled: bool,
}
fn endpoint_with_errors_intake_path(
mut endpoint: Endpoint,
direct_submission_enabled: bool,
) -> anyhow::Result<Endpoint> {
let mut uri_parts = endpoint.url.into_parts();
if uri_parts
.scheme
.as_ref()
.is_some_and(|scheme| scheme.as_str() != "file")
{
uri_parts.path_and_query = Some(PathAndQuery::from_static(
if endpoint.api_key.is_some() && direct_submission_enabled {
DIRECT_ERRORS_INTAKE_URL_PATH
} else {
AGENT_ERRORS_INTAKE_URL_PATH
},
));
}
endpoint.url = Uri::from_parts(uri_parts)?;
Ok(endpoint)
}
#[derive(Debug, Default)]
pub struct ErrorsIntakeSettings {
pub agent_host: Option<String>,
pub trace_agent_port: Option<u16>,
pub trace_agent_url: Option<String>,
pub trace_pipe_name: Option<String>,
pub direct_submission_enabled: bool,
pub api_key: Option<String>,
pub site: Option<String>,
pub errors_intake_dd_url: Option<String>,
pub shared_lib_debug: bool,
pub errors_intake_enabled: bool,
pub agent_uds_socket_found: bool,
}
impl ErrorsIntakeSettings {
const DD_TRACE_AGENT_URL: &'static str = "DD_TRACE_AGENT_URL";
const DD_AGENT_HOST: &'static str = "DD_AGENT_HOST";
const DD_TRACE_AGENT_PORT: &'static str = "DD_TRACE_AGENT_PORT";
const DD_TRACE_PIPE_NAME: &'static str = "DD_TRACE_PIPE_NAME";
const _DD_DIRECT_SUBMISSION_ENABLED: &'static str = "_DD_DIRECT_SUBMISSION_ENABLED";
const DD_API_KEY: &'static str = "DD_API_KEY";
const DD_SITE: &'static str = "DD_SITE";
const DD_ERRORS_INTAKE_DD_URL: &'static str = "DD_ERRORS_INTAKE_DD_URL";
const _DD_SHARED_LIB_DEBUG: &'static str = "_DD_SHARED_LIB_DEBUG";
const DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED: &'static str =
"DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED";
pub fn from_env() -> Self {
let default = Self::default();
Self {
agent_host: parse_env::str_not_empty(Self::DD_AGENT_HOST),
trace_agent_port: parse_env::int(Self::DD_TRACE_AGENT_PORT),
trace_agent_url: parse_env::str_not_empty(Self::DD_TRACE_AGENT_URL)
.or(default.trace_agent_url),
trace_pipe_name: parse_env::str_not_empty(Self::DD_TRACE_PIPE_NAME)
.or(default.trace_pipe_name),
direct_submission_enabled: parse_env::bool(Self::_DD_DIRECT_SUBMISSION_ENABLED)
.unwrap_or(default.direct_submission_enabled),
api_key: parse_env::str_not_empty(Self::DD_API_KEY),
site: parse_env::str_not_empty(Self::DD_SITE),
errors_intake_dd_url: parse_env::str_not_empty(Self::DD_ERRORS_INTAKE_DD_URL),
shared_lib_debug: parse_env::bool(Self::_DD_SHARED_LIB_DEBUG).unwrap_or(false),
errors_intake_enabled: parse_env::bool(Self::DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED)
.unwrap_or(true),
agent_uds_socket_found: (|| {
#[cfg(unix)]
return std::fs::metadata("/var/run/datadog/apm.socket").is_ok();
#[cfg(not(unix))]
return false;
})(),
}
}
}
impl ErrorsIntakeConfig {
fn trace_agent_url_from_setting(settings: &ErrorsIntakeSettings) -> String {
None.or_else(|| {
settings
.trace_agent_url
.as_deref()
.filter(|u| {
u.starts_with("unix://")
|| u.starts_with("http://")
|| u.starts_with("https://")
})
.map(ToString::to_string)
})
.or_else(|| {
#[cfg(windows)]
return settings
.trace_pipe_name
.as_ref()
.map(|pipe_name| format!("windows:{pipe_name}"));
#[cfg(not(windows))]
return None;
})
.or_else(|| {
#[cfg(unix)]
return settings
.agent_uds_socket_found
.then(|| "unix:///var/run/datadog/apm.socket".to_string());
#[cfg(not(unix))]
return None;
})
.or_else(|| match (&settings.agent_host, settings.trace_agent_port) {
(None, None) => None,
_ => Some(format!(
"http://{}:{}",
settings.agent_host.as_deref().unwrap_or(DEFAULT_AGENT_HOST),
settings.trace_agent_port.unwrap_or(DEFAULT_AGENT_PORT),
)),
})
.unwrap_or_else(|| format!("http://{DEFAULT_AGENT_HOST}:{DEFAULT_AGENT_PORT}"))
}
fn api_key_from_settings(settings: &ErrorsIntakeSettings) -> Option<Cow<'static, str>> {
if !settings.direct_submission_enabled {
return None;
}
settings.api_key.clone().map(Cow::Owned)
}
pub fn endpoint(&self) -> Option<&Endpoint> {
self.endpoint.as_ref()
}
pub fn is_errors_intake_enabled(&self) -> bool {
self.errors_intake_enabled
}
pub fn set_endpoint(&mut self, endpoint: Endpoint) -> anyhow::Result<()> {
self.endpoint = Some(endpoint_with_errors_intake_path(
endpoint,
self.direct_submission_enabled,
)?);
Ok(())
}
pub fn from_settings(settings: &ErrorsIntakeSettings) -> Self {
let api_key = Self::api_key_from_settings(settings);
let mut this = Self {
endpoint: None,
direct_submission_enabled: settings.direct_submission_enabled,
debug_enabled: settings.shared_lib_debug,
errors_intake_enabled: settings.errors_intake_enabled,
};
let url = if settings.direct_submission_enabled && settings.api_key.is_some() {
if let Some(ref errors_intake_url) = settings.errors_intake_dd_url {
errors_intake_url.clone()
} else {
let site = settings.site.as_deref().unwrap_or(DEFAULT_DD_SITE);
format!("https://{}.{}", PROD_ERRORS_INTAKE_SUBDOMAIN, site)
}
} else {
Self::trace_agent_url_from_setting(settings)
};
if let Ok(parsed_url) = parse_uri(&url) {
let _res = this.set_endpoint(Endpoint {
url: parsed_url,
api_key,
..Default::default()
});
}
this
}
pub fn from_env() -> Self {
let settings = ErrorsIntakeSettings::from_env();
Self::from_settings(&settings)
}
pub fn set_host_from_url(&mut self, host_url: &str) -> anyhow::Result<()> {
let endpoint = self.endpoint.take().unwrap_or_default();
self.set_endpoint(Endpoint {
url: parse_uri(host_url)?,
..endpoint
})
}
}
#[derive(serde::Serialize, Debug)]
pub struct ErrorObject {
#[serde(rename = "type", skip_serializing_if = "Option::is_none")]
pub error_type: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub experimental: Option<Experimental>,
#[serde(skip_serializing_if = "Option::is_none")]
pub is_crash: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub message: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub source_type: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub stack: Option<StackTrace>,
#[serde(skip_serializing_if = "Option::is_none")]
pub thread_name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub threads: Option<Vec<ThreadData>>,
}
#[derive(serde::Serialize, Debug)]
pub struct ErrorsIntakePayload {
pub ddsource: String,
pub ddtags: String,
pub error: ErrorObject,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub files: HashMap<String, Vec<String>>,
pub os_info: OsInfo,
#[serde(skip_serializing_if = "Option::is_none")]
pub proc_info: Option<ProcInfo>,
#[serde(skip_serializing_if = "Option::is_none")]
pub sig_info: Option<SigInfo>,
#[serde(skip_serializing_if = "Option::is_none")]
pub ucontext: Option<Ucontext>,
pub timestamp: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub trace_id: Option<String>,
}
#[derive(Debug, Default)]
struct ExtractedMetadata {
env: Option<String>,
language_name: Option<String>,
language_version: Option<String>,
service_name: String,
service_version: Option<String>,
tracer_version: Option<String>,
}
impl ExtractedMetadata {
fn from_metadata(metadata: &Metadata) -> Self {
let mut result = Self {
service_name: "unknown".to_string(),
..Default::default()
};
for tag in &metadata.tags {
if let Some((key, value)) = tag.split_once(':') {
match key {
"service" => result.service_name = value.to_string(),
"env" => result.env = Some(value.to_string()),
"version" | "service_version" => {
result.service_version = Some(value.to_string())
}
"language" => result.language_name = Some(value.to_string()),
"language_version" | "runtime_version" => {
result.language_version = Some(value.to_string())
}
"library_version" | "profiler_version" => {
result.tracer_version = Some(value.to_string())
}
_ => {}
}
}
}
result
}
fn append_base_tags(&self, tags: &mut String) {
tags.push_str(&format!("service:{}", self.service_name));
if let Some(env) = &self.env {
tags.push_str(&format!(",env:{env}"));
}
if let Some(version) = &self.service_version {
tags.push_str(&format!(",version:{version}"));
}
}
fn append_runtime_tags(&self, tags: &mut String) {
if let Some(language_name) = &self.language_name {
tags.push_str(&format!(",language_name:{language_name}"));
}
if let Some(language_version) = &self.language_version {
tags.push_str(&format!(",language_version:{language_version}"));
}
if let Some(tracer_version) = &self.tracer_version {
tags.push_str(&format!(",tracer_version:{tracer_version}"));
}
}
}
fn append_signal_tags(tags: &mut String, sig_info: &SigInfo) {
tags.push_str(&format!(
",si_code_human_readable:{:?}",
sig_info.si_code_human_readable
));
tags.push_str(&format!(",si_signo:{}", sig_info.si_signo));
tags.push_str(&format!(
",si_signo_human_readable:{:?}",
sig_info.si_signo_human_readable
));
}
fn build_crash_info_tags(crash_info: &CrashInfo) -> String {
let mut tags = format!("data_schema_version:{}", crash_info.data_schema_version);
if let Some(fingerprint) = &crash_info.fingerprint {
tags.push_str(&format!(",fingerprint:{fingerprint}"));
}
tags.push_str(&format!(",incomplete:{}", crash_info.incomplete));
tags.push_str(&format!(",is_crash:{}", crash_info.error.is_crash));
tags.push_str(&format!(",uuid:{}", crash_info.uuid));
for (counter, value) in &crash_info.counters {
tags.push_str(&format!(",{counter}:{value}"));
}
if let Some(siginfo) = &crash_info.sig_info {
if let Some(si_addr) = &siginfo.si_addr {
tags.push_str(&format!(",si_addr:{si_addr}"));
}
tags.push_str(&format!(",si_code:{}", siginfo.si_code));
append_signal_tags(&mut tags, siginfo);
}
tags.push_str(&format!(",runtime_platform:{TARGET_TRIPLE}"));
tags
}
impl ErrorsIntakePayload {
pub fn from_crash_info(crash_info: &CrashInfo) -> anyhow::Result<Self> {
let timestamp = crash_info.timestamp.parse::<DateTime<Utc>>().map_or_else(
|_| {
SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
},
|ts| ts.timestamp_millis() as u64,
);
let metadata = ExtractedMetadata::from_metadata(&crash_info.metadata);
let mut ddtags = String::new();
metadata.append_base_tags(&mut ddtags);
metadata.append_runtime_tags(&mut ddtags);
let crash_tags = build_crash_info_tags(crash_info);
ddtags.push_str(&format!(",{crash_tags}"));
let error_type = if let Some(sig_info) = &crash_info.sig_info {
Some(format!("{:?}", sig_info.si_signo_human_readable))
} else {
Some(format!("{:?}", crash_info.error.kind))
};
let error_message = crash_info.error.message.clone().or_else(|| {
crash_info.sig_info.as_ref().map(|sig_info| {
format!(
"Process terminated with {:?} ({:?})",
sig_info.si_code_human_readable, sig_info.si_signo_human_readable
)
.to_string()
})
});
let error_stack = if !crash_info.error.stack.frames.is_empty() {
Some(crash_info.error.stack.clone())
} else {
None
};
Ok(Self {
timestamp,
ddsource: "crashtracker".to_string(),
ddtags,
error: ErrorObject {
error_type,
message: error_message,
thread_name: crash_info.error.thread_name.clone(),
stack: error_stack,
is_crash: Some(true),
source_type: Some("Crashtracking".to_string()),
experimental: crash_info.experimental.clone(),
threads: crash_info.error.threads.clone(),
},
trace_id: None,
ucontext: crash_info.ucontext.clone(),
os_info: crash_info.os_info.clone(),
sig_info: crash_info.sig_info.clone(),
proc_info: crash_info.proc_info.clone(),
files: crash_info.files.clone(),
})
}
pub fn from_crash_ping(crash_ping: &CrashPing) -> anyhow::Result<Self> {
let timestamp = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let crash_uuid = crash_ping.crash_uuid();
let sig_info = crash_ping.siginfo();
let metadata = crash_ping.metadata();
let extracted_metadata = ExtractedMetadata::from_metadata(metadata);
let mut ddtags = format!(
"uuid:{},is_crash_ping:true,service:{}",
crash_uuid, extracted_metadata.service_name
);
extracted_metadata.append_runtime_tags(&mut ddtags);
if let Some(env) = &extracted_metadata.env {
ddtags.push_str(&format!(",env:{env}"));
}
if let Some(version) = &extracted_metadata.service_version {
ddtags.push_str(&format!(",version:{version}"));
}
if let Some(sig_info) = sig_info {
append_signal_tags(&mut ddtags, sig_info);
}
ddtags.push_str(&format!(",runtime_platform:{TARGET_TRIPLE}"));
let error_type = Some(
sig_info
.map(|s| format!("{:?}", s.si_signo_human_readable))
.unwrap_or_else(|| format!("{:?}", crash_ping.kind())),
);
let message = Some(crash_ping.message().to_string());
Ok(Self {
timestamp,
ddsource: "crashtracker".to_string(),
ddtags,
error: ErrorObject {
error_type,
message,
thread_name: None,
stack: None,
is_crash: Some(false),
source_type: Some("Crashtracking".to_string()),
experimental: None,
threads: None,
},
sig_info: sig_info.cloned(),
trace_id: None,
os_info: ::os_info::get().into(),
ucontext: None,
proc_info: None,
files: HashMap::new(),
})
}
}
pub struct ErrorsIntakeUploader {
cfg: ErrorsIntakeConfig,
}
impl ErrorsIntakeUploader {
pub fn new(endpoint: &Option<Endpoint>) -> anyhow::Result<Self> {
let mut cfg = ErrorsIntakeConfig::from_env();
if let Some(endpoint) = endpoint {
cfg.set_endpoint(endpoint.clone())?;
}
Ok(Self { cfg })
}
pub fn is_enabled(&self) -> bool {
self.cfg.is_errors_intake_enabled()
}
pub async fn upload_crash_ping(&self, crash_ping: &CrashPing) -> anyhow::Result<()> {
let payload = ErrorsIntakePayload::from_crash_ping(crash_ping)?;
self.send_payload(&payload).await
}
pub async fn upload_crash_info(&self, crash_info: &CrashInfo) -> anyhow::Result<()> {
let payload = ErrorsIntakePayload::from_crash_info(crash_info)?;
self.send_payload(&payload).await
}
async fn send_payload(&self, payload: &ErrorsIntakePayload) -> anyhow::Result<()> {
let Some(endpoint) = self.cfg.endpoint() else {
return Ok(());
};
if endpoint.url.scheme_str() == Some("file") {
let path = libdd_common::decode_uri_path_in_authority(&endpoint.url)
.context("errors intake file path is not valid")?;
let file_path = path.with_extension("errors");
let file = std::fs::File::create(&file_path).with_context(|| {
format!(
"Failed to create errors intake file {}",
file_path.display()
)
})?;
serde_json::to_writer_pretty(file, payload).with_context(|| {
format!(
"Failed to write errors intake JSON to {}",
file_path.display()
)
})?;
return Ok(());
}
let mut req_builder =
endpoint.to_request_builder(concat!("crashtracker/", env!("CARGO_PKG_VERSION")))?;
if endpoint.api_key.is_some() {
} else {
req_builder =
req_builder.header("X-Datadog-EVP-Subdomain", PROD_ERRORS_INTAKE_SUBDOMAIN);
}
let req = req_builder
.method(http::Method::POST)
.header(
http::header::CONTENT_TYPE,
libdd_common::header::APPLICATION_JSON,
)
.body(serde_json::to_string(payload)?.into())?;
let client = libdd_common::http_common::new_client_periodic();
tokio::time::timeout(
Duration::from_millis(endpoint.timeout_ms),
client.request(req),
)
.await??;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::crash_info::test_utils::TestInstance;
use std::sync::Mutex;
static ENV_TEST_LOCK: Mutex<()> = Mutex::new(());
fn clear_errors_intake_env() {
std::env::remove_var("DD_TRACE_AGENT_URL");
std::env::remove_var("DD_AGENT_HOST");
std::env::remove_var("DD_TRACE_AGENT_PORT");
std::env::remove_var("DD_TRACE_PIPE_NAME");
std::env::remove_var("_DD_DIRECT_SUBMISSION_ENABLED");
std::env::remove_var("DD_API_KEY");
std::env::remove_var("DD_SITE");
std::env::remove_var("DD_ERRORS_INTAKE_DD_URL");
std::env::remove_var("_DD_SHARED_LIB_DEBUG");
std::env::remove_var("DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED");
}
#[cfg_attr(miri, ignore)]
#[test]
fn test_errors_payload_from_crash_info() {
let crash_info = CrashInfo::test_instance(1);
let payload = ErrorsIntakePayload::from_crash_info(&crash_info).unwrap();
assert_eq!(payload.ddsource, "crashtracker");
assert_eq!(payload.error.source_type, Some("Crashtracking".to_string()));
assert_eq!(payload.error.is_crash, Some(true));
assert_eq!(payload.error.message, crash_info.error.message);
assert_eq!(payload.error.thread_name, crash_info.error.thread_name);
assert_eq!(payload.error.stack, Some(crash_info.error.stack.clone()));
assert_eq!(
payload.error.error_type,
Some(format!(
"{:?}",
crash_info
.sig_info
.as_ref()
.unwrap()
.si_signo_human_readable
))
);
assert_eq!(payload.error.experimental, crash_info.experimental);
assert_eq!(payload.os_info, crash_info.os_info);
assert_eq!(payload.sig_info, crash_info.sig_info);
assert_eq!(payload.proc_info, crash_info.proc_info);
assert_eq!(payload.files, crash_info.files);
let ddtags = &payload.ddtags;
assert!(ddtags.contains("service:foo"));
assert!(ddtags.contains("version:bar"));
assert!(ddtags.contains("language_name:native"));
assert!(ddtags.contains("data_schema_version:1.8"));
assert!(ddtags.contains("incomplete:true"));
assert!(ddtags.contains("is_crash:true"));
assert!(ddtags.contains("uuid:1d6b97cb-968c-40c9-af6e-e4b4d71e8781"));
assert!(ddtags.contains("collecting_sample:1"));
assert!(ddtags.contains("not_profiling:0"));
assert!(ddtags.contains("si_addr:0x0000000000001234"));
assert!(ddtags.contains("si_code:1"));
assert!(ddtags.contains("si_code_human_readable:SEGV_BNDERR"));
assert!(ddtags.contains("si_signo:11"));
assert!(ddtags.contains("si_signo_human_readable:SIGSEGV"));
}
#[cfg_attr(miri, ignore)]
#[test]
fn test_errors_payload_from_crash_ping() {
let metadata = Metadata::test_instance(1);
let sig_info = crate::SigInfo::test_instance(42);
let crash_uuid = uuid::Uuid::from_u128(0x01);
let crash_ping = crate::CrashPingBuilder::new(crash_uuid)
.with_metadata(metadata.clone())
.with_kind(crate::ErrorKind::UnixSignal)
.with_sig_info(sig_info.clone())
.build()
.unwrap();
let payload = ErrorsIntakePayload::from_crash_ping(&crash_ping).unwrap();
assert_eq!(payload.ddsource, "crashtracker");
assert_eq!(payload.error.source_type, Some("Crashtracking".to_string()));
assert_eq!(payload.error.is_crash, Some(false));
assert!(payload.error.stack.is_none());
assert_eq!(payload.error.message.as_deref(), Some(crash_ping.message()));
assert_eq!(
payload.error.error_type,
Some(format!("{:?}", sig_info.si_signo_human_readable))
);
let ddtags = &payload.ddtags;
assert!(ddtags.contains(&format!("uuid:{crash_uuid}")));
assert!(ddtags.contains("is_crash_ping:true"));
assert!(ddtags.contains("service:foo"));
assert!(ddtags.contains("language_name:native"));
assert!(ddtags.contains("version:bar"));
assert!(ddtags.contains("si_code_human_readable:SEGV_BNDERR"));
assert!(ddtags.contains("si_signo:11"));
assert!(ddtags.contains("si_signo_human_readable:SIGSEGV"));
}
#[cfg_attr(miri, ignore)]
#[test]
fn test_errors_intake_has_all_telemetry_tags() {
let crash_info = CrashInfo::test_instance(1);
let payload = ErrorsIntakePayload::from_crash_info(&crash_info).unwrap();
let expected_crash_tags = [
"data_schema_version:1.8",
"incomplete:true",
"is_crash:true",
"uuid:1d6b97cb-968c-40c9-af6e-e4b4d71e8781",
"collecting_sample:1",
"not_profiling:0",
"si_addr:0x0000000000001234",
"si_code:1",
"si_code_human_readable:SEGV_BNDERR",
"si_signo:11",
"si_signo_human_readable:SIGSEGV",
&format!("runtime_platform:{}", super::super::TARGET_TRIPLE),
];
let expected_metadata_tags = ["service:foo", "version:bar", "language_name:native"];
for tag in expected_crash_tags
.iter()
.chain(expected_metadata_tags.iter())
{
assert!(
payload.ddtags.contains(tag),
"Missing expected tag: {} in ddtags: {}",
tag,
payload.ddtags
);
}
}
#[cfg_attr(miri, ignore)]
#[test]
fn test_crash_ping_has_all_telemetry_tags() {
let metadata = Metadata::test_instance(1);
let sig_info = crate::SigInfo::test_instance(42);
let crash_uuid = uuid::Uuid::from_u128(0x02);
let crash_ping = crate::CrashPingBuilder::new(crash_uuid)
.with_metadata(metadata.clone())
.with_kind(crate::ErrorKind::UnixSignal)
.with_sig_info(sig_info.clone())
.build()
.unwrap();
let payload = ErrorsIntakePayload::from_crash_ping(&crash_ping).unwrap();
assert!(
payload.ddtags.contains(&format!("uuid:{crash_uuid}")),
"Missing uuid tag in ddtags: {}",
payload.ddtags
);
let expected_tags = [
"is_crash_ping:true",
"service:foo",
"language_name:native",
"version:bar",
"si_code_human_readable:SEGV_BNDERR",
"si_signo:11",
"si_signo_human_readable:SIGSEGV",
&format!("runtime_platform:{}", super::super::TARGET_TRIPLE),
];
for tag in expected_tags {
assert!(
payload.ddtags.contains(tag),
"Missing expected tag: {} in ddtags: {}",
tag,
payload.ddtags
);
}
}
#[test]
fn test_errors_intake_config_from_env() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
std::env::set_var("DD_API_KEY", "test-key");
std::env::set_var("_DD_DIRECT_SUBMISSION_ENABLED", "true");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(
endpoint.url.host(),
Some("error-tracking-intake.datadoghq.com")
);
assert_eq!(endpoint.url.scheme_str(), Some("https"));
assert!(endpoint.api_key.is_some());
assert_eq!(endpoint.url.path(), DIRECT_ERRORS_INTAKE_URL_PATH);
}
#[test]
fn test_errors_intake_config_custom_site() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
std::env::set_var("DD_API_KEY", "test-key");
std::env::set_var("_DD_DIRECT_SUBMISSION_ENABLED", "true");
std::env::set_var("DD_SITE", "us3.datadoghq.com");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(
endpoint.url.host(),
Some("error-tracking-intake.us3.datadoghq.com")
);
assert_eq!(endpoint.url.scheme_str(), Some("https"));
assert!(endpoint.api_key.is_some());
assert_eq!(endpoint.url.path(), DIRECT_ERRORS_INTAKE_URL_PATH);
}
#[test]
fn test_errors_intake_config_agent_proxy() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
std::env::set_var("DD_TRACE_AGENT_URL", "http://localhost:9126");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.host(), Some("localhost"));
assert_eq!(endpoint.url.port_u16(), Some(9126));
assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
}
#[test]
fn test_errors_intake_config_agent_with_api_key_but_no_direct() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
std::env::set_var("DD_TRACE_AGENT_URL", "http://localhost:9126");
std::env::set_var("DD_API_KEY", "test-key");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.host(), Some("localhost"));
assert_eq!(endpoint.url.port_u16(), Some(9126));
assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
assert!(endpoint.api_key.is_none());
}
#[test]
fn test_errors_intake_disabled_flag() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
let cfg = ErrorsIntakeConfig::from_env();
assert!(cfg.is_errors_intake_enabled());
let uploader = ErrorsIntakeUploader::new(&None).unwrap();
assert!(uploader.is_enabled());
std::env::set_var("DD_CRASHTRACKING_ERRORS_INTAKE_ENABLED", "false");
let uploader = ErrorsIntakeUploader::new(&None).unwrap();
assert!(!uploader.is_enabled());
}
#[test]
#[cfg(unix)]
fn test_errors_intake_config_uds_socket() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
let settings = ErrorsIntakeSettings {
agent_uds_socket_found: true,
..Default::default()
};
let cfg = ErrorsIntakeConfig::from_settings(&settings);
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.scheme_str(), Some("unix"));
let decoded_path = libdd_common::decode_uri_path_in_authority(&endpoint.url).unwrap();
assert_eq!(
decoded_path.to_string_lossy(),
"/var/run/datadog/apm.socket"
);
assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
assert!(endpoint.api_key.is_none());
}
#[test]
#[cfg(windows)]
fn test_errors_intake_config_named_pipe() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
std::env::set_var("DD_TRACE_PIPE_NAME", "my_custom_pipe");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.scheme_str(), Some("windows"));
let decoded_path = libdd_common::decode_uri_path_in_authority(&endpoint.url).unwrap();
assert_eq!(decoded_path.to_string_lossy(), "my_custom_pipe");
assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
assert!(endpoint.api_key.is_none());
}
#[test]
fn test_errors_intake_config_unix_url() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
std::env::set_var("DD_TRACE_AGENT_URL", "unix:///tmp/custom.socket");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.scheme_str(), Some("unix"));
let decoded_path = libdd_common::decode_uri_path_in_authority(&endpoint.url).unwrap();
assert_eq!(decoded_path.to_string_lossy(), "/tmp/custom.socket");
assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
assert!(endpoint.api_key.is_none());
}
#[test]
fn test_errors_intake_endpoint_priority_order() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
std::env::set_var("DD_TRACE_AGENT_URL", "http://priority-url:9999");
std::env::set_var("DD_AGENT_HOST", "ignored-host");
std::env::set_var("DD_TRACE_AGENT_PORT", "1111");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.host(), Some("priority-url"));
assert_eq!(endpoint.url.port_u16(), Some(9999));
clear_errors_intake_env();
std::env::set_var("DD_AGENT_HOST", "custom-host");
std::env::set_var("DD_TRACE_AGENT_PORT", "7777");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.host(), Some("custom-host"));
assert_eq!(endpoint.url.port_u16(), Some(7777));
clear_errors_intake_env();
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.host(), Some(DEFAULT_AGENT_HOST));
assert_eq!(endpoint.url.port_u16(), Some(DEFAULT_AGENT_PORT));
}
#[test]
fn test_errors_intake_direct_submission_vs_agent_priority() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
std::env::set_var("DD_TRACE_AGENT_URL", "http://agent-host:8888");
std::env::set_var("DD_API_KEY", "test-key");
std::env::set_var("_DD_DIRECT_SUBMISSION_ENABLED", "true");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(
endpoint.url.host(),
Some("error-tracking-intake.datadoghq.com")
);
assert_eq!(endpoint.url.scheme_str(), Some("https"));
assert!(endpoint.api_key.is_some());
assert_eq!(endpoint.url.path(), DIRECT_ERRORS_INTAKE_URL_PATH);
clear_errors_intake_env();
std::env::set_var("DD_TRACE_AGENT_URL", "http://agent-host:8888");
std::env::set_var("DD_API_KEY", "test-key");
let cfg = ErrorsIntakeConfig::from_env();
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.host(), Some("agent-host"));
assert_eq!(endpoint.url.port_u16(), Some(8888));
assert!(endpoint.api_key.is_none());
assert_eq!(endpoint.url.path(), AGENT_ERRORS_INTAKE_URL_PATH);
}
#[test]
#[cfg(unix)]
fn test_errors_intake_uds_priority_over_host_port() {
let _lock = ENV_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_errors_intake_env();
std::env::set_var("DD_AGENT_HOST", "ignored-host");
std::env::set_var("DD_TRACE_AGENT_PORT", "9999");
let settings = ErrorsIntakeSettings {
agent_host: Some("ignored-host".to_string()),
trace_agent_port: Some(9999),
agent_uds_socket_found: true,
..Default::default()
};
let cfg = ErrorsIntakeConfig::from_settings(&settings);
let endpoint = cfg.endpoint().unwrap();
assert_eq!(endpoint.url.scheme_str(), Some("unix"));
let decoded_path = libdd_common::decode_uri_path_in_authority(&endpoint.url).unwrap();
assert_eq!(
decoded_path.to_string_lossy(),
"/var/run/datadog/apm.socket"
);
}
}