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 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    /// Crash pings should only be initalized and built by the CrashInfoBuilder
39    /// We require the crash uuid to be passed in because a CrashPing is always
40    /// associated with a specific crash.
41    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    /// Sends this crash ping telemetry event to indicate that crash processing has started.
140    /// We no-op on file endpoints because unlike production environments, we know if
141    /// a crash report failed to send when file debugging.
142    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            // TODO: This changes the path part of the query to target the agent.
197            // What about if the crashtracker is sending directly to the intake?
198            // We probably need to remap the host from intake.profile.{site} to
199            // instrumentation-telemetry-intake.{site}?
200            // But do we want to support direct submission to the intake?
201
202            // ignore result because what are we going to do?
203            if endpoint.url.scheme_str() == Some("file") {
204                // Write telemetry alongside the crash file. Rebuild from the
205                // decoded path so the file path is re-encoded exactly once.
206                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    /// Send a general (non crash report, non crash ping) log message to the telemetry log intake.
265    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    /// Send a crash ping telemetry log to indicate that crash processing has started
281    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, // is_sensitive
295            false, // is_crash
296        )
297        .await
298    }
299
300    /// Send a crash info telemetry log to indicate that crash processing has completed
301    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, // is_sensitive
320            true, // is_crash
321        )
322        .await
323    }
324
325    /// Shared helper that builds `data::Telemetry` payload and calls `send_telemetry_payload`
326    /// to send the payload to the telemetry log intake.
327    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    /// Helper to perform actual HTTP submission via the native HTTP capability.
361    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        // This keeps alive for scope
527        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        // This keeps alive for scope
597        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        // Crash ping properties
636        assert_eq!(log_entry["is_sensitive"], false);
637        assert_eq!(log_entry["level"], "DEBUG");
638
639        // Structured message format
640        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        // This keeps alive for scope
676        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        // Crash ping properties
715        assert_eq!(log_entry["is_crash"], false);
716        assert_eq!(log_entry["is_sensitive"], false);
717        assert_eq!(log_entry["level"], "DEBUG");
718
719        // Structured message format
720        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        // Build crash ping through CrashInfoBuilder
779        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        // Use TelemetryCrashUploader to upload the crash ping
796        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        // Verify the .telemetry file was created with correct content
811        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        // Test that crash ping can be built with only metadata (no sig_info)
838        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        // Test that crash ping fails without metadata
852        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        // Test successful build with both fields
865        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        // Build crash ping through CrashInfoBuilder
886        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        // Build crash ping through CrashInfoBuilder
906        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        // Build crash ping through CrashInfoBuilder
929        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        // Test that when crash ping is uploaded via telemetry, all fields are preserved
951        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        // Build crash ping through CrashInfoBuilder
972        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        // Verify telemetry structure
984        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}