use core::fmt::Write;
use std::time::SystemTime;
use crate::{ErrorKind, SigInfo};
use super::{CrashInfo, Metadata, TARGET_TRIPLE};
use anyhow::Context;
use chrono::{DateTime, Utc};
use libdd_capabilities::HttpClientCapability;
use libdd_capabilities_impl::NativeCapabilities;
use libdd_common::Endpoint;
use libdd_telemetry::{
build_host,
data::{self, Application, LogLevel},
worker::http_client::request_builder,
};
use serde::Serialize;
use uuid::Uuid;
#[derive(Debug)]
struct TelemetryMetadata {
application: libdd_telemetry::data::Application,
host: libdd_telemetry::data::Host,
runtime_id: String,
}
pub struct CrashPingBuilder {
crash_uuid: Uuid,
custom_message: Option<String>,
kind: Option<ErrorKind>,
metadata: Option<Metadata>,
sig_info: Option<SigInfo>,
}
impl CrashPingBuilder {
pub fn new(crash_uuid: Uuid) -> Self {
Self {
crash_uuid,
custom_message: None,
kind: None,
metadata: None,
sig_info: None,
}
}
pub fn with_sig_info(mut self, sig_info: SigInfo) -> Self {
self.sig_info = Some(sig_info);
self
}
pub fn with_custom_message(mut self, message: String) -> Self {
self.custom_message = Some(message);
self
}
pub fn with_kind(mut self, kind: ErrorKind) -> Self {
self.kind = Some(kind);
self
}
pub fn with_metadata(mut self, metadata: Metadata) -> Self {
self.metadata = Some(metadata);
self
}
pub fn build(self) -> anyhow::Result<CrashPing> {
let crash_uuid = self.crash_uuid;
let sig_info = self.sig_info;
let metadata = self.metadata.context("metadata is required")?;
let kind = self.kind.context("kind is required")?;
let message = if let Some(custom_message) = self.custom_message {
format!("Crashtracker crash ping: crash processing started - {custom_message}")
} else if let Some(ref sig_info) = sig_info {
format!(
"Crashtracker crash ping: crash processing started - Process terminated with {:?} ({:?})",
sig_info.si_code_human_readable, sig_info.si_signo_human_readable
)
} else {
format!("Crashtracker crash ping: crash processing started - Process terminated due to {:?}", kind)
};
Ok(CrashPing {
crash_uuid: crash_uuid.to_string(),
message,
kind,
metadata,
siginfo: sig_info,
version: CrashPing::current_schema_version(),
})
}
}
#[derive(Debug, Serialize)]
pub struct CrashPing {
crash_uuid: String,
kind: ErrorKind,
message: String,
#[serde(skip_serializing_if = "Option::is_none")]
siginfo: Option<SigInfo>,
version: String,
metadata: Metadata,
}
impl CrashPing {
pub fn crash_uuid(&self) -> &str {
&self.crash_uuid
}
pub fn message(&self) -> &str {
&self.message
}
pub fn metadata(&self) -> &Metadata {
&self.metadata
}
pub fn siginfo(&self) -> Option<&SigInfo> {
self.siginfo.as_ref()
}
pub fn kind(&self) -> ErrorKind {
self.kind.clone()
}
pub fn upload_to_endpoint(&self, endpoint: &Option<Endpoint>) -> anyhow::Result<()> {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
rt.block_on(async { self.upload_to_endpoint_async(endpoint).await })
}
pub async fn upload_to_endpoint_async(
&self,
endpoint: &Option<Endpoint>,
) -> anyhow::Result<()> {
let telemetry_uploader = crate::TelemetryCrashUploader::new(self.metadata(), endpoint)?;
let errors_intake_uploader = crate::ErrorsIntakeUploader::new(endpoint)?;
let telemetry_future = telemetry_uploader.upload_crash_ping(self);
if errors_intake_uploader.is_enabled() {
let errors_intake_future = errors_intake_uploader.upload_crash_ping(self);
let (_telemetry_result, _errors_intake_result) =
tokio::join!(telemetry_future, errors_intake_future);
} else {
let _telemetry_result = telemetry_future.await;
}
Ok(())
}
fn current_schema_version() -> String {
"1.0".to_string()
}
}
macro_rules! parse_tags {
( $tag_iterator:expr,
$($tag_name:literal => $var:ident),* $(,)?) => {
$(
let mut $var: Option<&str> = None;
)*
for tag in $tag_iterator {
let Some((name, value)) = tag.split_once(':') else {
continue;
};
match name {
$($tag_name => {$var = Some(value);}, )*
_ => {},
}
}
};
}
pub struct TelemetryCrashUploader {
metadata: TelemetryMetadata,
cfg: libdd_telemetry::config::Config,
}
impl TelemetryCrashUploader {
pub fn new(
crashtracker_metadata: &Metadata,
endpoint: &Option<Endpoint>,
) -> anyhow::Result<Self> {
let mut cfg = libdd_telemetry::config::Config::from_env();
if let Some(endpoint) = endpoint {
if endpoint.url.scheme_str() == Some("file") {
let path = libdd_common::decode_uri_path_in_authority(&endpoint.url)
.context("file path is not valid")?;
let _ = cfg.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
url: Some(format!("file://{}.telemetry", path.display())),
..Default::default()
});
} else {
let _ = cfg.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
api_key: endpoint.api_key.as_deref().map(str::to_owned),
test_token: endpoint.test_token.as_deref().map(str::to_owned),
timeout_ms: endpoint.timeout_ms,
use_system_resolver: endpoint.use_system_resolver,
..Default::default()
});
let _ = cfg.set_endpoint_uri(endpoint.url.clone());
}
}
parse_tags!(
crashtracker_metadata.tags.iter(),
"env" => env,
"language" => language_name,
"library_version" => library_version,
"profiler_version" => profiler_version,
"runtime_version" => language_version,
"runtime-id" => runtime_id,
"service_version" => service_version,
"service" => service_name,
"process_tags" => process_tags,
);
let application = Application {
service_name: service_name.unwrap_or("unknown").to_owned(),
language_name: language_name.unwrap_or("unknown").to_owned(),
language_version: language_version.unwrap_or("unknown").to_owned(),
tracer_version: library_version
.or(profiler_version)
.unwrap_or("unknown")
.to_owned(),
env: env.map(ToOwned::to_owned),
service_version: service_version.map(ToOwned::to_owned),
process_tags: process_tags.map(ToOwned::to_owned),
..Default::default()
};
let host = build_host();
let s = Self {
metadata: TelemetryMetadata {
host,
application,
runtime_id: runtime_id.unwrap_or("unknown").to_owned(),
},
cfg,
};
Ok(s)
}
pub async fn upload_general_log(
&self,
message: String,
tags: String,
level: LogLevel,
) -> anyhow::Result<()> {
let tracer_time = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
self.send_log_payload(message, tags, tracer_time, level, false, false)
.await
}
pub async fn upload_crash_ping(&self, crash_ping: &CrashPing) -> anyhow::Result<()> {
let tags = self.build_crash_ping_tags(crash_ping.crash_uuid(), crash_ping.siginfo());
let tracer_time = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let message = serde_json::to_string(crash_ping)?;
self.send_log_payload(
message,
tags,
tracer_time,
LogLevel::Debug,
false, false, )
.await
}
pub async fn upload_crash_info(&self, crash_info: &CrashInfo) -> anyhow::Result<()> {
let message = serde_json::to_string(crash_info)?;
let tags = extract_crash_info_tags(crash_info).unwrap_or_default();
let tracer_time = crash_info.timestamp.parse::<DateTime<Utc>>().map_or_else(
|_| {
SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
},
|ts| ts.timestamp() as u64,
);
self.send_log_payload(
message,
tags,
tracer_time,
LogLevel::Error,
true, true, )
.await
}
async fn send_log_payload(
&self,
message: String,
tags: String,
tracer_time: u64,
level: LogLevel,
is_sensitive: bool,
is_crash: bool,
) -> anyhow::Result<()> {
let payload = data::Telemetry {
tracer_time,
api_version: libdd_telemetry::data::ApiVersion::V2,
runtime_id: &self.metadata.runtime_id,
seq_id: 1,
application: &self.metadata.application,
host: &self.metadata.host,
payload: &data::Payload::Logs(data::Logs {
logs: vec![data::Log {
message,
level,
stack_trace: None,
tags,
is_sensitive,
count: 1,
is_crash,
}],
}),
origin: Some("Crashtracker"),
};
self.send_telemetry_payload(&payload).await
}
async fn send_telemetry_payload(&self, payload: &data::Telemetry<'_>) -> anyhow::Result<()> {
let client = NativeCapabilities::new_client();
let req = request_builder(&self.cfg)?
.method(http::Method::POST)
.header(
http::header::CONTENT_TYPE,
libdd_common::header::APPLICATION_JSON,
)
.header(
libdd_telemetry::worker::http_client::header::API_VERSION,
libdd_telemetry::data::ApiVersion::V2.to_str(),
)
.header(
libdd_telemetry::worker::http_client::header::REQUEST_TYPE,
"logs",
)
.body(libdd_capabilities::Bytes::from(serde_json::to_vec(
&payload,
)?))?;
let timeout = core::time::Duration::from_millis({
if let Some(endp) = self.cfg.endpoint() {
endp.timeout_ms
} else {
Endpoint::DEFAULT_TIMEOUT
}
});
tokio::time::timeout(timeout, client.request(req))
.await
.map_err(|_| anyhow::anyhow!("Telemetry crash report timed out"))??;
Ok(())
}
fn build_crash_ping_tags(&self, crash_uuid: &str, sig_info: Option<&SigInfo>) -> String {
let metadata = &self.metadata;
let mut tags = format!(
"uuid:{},is_crash_ping:true,service:{},language_name:{},language_version:{},tracer_version:{}",
crash_uuid,
metadata.application.service_name,
metadata.application.language_name,
metadata.application.language_version,
metadata.application.tracer_version
);
if let Some(sig_info) = sig_info {
tags.push_str(&format!(
",si_code_human_readable:{:?},si_signo:{},si_signo_human_readable:{:?}",
sig_info.si_code_human_readable,
sig_info.si_signo,
sig_info.si_signo_human_readable
));
}
write!(tags, ",runtime_platform:{TARGET_TRIPLE}").ok();
self.append_optional_tags(&mut tags);
tags
}
fn append_optional_tags(&self, tags: &mut String) {
let metadata = &self.metadata;
if let Some(env) = &metadata.application.env {
tags.push_str(&format!(",env:{env}"));
}
if let Some(runtime_name) = &metadata.application.runtime_name {
tags.push_str(&format!(",runtime_name:{runtime_name}"));
}
if let Some(runtime_version) = &metadata.application.runtime_version {
tags.push_str(&format!(",runtime_version:{runtime_version}"));
}
}
}
fn extract_crash_info_tags(crash_info: &CrashInfo) -> anyhow::Result<String> {
let mut tags = String::new();
write!(
&mut tags,
"data_schema_version:{}",
crash_info.data_schema_version
)?;
if let Some(fingerprint) = &crash_info.fingerprint {
write!(&mut tags, ",fingerprint:{fingerprint}")?;
}
write!(&mut tags, ",incomplete:{}", crash_info.incomplete)?;
write!(&mut tags, ",is_crash:{}", crash_info.error.is_crash)?;
write!(&mut tags, ",uuid:{}", crash_info.uuid)?;
for (counter, value) in &crash_info.counters {
write!(&mut tags, ",{counter}:{value}")?;
}
if let Some(siginfo) = &crash_info.sig_info {
if let Some(si_addr) = &siginfo.si_addr {
write!(&mut tags, ",si_addr:{si_addr}")?;
}
write!(&mut tags, ",si_code:{}", siginfo.si_code)?;
write!(
&mut tags,
",si_code_human_readable:{:?}",
siginfo.si_code_human_readable
)?;
write!(&mut tags, ",si_signo:{}", siginfo.si_signo)?;
write!(
&mut tags,
",si_signo_human_readable:{:?}",
siginfo.si_signo_human_readable
)?;
}
write!(&mut tags, ",runtime_platform:{TARGET_TRIPLE}")?;
Ok(tags)
}
#[cfg(test)]
mod tests {
use super::TelemetryCrashUploader;
use crate::{
crash_info::{test_utils::TestInstance, CrashInfo, CrashInfoBuilder, Metadata},
ErrorKind,
};
use libdd_common::Endpoint;
use libdd_telemetry::data::LogLevel;
use std::{collections::HashSet, fs};
use uuid::Uuid;
fn new_test_uploader(seed: u64) -> TelemetryCrashUploader {
TelemetryCrashUploader::new(
&Metadata::test_instance(seed),
&Some(Endpoint::from_slice("http://localhost:8126")),
)
.unwrap()
}
fn new_test_uploader_with_process_tags(
seed: u64,
process_tags: &str,
) -> TelemetryCrashUploader {
let mut metadata = Metadata::test_instance(seed);
metadata.tags.push(format!("process_tags:{process_tags}"));
TelemetryCrashUploader::new(
&metadata,
&Some(Endpoint::from_slice("http://localhost:8126")),
)
.unwrap()
}
#[test]
#[cfg_attr(miri, ignore)]
fn test_profiler_config_extraction() {
let t = new_test_uploader(1);
let metadata = t.metadata;
assert_eq!(metadata.application.service_name, "foo");
assert_eq!(metadata.application.service_version.as_deref(), Some("bar"));
assert_eq!(metadata.application.language_name, "native");
assert_eq!(metadata.application.process_tags, None);
assert_eq!(metadata.runtime_id, "xyz");
let cfg = t.cfg;
assert_eq!(
cfg.endpoint().unwrap().url.to_string(),
"http://localhost:8126/telemetry/proxy/api/v2/apmtelemetry"
);
}
#[tokio::test]
#[cfg_attr(miri, ignore)]
async fn test_crash_request_content() -> anyhow::Result<()> {
let tmp = tempfile::tempdir().unwrap();
let output_filename = {
let mut p = tmp.keep();
p.push("crash_info");
p
};
let seed = 1;
let mut t =
new_test_uploader_with_process_tags(seed, "entrypoint.name:cli,entrypoint.type:script");
t.cfg
.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
url: Some(format!("file://{}", output_filename.to_str().unwrap())),
..Default::default()
})
.unwrap();
let test_instance = super::CrashInfo::test_instance(seed);
t.upload_crash_info(&test_instance).await.unwrap();
let payload: serde_json::value::Value =
serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap()).unwrap();
assert_eq!(payload["api_version"], "v2");
assert_eq!(payload["application"]["language_name"], "native");
assert_eq!(payload["application"]["service_name"], "foo");
assert_eq!(payload["application"]["service_version"], "bar");
assert_eq!(
payload["application"]["process_tags"],
"entrypoint.name:cli,entrypoint.type:script"
);
assert_eq!(payload["request_type"], "logs");
assert_eq!(payload["tracer_time"], 1568898000);
assert_eq!(payload["origin"], "Crashtracker");
assert_eq!(payload["payload"]["logs"].as_array().unwrap().len(), 1);
let tags = payload["payload"]["logs"][0]["tags"]
.as_str()
.unwrap()
.split(',')
.collect::<HashSet<_>>();
assert_eq!(
HashSet::from_iter([
"collecting_sample:1",
"data_schema_version:1.8",
"incomplete:true",
"is_crash:true",
"not_profiling:0",
"si_addr:0x0000000000001234",
"si_code_human_readable:SEGV_BNDERR",
"si_code:1",
"si_signo_human_readable:SIGSEGV",
"si_signo:11",
"uuid:1d6b97cb-968c-40c9-af6e-e4b4d71e8781",
&format!("runtime_platform:{}", super::super::TARGET_TRIPLE),
]),
tags
);
assert_eq!(payload["payload"]["logs"][0]["is_sensitive"], true);
assert_eq!(payload["payload"]["logs"][0]["level"], "ERROR");
let body: CrashInfo =
serde_json::from_str(payload["payload"]["logs"][0]["message"].as_str().unwrap())?;
assert_eq!(body, test_instance);
assert_eq!(payload["payload"]["logs"][0]["is_crash"], true);
Ok(())
}
#[tokio::test]
#[cfg_attr(miri, ignore)]
async fn test_crash_ping_content() -> anyhow::Result<()> {
let tmp = tempfile::tempdir().unwrap();
let output_filename = {
let mut p = tmp.keep();
p.push("crash_ping_info");
p
};
let seed = 1;
let mut t = new_test_uploader(seed);
t.cfg
.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
url: Some(format!("file://{}", output_filename.to_str().unwrap())),
..Default::default()
})
.unwrap();
let sig_info = crate::SigInfo::test_instance(42);
let metadata = Metadata::test_instance(1);
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
crash_info_builder.with_metadata(metadata.clone()).unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let crash_ping = crash_info_builder.build_crash_ping().unwrap();
t.upload_crash_ping(&crash_ping).await.unwrap();
let payload: serde_json::value::Value =
serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap()).unwrap();
assert_eq!(payload["api_version"], "v2");
assert_eq!(payload["application"]["language_name"], "native");
assert_eq!(payload["application"]["service_name"], "foo");
assert_eq!(payload["application"]["service_version"], "bar");
assert_eq!(payload["request_type"], "logs");
assert_eq!(payload["origin"], "Crashtracker");
assert_eq!(payload["payload"]["logs"].as_array().unwrap().len(), 1);
let log_entry = &payload["payload"]["logs"][0];
assert_eq!(log_entry["is_sensitive"], false);
assert_eq!(log_entry["level"], "DEBUG");
let message_json: serde_json::Value =
serde_json::from_str(log_entry["message"].as_str().unwrap())?;
assert_eq!(message_json["siginfo"], serde_json::to_value(&sig_info)?);
assert!(message_json["crash_uuid"].is_string());
assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
assert_eq!(message_json["version"], "1.0");
assert_eq!(message_json["kind"], "UnixSignal");
let metadata_in_message = &message_json["metadata"];
assert!(
metadata_in_message.is_object(),
"metadata should be an object"
);
let expected_metadata = serde_json::to_value(Metadata::test_instance(1))?;
assert_eq!(
metadata_in_message, &expected_metadata,
"metadata field should match expected structure"
);
let tags = log_entry["tags"].as_str().unwrap();
let uuid_str = message_json["crash_uuid"].as_str().unwrap();
assert!(tags.contains(&format!("uuid:{uuid_str}")));
assert!(tags.contains("is_crash_ping:true"));
assert!(tags.contains("service:foo"));
assert!(tags.contains("language_name:native"));
assert!(tags.contains("language_version:"));
assert!(tags.contains("tracer_version:"));
Ok(())
}
#[tokio::test]
#[cfg_attr(miri, ignore)]
async fn test_crash_ping_with_different_config() -> anyhow::Result<()> {
let tmp = tempfile::tempdir().unwrap();
let output_filename = {
let mut p = tmp.keep();
p.push("enhanced_crash_ping_info");
p
};
let seed = 1;
let mut t = new_test_uploader(seed);
t.cfg
.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
url: Some(format!("file://{}", output_filename.to_str().unwrap())),
..Default::default()
})
.unwrap();
let sig_info = crate::SigInfo::test_instance(123);
let metadata = Metadata::test_instance(1);
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
crash_info_builder.with_metadata(metadata.clone()).unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let crash_ping = crash_info_builder.build_crash_ping().unwrap();
t.upload_crash_ping(&crash_ping).await.unwrap();
let payload: serde_json::value::Value =
serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap()).unwrap();
assert_eq!(payload["api_version"], "v2");
assert_eq!(payload["application"]["language_name"], "native");
assert_eq!(payload["application"]["service_name"], "foo");
assert_eq!(payload["application"]["service_version"], "bar");
assert_eq!(payload["request_type"], "logs");
assert_eq!(payload["origin"], "Crashtracker");
assert_eq!(payload["payload"]["logs"].as_array().unwrap().len(), 1);
let log_entry = &payload["payload"]["logs"][0];
assert_eq!(log_entry["is_crash"], false);
assert_eq!(log_entry["is_sensitive"], false);
assert_eq!(log_entry["level"], "DEBUG");
let message_json: serde_json::Value =
serde_json::from_str(log_entry["message"].as_str().unwrap())?;
assert!(message_json["crash_uuid"].is_string());
assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
assert_eq!(
message_json["message"],
format!(
"Crashtracker crash ping: crash processing started - Process terminated with {:?} ({:?})",
sig_info.si_code_human_readable, sig_info.si_signo_human_readable
)
);
let metadata_in_message = &message_json["metadata"];
assert!(
metadata_in_message.is_object(),
"metadata should be an object"
);
let expected_metadata = serde_json::to_value(Metadata::test_instance(1))?;
assert_eq!(
metadata_in_message, &expected_metadata,
"metadata field should match expected structure"
);
let siginfo_in_message = &message_json["siginfo"];
let expected_siginfo = serde_json::to_value(&sig_info)?;
assert_eq!(
siginfo_in_message, &expected_siginfo,
"siginfo field should match expected structure"
);
assert_eq!(message_json["version"], "1.0");
assert_eq!(message_json["kind"], "UnixSignal");
let tags = log_entry["tags"].as_str().unwrap();
let uuid_str = message_json["crash_uuid"].as_str().unwrap();
assert!(tags.contains(&format!("uuid:{uuid_str}")));
assert!(tags.contains("is_crash_ping:true"));
assert!(tags.contains("service:foo"));
assert!(tags.contains("language_name:native"));
assert!(tags.contains("language_version:"));
assert!(tags.contains("tracer_version:"));
Ok(())
}
#[tokio::test]
#[cfg_attr(miri, ignore)]
async fn test_crash_ping_builder_basic() -> anyhow::Result<()> {
let tmp = tempfile::tempdir().unwrap();
let output_filename = {
let mut p = tmp.keep();
p.push("crash_ping_builder_test");
p
};
let sig_info = crate::SigInfo::test_instance(42);
let metadata = Metadata::test_instance(1);
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
crash_info_builder.with_metadata(metadata.clone()).unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let crash_ping = crash_info_builder.build_crash_ping()?;
let endpoint = Some(Endpoint::from_slice(&format!(
"file://{}",
output_filename.to_str().unwrap()
)));
assert!(!crash_ping.crash_uuid().is_empty());
assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
assert!(crash_ping.message().contains("crash processing started"));
assert_eq!(crash_ping.metadata(), &metadata);
let mut uploader = TelemetryCrashUploader::new(&metadata, &endpoint)?;
uploader
.cfg
.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
url: Some(format!(
"file://{}.telemetry",
output_filename.to_str().unwrap()
)),
..Default::default()
})
.unwrap();
uploader.upload_crash_ping(&crash_ping).await?;
let telemetry_filename = format!("{}.telemetry", output_filename.to_str().unwrap());
let payload: serde_json::value::Value =
serde_json::de::from_str(&std::fs::read_to_string(&telemetry_filename)?)?;
assert_eq!(payload["api_version"], "v2");
assert_eq!(payload["request_type"], "logs");
assert_eq!(payload["origin"], "Crashtracker");
let log_entry = &payload["payload"]["logs"][0];
assert_eq!(log_entry["level"], "DEBUG");
assert_eq!(log_entry["is_sensitive"], false);
assert_eq!(log_entry["is_crash"], false);
let message_json: serde_json::Value =
serde_json::from_str(log_entry["message"].as_str().unwrap())?;
assert!(message_json["crash_uuid"].is_string());
assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
assert_eq!(message_json["version"], "1.0");
assert_eq!(message_json["kind"], "UnixSignal");
Ok(())
}
#[test]
#[cfg_attr(miri, ignore)]
fn test_crash_ping_builder_validation() {
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder
.with_metadata(Metadata::test_instance(1))
.unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let result = crash_info_builder.build_crash_ping();
assert!(result.is_ok());
let crash_ping = result.unwrap();
assert!(crash_ping.siginfo().is_none());
assert!(crash_ping
.message()
.contains("Crashtracker crash ping: crash processing started - Process terminated"));
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder
.with_sig_info(crate::SigInfo::test_instance(1))
.unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let result = crash_info_builder.build_crash_ping();
assert!(result.is_err());
assert!(result
.unwrap_err()
.to_string()
.contains("metadata is required"));
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder
.with_sig_info(crate::SigInfo::test_instance(1))
.unwrap();
crash_info_builder
.with_metadata(Metadata::test_instance(1))
.unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let result = crash_info_builder.build_crash_ping();
assert!(result.is_ok());
let crash_ping = result.unwrap();
assert!(crash_ping.siginfo().is_some());
}
#[test]
#[cfg_attr(miri, ignore)]
fn test_crash_ping_all_fields_present() {
let sig_info = crate::SigInfo::test_instance(99);
let metadata = Metadata::test_instance(2);
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
crash_info_builder.with_metadata(metadata.clone()).unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let crash_ping = crash_info_builder.build_crash_ping().unwrap();
assert!(!crash_ping.crash_uuid().is_empty());
assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
assert!(crash_ping.message().contains("crash processing started"));
assert_eq!(crash_ping.metadata(), &metadata);
assert_eq!(crash_ping.siginfo(), Some(&sig_info));
}
#[test]
#[cfg_attr(miri, ignore)]
fn test_crash_ping_with_message_generated_from_sig_info() {
let sig_info = crate::SigInfo::test_instance(99);
let metadata = Metadata::test_instance(2);
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
crash_info_builder.with_metadata(metadata.clone()).unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let crash_ping = crash_info_builder.build_crash_ping().unwrap();
assert!(!crash_ping.crash_uuid().is_empty());
assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
assert_eq!(crash_ping.message(), format!(
"Crashtracker crash ping: crash processing started - Process terminated with {:?} ({:?})",
sig_info.si_code_human_readable, sig_info.si_signo_human_readable
));
assert_eq!(crash_ping.metadata(), &metadata);
assert_eq!(crash_ping.siginfo(), Some(&sig_info));
}
#[test]
#[cfg_attr(miri, ignore)]
fn test_crash_ping_with_custom_message() {
let sig_info = crate::SigInfo::test_instance(99);
let metadata = Metadata::test_instance(2);
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
crash_info_builder.with_metadata(metadata.clone()).unwrap();
crash_info_builder
.with_message("my process panicked".to_string())
.unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let crash_ping = crash_info_builder.build_crash_ping().unwrap();
assert!(!crash_ping.crash_uuid().is_empty());
assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
assert!(crash_ping
.message()
.contains("crash processing started - my process panicked"));
assert_eq!(crash_ping.metadata(), &metadata);
assert_eq!(crash_ping.siginfo(), Some(&sig_info));
}
#[tokio::test]
#[cfg_attr(miri, ignore)]
async fn test_crash_ping_telemetry_upload_all_fields() -> anyhow::Result<()> {
let tmp = tempfile::tempdir().unwrap();
let output_filename = {
let mut p = tmp.keep();
p.push("crash_ping_all_fields_upload");
p
};
let seed = 3;
let mut uploader = new_test_uploader(seed);
uploader
.cfg
.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
url: Some(format!("file://{}", output_filename.to_str().unwrap())),
..Default::default()
})
.unwrap();
let sig_info = crate::SigInfo::test_instance(150);
let metadata = Metadata::test_instance(3);
let mut crash_info_builder = CrashInfoBuilder::new();
crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
crash_info_builder.with_metadata(metadata.clone()).unwrap();
crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
let crash_ping = crash_info_builder.build_crash_ping().unwrap();
uploader.upload_crash_ping(&crash_ping).await?;
let payload: serde_json::value::Value =
serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap())?;
assert_eq!(payload["api_version"], "v2");
assert_eq!(payload["request_type"], "logs");
assert_eq!(payload["origin"], "Crashtracker");
let log_entry = &payload["payload"]["logs"][0];
assert_eq!(log_entry["level"], "DEBUG");
assert_eq!(log_entry["is_sensitive"], false);
assert_eq!(log_entry["is_crash"], false);
let message_json: serde_json::Value =
serde_json::from_str(log_entry["message"].as_str().unwrap())?;
assert!(message_json["crash_uuid"].is_string());
assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
assert_eq!(message_json["version"], "1.0");
assert_eq!(message_json["kind"], "UnixSignal");
let uploaded_siginfo = &message_json["siginfo"];
assert_eq!(uploaded_siginfo["si_signo"], sig_info.si_signo);
assert_eq!(uploaded_siginfo["si_code"], sig_info.si_code);
assert_eq!(
uploaded_siginfo["si_code_human_readable"],
serde_json::to_value(&sig_info.si_code_human_readable)?
);
assert_eq!(
uploaded_siginfo["si_signo_human_readable"],
serde_json::to_value(&sig_info.si_signo_human_readable)?
);
let uploaded_metadata = &message_json["metadata"];
assert!(uploaded_metadata.is_object());
let expected_metadata_json = serde_json::to_value(&metadata)?;
assert_eq!(uploaded_metadata, &expected_metadata_json);
assert!(message_json["message"].is_string());
assert!(message_json["message"]
.as_str()
.unwrap()
.contains("crash processing started"));
Ok(())
}
#[tokio::test]
#[cfg_attr(miri, ignore)]
async fn test_general_log_upload() -> anyhow::Result<()> {
let tmp = tempfile::tempdir().unwrap();
let output_filename = {
let mut p = tmp.keep();
p.push("general_log_upload");
p
};
let mut uploader = new_test_uploader(7);
uploader
.cfg
.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
url: Some(format!("file://{}", output_filename.to_str().unwrap())),
..Default::default()
})?;
uploader
.upload_general_log(
"hello general log".to_string(),
"service:foo,env:bar,crash_uuid:1234567890".to_string(),
LogLevel::Warn,
)
.await?;
let payload: serde_json::value::Value =
serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap())?;
println!("payload: {:?}", payload.to_string());
assert_eq!(payload["api_version"], "v2");
assert_eq!(payload["request_type"], "logs");
assert_eq!(payload["origin"], "Crashtracker");
let log_entry = &payload["payload"]["logs"][0];
assert_eq!(log_entry["level"], "WARN");
assert_eq!(log_entry["is_sensitive"], false);
assert_eq!(log_entry["is_crash"], false);
assert_eq!(log_entry["message"], "hello general log");
let tags = log_entry["tags"].as_str().unwrap();
assert!(tags.contains("service:foo"));
assert!(tags.contains("env:bar"));
assert!(tags.contains("crash_uuid:1234567890"));
Ok(())
}
}