1use std::{
2 fs, io,
3 path::{Path, PathBuf},
4 time::{Duration, SystemTime},
5};
6
7use detritus_protocol::{
8 CrashAttachment, CrashEnvelope, CrashMetadata,
9 multipart::{EnvelopeEncodings, PartEncoding},
10};
11use reqwest::header::{AUTHORIZATION, CONTENT_TYPE, HeaderValue};
12use secrecy::{ExposeSecret, SecretString};
13use url::Url;
14
15use crate::{
16 compression::{DEFAULT_COMPRESSION_LEVEL, compress, should_compress_content_type},
17 install_default_crypto_provider,
18 panic_hook::StoredUploadConfig,
19 spool::SpoolLock,
20};
21
22#[cfg(test)]
23#[path = "shipper/tests.rs"]
24mod regression_tests;
25
26pub const DEFAULT_SENT_RETENTION_DAYS: u64 = 90;
28
29#[derive(Debug, Clone)]
35pub struct CrashShipper {
36 client: reqwest::Client,
37 token: SecretString,
38 config: ShipConfig,
39}
40
41impl CrashShipper {
42 pub fn new(token: SecretString) -> Result<Self, ShipError> {
52 install_default_crypto_provider();
53 let client = reqwest::Client::builder()
55 .connect_timeout(Duration::from_secs(10))
56 .read_timeout(Duration::from_secs(30))
57 .timeout(Duration::from_secs(60))
58 .retry(reqwest::retry::never())
59 .build()?;
60 Ok(Self::with_client(client, token))
61 }
62
63 #[must_use]
65 pub fn with_client(client: reqwest::Client, token: SecretString) -> Self {
66 Self {
67 client,
68 token,
69 config: ShipConfig::default(),
70 }
71 }
72
73 #[must_use]
75 pub fn with_config(mut self, config: ShipConfig) -> Self {
76 self.config = config;
77 self
78 }
79
80 pub async fn ship_pending(
86 &self,
87 spool_dir: impl AsRef<Path>,
88 endpoint: Url,
89 ) -> Result<usize, ShipError> {
90 ship_pending_with_client(
91 spool_dir.as_ref(),
92 &endpoint,
93 &self.token,
94 &self.config,
95 &self.client,
96 )
97 .await
98 }
99
100 pub async fn ship_using_stored_config(
106 &self,
107 spool_dir: impl AsRef<Path>,
108 ) -> Result<usize, ShipError> {
109 ship_stored_with_client(spool_dir.as_ref(), &self.token, &self.config, &self.client).await
110 }
111}
112
113#[derive(Debug, Clone)]
119pub struct ShipConfig {
120 pub dump_compression_level: i32,
123 pub attachment_compression: bool,
126}
127
128impl Default for ShipConfig {
129 fn default() -> Self {
130 Self {
131 dump_compression_level: DEFAULT_COMPRESSION_LEVEL,
132 attachment_compression: true,
133 }
134 }
135}
136
137impl ShipConfig {
138 #[must_use]
140 pub fn with_dump_compression_level(mut self, level: i32) -> Self {
141 self.dump_compression_level = level;
142 self
143 }
144
145 #[must_use]
147 pub fn with_attachment_compression(mut self, enabled: bool) -> Self {
148 self.attachment_compression = enabled;
149 self
150 }
151}
152
153#[derive(Debug, thiserror::Error)]
155#[non_exhaustive]
156pub enum ShipError {
157 #[error("crash spool I/O error: {0}")]
159 Io(#[from] io::Error),
160 #[error("crash spool JSON error: {0}")]
162 Json(#[from] serde_json::Error),
163 #[error("crash upload HTTP error: {0}")]
165 Http(#[from] reqwest::Error),
166 #[error("crash upload bearer token is not a valid HTTP header value")]
168 InvalidToken,
169 #[error("crash upload failed with status {0}")]
171 Status(reqwest::StatusCode),
172 #[error("crash multipart encoding failed: {0}")]
174 Protocol(#[from] detritus_protocol::ProtocolError),
175 #[error("entry {} has no usable stored upload config", .0.display())]
178 MissingStoredConfig(PathBuf),
179}
180
181pub async fn ship_pending_crashes(
195 spool_dir: impl AsRef<Path>,
196 endpoint: Url,
197 token: SecretString,
198) -> Result<usize, ShipError> {
199 ship_pending_crashes_with_config(spool_dir, endpoint, token, ShipConfig::default()).await
200}
201
202pub async fn ship_pending_crashes_with_config(
212 spool_dir: impl AsRef<Path>,
213 endpoint: Url,
214 token: SecretString,
215 config: ShipConfig,
216) -> Result<usize, ShipError> {
217 CrashShipper::new(token)?
218 .with_config(config)
219 .ship_pending(spool_dir, endpoint)
220 .await
221}
222
223async fn ship_pending_with_client(
224 spool_dir: &Path,
225 endpoint: &Url,
226 token: &SecretString,
227 config: &ShipConfig,
228 client: &reqwest::Client,
229) -> Result<usize, ShipError> {
230 let _lock = SpoolLock::acquire(spool_dir)?;
231 let pending = spool_dir.join("pending");
232 let sent = spool_dir.join("sent");
233 fs::create_dir_all(&pending)?;
234 fs::create_dir_all(&sent)?;
235 cleanup_sent(&sent, DEFAULT_SENT_RETENTION_DAYS)?;
236
237 let mut shipped = 0;
238 for entry in pending_entries(&pending)? {
239 ship_one(&entry, &sent, endpoint, token, config, client).await?;
240 shipped += 1;
241 }
242 Ok(shipped)
243}
244
245pub async fn ship_pending_crashes_using_stored_config(
265 spool_dir: impl AsRef<Path>,
266 token: SecretString,
267) -> Result<usize, ShipError> {
268 ship_pending_crashes_using_stored_config_with_config(spool_dir, token, ShipConfig::default())
269 .await
270}
271
272pub async fn ship_pending_crashes_using_stored_config_with_config(
290 spool_dir: impl AsRef<Path>,
291 token: SecretString,
292 config: ShipConfig,
293) -> Result<usize, ShipError> {
294 CrashShipper::new(token)?
295 .with_config(config)
296 .ship_using_stored_config(spool_dir)
297 .await
298}
299
300async fn ship_stored_with_client(
301 spool_dir: &Path,
302 token: &SecretString,
303 config: &ShipConfig,
304 client: &reqwest::Client,
305) -> Result<usize, ShipError> {
306 let _lock = SpoolLock::acquire(spool_dir)?;
307 let pending = spool_dir.join("pending");
308 let sent = spool_dir.join("sent");
309 fs::create_dir_all(&pending)?;
310 fs::create_dir_all(&sent)?;
311
312 let mut shipped = 0;
313 let mut max_retention_days = None;
314 for entry in pending_entries(&pending)? {
315 let (endpoint, retention_days) = resolve_stored_endpoint(&entry)?;
316 ship_one(&entry, &sent, &endpoint, token, config, client).await?;
317 record_retention(&mut max_retention_days, retention_days);
318 shipped += 1;
319 }
320 cleanup_sent(&sent, retention_or_default(max_retention_days))?;
321 Ok(shipped)
322}
323
324fn pending_entries(pending: &Path) -> io::Result<Vec<PathBuf>> {
325 let mut entries = match fs::read_dir(pending) {
326 Ok(entries) => entries
327 .filter_map(Result::ok)
328 .map(|entry| entry.path())
329 .filter(|path| path.is_dir())
330 .collect::<Vec<_>>(),
331 Err(error) if error.kind() == io::ErrorKind::NotFound => Vec::new(),
332 Err(error) => return Err(error),
333 };
334 entries.sort();
335 Ok(entries)
336}
337
338async fn ship_one(
339 entry: &Path,
340 sent: &Path,
341 endpoint: &Url,
342 token: &SecretString,
343 config: &ShipConfig,
344 client: &reqwest::Client,
345) -> Result<(), ShipError> {
346 let envelope = read_envelope(entry)?;
347 post_envelope(endpoint, token, &envelope, config, client).await?;
348 let destination = sent.join(
349 entry
350 .file_name()
351 .ok_or_else(|| io::Error::other("pending entry has no file name"))?,
352 );
353 if destination.exists() {
354 fs::remove_dir_all(&destination)?;
355 }
356 fs::rename(entry, destination)?;
357 Ok(())
358}
359
360fn read_envelope(entry: &Path) -> Result<CrashEnvelope, ShipError> {
361 let metadata: CrashMetadata = serde_json::from_slice(&fs::read(entry.join("metadata.json"))?)?;
362 let dump = fs::read(entry.join("dump.bin"))?;
363 let attachments = metadata
364 .attachments
365 .iter()
366 .filter_map(|manifest| {
367 let filename = manifest.filename.as_ref()?;
368 let path = entry.join(filename);
369 Some((manifest, path))
370 })
371 .map(|(manifest, path)| {
372 Ok(CrashAttachment {
373 key: manifest.key.clone(),
374 content_type: manifest.content_type.clone(),
375 bytes: fs::read(path)?,
376 })
377 })
378 .collect::<io::Result<Vec<_>>>()?;
379 Ok(CrashEnvelope {
380 metadata,
381 dump,
382 attachments,
383 })
384}
385
386async fn post_envelope(
392 endpoint: &Url,
393 token: &SecretString,
394 envelope: &CrashEnvelope,
395 config: &ShipConfig,
396 client: &reqwest::Client,
397) -> Result<(), ShipError> {
398 install_default_crypto_provider();
399
400 let mut compressed_envelope = envelope.clone();
403 let mut encodings = EnvelopeEncodings {
404 dump: PartEncoding::default(),
405 attachments: Vec::with_capacity(envelope.attachments.len()),
406 };
407
408 let compressed_dump = compress(&envelope.dump, config.dump_compression_level)?;
410 compressed_envelope.dump = compressed_dump;
411 encodings.dump = PartEncoding {
412 content_encoding: Some("zstd".to_owned()),
413 };
414
415 for attachment in &envelope.attachments {
417 let should =
418 config.attachment_compression && should_compress_content_type(&attachment.content_type);
419 if should {
420 let compressed = compress(&attachment.bytes, config.dump_compression_level)?;
421 if let Some(local) = compressed_envelope
423 .attachments
424 .iter_mut()
425 .find(|a| a.key == attachment.key)
426 {
427 local.bytes = compressed;
428 }
429 encodings.attachments.push(PartEncoding {
430 content_encoding: Some("zstd".to_owned()),
431 });
432 } else {
433 encodings.attachments.push(PartEncoding::default());
434 }
435 }
436
437 let mut body = Vec::new();
438 compressed_envelope
439 .write_to_with_boundary_and_encodings(
440 &mut body,
441 detritus_protocol::multipart::DEFAULT_BOUNDARY,
442 &encodings,
443 )
444 .await?;
445
446 let url = crash_url(endpoint);
447 let mut authorization = HeaderValue::from_str(&format!("Bearer {}", token.expose_secret()))
448 .map_err(|_| ShipError::InvalidToken)?;
449 authorization.set_sensitive(true);
450 let mut response = client
451 .post(url)
452 .header(AUTHORIZATION, authorization)
453 .header(
454 CONTENT_TYPE,
455 format!(
456 "multipart/form-data; boundary={}",
457 detritus_protocol::multipart::DEFAULT_BOUNDARY
458 ),
459 )
460 .body(body)
461 .send()
462 .await?;
463 if response.status().is_success() {
464 while response.chunk().await?.is_some() {}
468 Ok(())
469 } else {
470 Err(ShipError::Status(response.status()))
471 }
472}
473
474fn crash_url(endpoint: &Url) -> Url {
475 if endpoint.path().ends_with("/v1/crashes") {
476 endpoint.clone()
477 } else {
478 endpoint
479 .join("/v1/crashes")
480 .unwrap_or_else(|_| endpoint.clone())
481 }
482}
483
484fn cleanup_sent(sent: &Path, retention_days: u64) -> io::Result<()> {
485 let cutoff = SystemTime::now()
489 .checked_sub(Duration::from_secs(
490 retention_days.saturating_mul(24 * 60 * 60),
491 ))
492 .unwrap_or(SystemTime::UNIX_EPOCH);
493 for entry in pending_entries(sent)? {
494 let modified = entry.metadata()?.modified().unwrap_or(SystemTime::now());
495 if modified < cutoff {
496 fs::remove_dir_all(entry)?;
497 }
498 }
499 Ok(())
500}
501
502fn record_retention(max_retention_days: &mut Option<u64>, retention_days: u64) {
503 *max_retention_days = Some(
504 max_retention_days
505 .map(|current| current.max(retention_days))
506 .unwrap_or(retention_days),
507 );
508}
509
510fn retention_or_default(max_retention_days: Option<u64>) -> u64 {
511 max_retention_days.unwrap_or(DEFAULT_SENT_RETENTION_DAYS)
512}
513
514fn read_stored_upload_config(entry: &Path) -> Result<Option<StoredUploadConfig>, ShipError> {
515 let path = entry.join("sdk-config.json");
516 if path.exists() {
517 Ok(Some(serde_json::from_slice(&fs::read(path)?)?))
518 } else {
519 Ok(None)
520 }
521}
522
523fn resolve_stored_endpoint(entry: &Path) -> Result<(Url, u64), ShipError> {
528 let config = read_stored_upload_config(entry)?
529 .ok_or_else(|| ShipError::MissingStoredConfig(entry.to_path_buf()))?;
530 let endpoint = Url::parse(&config.endpoint)
531 .map_err(|_| ShipError::MissingStoredConfig(entry.to_path_buf()))?;
532 Ok((endpoint, config.sent_retention_days))
533}
534
535#[cfg(test)]
536mod tests {
537 use super::*;
538
539 #[cfg_attr(coverage_nightly, coverage(off))]
540 fn write_stored_config(
541 entry: &Path,
542 endpoint: &str,
543 sent_retention_days: u64,
544 ) -> Result<(), ShipError> {
545 fs::write(
546 entry.join("sdk-config.json"),
547 serde_json::to_vec_pretty(&StoredUploadConfig {
548 endpoint: endpoint.to_owned(),
549 sent_retention_days,
550 })?,
551 )?;
552 Ok(())
553 }
554
555 #[test]
556 #[cfg_attr(coverage_nightly, coverage(off))]
557 fn read_stored_upload_config_round_trips_producer_format() {
558 let temp = tempfile::tempdir().expect("temp dir is created");
559 write_stored_config(temp.path(), "https://crashes.example.test/upload", 14)
560 .expect("stored config is written");
561
562 let config = read_stored_upload_config(temp.path())
563 .expect("stored config is readable")
564 .expect("stored config is present");
565
566 assert_eq!(config.endpoint, "https://crashes.example.test/upload");
567 assert_eq!(config.sent_retention_days, 14);
568 }
569
570 #[test]
571 #[cfg_attr(coverage_nightly, coverage(off))]
572 fn read_stored_upload_config_returns_none_when_missing() {
573 let temp = tempfile::tempdir().expect("temp dir is created");
574
575 let config =
576 read_stored_upload_config(temp.path()).expect("missing config is not an error");
577
578 assert!(config.is_none());
579 }
580
581 #[test]
582 #[cfg_attr(coverage_nightly, coverage(off))]
583 fn read_stored_upload_config_reports_malformed_json() {
584 let temp = tempfile::tempdir().expect("temp dir is created");
585 fs::write(temp.path().join("sdk-config.json"), b"{not-json")
586 .expect("malformed config is written");
587
588 assert!(matches!(
589 read_stored_upload_config(temp.path()),
590 Err(ShipError::Json(_))
591 ));
592 }
593
594 #[test]
595 #[cfg_attr(coverage_nightly, coverage(off))]
596 fn resolve_stored_endpoint_reports_missing_config() {
597 let temp = tempfile::tempdir().expect("temp dir is created");
598
599 assert!(matches!(
600 resolve_stored_endpoint(temp.path()),
601 Err(ShipError::MissingStoredConfig(path)) if path == temp.path()
602 ));
603 }
604
605 #[test]
606 #[cfg_attr(coverage_nightly, coverage(off))]
607 fn resolve_stored_endpoint_returns_url_and_retention() {
608 let temp = tempfile::tempdir().expect("temp dir is created");
609 write_stored_config(temp.path(), "https://crashes.example.test/base", 30)
610 .expect("stored config is written");
611
612 let (endpoint, retention_days) =
613 resolve_stored_endpoint(temp.path()).expect("stored endpoint is usable");
614
615 assert_eq!(endpoint.as_str(), "https://crashes.example.test/base");
616 assert_eq!(retention_days, 30);
617 }
618
619 #[test]
620 #[cfg_attr(coverage_nightly, coverage(off))]
621 fn resolve_stored_endpoint_reports_invalid_url() {
622 let temp = tempfile::tempdir().expect("temp dir is created");
623 write_stored_config(temp.path(), "not a url", 30).expect("stored config is written");
624
625 assert!(matches!(
626 resolve_stored_endpoint(temp.path()),
627 Err(ShipError::MissingStoredConfig(path)) if path == temp.path()
628 ));
629 }
630
631 #[test]
632 #[cfg_attr(coverage_nightly, coverage(off))]
633 fn record_retention_keeps_maximum_value() {
634 let mut retention = None;
635
636 record_retention(&mut retention, 7);
637 record_retention(&mut retention, 30);
638 record_retention(&mut retention, 0);
639
640 assert_eq!(retention_or_default(retention), 30);
641 }
642
643 #[test]
644 #[cfg_attr(coverage_nightly, coverage(off))]
645 fn retention_or_default_falls_back_when_no_entry_was_shipped() {
646 assert_eq!(retention_or_default(None), DEFAULT_SENT_RETENTION_DAYS);
647 }
648}