1use std::fmt;
46use std::io::{self, Write};
47use std::path::Path;
48use std::sync::atomic::{AtomicU64, Ordering};
49use std::sync::{mpsc, Arc};
50use std::time::{SystemTime, UNIX_EPOCH};
51
52use ant_protocol::transport::PeerId;
53use serde::Serialize;
54use tracing::{error, warn};
55
56pub const DIAGNOSTICS_SCHEMA_VERSION: u8 = 4;
58
59#[derive(Debug, Clone, PartialEq, Eq)]
62pub(crate) struct DownloadRequestCorrelation {
63 pub(crate) request_id: u64,
64 pub(crate) local_peer_id: String,
65}
66
67impl DownloadRequestCorrelation {
68 pub(crate) fn new(request_id: u64, local_peer_id: &PeerId) -> Self {
69 Self {
70 request_id,
71 local_peer_id: local_peer_id.to_string(),
72 }
73 }
74}
75
76const DIAGNOSTICS_CHANNEL_CAPACITY: usize = 1024;
80
81const DIAGNOSTICS_ERROR_MAX_CHARS: usize = 240;
85
86#[derive(Debug, Clone, Copy, PartialEq, Eq)]
92pub enum DownloadDiagnosticsOutcome {
93 Found,
95 NotFound,
97 Timeout,
99 NetworkError,
101 ProtocolError,
105 CacheHit,
107 LookupError,
109 Exhausted,
111}
112
113impl DownloadDiagnosticsOutcome {
114 #[must_use]
116 pub const fn as_str(self) -> &'static str {
117 match self {
118 Self::Found => "found",
119 Self::NotFound => "not_found",
120 Self::Timeout => "timeout",
121 Self::NetworkError => "network_error",
122 Self::ProtocolError => "protocol_error",
123 Self::CacheHit => "cache_hit",
124 Self::LookupError => "lookup_error",
125 Self::Exhausted => "exhausted",
126 }
127 }
128}
129
130impl fmt::Display for DownloadDiagnosticsOutcome {
131 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
132 f.write_str(self.as_str())
133 }
134}
135
136impl Serialize for DownloadDiagnosticsOutcome {
137 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
138 where
139 S: serde::Serializer,
140 {
141 serializer.serialize_str(self.as_str())
142 }
143}
144
145#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
152pub struct DownloadDiagnosticsRecord {
153 pub schema_version: u8,
155 pub timestamp: String,
157 pub request_started_unix_ms: Option<u64>,
161 pub request_completed_unix_ms: Option<u64>,
164 pub request_id: Option<u64>,
168 pub local_peer_id: Option<String>,
171 pub file_attempt: usize,
173 pub chunk_index: usize,
175 pub chunk_address: String,
177 pub sweep: String,
179 pub peer_attempt: Option<usize>,
182 pub lookup_duration_ms: Option<u64>,
185 pub lookup_correlation_id: Option<String>,
187 pub selected_peer_ordinal: Option<usize>,
189 pub expected_peer: Option<String>,
192 pub selected_peer_addresses: Option<Vec<String>>,
194 pub selected_peer_address_types: Option<Vec<String>>,
196 pub local_last_seen_age_ms: Option<u64>,
199 pub publisher_address_set_age_ms: Option<u64>,
202 pub publisher_address_set_unix_ns: Option<u64>,
204 pub source_peer: Option<String>,
208 pub transport_source: Option<String>,
212 pub route: String,
217 pub route_note: Option<String>,
220 pub peer_connected_before_request: Option<bool>,
223 pub active_requests_at_start: Option<usize>,
226 pub fetch_cap: Option<usize>,
229 pub response_elapsed_ms: Option<u64>,
232 pub ttfb_ms: Option<u64>,
234 pub ttfb_available: bool,
237 pub ttfb_unavailable_reason: String,
239 pub bytes: u64,
241 pub outcome: DownloadDiagnosticsOutcome,
243 pub error: Option<String>,
246}
247
248impl DownloadDiagnosticsRecord {
249 pub const TTFB_UNAVAILABLE_REASON: &'static str =
251 "protocol exposes only a complete-response event; first-byte/first-frame \
252 timing is not available";
253
254 pub const ROUTE_UNKNOWN_NOTE: &'static str =
258 "transport_source absent or did not match any known typed peer dial address; \
259 route could not be classified from the observed response";
260
261 #[allow(clippy::too_many_arguments)]
265 pub(crate) fn peer_attempt(
266 file_attempt: usize,
267 chunk_index: usize,
268 chunk_address: &[u8; 32],
269 sweep: &'static str,
270 peer_attempt: usize,
271 lookup_duration_ms: Option<u64>,
272 lookup_correlation_id: &str,
273 expected_peer: &str,
274 selected_peer_addresses: Vec<String>,
275 selected_peer_address_types: Vec<String>,
276 local_last_seen_age_ms: Option<u64>,
277 publisher_address_set_age_ms: Option<u64>,
278 publisher_address_set_unix_ns: Option<u64>,
279 source_peer: Option<&str>,
280 transport_source: Option<&str>,
281 route: &str,
282 route_note: Option<&str>,
283 peer_connected_before_request: Option<bool>,
284 active_requests_at_start: Option<usize>,
285 fetch_cap: Option<usize>,
286 request_started_unix_ms: u64,
287 request_completed_unix_ms: u64,
288 correlation: &DownloadRequestCorrelation,
289 response_elapsed_ms: u64,
290 bytes: u64,
291 outcome: DownloadDiagnosticsOutcome,
292 error: Option<String>,
293 ) -> Self {
294 let lookup = if peer_attempt == 1 {
295 lookup_duration_ms
296 } else {
297 None
298 };
299 Self {
300 schema_version: DIAGNOSTICS_SCHEMA_VERSION,
301 timestamp: utc_now_rfc3339(),
302 request_started_unix_ms: Some(request_started_unix_ms),
303 request_completed_unix_ms: Some(request_completed_unix_ms),
304 request_id: Some(correlation.request_id),
305 local_peer_id: Some(correlation.local_peer_id.clone()),
306 file_attempt,
307 chunk_index,
308 chunk_address: hex::encode(chunk_address),
309 sweep: sweep.to_string(),
310 peer_attempt: Some(peer_attempt),
311 lookup_duration_ms: lookup,
312 lookup_correlation_id: Some(lookup_correlation_id.to_string()),
313 selected_peer_ordinal: Some(peer_attempt),
314 expected_peer: Some(expected_peer.to_string()),
315 selected_peer_addresses: Some(selected_peer_addresses),
316 selected_peer_address_types: Some(selected_peer_address_types),
317 local_last_seen_age_ms,
318 publisher_address_set_age_ms,
319 publisher_address_set_unix_ns,
320 source_peer: source_peer.map(str::to_string),
321 transport_source: transport_source.map(str::to_string),
322 route: route.to_string(),
323 route_note: route_note.map(str::to_string),
324 peer_connected_before_request,
325 active_requests_at_start,
326 fetch_cap,
327 response_elapsed_ms: Some(response_elapsed_ms),
328 ttfb_ms: None,
329 ttfb_available: false,
330 ttfb_unavailable_reason: Self::TTFB_UNAVAILABLE_REASON.to_string(),
331 bytes,
332 outcome,
333 error,
334 }
335 }
336
337 #[allow(clippy::too_many_arguments)]
340 pub fn chunk_level(
341 file_attempt: usize,
342 chunk_index: usize,
343 chunk_address: &[u8; 32],
344 sweep: &'static str,
345 fetch_cap: Option<usize>,
346 bytes: u64,
347 outcome: DownloadDiagnosticsOutcome,
348 error: Option<String>,
349 ) -> Self {
350 Self {
351 schema_version: DIAGNOSTICS_SCHEMA_VERSION,
352 timestamp: utc_now_rfc3339(),
353 request_started_unix_ms: None,
354 request_completed_unix_ms: None,
355 request_id: None,
356 local_peer_id: None,
357 file_attempt,
358 chunk_index,
359 chunk_address: hex::encode(chunk_address),
360 sweep: sweep.to_string(),
361 peer_attempt: None,
362 lookup_duration_ms: None,
363 lookup_correlation_id: None,
364 selected_peer_ordinal: None,
365 expected_peer: None,
366 selected_peer_addresses: None,
367 selected_peer_address_types: None,
368 local_last_seen_age_ms: None,
369 publisher_address_set_age_ms: None,
370 publisher_address_set_unix_ns: None,
371 source_peer: None,
372 transport_source: None,
373 route: "unknown".to_string(),
374 route_note: None,
375 peer_connected_before_request: None,
376 active_requests_at_start: None,
377 fetch_cap,
378 response_elapsed_ms: None,
379 ttfb_ms: None,
380 ttfb_available: false,
381 ttfb_unavailable_reason: Self::TTFB_UNAVAILABLE_REASON.to_string(),
382 bytes,
383 outcome,
384 error,
385 }
386 }
387}
388
389#[derive(Clone)]
396pub struct DownloadDiagnosticsSender {
397 tx: mpsc::SyncSender<DownloadDiagnosticsRecord>,
398 dropped: Arc<AtomicU64>,
399}
400
401impl DownloadDiagnosticsSender {
402 pub fn try_emit(&self, record: DownloadDiagnosticsRecord) {
405 if self.tx.try_send(record).is_err() {
406 self.dropped.fetch_add(1, Ordering::Relaxed);
407 }
408 }
409}
410
411pub fn spawn_download_diagnostics_writer(
426 path: &Path,
427) -> io::Result<(DownloadDiagnosticsSender, std::thread::JoinHandle<()>)> {
428 let file = std::fs::OpenOptions::new()
429 .create(true)
430 .write(true)
431 .truncate(true)
432 .open(path)?;
433 let (tx, rx) = mpsc::sync_channel::<DownloadDiagnosticsRecord>(DIAGNOSTICS_CHANNEL_CAPACITY);
434 let dropped = Arc::new(AtomicU64::new(0));
435 let writer_dropped = Arc::clone(&dropped);
436 let builder = std::thread::Builder::new().name("ant-download-diagnostics-writer".to_string());
437 let handle = builder
438 .spawn(move || {
439 let mut writer = io::BufWriter::new(file);
440 for record in rx.iter() {
441 match serde_json::to_string(&record) {
442 Ok(line) => {
443 if let Err(err) = writeln!(writer, "{line}") {
444 error!(%err, "download diagnostics sidecar write failed");
445 break;
446 }
447 }
448 Err(err) => {
449 error!(%err, "download diagnostics record serialization failed");
453 continue;
454 }
455 }
456 }
457 if let Err(err) = writer.flush() {
458 error!(%err, "download diagnostics sidecar flush failed");
459 }
460 let dropped = writer_dropped.load(Ordering::Relaxed);
461 if dropped > 0 {
462 warn!(dropped, "download diagnostics records were dropped");
463 }
464 })
465 .map_err(|e| io::Error::other(format!("failed to spawn diagnostics writer thread: {e}")))?;
466 Ok((DownloadDiagnosticsSender { tx, dropped }, handle))
467}
468
469pub fn bounded_error(category: &str, detail: &str) -> String {
476 let prefix = if category.is_empty() {
477 String::new()
478 } else {
479 format!("{category}: ")
480 };
481 if prefix.len() + detail.len() <= DIAGNOSTICS_ERROR_MAX_CHARS {
482 return format!("{prefix}{detail}");
483 }
484 let remaining = DIAGNOSTICS_ERROR_MAX_CHARS.saturating_sub(prefix.len());
485 let mut truncated: String = detail.chars().take(remaining.saturating_sub(1)).collect();
486 truncated.push('…');
487 format!("{prefix}{truncated}")
488}
489
490fn utc_now_rfc3339() -> String {
493 let now = SystemTime::now()
494 .duration_since(UNIX_EPOCH)
495 .unwrap_or_default();
496 rfc3339_from_unix_secs(now.as_secs())
497}
498
499pub(crate) fn unix_now_ms() -> u64 {
503 let millis = SystemTime::now()
504 .duration_since(UNIX_EPOCH)
505 .unwrap_or_default()
506 .as_millis();
507 u64::try_from(millis).unwrap_or(u64::MAX)
508}
509
510fn rfc3339_from_unix_secs(secs: u64) -> String {
515 let days = (secs / 86_400) as i64;
516 let secs_of_day = secs % 86_400;
517 let hour = secs_of_day / 3600;
518 let minute = (secs_of_day % 3600) / 60;
519 let second = secs_of_day % 60;
520
521 let z = days + 719_468;
523 let era = if z >= 0 { z } else { z - 146_096 } / 146_097;
524 let doe = z - era * 146_097;
525 let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
526 let y = yoe + era * 400;
527 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
528 let mp = (5 * doy + 2) / 153;
529 let d = doy - (153 * mp + 2) / 5 + 1;
530 let m = if mp < 10 { mp + 3 } else { mp - 9 };
531 let year = if m <= 2 { y + 1 } else { y };
532
533 format!("{year:04}-{m:02}-{d:02}T{hour:02}:{minute:02}:{second:02}Z")
534}
535
536#[cfg(test)]
537mod tests {
538 use super::*;
539
540 fn test_correlation(request_id: u64) -> DownloadRequestCorrelation {
541 DownloadRequestCorrelation::new(request_id, &PeerId::from_bytes([42; 32]))
542 }
543
544 #[test]
545 fn outcome_as_str_is_stable_lowercase() {
546 assert_eq!(DownloadDiagnosticsOutcome::Found.as_str(), "found");
547 assert_eq!(DownloadDiagnosticsOutcome::NotFound.as_str(), "not_found");
548 assert_eq!(DownloadDiagnosticsOutcome::Timeout.as_str(), "timeout");
549 assert_eq!(
550 DownloadDiagnosticsOutcome::NetworkError.as_str(),
551 "network_error"
552 );
553 assert_eq!(
554 DownloadDiagnosticsOutcome::ProtocolError.as_str(),
555 "protocol_error"
556 );
557 assert_eq!(DownloadDiagnosticsOutcome::CacheHit.as_str(), "cache_hit");
558 assert_eq!(
559 DownloadDiagnosticsOutcome::LookupError.as_str(),
560 "lookup_error"
561 );
562 assert_eq!(DownloadDiagnosticsOutcome::Exhausted.as_str(), "exhausted");
563 }
564
565 #[test]
566 fn outcome_serializes_as_lowercase_string() {
567 let v = serde_json::to_string(&DownloadDiagnosticsOutcome::NetworkError).unwrap();
568 assert_eq!(v, "\"network_error\"");
569 }
570
571 #[test]
572 fn peer_attempt_record_pins_field_names_and_null_ttfb() {
573 let addr = [7u8; 32];
574 let correlation = test_correlation(9_001);
575 let record = DownloadDiagnosticsRecord::peer_attempt(
576 1,
577 3,
578 &addr,
579 "initial",
580 1,
581 Some(42),
582 "lookup-1",
583 "peer-abc",
584 vec!["/ip4/1.2.3.4/udp/9000/quic".to_string()],
585 vec!["direct".to_string()],
586 Some(1_500),
587 Some(2_000),
588 Some(1_234_000_000),
589 Some("peer-abc"),
590 Some("/ip4/1.2.3.4/udp/9000/quic"),
591 "direct",
592 None,
593 Some(true),
594 Some(8),
595 Some(8),
596 120,
597 1024,
598 &correlation,
599 904,
600 1024,
601 DownloadDiagnosticsOutcome::Found,
602 None,
603 );
604 let json = serde_json::to_value(&record).unwrap();
605 let obj = json.as_object().unwrap();
606 for field in [
608 "schema_version",
609 "timestamp",
610 "request_started_unix_ms",
611 "request_completed_unix_ms",
612 "request_id",
613 "local_peer_id",
614 "file_attempt",
615 "chunk_index",
616 "chunk_address",
617 "sweep",
618 "peer_attempt",
619 "lookup_duration_ms",
620 "lookup_correlation_id",
621 "selected_peer_ordinal",
622 "expected_peer",
623 "selected_peer_addresses",
624 "selected_peer_address_types",
625 "local_last_seen_age_ms",
626 "publisher_address_set_age_ms",
627 "publisher_address_set_unix_ns",
628 "source_peer",
629 "transport_source",
630 "route",
631 "route_note",
632 "peer_connected_before_request",
633 "active_requests_at_start",
634 "fetch_cap",
635 "response_elapsed_ms",
636 "ttfb_ms",
637 "ttfb_available",
638 "ttfb_unavailable_reason",
639 "bytes",
640 "outcome",
641 "error",
642 ] {
643 assert!(obj.contains_key(field), "missing field {field}");
644 }
645 assert_eq!(obj["ttfb_ms"], serde_json::Value::Null);
647 assert_eq!(obj["ttfb_available"], serde_json::Value::Bool(false));
648 assert!(
649 obj["ttfb_unavailable_reason"]
650 .as_str()
651 .unwrap()
652 .contains("first-byte"),
653 "ttfb reason must mention first-byte"
654 );
655 assert_eq!(obj["route"], serde_json::Value::String("direct".into()));
657 assert_eq!(obj["route_note"], serde_json::Value::Null);
658 assert_eq!(obj["source_peer"], serde_json::json!("peer-abc"));
660 assert_eq!(
661 obj["transport_source"],
662 serde_json::json!("/ip4/1.2.3.4/udp/9000/quic")
663 );
664 assert_eq!(obj["request_started_unix_ms"], serde_json::json!(120u64));
666 assert_eq!(obj["request_completed_unix_ms"], serde_json::json!(1024u64));
667 assert_eq!(obj["active_requests_at_start"], serde_json::json!(8usize));
668 assert_eq!(
669 obj["peer_connected_before_request"],
670 serde_json::json!(true)
671 );
672 assert_eq!(obj["fetch_cap"], serde_json::json!(8usize));
673 assert_eq!(obj["lookup_duration_ms"], serde_json::json!(42u64));
675 assert_eq!(obj["bytes"], serde_json::json!(1024u64));
676 assert_eq!(obj["outcome"], serde_json::json!("found"));
677 assert_eq!(obj["request_id"], serde_json::json!(9_001u64));
678 assert_eq!(
679 obj["local_peer_id"],
680 serde_json::json!(correlation.local_peer_id)
681 );
682 assert_eq!(obj["schema_version"], serde_json::json!(4u8));
683 assert_eq!(obj["lookup_correlation_id"], serde_json::json!("lookup-1"));
684 assert_eq!(obj["selected_peer_ordinal"], serde_json::json!(1usize));
685 assert_eq!(obj["local_last_seen_age_ms"], serde_json::json!(1_500u64));
686 assert_eq!(
687 obj["publisher_address_set_age_ms"],
688 serde_json::json!(2_000u64)
689 );
690 assert_eq!(
691 obj["chunk_address"],
692 serde_json::Value::String(hex::encode(addr))
693 );
694 }
695
696 #[test]
697 fn later_peer_attempt_omits_lookup_duration_and_carries_route_note_when_unknown() {
698 let addr = [9u8; 32];
699 let record = DownloadDiagnosticsRecord::peer_attempt(
700 1,
701 1,
702 &addr,
703 "retry",
704 3,
705 Some(10),
706 "lookup-2",
707 "peer-x",
708 vec!["/ip6/2001:db8::1/udp/9000/quic".to_string()],
709 vec!["unverified".to_string()],
710 None,
711 None,
712 None,
713 None,
715 None,
716 "unknown",
717 Some(DownloadDiagnosticsRecord::ROUTE_UNKNOWN_NOTE),
718 Some(false),
719 Some(4),
720 Some(4),
721 500,
722 1_000,
723 &test_correlation(9_002),
724 500,
725 0,
726 DownloadDiagnosticsOutcome::Timeout,
727 Some(bounded_error("timeout", "no response")),
728 );
729 let json = serde_json::to_value(&record).unwrap();
730 let obj = json.as_object().unwrap();
731 assert_eq!(obj["lookup_duration_ms"], serde_json::Value::Null);
732 assert_eq!(obj["lookup_correlation_id"], serde_json::json!("lookup-2"));
733 assert_eq!(obj["selected_peer_ordinal"], serde_json::json!(3usize));
734 assert_eq!(
735 obj["route_note"],
736 serde_json::json!(DownloadDiagnosticsRecord::ROUTE_UNKNOWN_NOTE)
737 );
738 assert_eq!(obj["source_peer"], serde_json::Value::Null);
739 assert_eq!(obj["transport_source"], serde_json::Value::Null);
740 assert_eq!(obj["route"], serde_json::json!("unknown"));
741 assert_eq!(
742 obj["peer_connected_before_request"],
743 serde_json::json!(false)
744 );
745 assert_eq!(obj["active_requests_at_start"], serde_json::json!(4usize));
746 assert_eq!(obj["fetch_cap"], serde_json::json!(4usize));
747 assert_eq!(obj["outcome"], serde_json::json!("timeout"));
748 assert_eq!(obj["bytes"], serde_json::json!(0u64));
749 assert!(obj["error"].as_str().unwrap().starts_with("timeout: "));
750 }
751
752 #[test]
753 fn cache_hit_record_has_no_peer_fields() {
754 let addr = [1u8; 32];
755 let record = DownloadDiagnosticsRecord::chunk_level(
756 1,
757 2,
758 &addr,
759 "initial",
760 Some(8),
761 4096,
762 DownloadDiagnosticsOutcome::CacheHit,
763 None,
764 );
765 let json = serde_json::to_value(&record).unwrap();
766 let obj = json.as_object().unwrap();
767 assert_eq!(obj["peer_attempt"], serde_json::Value::Null);
768 assert_eq!(obj["expected_peer"], serde_json::Value::Null);
769 assert_eq!(obj["source_peer"], serde_json::Value::Null);
770 assert_eq!(obj["transport_source"], serde_json::Value::Null);
771 assert_eq!(obj["lookup_duration_ms"], serde_json::Value::Null);
772 assert_eq!(obj["lookup_correlation_id"], serde_json::Value::Null);
773 assert_eq!(obj["selected_peer_addresses"], serde_json::Value::Null);
774 assert_eq!(obj["response_elapsed_ms"], serde_json::Value::Null);
775 assert_eq!(
776 obj["peer_connected_before_request"],
777 serde_json::Value::Null
778 );
779 assert_eq!(obj["request_started_unix_ms"], serde_json::Value::Null);
780 assert_eq!(obj["request_completed_unix_ms"], serde_json::Value::Null);
781 assert_eq!(obj["request_id"], serde_json::Value::Null);
782 assert_eq!(obj["local_peer_id"], serde_json::Value::Null);
783 assert_eq!(obj["active_requests_at_start"], serde_json::Value::Null);
784 assert_eq!(obj["fetch_cap"], serde_json::json!(8usize));
785 assert_eq!(obj["route"], serde_json::json!("unknown"));
786 assert_eq!(obj["route_note"], serde_json::Value::Null);
787 assert_eq!(obj["outcome"], serde_json::json!("cache_hit"));
788 assert_eq!(obj["bytes"], serde_json::json!(4096u64));
789 }
790
791 #[test]
792 fn exhausted_and_lookup_error_records_classify_correctly() {
793 let addr = [2u8; 32];
794 let exhausted = DownloadDiagnosticsRecord::chunk_level(
795 2,
796 5,
797 &addr,
798 "retry",
799 Some(2),
800 0,
801 DownloadDiagnosticsOutcome::Exhausted,
802 None,
803 );
804 assert_eq!(
805 serde_json::to_value(&exhausted).unwrap()["outcome"],
806 serde_json::json!("exhausted")
807 );
808
809 let lookup_err = DownloadDiagnosticsRecord::chunk_level(
810 2,
811 5,
812 &addr,
813 "initial",
814 Some(2),
815 0,
816 DownloadDiagnosticsOutcome::LookupError,
817 Some(bounded_error("lookup", "DHT returned no peers")),
818 );
819 let v = serde_json::to_value(&lookup_err).unwrap();
820 assert_eq!(v["outcome"], serde_json::json!("lookup_error"));
821 assert!(v["error"].as_str().unwrap().starts_with("lookup: "));
822 }
823
824 #[test]
825 fn bounded_error_truncates_long_detail() {
826 let long = "x".repeat(10_000);
827 let s = bounded_error("network", &long);
828 assert!(s.starts_with("network: "));
829 assert!(s.chars().count() <= DIAGNOSTICS_ERROR_MAX_CHARS);
831 assert!(s.ends_with('…'));
832 }
833
834 #[test]
835 fn bounded_error_preserves_short_detail_intact() {
836 let s = bounded_error("protocol", "mismatched address");
837 assert_eq!(s, "protocol: mismatched address");
838 }
839
840 #[test]
841 fn rfc3339_formatter_is_valid_for_known_epoch() {
842 let s = rfc3339_from_unix_secs(1_609_459_200);
844 assert_eq!(s, "2021-01-01T00:00:00Z");
845 assert_eq!(rfc3339_from_unix_secs(0), "1970-01-01T00:00:00Z");
847 assert_eq!(
849 rfc3339_from_unix_secs(1_709_164_800),
850 "2024-02-29T00:00:00Z"
851 );
852 }
853
854 #[test]
855 fn disabled_diagnostics_does_not_create_sidecar() {
856 let dir = tempfile::tempdir().unwrap();
857 let path = dir.path().join("disabled.jsonl");
858 let diagnostics: Option<DownloadDiagnosticsSender> = None;
859
860 assert!(diagnostics.is_none());
861 assert!(!path.exists());
862 }
863
864 #[test]
865 fn writer_emits_one_json_line_per_record_and_flushes_on_drop() {
866 let dir = std::env::temp_dir();
867 let path = dir.join(format!(
868 "ant-dl-diag-{}-{}.jsonl",
869 std::process::id(),
870 std::time::SystemTime::now()
871 .duration_since(UNIX_EPOCH)
872 .unwrap()
873 .as_nanos()
874 ));
875 let (sender, writer) = spawn_download_diagnostics_writer(&path).unwrap();
876 let addr = [3u8; 32];
877 sender.try_emit(DownloadDiagnosticsRecord::chunk_level(
878 1,
879 1,
880 &addr,
881 "initial",
882 Some(8),
883 128,
884 DownloadDiagnosticsOutcome::CacheHit,
885 None,
886 ));
887 sender.try_emit(DownloadDiagnosticsRecord::peer_attempt(
888 1,
889 1,
890 &addr,
891 "initial",
892 1,
893 Some(5),
894 "lookup-writer",
895 "peer-z",
896 vec!["/ip4/1.2.3.4/udp/9000/quic".to_string()],
897 vec!["direct".to_string()],
898 Some(50),
899 Some(100),
900 Some(1_234_000_000),
901 Some("peer-z"),
902 Some("/ip4/1.2.3.4/udp/9000/quic"),
903 "direct",
904 None,
905 Some(true),
906 Some(8),
907 Some(8),
908 30,
909 60,
910 &test_correlation(9_003),
911 30,
912 128,
913 DownloadDiagnosticsOutcome::Found,
914 None,
915 ));
916 drop(sender);
918 writer.join().unwrap();
919 let contents = std::fs::read_to_string(&path).unwrap();
920 let lines: Vec<&str> = contents.lines().filter(|l| !l.is_empty()).collect();
921 assert_eq!(
922 lines.len(),
923 2,
924 "expected 2 JSONL records, got: {contents:?}"
925 );
926 let first = serde_json::from_str::<serde_json::Value>(lines[0]).unwrap();
927 assert_eq!(first["outcome"], serde_json::json!("cache_hit"));
928 assert_eq!(first["schema_version"], serde_json::json!(4u8));
929 let second = serde_json::from_str::<serde_json::Value>(lines[1]).unwrap();
930 assert_eq!(second["outcome"], serde_json::json!("found"));
931 assert_eq!(second["route"], serde_json::json!("direct"));
932 assert_eq!(second["route_note"], serde_json::Value::Null);
933 assert_eq!(
934 second["transport_source"],
935 serde_json::json!("/ip4/1.2.3.4/udp/9000/quic")
936 );
937 let _ = std::fs::remove_file(&path);
938 }
939}