Skip to main content

detritus_protocol/
multipart.rs

1//! Feature-gated multipart helpers for crash envelopes.
2
3use std::fmt::Write as _;
4
5use bytes::Bytes;
6use futures_util::stream;
7use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
8
9use crate::{
10    PROTOCOL_VERSION,
11    crash::{CrashAttachment, CrashEnvelope, CrashMetadata, ProtocolError},
12};
13
14/// Default multipart boundary used by tests and simple SDK callers.
15pub const DEFAULT_BOUNDARY: &str = "detritus-boundary-v1";
16const CRLF: &str = "\r\n";
17
18/// Per-part encoding options for multipart crash upload.
19#[derive(Debug, Clone, Default)]
20pub struct PartEncoding {
21    /// Value of the `Content-Encoding` header to add to this part, e.g. `"zstd"`.
22    ///
23    /// `None` means no `Content-Encoding` header is emitted (plain bytes).
24    pub content_encoding: Option<String>,
25}
26
27impl CrashEnvelope {
28    /// Writes the envelope as an RFC 7578 multipart body with [`DEFAULT_BOUNDARY`].
29    ///
30    /// # Errors
31    ///
32    /// Returns [`ProtocolError::Json`] if the metadata fails to serialize, or
33    /// [`ProtocolError::Io`] if writing to `writer` fails.
34    pub async fn write_to<W>(&self, writer: &mut W) -> Result<(), ProtocolError>
35    where
36        W: AsyncWrite + Unpin,
37    {
38        self.write_to_with_boundary(writer, DEFAULT_BOUNDARY).await
39    }
40
41    /// Writes the envelope as an RFC 7578 multipart body.
42    ///
43    /// # Errors
44    ///
45    /// As [`write_to`](Self::write_to): [`ProtocolError::Json`] on metadata
46    /// serialization failure, [`ProtocolError::Io`] on writer failure.
47    pub async fn write_to_with_boundary<W>(
48        &self,
49        writer: &mut W,
50        boundary: &str,
51    ) -> Result<(), ProtocolError>
52    where
53        W: AsyncWrite + Unpin,
54    {
55        let encodings = EnvelopeEncodings::none();
56        self.write_to_with_boundary_and_encodings(writer, boundary, &encodings)
57            .await
58    }
59
60    /// Writes the envelope as an RFC 7578 multipart body, applying per-part
61    /// `Content-Encoding` headers as specified by `encodings`.
62    ///
63    /// # Errors
64    ///
65    /// Returns [`ProtocolError::Json`] if the metadata fails to serialize, or
66    /// [`ProtocolError::Io`] if writing to `writer` fails.
67    pub async fn write_to_with_boundary_and_encodings<W>(
68        &self,
69        writer: &mut W,
70        boundary: &str,
71        encodings: &EnvelopeEncodings,
72    ) -> Result<(), ProtocolError>
73    where
74        W: AsyncWrite + Unpin,
75    {
76        let metadata = serde_json::to_vec(&self.metadata)?;
77        write_part(
78            writer,
79            boundary,
80            "metadata",
81            "application/json",
82            None,
83            &metadata,
84        )
85        .await?;
86        write_part(
87            writer,
88            boundary,
89            "dump",
90            "application/octet-stream",
91            encodings.dump.content_encoding.as_deref(),
92            &self.dump,
93        )
94        .await?;
95        for (idx, attachment) in self.attachments.iter().enumerate() {
96            let name = format!("attach:{}", attachment.key);
97            let encoding = encodings
98                .attachments
99                .get(idx)
100                .and_then(|e| e.content_encoding.as_deref());
101            write_part(
102                writer,
103                boundary,
104                &name,
105                &attachment.content_type,
106                encoding,
107                &attachment.bytes,
108            )
109            .await?;
110        }
111        writer
112            .write_all(format!("--{boundary}--{CRLF}").as_bytes())
113            .await?;
114        writer.flush().await?;
115        Ok(())
116    }
117
118    /// Parses an envelope from a body using [`DEFAULT_BOUNDARY`].
119    ///
120    /// # Errors
121    ///
122    /// Returns [`ProtocolError::Multipart`] if the body is not valid multipart,
123    /// [`ProtocolError::MissingPart`] if the required `metadata` or `dump` part
124    /// is absent, [`ProtocolError::InvalidPartName`] or
125    /// [`ProtocolError::InvalidMultipart`] for a malformed part name or payload,
126    /// or [`ProtocolError::Json`] if the metadata part is not valid JSON.
127    pub async fn read_from<R>(reader: &mut R) -> Result<Self, ProtocolError>
128    where
129        R: AsyncRead + Unpin,
130    {
131        Self::read_from_with_boundary(reader, DEFAULT_BOUNDARY).await
132    }
133
134    /// Parses an envelope from a body with the supplied multipart boundary.
135    ///
136    /// # Errors
137    ///
138    /// Same as [`read_from`](Self::read_from).
139    pub async fn read_from_with_boundary<R>(
140        reader: &mut R,
141        boundary: &str,
142    ) -> Result<Self, ProtocolError>
143    where
144        R: AsyncRead + Unpin,
145    {
146        let mut body = Vec::new();
147        reader.read_to_end(&mut body).await?;
148        let body = Bytes::from(body);
149        let stream = stream::once(async move { Ok::<Bytes, std::io::Error>(body) });
150        let mut multipart = multer::Multipart::new(stream, boundary);
151        let mut metadata = None;
152        let mut dump = None;
153        let mut attachments = Vec::new();
154
155        while let Some(field) = multipart.next_field().await? {
156            let name = field
157                .name()
158                .ok_or(ProtocolError::InvalidPartName)?
159                .to_owned();
160            let content_type = field.content_type().map_or_else(
161                || "application/octet-stream".to_owned(),
162                ToString::to_string,
163            );
164            let bytes = field.bytes().await?.to_vec();
165            match name.as_str() {
166                "metadata" => metadata = Some(serde_json::from_slice::<CrashMetadata>(&bytes)?),
167                "dump" => dump = Some(bytes),
168                name if name.starts_with("attach:") => {
169                    let key = name
170                        .strip_prefix("attach:")
171                        .ok_or_else(|| ProtocolError::InvalidMultipart(name.to_owned()))?
172                        .to_owned();
173                    attachments.push(CrashAttachment {
174                        key,
175                        content_type,
176                        bytes,
177                    });
178                }
179                _ => return Err(ProtocolError::InvalidMultipart(name)),
180            }
181        }
182
183        let metadata = metadata.ok_or(ProtocolError::MissingPart("metadata"))?;
184        if metadata.schema_version != PROTOCOL_VERSION {
185            return Err(ProtocolError::InvalidMultipart(format!(
186                "schema version {} does not match protocol version {}",
187                metadata.schema_version, PROTOCOL_VERSION
188            )));
189        }
190        let dump = dump.ok_or(ProtocolError::MissingPart("dump"))?;
191        Ok(Self {
192            metadata,
193            dump,
194            attachments,
195        })
196    }
197}
198
199/// Per-part encoding settings for an entire [`CrashEnvelope`].
200///
201/// `attachments[i]` corresponds to `envelope.attachments[i]`. If the slice is
202/// shorter than the attachment list the remaining attachments are written with
203/// no `Content-Encoding` header.
204#[derive(Debug, Clone, Default)]
205pub struct EnvelopeEncodings {
206    /// Encoding for the `dump` part.
207    pub dump: PartEncoding,
208    /// Encodings for each `attach:<key>` part, in the same order as
209    /// [`CrashEnvelope::attachments`].
210    pub attachments: Vec<PartEncoding>,
211}
212
213impl EnvelopeEncodings {
214    /// Returns an `EnvelopeEncodings` with no `Content-Encoding` set on any part.
215    #[must_use]
216    pub fn none() -> Self {
217        Self::default()
218    }
219}
220
221async fn write_part<W>(
222    writer: &mut W,
223    boundary: &str,
224    name: &str,
225    content_type: &str,
226    content_encoding: Option<&str>,
227    bytes: &[u8],
228) -> Result<(), ProtocolError>
229where
230    W: AsyncWrite + Unpin,
231{
232    writer
233        .write_all(format!("--{boundary}{CRLF}").as_bytes())
234        .await?;
235    let mut headers = format!(
236        "Content-Disposition: form-data; name=\"{name}\"{CRLF}Content-Type: {content_type}{CRLF}"
237    );
238    if let Some(encoding) = content_encoding {
239        let _ = write!(headers, "Content-Encoding: {encoding}{CRLF}");
240    }
241    headers.push_str(CRLF);
242    writer.write_all(headers.as_bytes()).await?;
243    writer.write_all(bytes).await?;
244    writer.write_all(CRLF.as_bytes()).await?;
245    Ok(())
246}