toolkit/http/multipart.rs
1//! `multipart/mixed` streaming responses — the server counterpart to
2//! `toolkit_contract::runtime::multipart`.
3//!
4//! Sits beside [`super::sse`] as the second wire framing for a long-lived
5//! streaming endpoint: one JSON item per body part, for server-to-server
6//! consumers, where SSE serves browser-direct `EventSource`-style ones.
7//!
8//! This is streaming *response* framing only. `multipart/form-data` requests,
9//! mixed per-part content types and `Content-Disposition` are all out of scope,
10//! and the `boundary=` parameter is generated here rather than declared by the
11//! author.
12
13use axum::body::{Body, Bytes};
14use axum::response::{IntoResponse, Response};
15use futures_core::Stream;
16use futures_util::StreamExt as _;
17use http::HeaderValue;
18use http::header::CONTENT_TYPE;
19use serde::Serialize;
20use toolkit_canonical_errors::Problem;
21use toolkit_contract::runtime::multipart::MAX_ACCUMULATED_BYTES;
22
23/// RFC 2046 caps a boundary at 70 characters.
24const MAX_BOUNDARY_LEN: usize = 70;
25
26/// Media type of a typed **error part**. A post-open domain error is framed as
27/// a normal part whose `Content-Type` is this rather than `application/json`,
28/// so the reader can decode its body as an RFC 9457 [`Problem`] and surface a
29/// typed `Err`, instead of the whole body aborting. Kept byte-for-byte in sync
30/// with the reader's matcher in `toolkit_contract::runtime::multipart`.
31const PROBLEM_CONTENT_TYPE: &str = "application/problem+json";
32
33/// Largest serialized item this framer will emit as a single part.
34///
35/// Bound to the paired reader's accumulation ceiling
36/// ([`toolkit_contract::runtime::multipart::MAX_ACCUMULATED_BYTES`]) — the
37/// reader is the source of truth. The reader rejects any part whose body
38/// exceeds it, so emitting a larger part would produce a stream the toolkit
39/// client cannot consume; the framer aborts the body instead (see [`frame`]).
40const MAX_PART_BYTES: usize = MAX_ACCUMULATED_BYTES;
41
42/// A `multipart/mixed` streaming response, one JSON part per item.
43///
44/// Emits, per item: `--<boundary>CRLF`, `Content-Type: application/json`,
45/// `Content-Length: <n>`, `CRLF`, the JSON body, `CRLF`; and
46/// `--<boundary>--CRLF` on a clean end of the source stream.
47///
48/// The per-part `Content-Length` is what lets a reader emit each part from that
49/// part's own bytes instead of waiting for the following delimiter — so a live
50/// stream is never delivered a part behind. It is a fast path, not a protocol
51/// fork: RFC 2046 §5.1 makes the delimiter authoritative, so a strict parser
52/// that ignores the header still reads the stream correctly.
53///
54/// # Post-open domain error
55///
56/// The wrapped stream yields `Result<T, E>`. Once the `200` is on the wire a
57/// domain failure can no longer become an HTTP status, so an `Err(e)` is framed
58/// as a typed **error part**: one `application/problem+json` part carrying
59/// `Into::<Problem>::into(e)`, immediately followed by the close delimiter. The
60/// stream therefore ends *cleanly* and the reader surfaces the error as a typed
61/// `Err` item, not a truncation — an `Err` is terminal, so any later source
62/// items are not framed.
63///
64/// If the error `Problem` itself would exceed [`MAX_PART_BYTES`], its unbounded
65/// fields (`detail`, `context`) are dropped so the typed error — status, type,
66/// `error_code` — still reaches the client; only a `Problem` that cannot fit or
67/// serialize even trimmed falls back to an abort.
68///
69/// # Mid-stream serialization failure (genuine transport fault)
70///
71/// Distinct from the above: an `Ok(item)` that will not *serialize*, or one
72/// larger than [`MAX_PART_BYTES`], is not a domain error — it is a value the
73/// framer cannot put on the wire at all. There is no status left to send, so
74/// the response **body is aborted** and the failure is logged at `error`.
75///
76/// "Abort" specifically means yielding an error into the body stream, which
77/// truncates the chunked encoding and makes the reader report a transport
78/// error — *not* ending the stream, which would be a graceful EOF and which
79/// the reader treats as a clean end. A silently short stream is worse than a
80/// broken connection: the consumer's reopen path already handles an ungraceful
81/// close, but it cannot detect a stream that simply stopped early.
82pub struct MultipartJsonStream<S> {
83 stream: S,
84 boundary: String,
85 /// Prebuilt `Content-Type`, so rendering the response has no fallible
86 /// step — the boundary is validated once, on the way in.
87 content_type: HeaderValue,
88}
89
90/// A caller-supplied boundary was rejected by
91/// [`MultipartJsonStream::with_boundary`].
92///
93/// Returned rather than silently substituted so the caller decides whether to
94/// fall back to a generated boundary ([`MultipartJsonStream::new`]) or surface
95/// the failure.
96#[derive(Debug, Clone, PartialEq, Eq)]
97pub struct BoundaryError {
98 boundary: String,
99 reason: &'static str,
100}
101
102impl BoundaryError {
103 /// Why the boundary was rejected.
104 #[must_use]
105 pub fn reason(&self) -> &str {
106 self.reason
107 }
108}
109
110impl std::fmt::Display for BoundaryError {
111 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
112 // `escape_debug` so control characters in a rejected value are escaped
113 // rather than written raw into the message.
114 write!(
115 f,
116 "invalid multipart/mixed boundary \"{}\": {}",
117 self.boundary.escape_debug(),
118 self.reason
119 )
120 }
121}
122
123impl std::error::Error for BoundaryError {}
124
125impl<S> MultipartJsonStream<S> {
126 /// Wrap a stream of serializable items, with a generated boundary.
127 #[must_use]
128 pub fn new(stream: S) -> Self {
129 Self::with_validated_boundary(stream, generated_boundary())
130 }
131
132 /// Wrap a stream of serializable items, with a caller-chosen boundary.
133 ///
134 /// The boundary must be RFC 2046-legal *and* an HTTP token: 1–70 characters
135 /// drawn from `A-Z a-z 0-9 ' + - . _` (see [`validate_boundary`]).
136 ///
137 /// Returns [`BoundaryError`] if the boundary is rejected, leaving the
138 /// fallback decision to the caller — typically [`MultipartJsonStream::new`]
139 /// for a generated boundary — rather than silently substituting one, which
140 /// would hide the caller's mistake behind a boundary they never chose.
141 ///
142 /// The boundary must also not occur inside any item's JSON. A generated
143 /// boundary makes that impossible in practice; a caller-chosen one makes it
144 /// the caller's responsibility.
145 ///
146 /// # Errors
147 ///
148 /// Returns [`BoundaryError`] when `boundary` is empty, exceeds 70
149 /// characters, or contains a character outside the HTTP token set.
150 pub fn with_boundary(stream: S, boundary: impl Into<String>) -> Result<Self, BoundaryError> {
151 let boundary = boundary.into();
152 match validate_boundary(&boundary) {
153 Ok(()) => Ok(Self::with_validated_boundary(stream, boundary)),
154 Err(reason) => Err(BoundaryError { boundary, reason }),
155 }
156 }
157
158 fn with_validated_boundary(stream: S, boundary: String) -> Self {
159 if let Some(content_type) = content_type_for(&boundary) {
160 return Self {
161 stream,
162 boundary,
163 content_type,
164 };
165 }
166 // Unreachable: a validated boundary is an HTTP token, and `new` uses a
167 // hex boundary — both are legal header values. If that invariant is ever
168 // broken we must NOT emit the old silent fallback (a bare
169 // `multipart/mixed` with no `boundary=`), which disagreed with the
170 // `--<boundary>` delimiters the framer still writes and produced an
171 // unparseable body. Substitute a known-good boundary and its matching
172 // header together so the two never desync. Panic-free: `expect`/`unwrap`
173 // are denied.
174 tracing::error!(
175 boundary = ?boundary,
176 "validated multipart/mixed boundary did not form a legal header value; substituting a known-good boundary"
177 );
178 Self {
179 stream,
180 boundary: FALLBACK_BOUNDARY.to_owned(),
181 content_type: HeaderValue::from_static(FALLBACK_CONTENT_TYPE),
182 }
183 }
184
185 /// The boundary this response will frame its parts with.
186 #[must_use]
187 pub fn boundary(&self) -> &str {
188 &self.boundary
189 }
190}
191
192// The error bound is `Into<Problem>` by deliberate design, not incidentally: the
193// paired reader (`toolkit_contract::runtime::multipart`) decodes an error part
194// as an RFC 9457 `Problem`, and `Problem` is the platform's single canonical
195// wire error across REST, SSE and gRPC. This framer and that reader are two
196// halves of one protocol and must agree on the error envelope, so the framer
197// speaks `Problem` rather than being generic over the error format.
198impl<S, T, E> IntoResponse for MultipartJsonStream<S>
199where
200 S: Stream<Item = Result<T, E>> + Send + 'static,
201 T: Serialize,
202 E: Into<Problem>,
203{
204 fn into_response(self) -> Response {
205 let mut response = Response::new(Body::from_stream(frame(self.stream, self.boundary)));
206 response
207 .headers_mut()
208 .insert(CONTENT_TYPE, self.content_type);
209 response
210 }
211}
212
213/// Where the framer is in the response body.
214struct FramerState<S> {
215 stream: S,
216 boundary: String,
217 /// Set once the close delimiter has been written, or once an item failed
218 /// to serialize — either way there is nothing further to emit.
219 finished: bool,
220}
221
222/// Turn a stream of `Result<T, E>` items into the `multipart/mixed` body byte
223/// stream. `Ok` items become `application/json` data parts; an `Err` becomes a
224/// terminal `application/problem+json` error part (see [`MultipartJsonStream`]).
225fn frame<S, T, E>(
226 stream: S,
227 boundary: String,
228) -> impl Stream<Item = Result<Bytes, std::io::Error>> + Send + 'static
229where
230 S: Stream<Item = Result<T, E>> + Send + 'static,
231 T: Serialize,
232 E: Into<Problem>,
233{
234 let state = FramerState {
235 stream: Box::pin(stream),
236 boundary,
237 finished: false,
238 };
239
240 futures_util::stream::unfold(state, |mut state| async move {
241 if state.finished {
242 return None;
243 }
244 let Some(item) = state.stream.next().await else {
245 // Clean end of the source stream: write the close delimiter, then
246 // let the body end gracefully on the next poll.
247 state.finished = true;
248 let close = format!("--{}--\r\n", state.boundary);
249 return Some((Ok(Bytes::from(close)), state));
250 };
251 let value = match item {
252 Ok(value) => value,
253 Err(e) => {
254 // A post-open domain error. The status is long gone, so it
255 // cannot be an HTTP error — but it must not be an abort either:
256 // deliver it as a typed `application/problem+json` error part
257 // followed by the close delimiter, so the reader surfaces
258 // `Err(typed)` and a *clean* end. An error is terminal.
259 state.finished = true;
260 let mut problem: Problem = e.into();
261 // A mid-stream failure after the `200` is invisible to the
262 // status line; log it so it stays observable server-side.
263 tracing::warn!(
264 status = ?problem.status,
265 error_code = ?problem.error_code,
266 "multipart/mixed stream ended with a domain error; framing it as a typed error part"
267 );
268 let Some(bytes) = encode_error_part_within_limit(&state.boundary, &mut problem)
269 else {
270 // Even trimmed, the `Problem` will not fit or will not
271 // serialize — there is nothing well-formed left to frame, so
272 // abort the body rather than emit an unreadable part.
273 tracing::error!(
274 status = ?problem.status,
275 "multipart/mixed error part exceeds the maximum part size even after trimming, or would not serialize; aborting the response body"
276 );
277 return Some((
278 Err(std::io::Error::other(
279 "multipart/mixed error part exceeds the maximum part size even after trimming",
280 )),
281 state,
282 ));
283 };
284 return Some((Ok(bytes), state));
285 }
286 };
287 match serde_json::to_vec(&value) {
288 Ok(json) if json.len() > MAX_PART_BYTES => {
289 // The reader rejects any part whose body exceeds its
290 // accumulation ceiling, so emitting a larger part would produce
291 // a stream the toolkit client cannot consume. The status is long
292 // gone, so abort the body — exactly as the serialization-failure
293 // arm below does — rather than write an unreadable part.
294 tracing::error!(
295 item_type = std::any::type_name::<T>(),
296 item_bytes = json.len(),
297 max_bytes = MAX_PART_BYTES,
298 "multipart/mixed stream item exceeds the maximum part size; aborting the response body"
299 );
300 state.finished = true;
301 Some((
302 Err(std::io::Error::other(format!(
303 "multipart/mixed stream item is {} bytes, exceeding the maximum part size of {MAX_PART_BYTES} bytes",
304 json.len()
305 ))),
306 state,
307 ))
308 }
309 Ok(json) => {
310 let bytes = encode_part(&state.boundary, &json);
311 Some((Ok(bytes), state))
312 }
313 Err(e) => {
314 // Q5: the status is long gone, so abort the body rather than
315 // short it. Yielding an `Err` truncates the chunked encoding,
316 // which the reader surfaces as a transport error; returning
317 // `None` here would be a graceful EOF and would read as a
318 // clean, complete stream.
319 tracing::error!(
320 item_type = std::any::type_name::<T>(),
321 error = %e,
322 "failed to serialize a multipart/mixed stream item; aborting the response body"
323 );
324 state.finished = true;
325 Some((
326 Err(std::io::Error::other(format!(
327 "failed to serialize a multipart/mixed stream item: {e}"
328 ))),
329 state,
330 ))
331 }
332 }
333 })
334}
335
336fn encode_part(boundary: &str, json: &[u8]) -> Bytes {
337 let header = format!(
338 "--{boundary}\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n",
339 json.len()
340 );
341 let mut out = Vec::with_capacity(header.len() + json.len() + 2);
342 out.extend_from_slice(header.as_bytes());
343 out.extend_from_slice(json);
344 out.extend_from_slice(b"\r\n");
345 Bytes::from(out)
346}
347
348/// Encode a terminal error part — one `application/problem+json` part carrying
349/// the serialized [`Problem`] — immediately followed by the close delimiter.
350///
351/// Emitting the close in the same chunk is what makes an error a *clean* end:
352/// the reader decodes the part into a typed `Err`, then sees the closing
353/// delimiter and reports a graceful close, never a truncation or a manufactured
354/// missing-terminator error.
355fn encode_error_part_and_close(boundary: &str, problem_json: &[u8]) -> Bytes {
356 let header = format!(
357 "--{boundary}\r\nContent-Type: {PROBLEM_CONTENT_TYPE}\r\nContent-Length: {}\r\n\r\n",
358 problem_json.len()
359 );
360 let close = format!("\r\n--{boundary}--\r\n");
361 let mut out = Vec::with_capacity(header.len() + problem_json.len() + close.len());
362 out.extend_from_slice(header.as_bytes());
363 out.extend_from_slice(problem_json);
364 out.extend_from_slice(close.as_bytes());
365 Bytes::from(out)
366}
367
368/// Serialize `problem` as a terminal error part (+ close), keeping it within
369/// [`MAX_PART_BYTES`].
370///
371/// If the full `Problem` would exceed the reader's ceiling, its two unbounded
372/// fields (`detail`, `context`) are dropped and it is retried, so the *typed*
373/// error — `status`, `type`, `error_code`/`error_domain` — still reaches the
374/// client rather than degrading to a truncation. Returns `None` only when even
375/// the trimmed form will not fit or will not serialize; the caller aborts then.
376///
377/// This trimming is for a real domain `Err`. An oversized or unserializable
378/// `Ok` *data* item is a different case (a value the framer cannot render at
379/// all) and is deliberately aborted rather than replaced with a synthetic error.
380fn encode_error_part_within_limit(boundary: &str, problem: &mut Problem) -> Option<Bytes> {
381 let full = serde_json::to_vec(problem)
382 .ok()
383 .filter(|json| json.len() <= MAX_PART_BYTES);
384 if let Some(json) = full {
385 return Some(encode_error_part_and_close(boundary, &json));
386 }
387 problem.detail = String::new();
388 problem.context = serde_json::Value::Null;
389 let trimmed = serde_json::to_vec(problem)
390 .ok()
391 .filter(|json| json.len() <= MAX_PART_BYTES)?;
392 Some(encode_error_part_and_close(boundary, &trimmed))
393}
394
395fn generated_boundary() -> String {
396 // Hex only, so RFC 2046-legal by construction, and wide enough that it
397 // cannot collide with item payload bytes in practice.
398 uuid::Uuid::now_v7().simple().to_string()
399}
400
401/// Build the `Content-Type` for a boundary, or `None` if it does not form a
402/// legal header value. A validated boundary (an HTTP token) always yields
403/// `Some`; `None` is the unreachable guard handled in
404/// [`MultipartJsonStream::with_validated_boundary`].
405fn content_type_for(boundary: &str) -> Option<HeaderValue> {
406 HeaderValue::from_str(&format!("multipart/mixed; boundary={boundary}")).ok()
407}
408
409/// Known-good boundary substituted only if a *validated* boundary somehow fails
410/// to form a header value (unreachable). Hex, so it is both RFC 2046-legal and
411/// an HTTP token. Kept byte-for-byte consistent with [`FALLBACK_CONTENT_TYPE`]
412/// so the emitted `boundary=` always matches the `--<boundary>` delimiters.
413const FALLBACK_BOUNDARY: &str = "0f0f0f0f0f0f0f0f0f0f0f0f0f0f0f0f";
414const FALLBACK_CONTENT_TYPE: &str = "multipart/mixed; boundary=0f0f0f0f0f0f0f0f0f0f0f0f0f0f0f0f";
415
416/// Check a boundary against RFC 2046 §5.1.1 **and** the HTTP token grammar
417/// (RFC 7230 §3.2.6). `Err` carries a reason suitable for a log line.
418///
419/// RFC 2046 bchars additionally allow `( ) , / : = ?` and space, but those are
420/// HTTP `tspecials`: a boundary containing one would have to be quoted to
421/// survive a conformant `Content-Type` parser, yet this crate emits — and
422/// reads — the `boundary=` parameter unquoted (see [`with_validated_boundary`]
423/// and the client's `;`-splitting reader). Restricting to the intersection of
424/// the two grammars keeps the emitted header unambiguous, so we reject those
425/// characters rather than produce a `Content-Type` a conformant client would
426/// read differently from this crate.
427fn validate_boundary(boundary: &str) -> Result<(), &'static str> {
428 if boundary.is_empty() {
429 return Err("boundary must not be empty");
430 }
431 if boundary.len() > MAX_BOUNDARY_LEN {
432 return Err("boundary exceeds 70 characters");
433 }
434 // Intersection of RFC 2046 bchars and HTTP token chars:
435 // DIGIT / ALPHA / "'" / "+" / "-" / "." / "_"
436 // (space can't appear at all, so no trailing-space check is needed.)
437 if !boundary
438 .bytes()
439 .all(|b| b.is_ascii_alphanumeric() || b"'+-._".contains(&b))
440 {
441 return Err("boundary contains a character outside the HTTP token set");
442 }
443 Ok(())
444}
445
446#[cfg(test)]
447#[cfg_attr(coverage_nightly, coverage(off))]
448#[allow(clippy::unwrap_used)]
449mod tests {
450 use super::*;
451 use axum::body::to_bytes;
452 use serde::Serialize;
453 use toolkit_canonical_errors::CanonicalError;
454
455 #[derive(Serialize)]
456 struct Item {
457 id: u32,
458 }
459
460 /// An item whose `Serialize` impl fails, to exercise the abort path — a
461 /// value the framer cannot put on the wire at all (distinct from a domain
462 /// `Err`, which becomes a typed error part).
463 struct Unserializable;
464 impl Serialize for Unserializable {
465 fn serialize<S: serde::Serializer>(&self, _: S) -> Result<S::Ok, S::Error> {
466 Err(serde::ser::Error::custom("nope"))
467 }
468 }
469
470 #[tokio::test]
471 async fn frames_one_part_per_item_with_a_content_length() {
472 let items = futures_util::stream::iter(vec![
473 Ok::<_, CanonicalError>(Item { id: 1 }),
474 Ok(Item { id: 2 }),
475 ]);
476 let response = MultipartJsonStream::with_boundary(items, "BOUND")
477 .unwrap()
478 .into_response();
479 assert_eq!(
480 response.headers().get(CONTENT_TYPE).unwrap(),
481 "multipart/mixed; boundary=BOUND"
482 );
483 let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
484 assert_eq!(
485 String::from_utf8(body.to_vec()).unwrap(),
486 "--BOUND\r\nContent-Type: application/json\r\nContent-Length: 8\r\n\r\n{\"id\":1}\r\n\
487 --BOUND\r\nContent-Type: application/json\r\nContent-Length: 8\r\n\r\n{\"id\":2}\r\n\
488 --BOUND--\r\n"
489 );
490 }
491
492 #[tokio::test]
493 async fn an_empty_stream_is_just_the_close_delimiter() {
494 let items = futures_util::stream::iter(Vec::<Result<Item, CanonicalError>>::new());
495 let response = MultipartJsonStream::with_boundary(items, "BOUND")
496 .unwrap()
497 .into_response();
498 let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
499 assert_eq!(String::from_utf8(body.to_vec()).unwrap(), "--BOUND--\r\n");
500 }
501
502 #[tokio::test]
503 async fn an_unserializable_item_aborts_the_body_rather_than_ending_it() {
504 // Q5's implementation trap: the body must ERROR, not end. If this ever
505 // regresses to ending the stream, `to_bytes` succeeds and the reader
506 // sees a clean, complete — but silently short — stream.
507 let items = futures_util::stream::iter(vec![Ok::<_, CanonicalError>(Unserializable)]);
508 let response = MultipartJsonStream::with_boundary(items, "BOUND")
509 .unwrap()
510 .into_response();
511 let result = to_bytes(response.into_body(), usize::MAX).await;
512 assert!(
513 result.is_err(),
514 "the body must abort, not end gracefully; got {:?}",
515 result.map(|b| String::from_utf8_lossy(&b).into_owned())
516 );
517 }
518
519 #[tokio::test]
520 async fn an_err_item_becomes_a_typed_error_part_then_a_clean_close() {
521 // A domain `Err` is NOT an abort: it is framed as a terminal
522 // `application/problem+json` part followed by the close delimiter, so
523 // the body ends cleanly and the reader can surface `Err(typed)`.
524 let items = futures_util::stream::iter(vec![
525 Ok(Item { id: 1 }),
526 Err(CanonicalError::internal("mid-stream boom").create()),
527 // Anything after the error is not framed — an error is terminal.
528 Ok(Item { id: 2 }),
529 ]);
530 let response = MultipartJsonStream::with_boundary(items, "BOUND")
531 .unwrap()
532 .into_response();
533 // The body must END (not abort): `to_bytes` succeeds.
534 let body = to_bytes(response.into_body(), usize::MAX)
535 .await
536 .expect("an error item ends the body cleanly, it must not abort");
537 let text = String::from_utf8(body.to_vec()).unwrap();
538 // One data part, then the problem+json error part, then the close —
539 // and no second data part (the error is terminal).
540 assert!(
541 text.starts_with(
542 "--BOUND\r\nContent-Type: application/json\r\nContent-Length: 8\r\n\r\n{\"id\":1}\r\n"
543 ),
544 "expected the data part first; got:\n{text}"
545 );
546 assert!(
547 text.contains("Content-Type: application/problem+json"),
548 "expected a problem+json error part; got:\n{text}"
549 );
550 assert!(
551 text.ends_with("--BOUND--\r\n"),
552 "expected a clean close; got:\n{text}"
553 );
554 assert!(
555 !text.contains("{\"id\":2}"),
556 "an error is terminal; later items must not be framed; got:\n{text}"
557 );
558 }
559
560 #[tokio::test]
561 async fn an_oversized_error_problem_is_trimmed_not_aborted() {
562 // A domain error whose `Problem` exceeds `MAX_PART_BYTES` must still
563 // reach the client as a typed error part (with its unbounded fields
564 // trimmed), NOT degrade to a truncation/abort. Contrast with an oversized
565 // `Ok` data item, which does abort.
566 let huge = "a".repeat(MAX_PART_BYTES + 4096);
567 let items = futures_util::stream::iter(vec![Err::<Item, CanonicalError>(
568 CanonicalError::internal(huge).create(),
569 )]);
570 let response = MultipartJsonStream::with_boundary(items, "BOUND")
571 .unwrap()
572 .into_response();
573 // The body must END cleanly (trimmed part + close), not abort.
574 let body = to_bytes(response.into_body(), usize::MAX)
575 .await
576 .expect("an oversized error Problem must be trimmed and framed, not aborted");
577 let text = String::from_utf8_lossy(&body);
578 assert!(
579 text.contains("Content-Type: application/problem+json"),
580 "expected a problem+json error part; got {} bytes",
581 body.len()
582 );
583 assert!(
584 text.ends_with("--BOUND--\r\n"),
585 "expected a clean close; got {} bytes",
586 body.len()
587 );
588 // Trimmed: the ~16 MiB detail was dropped, so the whole body is small.
589 assert!(
590 body.len() < MAX_PART_BYTES,
591 "the oversized detail must be trimmed away; body was {} bytes",
592 body.len()
593 );
594 }
595
596 #[tokio::test]
597 async fn an_oversized_item_aborts_the_body_rather_than_emitting_an_unreadable_part() {
598 // A part larger than the reader's accumulation ceiling would produce a
599 // stream the toolkit client cannot consume, so the framer must abort the
600 // body (error), not emit the part. `String` serializes to its own bytes
601 // plus two quotes, so this clears `MAX_PART_BYTES`.
602 let oversized = "a".repeat(MAX_PART_BYTES);
603 let items = futures_util::stream::iter(vec![Ok::<_, CanonicalError>(oversized)]);
604 let response = MultipartJsonStream::with_boundary(items, "BOUND")
605 .unwrap()
606 .into_response();
607 let result = to_bytes(response.into_body(), usize::MAX).await;
608 assert!(
609 result.is_err(),
610 "an oversized item must abort the body, not emit an unreadable part; got {:?}",
611 result.map(|b| b.len())
612 );
613 }
614
615 #[test]
616 fn the_emitted_content_type_carries_the_framing_boundary() {
617 // The invariant #19 is about: the header's `boundary=` must match the
618 // `--<boundary>` delimiters the framer writes, never degrade to a
619 // boundary-less `multipart/mixed`.
620 let framer = MultipartJsonStream::with_boundary(
621 futures_util::stream::iter(Vec::<Result<Item, CanonicalError>>::new()),
622 "abc123",
623 )
624 .unwrap();
625 let expected = format!("multipart/mixed; boundary={}", framer.boundary());
626 let response = framer.into_response();
627 assert_eq!(response.headers().get(CONTENT_TYPE).unwrap(), &expected);
628 }
629
630 #[test]
631 fn the_fallback_boundary_and_header_stay_consistent() {
632 // The unreachable substitution relies on these two constants agreeing.
633 assert_eq!(
634 FALLBACK_CONTENT_TYPE,
635 format!("multipart/mixed; boundary={FALLBACK_BOUNDARY}")
636 );
637 assert!(validate_boundary(FALLBACK_BOUNDARY).is_ok());
638 assert!(content_type_for(FALLBACK_BOUNDARY).is_some());
639 }
640
641 #[test]
642 fn a_generated_boundary_is_used_when_none_is_given() {
643 let items = futures_util::stream::iter(Vec::<Result<Item, CanonicalError>>::new());
644 let framer = MultipartJsonStream::new(items);
645 assert!(validate_boundary(framer.boundary()).is_ok());
646 assert!(!framer.boundary().is_empty());
647 }
648
649 #[test]
650 fn an_illegal_boundary_is_rejected_not_silently_substituted() {
651 let items = futures_util::stream::iter(Vec::<Result<Item, CanonicalError>>::new());
652 // `MultipartJsonStream` isn't `Debug`, so match rather than `unwrap_err`.
653 let Err(err) = MultipartJsonStream::with_boundary(items, "not\r\nlegal") else {
654 panic!("an illegal boundary must be rejected");
655 };
656 // The rejected value is escaped in the message, not written raw.
657 assert!(
658 err.to_string().contains("not\\r\\nlegal"),
659 "message should escape control chars: {err}"
660 );
661 // The caller can still opt into a generated boundary via `new`.
662 let framer = MultipartJsonStream::new(futures_util::stream::iter(Vec::<
663 Result<Item, CanonicalError>,
664 >::new()));
665 assert!(validate_boundary(framer.boundary()).is_ok());
666 }
667
668 #[test]
669 fn boundary_validation_requires_an_rfc_2046_http_token() {
670 // Characters in both RFC 2046 bchars and the HTTP token set.
671 assert!(validate_boundary("abcABC012'+-._").is_ok());
672 assert!(validate_boundary("").is_err());
673 assert!(validate_boundary(&"a".repeat(MAX_BOUNDARY_LEN + 1)).is_err());
674 // RFC 2046-legal but HTTP tspecials: rejected so the emitted
675 // Content-Type can't be misparsed by a conformant client.
676 assert!(validate_boundary("has space").is_err());
677 assert!(validate_boundary("a/b").is_err());
678 assert!(validate_boundary("a:b=c?").is_err());
679 assert!(validate_boundary("(paren)").is_err());
680 assert!(validate_boundary("comma,d").is_err());
681 // Not even RFC 2046-legal.
682 assert!(validate_boundary("semi;colon").is_err());
683 assert!(validate_boundary("quote\"d").is_err());
684 }
685}