Skip to main content

detritus/
shipper.rs

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
26/// Default number of days to keep successfully sent crash entries.
27pub const DEFAULT_SENT_RETENTION_DAYS: u64 = 90;
28
29/// Reusable crash uploader with pooled HTTP connections and an in-memory token.
30///
31/// Convenience shipping functions create one uploader per scan. Keep this value
32/// to reuse its pool across scans, or supply a Reqwest client for private CAs,
33/// mutual TLS, proxy settings, and custom timeouts.
34#[derive(Debug, Clone)]
35pub struct CrashShipper {
36    client: reqwest::Client,
37    token: SecretString,
38    config: ShipConfig,
39}
40
41impl CrashShipper {
42    /// Creates an uploader with a 10-second connect timeout, 30-second read
43    /// timeout, and 60-second total request timeout.
44    ///
45    /// The offline spool owns retries; HTTP-level retries are disabled to avoid
46    /// recording a crash index more than once during one upload attempt.
47    ///
48    /// # Errors
49    ///
50    /// Returns [`ShipError::Http`] if the HTTP client cannot be initialized.
51    pub fn new(token: SecretString) -> Result<Self, ShipError> {
52        install_default_crypto_provider();
53        // https://docs.rs/reqwest/0.13.5/reqwest/struct.ClientBuilder.html
54        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    /// Uses a caller-configured HTTP client. Its timeout and retry policy apply.
64    #[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    /// Sets dump and attachment compression options.
74    #[must_use]
75    pub fn with_config(mut self, config: ShipConfig) -> Self {
76        self.config = config;
77        self
78    }
79
80    /// Ships pending entries to an explicit endpoint using this client's connection pool.
81    ///
82    /// # Errors
83    ///
84    /// See [`ship_pending_crashes_with_config`]. Failed entries remain pending.
85    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    /// Ships pending entries using their stored endpoints and retention policies.
101    ///
102    /// # Errors
103    ///
104    /// See [`ship_pending_crashes_using_stored_config_with_config`].
105    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/// Configuration knobs for crash-dump upload compression.
114///
115/// Created with [`ShipConfig::default`] and customised with the builder
116/// methods.  The defaults compress dumps at zstd level 19 and also
117/// compress text-ish attachments.
118#[derive(Debug, Clone)]
119pub struct ShipConfig {
120    /// Zstd level used for the `dump` part. Valid range is `1..=22`.
121    /// Defaults to [`DEFAULT_COMPRESSION_LEVEL`] (19).
122    pub dump_compression_level: i32,
123    /// When `true`, text-ish attachments are compressed with zstd before
124    /// upload. Defaults to `true`.
125    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    /// Sets the zstd compression level for the `dump` part.
139    #[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    /// Enables or disables zstd compression for text-ish attachment parts.
146    #[must_use]
147    pub fn with_attachment_compression(mut self, enabled: bool) -> Self {
148        self.attachment_compression = enabled;
149        self
150    }
151}
152
153/// Errors returned while shipping pending crash reports.
154#[derive(Debug, thiserror::Error)]
155#[non_exhaustive]
156pub enum ShipError {
157    /// Filesystem operation failed.
158    #[error("crash spool I/O error: {0}")]
159    Io(#[from] io::Error),
160    /// JSON metadata failed to parse or encode.
161    #[error("crash spool JSON error: {0}")]
162    Json(#[from] serde_json::Error),
163    /// HTTP client failed.
164    #[error("crash upload HTTP error: {0}")]
165    Http(#[from] reqwest::Error),
166    /// Bearer token contains bytes that are not valid in an HTTP header.
167    #[error("crash upload bearer token is not a valid HTTP header value")]
168    InvalidToken,
169    /// Server rejected the upload.
170    #[error("crash upload failed with status {0}")]
171    Status(reqwest::StatusCode),
172    /// Multipart encoding failed.
173    #[error("crash multipart encoding failed: {0}")]
174    Protocol(#[from] detritus_protocol::ProtocolError),
175    /// A spool entry has no usable stored upload config: `sdk-config.json` is
176    /// absent when required or its stored endpoint is not a valid URL.
177    #[error("entry {} has no usable stored upload config", .0.display())]
178    MissingStoredConfig(PathBuf),
179}
180
181/// Ships pending crash artifacts and moves successful entries to `sent/`.
182///
183/// A filesystem lock on `spool_dir/.lock` prevents two processes from scanning
184/// the same pending directory concurrently. One spool directory per process is
185/// still the recommended default.
186///
187/// Dump bytes and text-ish attachments are compressed with zstd
188/// (level [`DEFAULT_COMPRESSION_LEVEL`]) **before** the SHA-256 is computed,
189/// so content-addressed dedup is based on the compressed bytes.
190///
191/// # Errors
192///
193/// See [`ship_pending_crashes_with_config`].
194pub 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
202/// Like [`ship_pending_crashes`] but with explicit compression settings.
203///
204/// # Errors
205///
206/// Returns [`ShipError::Io`] on spool filesystem errors, [`ShipError::Json`] if
207/// stored metadata cannot be parsed, [`ShipError::Protocol`] if multipart
208/// encoding fails, [`ShipError::InvalidToken`] for malformed bearer headers,
209/// [`ShipError::Http`] if the upload request fails, or
210/// [`ShipError::Status`] if the server rejects an upload.
211pub 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
245/// Ships pending crash artifacts using each entry's stored endpoint.
246///
247/// The caller supplies the bearer token. The token is never read from the spool;
248/// each pending entry only contributes its stored endpoint and sent-entry
249/// retention from `sdk-config.json`.
250///
251/// A filesystem lock on `spool_dir/.lock` prevents two processes from scanning
252/// the same pending directory concurrently. The lock is held for the whole
253/// scan-and-ship run.
254///
255/// # Errors
256///
257/// Returns [`ShipError::Io`] on spool filesystem errors, [`ShipError::Json`] if
258/// stored metadata or stored upload config cannot be parsed,
259/// [`ShipError::MissingStoredConfig`] if a pending entry lacks a usable
260/// `sdk-config.json`, [`ShipError::Protocol`] if multipart encoding fails,
261/// [`ShipError::InvalidToken`] for malformed bearer headers,
262/// [`ShipError::Http`] if an upload request fails, or [`ShipError::Status`] if
263/// the server rejects an upload. The first error aborts the run.
264pub 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
272/// Like [`ship_pending_crashes_using_stored_config`] but with explicit
273/// compression settings.
274///
275/// Each pending entry is posted to the endpoint recovered from that entry's
276/// `sdk-config.json`. Successfully shipped entries are moved to `sent/`, then
277/// `sent/` is cleaned using the maximum stored retention across the entries
278/// shipped by this run. If no entry was shipped, cleanup uses
279/// [`DEFAULT_SENT_RETENTION_DAYS`].
280///
281/// # Errors
282///
283/// Returns [`ShipError::Io`] on spool filesystem errors, [`ShipError::Json`] if
284/// stored metadata or stored upload config cannot be parsed,
285/// [`ShipError::MissingStoredConfig`] if a pending entry lacks a usable
286/// `sdk-config.json`, [`ShipError::Protocol`] if multipart encoding fails,
287/// [`ShipError::Http`] if an upload request fails, or [`ShipError::Status`] if
288/// the server rejects an upload. The first error aborts the run.
289pub 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
386/// Applies zstd compression to the envelope's dump (and optionally its
387/// attachments), then serialises to multipart and POSTs it.
388///
389/// Compressed bytes replace the originals in-place inside a locally-owned
390/// clone so the original `CrashEnvelope` is not mutated.
391async 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    // Build a locally-owned envelope with compressed bytes plus the matching
401    // EnvelopeEncodings that records which parts are zstd-encoded.
402    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    // Compress the dump unconditionally.
409    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    // Optionally compress text-ish attachments.
416    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            // Find the matching attachment in our local clone and replace bytes.
422            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        // HTTP/1 connections return to the pool only after the response body is
465        // consumed. Stream and discard it to keep memory bounded; the client's
466        // read and total timeouts also apply while draining.
467        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    // saturating_mul avoids an overflow panic if retention_days is misconfigured
486    // to an absurd value; the checked_sub below then saturates the cutoff to
487    // UNIX_EPOCH, i.e. "keep everything", which is the intended extreme behavior.
488    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
523/// Resolve a spool entry's stored upload endpoint and retention.
524///
525/// Returns the parsed endpoint URL and retention, or an error if the entry lacks
526/// a usable `sdk-config.json`.
527fn 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}