1use core::fmt::Write;
4use std::time::SystemTime;
5
6use crate::{ErrorKind, SigInfo};
7
8use super::{CrashInfo, Metadata, TARGET_TRIPLE};
9use anyhow::Context;
10use chrono::{DateTime, Utc};
11use libdd_capabilities::HttpClientCapability;
12use libdd_capabilities_impl::NativeCapabilities;
13use libdd_common::Endpoint;
14use libdd_telemetry::{
15 build_host,
16 data::{self, Application, LogLevel},
17 worker::http_client::request_builder,
18};
19use serde::Serialize;
20use uuid::Uuid;
21
22#[derive(Debug)]
23struct TelemetryMetadata {
24 application: libdd_telemetry::data::Application,
25 host: libdd_telemetry::data::Host,
26 runtime_id: String,
27}
28
29pub struct CrashPingBuilder {
30 crash_uuid: Uuid,
31 custom_message: Option<String>,
32 kind: Option<ErrorKind>,
33 metadata: Option<Metadata>,
34 sig_info: Option<SigInfo>,
35}
36
37impl CrashPingBuilder {
38 pub fn new(crash_uuid: Uuid) -> Self {
42 Self {
43 crash_uuid,
44 custom_message: None,
45 kind: None,
46 metadata: None,
47 sig_info: None,
48 }
49 }
50
51 pub fn with_sig_info(mut self, sig_info: SigInfo) -> Self {
52 self.sig_info = Some(sig_info);
53 self
54 }
55
56 pub fn with_custom_message(mut self, message: String) -> Self {
57 self.custom_message = Some(message);
58 self
59 }
60
61 pub fn with_kind(mut self, kind: ErrorKind) -> Self {
62 self.kind = Some(kind);
63 self
64 }
65
66 pub fn with_metadata(mut self, metadata: Metadata) -> Self {
67 self.metadata = Some(metadata);
68 self
69 }
70
71 pub fn build(self) -> anyhow::Result<CrashPing> {
72 let crash_uuid = self.crash_uuid;
73 let sig_info = self.sig_info;
74 let metadata = self.metadata.context("metadata is required")?;
75 let kind = self.kind.context("kind is required")?;
76
77 let message = if let Some(custom_message) = self.custom_message {
78 format!("Crashtracker crash ping: crash processing started - {custom_message}")
79 } else if let Some(ref sig_info) = sig_info {
80 format!(
81 "Crashtracker crash ping: crash processing started - Process terminated with {:?} ({:?})",
82 sig_info.si_code_human_readable, sig_info.si_signo_human_readable
83 )
84 } else {
85 format!("Crashtracker crash ping: crash processing started - Process terminated due to {:?}", kind)
86 };
87
88 Ok(CrashPing {
89 crash_uuid: crash_uuid.to_string(),
90 message,
91 kind,
92 metadata,
93 siginfo: sig_info,
94 version: CrashPing::current_schema_version(),
95 })
96 }
97}
98
99#[derive(Debug, Serialize)]
100pub struct CrashPing {
101 crash_uuid: String,
102 kind: ErrorKind,
103 message: String,
104 #[serde(skip_serializing_if = "Option::is_none")]
105 siginfo: Option<SigInfo>,
106 version: String,
107 metadata: Metadata,
108}
109
110impl CrashPing {
111 pub fn crash_uuid(&self) -> &str {
112 &self.crash_uuid
113 }
114
115 pub fn message(&self) -> &str {
116 &self.message
117 }
118
119 pub fn metadata(&self) -> &Metadata {
120 &self.metadata
121 }
122
123 pub fn siginfo(&self) -> Option<&SigInfo> {
124 self.siginfo.as_ref()
125 }
126
127 pub fn kind(&self) -> ErrorKind {
128 self.kind.clone()
129 }
130
131 pub fn upload_to_endpoint(&self, endpoint: &Option<Endpoint>) -> anyhow::Result<()> {
132 let rt = tokio::runtime::Builder::new_current_thread()
133 .enable_all()
134 .build()?;
135
136 rt.block_on(async { self.upload_to_endpoint_async(endpoint).await })
137 }
138
139 pub async fn upload_to_endpoint_async(
143 &self,
144 endpoint: &Option<Endpoint>,
145 ) -> anyhow::Result<()> {
146 let telemetry_uploader = crate::TelemetryCrashUploader::new(self.metadata(), endpoint)?;
147 let errors_intake_uploader = crate::ErrorsIntakeUploader::new(endpoint)?;
148 let telemetry_future = telemetry_uploader.upload_crash_ping(self);
149
150 if errors_intake_uploader.is_enabled() {
151 let errors_intake_future = errors_intake_uploader.upload_crash_ping(self);
152 let (_telemetry_result, _errors_intake_result) =
153 tokio::join!(telemetry_future, errors_intake_future);
154 } else {
155 let _telemetry_result = telemetry_future.await;
156 }
157 Ok(())
158 }
159
160 fn current_schema_version() -> String {
161 "1.0".to_string()
162 }
163}
164
165macro_rules! parse_tags {
166 ( $tag_iterator:expr,
167 $($tag_name:literal => $var:ident),* $(,)?) => {
168 $(
169 let mut $var: Option<&str> = None;
170 )*
171 for tag in $tag_iterator {
172 let Some((name, value)) = tag.split_once(':') else {
173 continue;
174 };
175 match name {
176 $($tag_name => {$var = Some(value);}, )*
177 _ => {},
178 }
179 }
180
181 };
182}
183
184pub struct TelemetryCrashUploader {
185 metadata: TelemetryMetadata,
186 cfg: libdd_telemetry::config::Config,
187}
188
189impl TelemetryCrashUploader {
190 pub fn new(
191 crashtracker_metadata: &Metadata,
192 endpoint: &Option<Endpoint>,
193 ) -> anyhow::Result<Self> {
194 let mut cfg = libdd_telemetry::config::Config::from_env();
195 if let Some(endpoint) = endpoint {
196 if endpoint.url.scheme_str() == Some("file") {
204 let path = libdd_common::decode_uri_path_in_authority(&endpoint.url)
207 .context("file path is not valid")?;
208 let _ = cfg.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
209 url: Some(format!("file://{}.telemetry", path.display())),
210 ..Default::default()
211 });
212 } else {
213 let _ = cfg.set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
214 api_key: endpoint.api_key.as_deref().map(str::to_owned),
215 test_token: endpoint.test_token.as_deref().map(str::to_owned),
216 timeout_ms: endpoint.timeout_ms,
217 use_system_resolver: endpoint.use_system_resolver,
218 ..Default::default()
219 });
220 let _ = cfg.set_endpoint_uri(endpoint.url.clone());
221 }
222 }
223
224 parse_tags!(
225 crashtracker_metadata.tags.iter(),
226 "env" => env,
227 "language" => language_name,
228 "library_version" => library_version,
229 "profiler_version" => profiler_version,
230 "runtime_version" => language_version,
231 "runtime-id" => runtime_id,
232 "service_version" => service_version,
233 "service" => service_name,
234 "process_tags" => process_tags,
235 );
236
237 let application = Application {
238 service_name: service_name.unwrap_or("unknown").to_owned(),
239 language_name: language_name.unwrap_or("unknown").to_owned(),
240 language_version: language_version.unwrap_or("unknown").to_owned(),
241 tracer_version: library_version
242 .or(profiler_version)
243 .unwrap_or("unknown")
244 .to_owned(),
245 env: env.map(ToOwned::to_owned),
246 service_version: service_version.map(ToOwned::to_owned),
247 process_tags: process_tags.map(ToOwned::to_owned),
248 ..Default::default()
249 };
250
251 let host = build_host();
252
253 let s = Self {
254 metadata: TelemetryMetadata {
255 host,
256 application,
257 runtime_id: runtime_id.unwrap_or("unknown").to_owned(),
258 },
259 cfg,
260 };
261 Ok(s)
262 }
263
264 pub async fn upload_general_log(
266 &self,
267 message: String,
268 tags: String,
269 level: LogLevel,
270 ) -> anyhow::Result<()> {
271 let tracer_time = SystemTime::now()
272 .duration_since(SystemTime::UNIX_EPOCH)
273 .map(|d| d.as_secs())
274 .unwrap_or(0);
275
276 self.send_log_payload(message, tags, tracer_time, level, false, false)
277 .await
278 }
279
280 pub async fn upload_crash_ping(&self, crash_ping: &CrashPing) -> anyhow::Result<()> {
282 let tags = self.build_crash_ping_tags(crash_ping.crash_uuid(), crash_ping.siginfo());
283 let tracer_time = SystemTime::now()
284 .duration_since(SystemTime::UNIX_EPOCH)
285 .map(|d| d.as_secs())
286 .unwrap_or(0);
287 let message = serde_json::to_string(crash_ping)?;
288
289 self.send_log_payload(
290 message,
291 tags,
292 tracer_time,
293 LogLevel::Debug,
294 false, false, )
297 .await
298 }
299
300 pub async fn upload_crash_info(&self, crash_info: &CrashInfo) -> anyhow::Result<()> {
302 let message = serde_json::to_string(crash_info)?;
303 let tags = extract_crash_info_tags(crash_info).unwrap_or_default();
304 let tracer_time = crash_info.timestamp.parse::<DateTime<Utc>>().map_or_else(
305 |_| {
306 SystemTime::now()
307 .duration_since(SystemTime::UNIX_EPOCH)
308 .map(|d| d.as_secs())
309 .unwrap_or(0)
310 },
311 |ts| ts.timestamp() as u64,
312 );
313
314 self.send_log_payload(
315 message,
316 tags,
317 tracer_time,
318 LogLevel::Error,
319 true, true, )
322 .await
323 }
324
325 async fn send_log_payload(
328 &self,
329 message: String,
330 tags: String,
331 tracer_time: u64,
332 level: LogLevel,
333 is_sensitive: bool,
334 is_crash: bool,
335 ) -> anyhow::Result<()> {
336 let payload = data::Telemetry {
337 tracer_time,
338 api_version: libdd_telemetry::data::ApiVersion::V2,
339 runtime_id: &self.metadata.runtime_id,
340 seq_id: 1,
341 application: &self.metadata.application,
342 host: &self.metadata.host,
343 payload: &data::Payload::Logs(data::Logs {
344 logs: vec![data::Log {
345 message,
346 level,
347 stack_trace: None,
348 tags,
349 is_sensitive,
350 count: 1,
351 is_crash,
352 }],
353 }),
354 origin: Some("Crashtracker"),
355 };
356
357 self.send_telemetry_payload(&payload).await
358 }
359
360 async fn send_telemetry_payload(&self, payload: &data::Telemetry<'_>) -> anyhow::Result<()> {
362 let client = NativeCapabilities::new_client();
363 let req = request_builder(&self.cfg)?
364 .method(http::Method::POST)
365 .header(
366 http::header::CONTENT_TYPE,
367 libdd_common::header::APPLICATION_JSON,
368 )
369 .header(
370 libdd_telemetry::worker::http_client::header::API_VERSION,
371 libdd_telemetry::data::ApiVersion::V2.to_str(),
372 )
373 .header(
374 libdd_telemetry::worker::http_client::header::REQUEST_TYPE,
375 "logs",
376 )
377 .body(libdd_capabilities::Bytes::from(serde_json::to_vec(
378 &payload,
379 )?))?;
380
381 let timeout = core::time::Duration::from_millis({
382 if let Some(endp) = self.cfg.endpoint() {
383 endp.timeout_ms
384 } else {
385 Endpoint::DEFAULT_TIMEOUT
386 }
387 });
388 tokio::time::timeout(timeout, client.request(req))
389 .await
390 .map_err(|_| anyhow::anyhow!("Telemetry crash report timed out"))??;
391
392 Ok(())
393 }
394
395 fn build_crash_ping_tags(&self, crash_uuid: &str, sig_info: Option<&SigInfo>) -> String {
396 let metadata = &self.metadata;
397 let mut tags = format!(
398 "uuid:{},is_crash_ping:true,service:{},language_name:{},language_version:{},tracer_version:{}",
399 crash_uuid,
400 metadata.application.service_name,
401 metadata.application.language_name,
402 metadata.application.language_version,
403 metadata.application.tracer_version
404 );
405
406 if let Some(sig_info) = sig_info {
407 tags.push_str(&format!(
408 ",si_code_human_readable:{:?},si_signo:{},si_signo_human_readable:{:?}",
409 sig_info.si_code_human_readable,
410 sig_info.si_signo,
411 sig_info.si_signo_human_readable
412 ));
413 }
414
415 write!(tags, ",runtime_platform:{TARGET_TRIPLE}").ok();
416 self.append_optional_tags(&mut tags);
417 tags
418 }
419
420 fn append_optional_tags(&self, tags: &mut String) {
421 let metadata = &self.metadata;
422 if let Some(env) = &metadata.application.env {
423 tags.push_str(&format!(",env:{env}"));
424 }
425 if let Some(runtime_name) = &metadata.application.runtime_name {
426 tags.push_str(&format!(",runtime_name:{runtime_name}"));
427 }
428 if let Some(runtime_version) = &metadata.application.runtime_version {
429 tags.push_str(&format!(",runtime_version:{runtime_version}"));
430 }
431 }
432}
433
434fn extract_crash_info_tags(crash_info: &CrashInfo) -> anyhow::Result<String> {
435 let mut tags = String::new();
436 write!(
437 &mut tags,
438 "data_schema_version:{}",
439 crash_info.data_schema_version
440 )?;
441 if let Some(fingerprint) = &crash_info.fingerprint {
442 write!(&mut tags, ",fingerprint:{fingerprint}")?;
443 }
444 write!(&mut tags, ",incomplete:{}", crash_info.incomplete)?;
445 write!(&mut tags, ",is_crash:{}", crash_info.error.is_crash)?;
446 write!(&mut tags, ",uuid:{}", crash_info.uuid)?;
447 for (counter, value) in &crash_info.counters {
448 write!(&mut tags, ",{counter}:{value}")?;
449 }
450
451 if let Some(siginfo) = &crash_info.sig_info {
452 if let Some(si_addr) = &siginfo.si_addr {
453 write!(&mut tags, ",si_addr:{si_addr}")?;
454 }
455 write!(&mut tags, ",si_code:{}", siginfo.si_code)?;
456 write!(
457 &mut tags,
458 ",si_code_human_readable:{:?}",
459 siginfo.si_code_human_readable
460 )?;
461 write!(&mut tags, ",si_signo:{}", siginfo.si_signo)?;
462 write!(
463 &mut tags,
464 ",si_signo_human_readable:{:?}",
465 siginfo.si_signo_human_readable
466 )?;
467 }
468 write!(&mut tags, ",runtime_platform:{TARGET_TRIPLE}")?;
469 Ok(tags)
470}
471
472#[cfg(test)]
473mod tests {
474 use super::TelemetryCrashUploader;
475 use crate::{
476 crash_info::{test_utils::TestInstance, CrashInfo, CrashInfoBuilder, Metadata},
477 ErrorKind,
478 };
479 use libdd_common::Endpoint;
480 use libdd_telemetry::data::LogLevel;
481 use std::{collections::HashSet, fs};
482 use uuid::Uuid;
483
484 fn new_test_uploader(seed: u64) -> TelemetryCrashUploader {
485 TelemetryCrashUploader::new(
486 &Metadata::test_instance(seed),
487 &Some(Endpoint::from_slice("http://localhost:8126")),
488 )
489 .unwrap()
490 }
491
492 fn new_test_uploader_with_process_tags(
493 seed: u64,
494 process_tags: &str,
495 ) -> TelemetryCrashUploader {
496 let mut metadata = Metadata::test_instance(seed);
497 metadata.tags.push(format!("process_tags:{process_tags}"));
498 TelemetryCrashUploader::new(
499 &metadata,
500 &Some(Endpoint::from_slice("http://localhost:8126")),
501 )
502 .unwrap()
503 }
504
505 #[test]
506 #[cfg_attr(miri, ignore)]
507 fn test_profiler_config_extraction() {
508 let t = new_test_uploader(1);
509
510 let metadata = t.metadata;
511 assert_eq!(metadata.application.service_name, "foo");
512 assert_eq!(metadata.application.service_version.as_deref(), Some("bar"));
513 assert_eq!(metadata.application.language_name, "native");
514 assert_eq!(metadata.application.process_tags, None);
515 assert_eq!(metadata.runtime_id, "xyz");
516 let cfg = t.cfg;
517 assert_eq!(
518 cfg.endpoint().unwrap().url.to_string(),
519 "http://localhost:8126/telemetry/proxy/api/v2/apmtelemetry"
520 );
521 }
522
523 #[tokio::test]
524 #[cfg_attr(miri, ignore)]
525 async fn test_crash_request_content() -> anyhow::Result<()> {
526 let tmp = tempfile::tempdir().unwrap();
528 let output_filename = {
529 let mut p = tmp.keep();
530 p.push("crash_info");
531 p
532 };
533 let seed = 1;
534 let mut t =
535 new_test_uploader_with_process_tags(seed, "entrypoint.name:cli,entrypoint.type:script");
536
537 t.cfg
538 .set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
539 url: Some(format!("file://{}", output_filename.to_str().unwrap())),
540 ..Default::default()
541 })
542 .unwrap();
543 let test_instance = super::CrashInfo::test_instance(seed);
544
545 t.upload_crash_info(&test_instance).await.unwrap();
546
547 let payload: serde_json::value::Value =
548 serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap()).unwrap();
549 assert_eq!(payload["api_version"], "v2");
550 assert_eq!(payload["application"]["language_name"], "native");
551 assert_eq!(payload["application"]["service_name"], "foo");
552 assert_eq!(payload["application"]["service_version"], "bar");
553 assert_eq!(
554 payload["application"]["process_tags"],
555 "entrypoint.name:cli,entrypoint.type:script"
556 );
557 assert_eq!(payload["request_type"], "logs");
558 assert_eq!(payload["tracer_time"], 1568898000);
559 assert_eq!(payload["origin"], "Crashtracker");
560
561 assert_eq!(payload["payload"]["logs"].as_array().unwrap().len(), 1);
562 let tags = payload["payload"]["logs"][0]["tags"]
563 .as_str()
564 .unwrap()
565 .split(',')
566 .collect::<HashSet<_>>();
567 assert_eq!(
568 HashSet::from_iter([
569 "collecting_sample:1",
570 "data_schema_version:1.8",
571 "incomplete:true",
572 "is_crash:true",
573 "not_profiling:0",
574 "si_addr:0x0000000000001234",
575 "si_code_human_readable:SEGV_BNDERR",
576 "si_code:1",
577 "si_signo_human_readable:SIGSEGV",
578 "si_signo:11",
579 "uuid:1d6b97cb-968c-40c9-af6e-e4b4d71e8781",
580 &format!("runtime_platform:{}", super::super::TARGET_TRIPLE),
581 ]),
582 tags
583 );
584 assert_eq!(payload["payload"]["logs"][0]["is_sensitive"], true);
585 assert_eq!(payload["payload"]["logs"][0]["level"], "ERROR");
586 let body: CrashInfo =
587 serde_json::from_str(payload["payload"]["logs"][0]["message"].as_str().unwrap())?;
588 assert_eq!(body, test_instance);
589 assert_eq!(payload["payload"]["logs"][0]["is_crash"], true);
590 Ok(())
591 }
592
593 #[tokio::test]
594 #[cfg_attr(miri, ignore)]
595 async fn test_crash_ping_content() -> anyhow::Result<()> {
596 let tmp = tempfile::tempdir().unwrap();
598 let output_filename = {
599 let mut p = tmp.keep();
600 p.push("crash_ping_info");
601 p
602 };
603 let seed = 1;
604 let mut t = new_test_uploader(seed);
605
606 t.cfg
607 .set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
608 url: Some(format!("file://{}", output_filename.to_str().unwrap())),
609 ..Default::default()
610 })
611 .unwrap();
612
613 let sig_info = crate::SigInfo::test_instance(42);
614 let metadata = Metadata::test_instance(1);
615
616 let mut crash_info_builder = CrashInfoBuilder::new();
617 crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
618 crash_info_builder.with_metadata(metadata.clone()).unwrap();
619 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
620 let crash_ping = crash_info_builder.build_crash_ping().unwrap();
621 t.upload_crash_ping(&crash_ping).await.unwrap();
622
623 let payload: serde_json::value::Value =
624 serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap()).unwrap();
625 assert_eq!(payload["api_version"], "v2");
626 assert_eq!(payload["application"]["language_name"], "native");
627 assert_eq!(payload["application"]["service_name"], "foo");
628 assert_eq!(payload["application"]["service_version"], "bar");
629 assert_eq!(payload["request_type"], "logs");
630 assert_eq!(payload["origin"], "Crashtracker");
631
632 assert_eq!(payload["payload"]["logs"].as_array().unwrap().len(), 1);
633 let log_entry = &payload["payload"]["logs"][0];
634
635 assert_eq!(log_entry["is_sensitive"], false);
637 assert_eq!(log_entry["level"], "DEBUG");
638
639 let message_json: serde_json::Value =
641 serde_json::from_str(log_entry["message"].as_str().unwrap())?;
642 assert_eq!(message_json["siginfo"], serde_json::to_value(&sig_info)?);
643 assert!(message_json["crash_uuid"].is_string());
644 assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
645
646 assert_eq!(message_json["version"], "1.0");
647 assert_eq!(message_json["kind"], "UnixSignal");
648
649 let metadata_in_message = &message_json["metadata"];
650 assert!(
651 metadata_in_message.is_object(),
652 "metadata should be an object"
653 );
654 let expected_metadata = serde_json::to_value(Metadata::test_instance(1))?;
655 assert_eq!(
656 metadata_in_message, &expected_metadata,
657 "metadata field should match expected structure"
658 );
659
660 let tags = log_entry["tags"].as_str().unwrap();
661 let uuid_str = message_json["crash_uuid"].as_str().unwrap();
662 assert!(tags.contains(&format!("uuid:{uuid_str}")));
663 assert!(tags.contains("is_crash_ping:true"));
664 assert!(tags.contains("service:foo"));
665 assert!(tags.contains("language_name:native"));
666 assert!(tags.contains("language_version:"));
667 assert!(tags.contains("tracer_version:"));
668
669 Ok(())
670 }
671
672 #[tokio::test]
673 #[cfg_attr(miri, ignore)]
674 async fn test_crash_ping_with_different_config() -> anyhow::Result<()> {
675 let tmp = tempfile::tempdir().unwrap();
677 let output_filename = {
678 let mut p = tmp.keep();
679 p.push("enhanced_crash_ping_info");
680 p
681 };
682 let seed = 1;
683 let mut t = new_test_uploader(seed);
684
685 t.cfg
686 .set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
687 url: Some(format!("file://{}", output_filename.to_str().unwrap())),
688 ..Default::default()
689 })
690 .unwrap();
691
692 let sig_info = crate::SigInfo::test_instance(123);
693 let metadata = Metadata::test_instance(1);
694
695 let mut crash_info_builder = CrashInfoBuilder::new();
696 crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
697 crash_info_builder.with_metadata(metadata.clone()).unwrap();
698 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
699 let crash_ping = crash_info_builder.build_crash_ping().unwrap();
700 t.upload_crash_ping(&crash_ping).await.unwrap();
701
702 let payload: serde_json::value::Value =
703 serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap()).unwrap();
704 assert_eq!(payload["api_version"], "v2");
705 assert_eq!(payload["application"]["language_name"], "native");
706 assert_eq!(payload["application"]["service_name"], "foo");
707 assert_eq!(payload["application"]["service_version"], "bar");
708 assert_eq!(payload["request_type"], "logs");
709 assert_eq!(payload["origin"], "Crashtracker");
710
711 assert_eq!(payload["payload"]["logs"].as_array().unwrap().len(), 1);
712 let log_entry = &payload["payload"]["logs"][0];
713
714 assert_eq!(log_entry["is_crash"], false);
716 assert_eq!(log_entry["is_sensitive"], false);
717 assert_eq!(log_entry["level"], "DEBUG");
718
719 let message_json: serde_json::Value =
721 serde_json::from_str(log_entry["message"].as_str().unwrap())?;
722 assert!(message_json["crash_uuid"].is_string());
723 assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
724 assert_eq!(
725 message_json["message"],
726 format!(
727 "Crashtracker crash ping: crash processing started - Process terminated with {:?} ({:?})",
728 sig_info.si_code_human_readable, sig_info.si_signo_human_readable
729 )
730 );
731
732 let metadata_in_message = &message_json["metadata"];
733 assert!(
734 metadata_in_message.is_object(),
735 "metadata should be an object"
736 );
737 let expected_metadata = serde_json::to_value(Metadata::test_instance(1))?;
738 assert_eq!(
739 metadata_in_message, &expected_metadata,
740 "metadata field should match expected structure"
741 );
742
743 let siginfo_in_message = &message_json["siginfo"];
744 let expected_siginfo = serde_json::to_value(&sig_info)?;
745 assert_eq!(
746 siginfo_in_message, &expected_siginfo,
747 "siginfo field should match expected structure"
748 );
749
750 assert_eq!(message_json["version"], "1.0");
751 assert_eq!(message_json["kind"], "UnixSignal");
752
753 let tags = log_entry["tags"].as_str().unwrap();
754 let uuid_str = message_json["crash_uuid"].as_str().unwrap();
755 assert!(tags.contains(&format!("uuid:{uuid_str}")));
756 assert!(tags.contains("is_crash_ping:true"));
757 assert!(tags.contains("service:foo"));
758 assert!(tags.contains("language_name:native"));
759 assert!(tags.contains("language_version:"));
760 assert!(tags.contains("tracer_version:"));
761
762 Ok(())
763 }
764
765 #[tokio::test]
766 #[cfg_attr(miri, ignore)]
767 async fn test_crash_ping_builder_basic() -> anyhow::Result<()> {
768 let tmp = tempfile::tempdir().unwrap();
769 let output_filename = {
770 let mut p = tmp.keep();
771 p.push("crash_ping_builder_test");
772 p
773 };
774
775 let sig_info = crate::SigInfo::test_instance(42);
776 let metadata = Metadata::test_instance(1);
777
778 let mut crash_info_builder = CrashInfoBuilder::new();
780 crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
781 crash_info_builder.with_metadata(metadata.clone()).unwrap();
782 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
783 let crash_ping = crash_info_builder.build_crash_ping()?;
784
785 let endpoint = Some(Endpoint::from_slice(&format!(
786 "file://{}",
787 output_filename.to_str().unwrap()
788 )));
789
790 assert!(!crash_ping.crash_uuid().is_empty());
791 assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
792 assert!(crash_ping.message().contains("crash processing started"));
793 assert_eq!(crash_ping.metadata(), &metadata);
794
795 let mut uploader = TelemetryCrashUploader::new(&metadata, &endpoint)?;
797 uploader
798 .cfg
799 .set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
800 url: Some(format!(
801 "file://{}.telemetry",
802 output_filename.to_str().unwrap()
803 )),
804 ..Default::default()
805 })
806 .unwrap();
807
808 uploader.upload_crash_ping(&crash_ping).await?;
809
810 let telemetry_filename = format!("{}.telemetry", output_filename.to_str().unwrap());
812 let payload: serde_json::value::Value =
813 serde_json::de::from_str(&std::fs::read_to_string(&telemetry_filename)?)?;
814
815 assert_eq!(payload["api_version"], "v2");
816 assert_eq!(payload["request_type"], "logs");
817 assert_eq!(payload["origin"], "Crashtracker");
818
819 let log_entry = &payload["payload"]["logs"][0];
820 assert_eq!(log_entry["level"], "DEBUG");
821 assert_eq!(log_entry["is_sensitive"], false);
822 assert_eq!(log_entry["is_crash"], false);
823
824 let message_json: serde_json::Value =
825 serde_json::from_str(log_entry["message"].as_str().unwrap())?;
826 assert!(message_json["crash_uuid"].is_string());
827 assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
828 assert_eq!(message_json["version"], "1.0");
829 assert_eq!(message_json["kind"], "UnixSignal");
830
831 Ok(())
832 }
833
834 #[test]
835 #[cfg_attr(miri, ignore)]
836 fn test_crash_ping_builder_validation() {
837 let mut crash_info_builder = CrashInfoBuilder::new();
839 crash_info_builder
840 .with_metadata(Metadata::test_instance(1))
841 .unwrap();
842 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
843 let result = crash_info_builder.build_crash_ping();
844 assert!(result.is_ok());
845 let crash_ping = result.unwrap();
846 assert!(crash_ping.siginfo().is_none());
847 assert!(crash_ping
848 .message()
849 .contains("Crashtracker crash ping: crash processing started - Process terminated"));
850
851 let mut crash_info_builder = CrashInfoBuilder::new();
853 crash_info_builder
854 .with_sig_info(crate::SigInfo::test_instance(1))
855 .unwrap();
856 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
857 let result = crash_info_builder.build_crash_ping();
858 assert!(result.is_err());
859 assert!(result
860 .unwrap_err()
861 .to_string()
862 .contains("metadata is required"));
863
864 let mut crash_info_builder = CrashInfoBuilder::new();
866 crash_info_builder
867 .with_sig_info(crate::SigInfo::test_instance(1))
868 .unwrap();
869 crash_info_builder
870 .with_metadata(Metadata::test_instance(1))
871 .unwrap();
872 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
873 let result = crash_info_builder.build_crash_ping();
874 assert!(result.is_ok());
875 let crash_ping = result.unwrap();
876 assert!(crash_ping.siginfo().is_some());
877 }
878
879 #[test]
880 #[cfg_attr(miri, ignore)]
881 fn test_crash_ping_all_fields_present() {
882 let sig_info = crate::SigInfo::test_instance(99);
883 let metadata = Metadata::test_instance(2);
884
885 let mut crash_info_builder = CrashInfoBuilder::new();
887 crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
888 crash_info_builder.with_metadata(metadata.clone()).unwrap();
889 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
890 let crash_ping = crash_info_builder.build_crash_ping().unwrap();
891
892 assert!(!crash_ping.crash_uuid().is_empty());
893 assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
894 assert!(crash_ping.message().contains("crash processing started"));
895 assert_eq!(crash_ping.metadata(), &metadata);
896 assert_eq!(crash_ping.siginfo(), Some(&sig_info));
897 }
898
899 #[test]
900 #[cfg_attr(miri, ignore)]
901 fn test_crash_ping_with_message_generated_from_sig_info() {
902 let sig_info = crate::SigInfo::test_instance(99);
903 let metadata = Metadata::test_instance(2);
904
905 let mut crash_info_builder = CrashInfoBuilder::new();
907 crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
908 crash_info_builder.with_metadata(metadata.clone()).unwrap();
909 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
910 let crash_ping = crash_info_builder.build_crash_ping().unwrap();
911
912 assert!(!crash_ping.crash_uuid().is_empty());
913 assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
914 assert_eq!(crash_ping.message(), format!(
915 "Crashtracker crash ping: crash processing started - Process terminated with {:?} ({:?})",
916 sig_info.si_code_human_readable, sig_info.si_signo_human_readable
917 ));
918 assert_eq!(crash_ping.metadata(), &metadata);
919 assert_eq!(crash_ping.siginfo(), Some(&sig_info));
920 }
921
922 #[test]
923 #[cfg_attr(miri, ignore)]
924 fn test_crash_ping_with_custom_message() {
925 let sig_info = crate::SigInfo::test_instance(99);
926 let metadata = Metadata::test_instance(2);
927
928 let mut crash_info_builder = CrashInfoBuilder::new();
930 crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
931 crash_info_builder.with_metadata(metadata.clone()).unwrap();
932 crash_info_builder
933 .with_message("my process panicked".to_string())
934 .unwrap();
935 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
936 let crash_ping = crash_info_builder.build_crash_ping().unwrap();
937
938 assert!(!crash_ping.crash_uuid().is_empty());
939 assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
940 assert!(crash_ping
941 .message()
942 .contains("crash processing started - my process panicked"));
943 assert_eq!(crash_ping.metadata(), &metadata);
944 assert_eq!(crash_ping.siginfo(), Some(&sig_info));
945 }
946
947 #[tokio::test]
948 #[cfg_attr(miri, ignore)]
949 async fn test_crash_ping_telemetry_upload_all_fields() -> anyhow::Result<()> {
950 let tmp = tempfile::tempdir().unwrap();
952 let output_filename = {
953 let mut p = tmp.keep();
954 p.push("crash_ping_all_fields_upload");
955 p
956 };
957 let seed = 3;
958 let mut uploader = new_test_uploader(seed);
959
960 uploader
961 .cfg
962 .set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
963 url: Some(format!("file://{}", output_filename.to_str().unwrap())),
964 ..Default::default()
965 })
966 .unwrap();
967
968 let sig_info = crate::SigInfo::test_instance(150);
969 let metadata = Metadata::test_instance(3);
970
971 let mut crash_info_builder = CrashInfoBuilder::new();
973 crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
974 crash_info_builder.with_metadata(metadata.clone()).unwrap();
975 crash_info_builder.with_kind(ErrorKind::UnixSignal).unwrap();
976 let crash_ping = crash_info_builder.build_crash_ping().unwrap();
977
978 uploader.upload_crash_ping(&crash_ping).await?;
979
980 let payload: serde_json::value::Value =
981 serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap())?;
982
983 assert_eq!(payload["api_version"], "v2");
985 assert_eq!(payload["request_type"], "logs");
986 assert_eq!(payload["origin"], "Crashtracker");
987
988 let log_entry = &payload["payload"]["logs"][0];
989 assert_eq!(log_entry["level"], "DEBUG");
990 assert_eq!(log_entry["is_sensitive"], false);
991 assert_eq!(log_entry["is_crash"], false);
992
993 let message_json: serde_json::Value =
994 serde_json::from_str(log_entry["message"].as_str().unwrap())?;
995
996 assert!(message_json["crash_uuid"].is_string());
997 assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
998 assert_eq!(message_json["version"], "1.0");
999 assert_eq!(message_json["kind"], "UnixSignal");
1000
1001 let uploaded_siginfo = &message_json["siginfo"];
1002 assert_eq!(uploaded_siginfo["si_signo"], sig_info.si_signo);
1003 assert_eq!(uploaded_siginfo["si_code"], sig_info.si_code);
1004 assert_eq!(
1005 uploaded_siginfo["si_code_human_readable"],
1006 serde_json::to_value(&sig_info.si_code_human_readable)?
1007 );
1008 assert_eq!(
1009 uploaded_siginfo["si_signo_human_readable"],
1010 serde_json::to_value(&sig_info.si_signo_human_readable)?
1011 );
1012
1013 let uploaded_metadata = &message_json["metadata"];
1014 assert!(uploaded_metadata.is_object());
1015
1016 let expected_metadata_json = serde_json::to_value(&metadata)?;
1017 assert_eq!(uploaded_metadata, &expected_metadata_json);
1018
1019 assert!(message_json["message"].is_string());
1020 assert!(message_json["message"]
1021 .as_str()
1022 .unwrap()
1023 .contains("crash processing started"));
1024 Ok(())
1025 }
1026
1027 #[tokio::test]
1028 #[cfg_attr(miri, ignore)]
1029 async fn test_general_log_upload() -> anyhow::Result<()> {
1030 let tmp = tempfile::tempdir().unwrap();
1031 let output_filename = {
1032 let mut p = tmp.keep();
1033 p.push("general_log_upload");
1034 p
1035 };
1036
1037 let mut uploader = new_test_uploader(7);
1038 uploader
1039 .cfg
1040 .set_endpoint(libdd_telemetry::config::TelemetryEndpoint {
1041 url: Some(format!("file://{}", output_filename.to_str().unwrap())),
1042 ..Default::default()
1043 })?;
1044
1045 uploader
1046 .upload_general_log(
1047 "hello general log".to_string(),
1048 "service:foo,env:bar,crash_uuid:1234567890".to_string(),
1049 LogLevel::Warn,
1050 )
1051 .await?;
1052
1053 let payload: serde_json::value::Value =
1054 serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap())?;
1055 println!("payload: {:?}", payload.to_string());
1056
1057 assert_eq!(payload["api_version"], "v2");
1058 assert_eq!(payload["request_type"], "logs");
1059 assert_eq!(payload["origin"], "Crashtracker");
1060
1061 let log_entry = &payload["payload"]["logs"][0];
1062 assert_eq!(log_entry["level"], "WARN");
1063 assert_eq!(log_entry["is_sensitive"], false);
1064 assert_eq!(log_entry["is_crash"], false);
1065 assert_eq!(log_entry["message"], "hello general log");
1066 let tags = log_entry["tags"].as_str().unwrap();
1067 assert!(tags.contains("service:foo"));
1068 assert!(tags.contains("env:bar"));
1069 assert!(tags.contains("crash_uuid:1234567890"));
1070
1071 Ok(())
1072 }
1073}