Skip to main content

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}