Skip to main content

libdd_crashtracker/crash_info/
telemetry.rs

1// Copyright 2024-Present Datadog, Inc. https://www.datadoghq.com/
2// SPDX-License-Identifier: Apache-2.0
3use std::{fmt::Write, time::SystemTime};
4
5use crate::SigInfo;
6
7use super::{CrashInfo, Metadata};
8use anyhow::Context;
9use chrono::{DateTime, Utc};
10use libdd_common::Endpoint;
11use libdd_telemetry::{
12    build_host,
13    data::{self, Application, LogLevel},
14    worker::http_client::request_builder,
15};
16use serde::Serialize;
17use uuid::Uuid;
18
19#[derive(Debug)]
20struct TelemetryMetadata {
21    application: libdd_telemetry::data::Application,
22    host: libdd_telemetry::data::Host,
23    runtime_id: String,
24}
25
26pub struct CrashPingBuilder {
27    crash_uuid: Uuid,
28    custom_message: Option<String>,
29    metadata: Option<Metadata>,
30    sig_info: Option<SigInfo>,
31}
32
33impl CrashPingBuilder {
34    /// Crash pings should only be initalized and built by the CrashInfoBuilder
35    /// We require the crash uuid to be passed in because a CrashPing is always
36    /// associated with a specific crash.
37    pub fn new(crash_uuid: Uuid) -> Self {
38        Self {
39            crash_uuid,
40            custom_message: None,
41            metadata: None,
42            sig_info: None,
43        }
44    }
45
46    pub fn with_sig_info(mut self, sig_info: SigInfo) -> Self {
47        self.sig_info = Some(sig_info);
48        self
49    }
50
51    pub fn with_custom_message(mut self, message: String) -> Self {
52        self.custom_message = Some(message);
53        self
54    }
55
56    pub fn with_metadata(mut self, metadata: Metadata) -> Self {
57        self.metadata = Some(metadata);
58        self
59    }
60
61    pub fn build(self) -> anyhow::Result<CrashPing> {
62        let crash_uuid = self.crash_uuid;
63        let sig_info = self.sig_info;
64        let metadata = self.metadata.context("metadata is required")?;
65
66        let message = self.custom_message.unwrap_or_else(|| {
67            if let Some(ref sig_info) = sig_info {
68                format!(
69                    "Crashtracker crash ping: crash processing started - Process terminated with {:?} ({:?})",
70                    sig_info.si_code_human_readable, sig_info.si_signo_human_readable
71                )
72            } else {
73                "Crashtracker crash ping: crash processing started - Process terminated".to_string()
74            }
75        });
76
77        Ok(CrashPing {
78            crash_uuid: crash_uuid.to_string(),
79            kind: "Crash ping".to_string(),
80            message,
81            metadata,
82            siginfo: sig_info,
83            version: CrashPing::current_schema_version(),
84        })
85    }
86}
87
88#[derive(Debug, Serialize)]
89pub struct CrashPing {
90    crash_uuid: String,
91    #[serde(skip_serializing_if = "Option::is_none")]
92    siginfo: Option<SigInfo>,
93    message: String,
94    version: String,
95    kind: String,
96    metadata: Metadata,
97}
98
99impl CrashPing {
100    pub fn crash_uuid(&self) -> &str {
101        &self.crash_uuid
102    }
103
104    pub fn message(&self) -> &str {
105        &self.message
106    }
107
108    pub fn metadata(&self) -> &Metadata {
109        &self.metadata
110    }
111
112    pub fn siginfo(&self) -> Option<&SigInfo> {
113        self.siginfo.as_ref()
114    }
115
116    pub fn upload_to_endpoint(&self, endpoint: &Option<Endpoint>) -> anyhow::Result<()> {
117        let rt = tokio::runtime::Builder::new_current_thread()
118            .enable_all()
119            .build()?;
120
121        rt.block_on(async { self.upload_to_endpoint_async(endpoint).await })
122    }
123
124    /// Sends this crash ping telemetry event to indicate that crash processing has started.
125    /// We no-op on file endpoints because unlike production environments, we know if
126    /// a crash report failed to send when file debugging.
127    pub async fn upload_to_endpoint_async(
128        &self,
129        endpoint: &Option<Endpoint>,
130    ) -> anyhow::Result<()> {
131        let telemetry_uploader = crate::TelemetryCrashUploader::new(self.metadata(), endpoint)?;
132        let errors_intake_uploader = crate::ErrorsIntakeUploader::new(endpoint)?;
133        let telemetry_future = telemetry_uploader.upload_crash_ping(self);
134
135        if errors_intake_uploader.is_enabled() {
136            let errors_intake_future = errors_intake_uploader.upload_crash_ping(
137                &self.crash_uuid,
138                self.siginfo.as_ref(),
139                self.metadata(),
140            );
141            let (_telemetry_result, _errors_intake_result) =
142                tokio::join!(telemetry_future, errors_intake_future);
143        } else {
144            let _telemetry_result = telemetry_future.await;
145        }
146        Ok(())
147    }
148
149    fn current_schema_version() -> String {
150        "1.0".to_string()
151    }
152}
153
154macro_rules! parse_tags {
155    (   $tag_iterator:expr,
156        $($tag_name:literal => $var:ident),* $(,)?)  => {
157        $(
158            let mut $var: Option<&str> = None;
159        )*
160        for tag in $tag_iterator {
161            let Some((name, value)) = tag.split_once(':') else {
162                continue;
163            };
164            match name {
165                $($tag_name => {$var = Some(value);}, )*
166                _ => {},
167            }
168        }
169
170    };
171}
172
173pub struct TelemetryCrashUploader {
174    metadata: TelemetryMetadata,
175    cfg: libdd_telemetry::config::Config,
176}
177
178impl TelemetryCrashUploader {
179    pub fn new(
180        crashtracker_metadata: &Metadata,
181        endpoint: &Option<Endpoint>,
182    ) -> anyhow::Result<Self> {
183        let mut cfg = libdd_telemetry::config::Config::from_env();
184        if let Some(endpoint) = endpoint {
185            // TODO: This changes the path part of the query to target the agent.
186            // What about if the crashtracker is sending directly to the intake?
187            // We probably need to remap the host from intake.profile.{site} to
188            // instrumentation-telemetry-intake.{site}?
189            // But do we want to support direct submission to the intake?
190
191            // ignore result because what are we going to do?
192            let _ = if endpoint.url.scheme_str() == Some("file") {
193                let path = libdd_common::decode_uri_path_in_authority(&endpoint.url)
194                    .context("file path is not valid")?;
195                cfg.set_host_from_url(&format!("file://{}.telemetry", path.display()))
196            } else {
197                cfg.set_endpoint(endpoint.clone())
198            };
199        }
200
201        parse_tags!(
202            crashtracker_metadata.tags.iter(),
203            "env" => env,
204            "language" => language_name,
205            "library_version" => library_version,
206            "profiler_version" => profiler_version,
207            "runtime_version" => language_version,
208            "runtime-id" => runtime_id,
209            "service_version" => service_version,
210            "service" => service_name,
211        );
212
213        let application = Application {
214            service_name: service_name.unwrap_or("unknown").to_owned(),
215            language_name: language_name.unwrap_or("unknown").to_owned(),
216            language_version: language_version.unwrap_or("unknown").to_owned(),
217            tracer_version: library_version
218                .or(profiler_version)
219                .unwrap_or("unknown")
220                .to_owned(),
221            env: env.map(ToOwned::to_owned),
222            service_version: service_version.map(ToOwned::to_owned),
223            ..Default::default()
224        };
225
226        let host = build_host();
227
228        let s = Self {
229            metadata: TelemetryMetadata {
230                host,
231                application,
232                runtime_id: runtime_id.unwrap_or("unknown").to_owned(),
233            },
234            cfg,
235        };
236        Ok(s)
237    }
238
239    pub async fn upload_crash_ping(&self, crash_ping: &CrashPing) -> anyhow::Result<()> {
240        let tags = self.build_crash_ping_tags(crash_ping.crash_uuid(), crash_ping.siginfo());
241        let tracer_time = SystemTime::now()
242            .duration_since(SystemTime::UNIX_EPOCH)
243            .map(|d| d.as_secs())
244            .unwrap_or(0);
245        let message = serde_json::to_string(crash_ping)?;
246
247        self.send_log_payload(
248            message,
249            tags,
250            tracer_time,
251            LogLevel::Debug,
252            false, // is_sensitive
253            false, // is_crash
254        )
255        .await
256    }
257
258    fn build_crash_ping_tags(&self, crash_uuid: &str, sig_info: Option<&SigInfo>) -> String {
259        let metadata = &self.metadata;
260        let mut tags = format!(
261            "uuid:{},is_crash_ping:true,service:{},language_name:{},language_version:{},tracer_version:{}",
262            crash_uuid,
263            metadata.application.service_name,
264            metadata.application.language_name,
265            metadata.application.language_version,
266            metadata.application.tracer_version
267        );
268
269        if let Some(sig_info) = sig_info {
270            tags.push_str(&format!(
271                ",si_code_human_readable:{:?},si_signo:{},si_signo_human_readable:{:?}",
272                sig_info.si_code_human_readable,
273                sig_info.si_signo,
274                sig_info.si_signo_human_readable
275            ));
276        }
277
278        self.append_optional_tags(&mut tags);
279        tags
280    }
281
282    fn append_optional_tags(&self, tags: &mut String) {
283        let metadata = &self.metadata;
284        if let Some(env) = &metadata.application.env {
285            tags.push_str(&format!(",env:{env}"));
286        }
287        if let Some(runtime_name) = &metadata.application.runtime_name {
288            tags.push_str(&format!(",runtime_name:{runtime_name}"));
289        }
290        if let Some(runtime_version) = &metadata.application.runtime_version {
291            tags.push_str(&format!(",runtime_version:{runtime_version}"));
292        }
293    }
294
295    pub async fn upload_to_telemetry(&self, crash_info: &CrashInfo) -> anyhow::Result<()> {
296        let message = serde_json::to_string(crash_info)?;
297        let tags = extract_crash_info_tags(crash_info).unwrap_or_default();
298        let tracer_time = crash_info.timestamp.parse::<DateTime<Utc>>().map_or_else(
299            |_| {
300                SystemTime::now()
301                    .duration_since(SystemTime::UNIX_EPOCH)
302                    .map(|d| d.as_secs())
303                    .unwrap_or(0)
304            },
305            |ts| ts.timestamp() as u64,
306        );
307
308        self.send_log_payload(
309            message,
310            tags,
311            tracer_time,
312            LogLevel::Error,
313            true, // is_sensitive
314            true, // is_crash
315        )
316        .await
317    }
318
319    async fn send_log_payload(
320        &self,
321        message: String,
322        tags: String,
323        tracer_time: u64,
324        level: LogLevel,
325        is_sensitive: bool,
326        is_crash: bool,
327    ) -> anyhow::Result<()> {
328        let payload = data::Telemetry {
329            tracer_time,
330            api_version: libdd_telemetry::data::ApiVersion::V2,
331            runtime_id: &self.metadata.runtime_id,
332            seq_id: 1,
333            application: &self.metadata.application,
334            host: &self.metadata.host,
335            payload: &data::Payload::Logs(vec![data::Log {
336                message,
337                level,
338                stack_trace: None,
339                tags,
340                is_sensitive,
341                count: 1,
342                is_crash,
343            }]),
344            origin: Some("Crashtracker"),
345        };
346
347        self.send_telemetry_payload(&payload).await
348    }
349
350    async fn send_telemetry_payload(&self, payload: &data::Telemetry<'_>) -> anyhow::Result<()> {
351        let client = libdd_telemetry::worker::http_client::from_config(&self.cfg);
352        let req = request_builder(&self.cfg)?
353            .method(http::Method::POST)
354            .header(
355                http::header::CONTENT_TYPE,
356                libdd_common::header::APPLICATION_JSON,
357            )
358            .header(
359                libdd_telemetry::worker::http_client::header::API_VERSION,
360                libdd_telemetry::data::ApiVersion::V2.to_str(),
361            )
362            .header(
363                libdd_telemetry::worker::http_client::header::REQUEST_TYPE,
364                "logs",
365            )
366            .body(serde_json::to_string(&payload)?.into())?;
367
368        tokio::time::timeout(
369            std::time::Duration::from_millis({
370                if let Some(endp) = self.cfg.endpoint() {
371                    endp.timeout_ms
372                } else {
373                    Endpoint::DEFAULT_TIMEOUT
374                }
375            }),
376            client.request(req),
377        )
378        .await??;
379
380        Ok(())
381    }
382}
383
384fn extract_crash_info_tags(crash_info: &CrashInfo) -> anyhow::Result<String> {
385    let mut tags = String::new();
386    write!(
387        &mut tags,
388        "data_schema_version:{}",
389        crash_info.data_schema_version
390    )?;
391    if let Some(fingerprint) = &crash_info.fingerprint {
392        write!(&mut tags, ",fingerprint:{fingerprint}")?;
393    }
394    write!(&mut tags, ",incomplete:{}", crash_info.incomplete)?;
395    write!(&mut tags, ",is_crash:{}", crash_info.error.is_crash)?;
396    write!(&mut tags, ",uuid:{}", crash_info.uuid)?;
397    for (counter, value) in &crash_info.counters {
398        write!(&mut tags, ",{counter}:{value}")?;
399    }
400
401    if let Some(siginfo) = &crash_info.sig_info {
402        if let Some(si_addr) = &siginfo.si_addr {
403            write!(&mut tags, ",si_addr:{si_addr}")?;
404        }
405        write!(&mut tags, ",si_code:{}", siginfo.si_code)?;
406        write!(
407            &mut tags,
408            ",si_code_human_readable:{:?}",
409            siginfo.si_code_human_readable
410        )?;
411        write!(&mut tags, ",si_signo:{}", siginfo.si_signo)?;
412        write!(
413            &mut tags,
414            ",si_signo_human_readable:{:?}",
415            siginfo.si_signo_human_readable
416        )?;
417    }
418    Ok(tags)
419}
420
421#[cfg(test)]
422mod tests {
423    use super::TelemetryCrashUploader;
424    use crate::crash_info::{test_utils::TestInstance, CrashInfo, CrashInfoBuilder, Metadata};
425    use libdd_common::Endpoint;
426    use std::{collections::HashSet, fs};
427    use uuid::Uuid;
428
429    fn new_test_uploader(seed: u64) -> TelemetryCrashUploader {
430        TelemetryCrashUploader::new(
431            &Metadata::test_instance(seed),
432            &Some(Endpoint::from_slice("http://localhost:8126")),
433        )
434        .unwrap()
435    }
436
437    #[test]
438    #[cfg_attr(miri, ignore)]
439    fn test_profiler_config_extraction() {
440        let t = new_test_uploader(1);
441
442        let metadata = t.metadata;
443        assert_eq!(metadata.application.service_name, "foo");
444        assert_eq!(metadata.application.service_version.as_deref(), Some("bar"));
445        assert_eq!(metadata.application.language_name, "native");
446        assert_eq!(metadata.runtime_id, "xyz");
447        let cfg = t.cfg;
448        assert_eq!(
449            cfg.endpoint().unwrap().url.to_string(),
450            "http://localhost:8126/telemetry/proxy/api/v2/apmtelemetry"
451        );
452    }
453
454    #[tokio::test]
455    #[cfg_attr(miri, ignore)]
456    async fn test_crash_request_content() -> anyhow::Result<()> {
457        // This keeps alive for scope
458        let tmp = tempfile::tempdir().unwrap();
459        let output_filename = {
460            let mut p = tmp.keep();
461            p.push("crash_info");
462            p
463        };
464        let seed = 1;
465        let mut t = new_test_uploader(seed);
466
467        t.cfg
468            .set_host_from_url(&format!("file://{}", output_filename.to_str().unwrap()))
469            .unwrap();
470        let test_instance = super::CrashInfo::test_instance(seed);
471
472        t.upload_to_telemetry(&test_instance).await.unwrap();
473
474        let payload: serde_json::value::Value =
475            serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap()).unwrap();
476        assert_eq!(payload["api_version"], "v2");
477        assert_eq!(payload["application"]["language_name"], "native");
478        assert_eq!(payload["application"]["service_name"], "foo");
479        assert_eq!(payload["application"]["service_version"], "bar");
480        assert_eq!(payload["request_type"], "logs");
481        assert_eq!(payload["tracer_time"], 1568898000);
482        assert_eq!(payload["origin"], "Crashtracker");
483
484        assert_eq!(payload["payload"].as_array().unwrap().len(), 1);
485        let tags = payload["payload"][0]["tags"]
486            .as_str()
487            .unwrap()
488            .split(',')
489            .collect::<HashSet<_>>();
490        assert_eq!(
491            HashSet::from_iter([
492                "collecting_sample:1",
493                "data_schema_version:1.4",
494                "incomplete:true",
495                "is_crash:true",
496                "not_profiling:0",
497                "si_addr:0x0000000000001234",
498                "si_code_human_readable:SEGV_BNDERR",
499                "si_code:1",
500                "si_signo_human_readable:SIGSEGV",
501                "si_signo:11",
502                "uuid:1d6b97cb-968c-40c9-af6e-e4b4d71e8781",
503            ]),
504            tags
505        );
506        assert_eq!(payload["payload"][0]["is_sensitive"], true);
507        assert_eq!(payload["payload"][0]["level"], "ERROR");
508        let body: CrashInfo =
509            serde_json::from_str(payload["payload"][0]["message"].as_str().unwrap())?;
510        assert_eq!(body, test_instance);
511        assert_eq!(payload["payload"][0]["is_crash"], true);
512        Ok(())
513    }
514
515    #[tokio::test]
516    #[cfg_attr(miri, ignore)]
517    async fn test_crash_ping_content() -> anyhow::Result<()> {
518        // This keeps alive for scope
519        let tmp = tempfile::tempdir().unwrap();
520        let output_filename = {
521            let mut p = tmp.keep();
522            p.push("crash_ping_info");
523            p
524        };
525        let seed = 1;
526        let mut t = new_test_uploader(seed);
527
528        t.cfg
529            .set_host_from_url(&format!("file://{}", output_filename.to_str().unwrap()))
530            .unwrap();
531
532        let sig_info = crate::SigInfo::test_instance(42);
533        let metadata = Metadata::test_instance(1);
534
535        let mut crash_info_builder = CrashInfoBuilder::new();
536        crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
537        crash_info_builder.with_metadata(metadata.clone()).unwrap();
538        let crash_ping = crash_info_builder.build_crash_ping().unwrap();
539        t.upload_crash_ping(&crash_ping).await.unwrap();
540
541        let payload: serde_json::value::Value =
542            serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap()).unwrap();
543        assert_eq!(payload["api_version"], "v2");
544        assert_eq!(payload["application"]["language_name"], "native");
545        assert_eq!(payload["application"]["service_name"], "foo");
546        assert_eq!(payload["application"]["service_version"], "bar");
547        assert_eq!(payload["request_type"], "logs");
548        assert_eq!(payload["origin"], "Crashtracker");
549
550        assert_eq!(payload["payload"].as_array().unwrap().len(), 1);
551        let log_entry = &payload["payload"][0];
552
553        // Crash ping properties
554        assert_eq!(log_entry["is_sensitive"], false);
555        assert_eq!(log_entry["level"], "DEBUG");
556
557        // Structured message format
558        let message_json: serde_json::Value =
559            serde_json::from_str(log_entry["message"].as_str().unwrap())?;
560        assert_eq!(message_json["siginfo"], serde_json::to_value(&sig_info)?);
561        assert!(message_json["crash_uuid"].is_string());
562        assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
563
564        assert_eq!(message_json["version"], "1.0");
565        assert_eq!(message_json["kind"], "Crash ping");
566
567        let metadata_in_message = &message_json["metadata"];
568        assert!(
569            metadata_in_message.is_object(),
570            "metadata should be an object"
571        );
572        let expected_metadata = serde_json::to_value(Metadata::test_instance(1))?;
573        assert_eq!(
574            metadata_in_message, &expected_metadata,
575            "metadata field should match expected structure"
576        );
577
578        let tags = log_entry["tags"].as_str().unwrap();
579        let uuid_str = message_json["crash_uuid"].as_str().unwrap();
580        assert!(tags.contains(&format!("uuid:{uuid_str}")));
581        assert!(tags.contains("is_crash_ping:true"));
582        assert!(tags.contains("service:foo"));
583        assert!(tags.contains("language_name:native"));
584        assert!(tags.contains("language_version:"));
585        assert!(tags.contains("tracer_version:"));
586
587        Ok(())
588    }
589
590    #[tokio::test]
591    #[cfg_attr(miri, ignore)]
592    async fn test_crash_ping_with_different_config() -> anyhow::Result<()> {
593        // This keeps alive for scope
594        let tmp = tempfile::tempdir().unwrap();
595        let output_filename = {
596            let mut p = tmp.keep();
597            p.push("enhanced_crash_ping_info");
598            p
599        };
600        let seed = 1;
601        let mut t = new_test_uploader(seed);
602
603        t.cfg
604            .set_host_from_url(&format!("file://{}", output_filename.to_str().unwrap()))
605            .unwrap();
606
607        let sig_info = crate::SigInfo::test_instance(123);
608        let metadata = Metadata::test_instance(1);
609
610        let mut crash_info_builder = CrashInfoBuilder::new();
611        crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
612        crash_info_builder.with_metadata(metadata.clone()).unwrap();
613        let crash_ping = crash_info_builder.build_crash_ping().unwrap();
614        t.upload_crash_ping(&crash_ping).await.unwrap();
615
616        let payload: serde_json::value::Value =
617            serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap()).unwrap();
618        assert_eq!(payload["api_version"], "v2");
619        assert_eq!(payload["application"]["language_name"], "native");
620        assert_eq!(payload["application"]["service_name"], "foo");
621        assert_eq!(payload["application"]["service_version"], "bar");
622        assert_eq!(payload["request_type"], "logs");
623        assert_eq!(payload["origin"], "Crashtracker");
624
625        assert_eq!(payload["payload"].as_array().unwrap().len(), 1);
626        let log_entry = &payload["payload"][0];
627
628        // Crash ping properties
629        assert_eq!(log_entry["is_crash"], false);
630        assert_eq!(log_entry["is_sensitive"], false);
631        assert_eq!(log_entry["level"], "DEBUG");
632
633        // Structured message format
634        let message_json: serde_json::Value =
635            serde_json::from_str(log_entry["message"].as_str().unwrap())?;
636        assert!(message_json["crash_uuid"].is_string());
637        assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
638        assert_eq!(
639            message_json["message"],
640            format!(
641                "Crashtracker crash ping: crash processing started - Process terminated with {:?} ({:?})",
642                sig_info.si_code_human_readable, sig_info.si_signo_human_readable
643            )
644        );
645
646        let metadata_in_message = &message_json["metadata"];
647        assert!(
648            metadata_in_message.is_object(),
649            "metadata should be an object"
650        );
651        let expected_metadata = serde_json::to_value(Metadata::test_instance(1))?;
652        assert_eq!(
653            metadata_in_message, &expected_metadata,
654            "metadata field should match expected structure"
655        );
656
657        let siginfo_in_message = &message_json["siginfo"];
658        let expected_siginfo = serde_json::to_value(&sig_info)?;
659        assert_eq!(
660            siginfo_in_message, &expected_siginfo,
661            "siginfo field should match expected structure"
662        );
663
664        assert_eq!(message_json["version"], "1.0");
665        assert_eq!(message_json["kind"], "Crash ping");
666
667        let tags = log_entry["tags"].as_str().unwrap();
668        let uuid_str = message_json["crash_uuid"].as_str().unwrap();
669        assert!(tags.contains(&format!("uuid:{uuid_str}")));
670        assert!(tags.contains("is_crash_ping:true"));
671        assert!(tags.contains("service:foo"));
672        assert!(tags.contains("language_name:native"));
673        assert!(tags.contains("language_version:"));
674        assert!(tags.contains("tracer_version:"));
675
676        Ok(())
677    }
678
679    #[tokio::test]
680    #[cfg_attr(miri, ignore)]
681    async fn test_crash_ping_builder_basic() -> anyhow::Result<()> {
682        let tmp = tempfile::tempdir().unwrap();
683        let output_filename = {
684            let mut p = tmp.keep();
685            p.push("crash_ping_builder_test");
686            p
687        };
688
689        let sig_info = crate::SigInfo::test_instance(42);
690        let metadata = Metadata::test_instance(1);
691
692        // Build crash ping through CrashInfoBuilder
693        let mut crash_info_builder = CrashInfoBuilder::new();
694        crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
695        crash_info_builder.with_metadata(metadata.clone()).unwrap();
696        let crash_ping = crash_info_builder.build_crash_ping()?;
697
698        let endpoint = Some(Endpoint::from_slice(&format!(
699            "file://{}",
700            output_filename.to_str().unwrap()
701        )));
702
703        assert!(!crash_ping.crash_uuid().is_empty());
704        assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
705        assert!(crash_ping.message().contains("crash processing started"));
706        assert_eq!(crash_ping.metadata(), &metadata);
707
708        // Use TelemetryCrashUploader to upload the crash ping
709        let mut uploader = TelemetryCrashUploader::new(&metadata, &endpoint)?;
710        uploader
711            .cfg
712            .set_host_from_url(&format!(
713                "file://{}.telemetry",
714                output_filename.to_str().unwrap()
715            ))
716            .unwrap();
717
718        uploader.upload_crash_ping(&crash_ping).await?;
719
720        // Verify the .telemetry file was created with correct content
721        let telemetry_filename = format!("{}.telemetry", output_filename.to_str().unwrap());
722        let payload: serde_json::value::Value =
723            serde_json::de::from_str(&std::fs::read_to_string(&telemetry_filename)?)?;
724
725        assert_eq!(payload["api_version"], "v2");
726        assert_eq!(payload["request_type"], "logs");
727        assert_eq!(payload["origin"], "Crashtracker");
728
729        let log_entry = &payload["payload"][0];
730        assert_eq!(log_entry["level"], "DEBUG");
731        assert_eq!(log_entry["is_sensitive"], false);
732        assert_eq!(log_entry["is_crash"], false);
733
734        let message_json: serde_json::Value =
735            serde_json::from_str(log_entry["message"].as_str().unwrap())?;
736        assert!(message_json["crash_uuid"].is_string());
737        assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
738        assert_eq!(message_json["version"], "1.0");
739        assert_eq!(message_json["kind"], "Crash ping");
740
741        Ok(())
742    }
743
744    #[test]
745    #[cfg_attr(miri, ignore)]
746    fn test_crash_ping_builder_validation() {
747        // Test that crash ping can be built with only metadata (no sig_info)
748        let mut crash_info_builder = CrashInfoBuilder::new();
749        crash_info_builder
750            .with_metadata(Metadata::test_instance(1))
751            .unwrap();
752        let result = crash_info_builder.build_crash_ping();
753        assert!(result.is_ok());
754        let crash_ping = result.unwrap();
755        assert!(crash_ping.siginfo().is_none());
756        assert!(crash_ping
757            .message()
758            .contains("Crashtracker crash ping: crash processing started - Process terminated"));
759
760        // Test that crash ping fails without metadata
761        let mut crash_info_builder = CrashInfoBuilder::new();
762        crash_info_builder
763            .with_sig_info(crate::SigInfo::test_instance(1))
764            .unwrap();
765        let result = crash_info_builder.build_crash_ping();
766        assert!(result.is_err());
767        assert!(result
768            .unwrap_err()
769            .to_string()
770            .contains("metadata is required"));
771
772        // Test successful build with both fields
773        let mut crash_info_builder = CrashInfoBuilder::new();
774        crash_info_builder
775            .with_sig_info(crate::SigInfo::test_instance(1))
776            .unwrap();
777        crash_info_builder
778            .with_metadata(Metadata::test_instance(1))
779            .unwrap();
780        let result = crash_info_builder.build_crash_ping();
781        assert!(result.is_ok());
782        let crash_ping = result.unwrap();
783        assert!(crash_ping.siginfo().is_some());
784    }
785
786    #[test]
787    #[cfg_attr(miri, ignore)]
788    fn test_crash_ping_all_fields_present() {
789        let sig_info = crate::SigInfo::test_instance(99);
790        let metadata = Metadata::test_instance(2);
791
792        // Build crash ping through CrashInfoBuilder
793        let mut crash_info_builder = CrashInfoBuilder::new();
794        crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
795        crash_info_builder.with_metadata(metadata.clone()).unwrap();
796        let crash_ping = crash_info_builder.build_crash_ping().unwrap();
797
798        assert!(!crash_ping.crash_uuid().is_empty());
799        assert!(Uuid::parse_str(crash_ping.crash_uuid()).is_ok());
800        assert!(crash_ping.message().contains("crash processing started"));
801        assert_eq!(crash_ping.metadata(), &metadata);
802        assert_eq!(crash_ping.siginfo(), Some(&sig_info));
803    }
804
805    #[tokio::test]
806    #[cfg_attr(miri, ignore)]
807    async fn test_crash_ping_telemetry_upload_all_fields() -> anyhow::Result<()> {
808        // Test that when crash ping is uploaded via telemetry, all fields are preserved
809        let tmp = tempfile::tempdir().unwrap();
810        let output_filename = {
811            let mut p = tmp.keep();
812            p.push("crash_ping_all_fields_upload");
813            p
814        };
815        let seed = 3;
816        let mut uploader = new_test_uploader(seed);
817
818        uploader
819            .cfg
820            .set_host_from_url(&format!("file://{}", output_filename.to_str().unwrap()))
821            .unwrap();
822
823        let sig_info = crate::SigInfo::test_instance(150);
824        let metadata = Metadata::test_instance(3);
825
826        // Build crash ping through CrashInfoBuilder
827        let mut crash_info_builder = CrashInfoBuilder::new();
828        crash_info_builder.with_sig_info(sig_info.clone()).unwrap();
829        crash_info_builder.with_metadata(metadata.clone()).unwrap();
830        let crash_ping = crash_info_builder.build_crash_ping().unwrap();
831
832        uploader.upload_crash_ping(&crash_ping).await?;
833
834        let payload: serde_json::value::Value =
835            serde_json::de::from_str(&fs::read_to_string(&output_filename).unwrap())?;
836
837        // Verify telemetry structure
838        assert_eq!(payload["api_version"], "v2");
839        assert_eq!(payload["request_type"], "logs");
840        assert_eq!(payload["origin"], "Crashtracker");
841
842        let log_entry = &payload["payload"][0];
843        assert_eq!(log_entry["level"], "DEBUG");
844        assert_eq!(log_entry["is_sensitive"], false);
845        assert_eq!(log_entry["is_crash"], false);
846
847        let message_json: serde_json::Value =
848            serde_json::from_str(log_entry["message"].as_str().unwrap())?;
849
850        assert!(message_json["crash_uuid"].is_string());
851        assert!(Uuid::parse_str(message_json["crash_uuid"].as_str().unwrap()).is_ok());
852        assert_eq!(message_json["version"], "1.0");
853        assert_eq!(message_json["kind"], "Crash ping");
854
855        let uploaded_siginfo = &message_json["siginfo"];
856        assert_eq!(uploaded_siginfo["si_signo"], sig_info.si_signo);
857        assert_eq!(uploaded_siginfo["si_code"], sig_info.si_code);
858        assert_eq!(
859            uploaded_siginfo["si_code_human_readable"],
860            serde_json::to_value(&sig_info.si_code_human_readable)?
861        );
862        assert_eq!(
863            uploaded_siginfo["si_signo_human_readable"],
864            serde_json::to_value(&sig_info.si_signo_human_readable)?
865        );
866
867        let uploaded_metadata = &message_json["metadata"];
868        assert!(uploaded_metadata.is_object());
869
870        let expected_metadata_json = serde_json::to_value(&metadata)?;
871        assert_eq!(uploaded_metadata, &expected_metadata_json);
872
873        assert!(message_json["message"].is_string());
874        assert!(message_json["message"]
875            .as_str()
876            .unwrap()
877            .contains("crash processing started"));
878        Ok(())
879    }
880}