detritus_protocol/
multipart.rs1use 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
14pub const DEFAULT_BOUNDARY: &str = "detritus-boundary-v1";
16const CRLF: &str = "\r\n";
17
18#[derive(Debug, Clone, Default)]
20pub struct PartEncoding {
21 pub content_encoding: Option<String>,
25}
26
27impl CrashEnvelope {
28 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 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 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 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 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#[derive(Debug, Clone, Default)]
205pub struct EnvelopeEncodings {
206 pub dump: PartEncoding,
208 pub attachments: Vec<PartEncoding>,
211}
212
213impl EnvelopeEncodings {
214 #[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}