Skip to main content

web_faith/
response.rs

1//! Responses, bodies, and request timing.
2
3pub use crate::timing::RequestTiming;
4
5// spec:RESP spec:TRL spec:BODY
6
7use std::{
8	fmt::Debug,
9	net::SocketAddr,
10	path::{Path, PathBuf},
11	pin::Pin,
12	sync::{
13		Arc,
14		atomic::{AtomicBool, Ordering},
15	},
16	task::{Context, Poll},
17	time::{Duration, Instant},
18};
19
20use bytes::Bytes;
21use futures::{Stream, StreamExt, stream};
22use http::header::{CONTENT_LENGTH, HeaderMap};
23use reqwest::{StatusCode, Url, Version};
24use serde::de::DeserializeOwned;
25use tokio::{io::AsyncWriteExt, sync::watch};
26
27pub use crate::body::BodyReader;
28use crate::{
29	body::Claim,
30	error::{FaithError, FaithErrorKind},
31	timing::TimingSlot,
32};
33
34use crate::integrity::{finish_integrity, integrity_checker, verify_integrity};
35
36/// The peer that sent a response.
37#[derive(Debug)]
38pub struct PeerInformation {
39	/// The peer's address and port, where the connection could report one.
40	pub address: Option<SocketAddr>,
41	/// The peer's DER-encoded leaf certificate, for a response that arrived over HTTPS.
42	pub certificate: Option<Vec<u8>>,
43}
44
45/// Where a response body is written.
46#[derive(Debug, Clone, Default)]
47pub struct FileDestination {
48	/// Truncate and replace an occupied destination. The default refuses one instead, leaving what
49	/// is there untouched.
50	pub overwrite: bool,
51	/// The permissions a newly created file is given. Ignored on platforms without Unix file modes.
52	pub mode: Option<u32>,
53}
54
55// Reporting every chunk would cross the surface boundary thousands of times for a large body,
56// which is the cost writing to a file directly exists to avoid. The final report always lands
57// regardless.
58pub(crate) const PROGRESS_INTERVAL: Duration = Duration::from_millis(50);
59
60/// Open the destination file for a body write, mapping filesystem refusals to Faith's errors.
61// spec:BODY#tofile
62pub(crate) async fn open_destination(
63	path: &Path,
64	options: &FileDestination,
65) -> Result<tokio::fs::File, FaithError> {
66	let mut open = tokio::fs::OpenOptions::new();
67	open.write(true);
68	if options.overwrite {
69		// An occupied destination is truncated and replaced.
70		open.create(true).truncate(true);
71	} else {
72		// The safe default refuses an occupied destination outright.
73		open.create_new(true);
74	}
75	#[cfg(unix)]
76	if let Some(mode) = options.mode {
77		open.mode(mode);
78	}
79
80	match open.open(path).await {
81		Ok(file) => Ok(file),
82		Err(err) => Err(classify_open_error(path, err).await),
83	}
84}
85
86/// Classify a failure to open the destination: `FileExists` for an occupied path, `FileWrite` for
87/// a directory or any other refusal, carrying the OS detail.
88pub(crate) async fn classify_open_error(path: &Path, err: std::io::Error) -> FaithError {
89	let kind = if err.kind() == std::io::ErrorKind::AlreadyExists {
90		match tokio::fs::symlink_metadata(path).await {
91			Ok(meta) if meta.is_dir() => FaithErrorKind::FileWrite,
92			_ => FaithErrorKind::FileExists,
93		}
94	} else {
95		FaithErrorKind::FileWrite
96	};
97	FaithError::new(kind, err.to_string())
98}
99
100/// The trailing headers a response carried, once its body has ended.
101#[derive(Clone, Debug, Default)]
102pub enum Trailers {
103	/// The body has not ended, so the question is still open.
104	#[default]
105	NotYet,
106	/// The body ended carrying no trailers.
107	None,
108	/// The trailers that arrived.
109	Some(HeaderMap),
110}
111
112/// Where the trailers land: written by whoever finishes the body, awaited by `trailers()`.
113///
114/// A watch channel, so waiting parks instead of spinning. The fetch standard's trailers proposal
115/// (<https://github.com/whatwg/fetch/pull/1940>) has this not resolving until the body is
116/// consumed, so the wait is unbounded by design.
117#[derive(Debug)]
118pub(crate) struct TrailersSlot(watch::Sender<Trailers>);
119
120impl Default for TrailersSlot {
121	fn default() -> Self {
122		Self(watch::channel(Trailers::NotYet).0)
123	}
124}
125
126impl TrailersSlot {
127	/// Record trailers that arrived, waking whoever is waiting.
128	pub fn arrived(&self, trailers: HeaderMap) {
129		self.0.send_replace(Trailers::Some(trailers));
130	}
131
132	/// Record that the body ended, if no trailers frame got there first.
133	///
134	/// `send_if_modified` keeps the read and write one step, and wakes waiters only from the call
135	/// that settled it.
136	pub fn ended(&self) {
137		self.0.send_if_modified(|state| {
138			if matches!(state, Trailers::NotYet) {
139				*state = Trailers::None;
140				true
141			} else {
142				false
143			}
144		});
145	}
146
147	/// Wait until the body has settled the question.
148	pub async fn settled(&self) -> Trailers {
149		let mut rx = self.0.subscribe();
150		// `wait_for` tests the current value before waiting, so trailers that already
151		// arrived return without yielding. Its error case is the sender being gone, which
152		// means the response was dropped and nothing can ever set this -- no trailers is
153		// the only answer left.
154		match rx
155			.wait_for(|state| !matches!(state, Trailers::NotYet))
156			.await
157		{
158			Ok(state) => state.clone(),
159			Err(_) => Trailers::None,
160		}
161	}
162}
163
164#[cfg(test)]
165mod tests {
166	use std::{
167		future::Future,
168		pin::pin,
169		sync::atomic::{AtomicUsize, Ordering},
170		task::{Context, Poll, Wake, Waker},
171	};
172
173	use super::*;
174
175	/// A waker that counts how many times the task asks to be polled again.
176	struct CountingWaker(AtomicUsize);
177
178	impl CountingWaker {
179		fn wakes(&self) -> usize {
180			self.0.load(Ordering::SeqCst)
181		}
182	}
183
184	impl Wake for CountingWaker {
185		fn wake(self: Arc<Self>) {
186			self.wake_by_ref();
187		}
188
189		fn wake_by_ref(self: &Arc<Self>) {
190			self.0.fetch_add(1, Ordering::SeqCst);
191		}
192	}
193
194	/// A response that cannot carry a body converts to an `http::Response` with an empty one,
195	/// carrying its status, version, and headers across.
196	#[test]
197	fn a_bodyless_response_converts_to_an_http_response() {
198		let mut headers = HeaderMap::new();
199		headers.insert("x-test", "yes".parse().expect("a valid header value"));
200
201		let response = Response {
202			claim: None,
203			disturbed: Arc::new(AtomicBool::new(false)),
204			headers,
205			integrity: None,
206			peer: Arc::new(PeerInformation {
207				address: None,
208				certificate: None,
209			}),
210			redirected: false,
211			status_code: StatusCode::NO_CONTENT,
212			timing: Arc::new(TimingSlot::new(
213				Instant::now(),
214				crate::timing::RequestTiming::default(),
215			)),
216			trailers: Arc::new(TrailersSlot::default()),
217			url: Url::parse("https://example.com/").expect("a valid url"),
218			version: Version::HTTP_2,
219		};
220
221		let http = response.into_http().expect("an undisturbed body converts");
222		assert_eq!(http.status(), StatusCode::NO_CONTENT);
223		assert_eq!(http.version(), Version::HTTP_2);
224		assert_eq!(
225			http.headers().get("x-test").map(|v| v.as_bytes()),
226			Some(&b"yes"[..])
227		);
228
229		// Draining it yields nothing, which is what a consumer of the body actually observes.
230		let collected =
231			futures::executor::block_on(http_body_util::BodyExt::collect(http.into_body()))
232				.expect("an empty body collects");
233		assert!(collected.to_bytes().is_empty());
234	}
235
236	/// Waiting for trailers parks until the body settles the question, rather than polling
237	/// for it.
238	///
239	/// The bug this guards against was a `yield_now` loop, which is visible here as the shape
240	/// of the wait rather than as a quantity of CPU: a spin re-arms its own waker on every
241	/// poll, so it is scheduled again immediately, while a parked wait asks for nothing until
242	/// something else moves. Asserting the wake count keeps this deterministic -- timing how
243	/// much CPU the process burns measures the machine as much as the code.
244	#[test]
245	fn waiting_for_trailers_parks_rather_than_spinning() {
246		let slot = TrailersSlot::default();
247		let counter = Arc::new(CountingWaker(AtomicUsize::new(0)));
248		let waker = Waker::from(counter.clone());
249		let mut cx = Context::from_waker(&waker);
250		let mut settled = pin!(slot.settled());
251
252		// Nothing has settled the question, so the wait parks...
253		assert!(matches!(settled.as_mut().poll(&mut cx), Poll::Pending));
254		// ...without scheduling itself to be polled again, which is what a spin does.
255		assert_eq!(counter.wakes(), 0, "a parked wait asks for no wake-up");
256
257		// Polling again changes nothing: still parked, still asking for nothing.
258		assert!(matches!(settled.as_mut().poll(&mut cx), Poll::Pending));
259		assert_eq!(counter.wakes(), 0, "polling again does not arm a wake-up");
260
261		// The body ending is what wakes it, and it resolves on the next poll.
262		slot.ended();
263		assert!(counter.wakes() >= 1, "the body ending wakes the waiter");
264		assert!(matches!(
265			settled.as_mut().poll(&mut cx),
266			Poll::Ready(Trailers::None)
267		));
268	}
269
270	/// Trailers that arrived before anyone asked resolve without parking at all.
271	#[test]
272	fn trailers_already_there_resolve_on_the_first_poll() {
273		let slot = TrailersSlot::default();
274		let mut headers = HeaderMap::new();
275		headers.insert("x-checksum", "abc123".parse().unwrap());
276		slot.arrived(headers);
277
278		let counter = Arc::new(CountingWaker(AtomicUsize::new(0)));
279		let waker = Waker::from(counter.clone());
280		let mut cx = Context::from_waker(&waker);
281		let mut settled = pin!(slot.settled());
282
283		assert!(matches!(
284			settled.as_mut().poll(&mut cx),
285			Poll::Ready(Trailers::Some(_))
286		));
287	}
288}
289
290/// A progress report from a body write in flight.
291#[derive(Debug, Clone, Copy, PartialEq, Eq)]
292#[non_exhaustive]
293pub struct FileProgress {
294	/// Bytes written to the file so far.
295	pub bytes_written: u64,
296	/// The `Content-Length` the response advertised. Absent for a chunked response, and for one
297	/// being decoded, where the final size is not known ahead of time.
298	pub content_length: Option<u64>,
299}
300
301/// The result of writing a body to a file.
302#[derive(Debug, Clone, PartialEq, Eq)]
303#[non_exhaustive]
304pub struct FileWritten {
305	/// The absolute filesystem path written to.
306	pub path: PathBuf,
307	/// The number of bytes that landed at the destination.
308	pub bytes_written: u64,
309}
310
311/// A response to a request.
312///
313/// Arrives from a request; it is not constructed directly. Reading the body consumes it, so a
314/// second read fails. [`Self::try_clone`] gets a copy that can be read separately.
315#[derive(Debug, Clone)]
316pub struct Response {
317	/// This response's claim on the body, shared by its in-process copies and given up when the
318	/// last of them goes. `None` for a response that cannot carry a body.
319	// spec:BODY#giving-up-the-body
320	pub(crate) claim: Option<Arc<Claim>>,
321	pub(crate) disturbed: Arc<AtomicBool>,
322	pub(crate) headers: HeaderMap,
323	pub(crate) integrity: Option<String>,
324	pub(crate) peer: Arc<PeerInformation>,
325	pub(crate) redirected: bool,
326	pub(crate) status_code: StatusCode,
327	pub(crate) timing: Arc<TimingSlot>,
328	pub(crate) trailers: Arc<TrailersSlot>,
329	pub(crate) url: Url,
330	pub(crate) version: Version,
331}
332
333impl Response {
334	/// The response's status.
335	pub fn status(&self) -> StatusCode {
336		self.status_code
337	}
338
339	/// The canonical reason phrase for the status, or empty for a code with no well-known one.
340	///
341	/// Always the canonical phrase. HTTP/1 lets a server send its own, which is not surfaced here;
342	/// HTTP/2 and HTTP/3 carry none at all.
343	pub fn status_text(&self) -> &'static str {
344		self.status_code.canonical_reason().unwrap_or_default()
345	}
346
347	/// Whether the status is in the 2xx range.
348	pub fn ok(&self) -> bool {
349		self.status_code.is_success()
350	}
351
352	/// The response's headers.
353	pub fn headers(&self) -> &HeaderMap {
354		&self.headers
355	}
356
357	/// The URL the response came from, after any redirects.
358	pub fn url(&self) -> &Url {
359		&self.url
360	}
361
362	/// Whether a redirect was followed to reach this response.
363	pub fn redirected(&self) -> bool {
364		self.redirected
365	}
366
367	/// The HTTP version the response arrived over.
368	pub fn version(&self) -> Version {
369		self.version
370	}
371
372	/// The peer that sent the response.
373	pub fn peer(&self) -> &PeerInformation {
374		&self.peer
375	}
376
377	/// Copy the response, so the body can be read twice.
378	///
379	/// Both copies read the same underlying body, and neither is disturbed by the other having
380	/// been cloned. Fails if the body has already been read.
381	pub fn try_clone(&self) -> Result<Self, FaithError> {
382		// A read, not `check_stream_disturbed`: that one swaps the flag, which would disturb the
383		// response being cloned and leave neither copy readable.
384		if self.body_used() {
385			return Err(FaithErrorKind::ResponseAlreadyDisturbed.into());
386		}
387
388		// The copy holds a claim of its own, so it keeps the transfer going after this one
389		// gives up. A discarded body has no claim left to copy.
390		let claim = match &self.claim {
391			None => None,
392			Some(claim) => Some(
393				claim
394					.duplicate()
395					.ok_or(FaithErrorKind::ResponseAlreadyDisturbed)?,
396			),
397		};
398
399		Ok(Self {
400			claim,
401			disturbed: Arc::new(AtomicBool::new(false)),
402			..Clone::clone(self)
403		})
404	}
405
406	/// Whether the body has been read, or handed out as a stream.
407	pub fn body_used(&self) -> bool {
408		self.disturbed.load(Ordering::SeqCst)
409	}
410
411	/// Read the whole body.
412	///
413	/// Consumes it, so a second read fails. A request's `integrity` is verified here, once the
414	/// whole body is in hand.
415	// spec:BODY
416	pub async fn bytes(&self) -> Result<Vec<u8>, FaithError> {
417		self.check_stream_disturbed()?;
418		self.gather_contiguous().await
419	}
420
421	/// Read the whole body as text.
422	///
423	/// Decoded as UTF-8, with invalid sequences replaced by U+FFFD.
424	pub async fn text(&self) -> Result<String, FaithError> {
425		let bytes = self.bytes().await?;
426		Ok(String::from_utf8(bytes)
427			.unwrap_or_else(|err| String::from_utf8_lossy(err.as_bytes()).into_owned()))
428	}
429
430	/// Read the whole body and deserialise it from JSON.
431	///
432	/// Reads into memory before parsing, which can cost twice the body's size; [`Self::body_stream`]
433	/// avoids that.
434	pub async fn json<T: DeserializeOwned>(&self) -> Result<T, FaithError> {
435		let bytes = self.bytes().await?;
436		serde_json::from_slice(&bytes)
437			.map_err(|err| FaithError::new(FaithErrorKind::JsonParse, err.to_string()))
438	}
439
440	/// The body as a stream of chunks, decoded under whichever coding was negotiated.
441	///
442	/// `None` for a response that cannot carry a body. Callable more than once: every stream a
443	/// response hands out reads from its one position in the body, so one taken after part of
444	/// the body was read continues from there. Dropping a stream gives the response's claim on
445	/// the body up, stopping the transfer once no clone still wants it. Fails if the body has
446	/// already been read to the end or discarded.
447	// spec:BODY
448	pub fn body_stream(&self) -> Result<Option<BodyReader>, FaithError> {
449		// The body counts as disturbed from here, though the stream itself stays re-readable.
450		let _ = self.check_stream_disturbed();
451
452		match &self.claim {
453			None => Ok(None),
454			Some(claim) => claim.reader().map(Some),
455		}
456	}
457
458	/// Give up on the body, releasing the connection back to the pool.
459	///
460	/// Worth doing when the body is not wanted: left unread, the claim on it is held until the
461	/// response drops. Gives up this response's claim only, so a clone still reading carries on;
462	/// when it was the last claim, resolves once the transfer has been stopped: the stream reset
463	/// on HTTP/2 and HTTP/3, the connection drained back to the pool or closed on HTTP/1.
464	// spec:BODY#discard spec:BODY#giving-up-the-body
465	pub async fn discard(&self) {
466		let Some(claim) = &self.claim else {
467			return;
468		};
469		claim.give_up();
470		if claim.body().claims_left() == 0 {
471			claim.body().settled().await;
472		}
473	}
474
475	/// The timing of the request that produced this response, once its body has ended.
476	// spec:RESP#request-timing
477	pub async fn timing(&self) -> crate::timing::RequestTiming {
478		self.timing.settled().await
479	}
480
481	/// The trailers, once the body has ended.
482	///
483	/// A body that is never read never ends, so this waits indefinitely; see [`Trailers`].
484	// spec:TRL
485	pub async fn trailers(&self) -> Trailers {
486		self.trailers.settled().await
487	}
488
489	pub(crate) fn check_stream_disturbed(&self) -> Result<(), FaithError> {
490		if self.disturbed.swap(true, Ordering::SeqCst) {
491			Err(FaithErrorKind::ResponseAlreadyDisturbed.into())
492		} else {
493			Ok(())
494		}
495	}
496
497	/// Read the whole body as the chunks it arrived in, without copying them.
498	///
499	/// [`Self::bytes`] and its siblings are built on this.
500	pub(crate) async fn gather(&self) -> Result<Arc<[Bytes]>, FaithError> {
501		let Some(claim) = &self.claim else {
502			return Ok(Default::default());
503		};
504
505		let mut stream = claim.reader()?;
506		let mut chunks = Vec::new();
507		while let Some(chunk) = stream.next().await {
508			chunks.push(chunk?);
509		}
510
511		Ok(Arc::from(chunks.into_boxed_slice()))
512	}
513
514	/// [`Self::gather`], then copy the chunks into one contiguous buffer.
515	pub(crate) async fn gather_contiguous(&self) -> Result<Vec<u8>, FaithError> {
516		let body = self.gather().await?;
517		let length = body.iter().map(|chunk| chunk.len()).sum();
518		let mut bytes = Vec::with_capacity(length);
519		for chunk in body.into_iter() {
520			bytes.extend_from_slice(chunk);
521		}
522
523		if let Some(ref integrity) = self.integrity {
524			verify_integrity(&bytes, integrity)?;
525		}
526
527		Ok(bytes)
528	}
529
530	/// Write the body to a file, reporting progress as the bytes land.
531	///
532	/// `on_progress` is called with the bytes written so far and the advertised length when it is
533	/// known, at most every 50ms, and once more when the last byte is written.
534	// spec:BODY#tofile
535	pub async fn write_to_file(
536		&self,
537		path: impl AsRef<Path>,
538		options: &FileDestination,
539		mut on_progress: impl FnMut(FileProgress),
540	) -> Result<FileWritten, FaithError> {
541		let path = path.as_ref();
542		// A response that cannot carry a body has nothing to write, and this is settled
543		// before any file is created (spec:BODY#tofile).
544		let Some(claim) = self.claim.clone() else {
545			return Err(FaithErrorKind::ResponseBodyNull.into());
546		};
547
548		// A body already read, discarded, or whose stream was handed out, has no second read to
549		// give. Checked without committing so an open failure below still leaves the body
550		// undisturbed and the caller free to retry to another path.
551		if self.disturbed.load(Ordering::SeqCst) || claim.is_given_up() {
552			return Err(FaithErrorKind::ResponseAlreadyDisturbed.into());
553		}
554
555		// Reject a malformed integrity value up front, before the body is touched, the same
556		// as the other verified reads reject it when the whole body is in hand.
557		let mut checker = integrity_checker(self.integrity.as_deref())?;
558
559		// The advertised length, when the server sent one. It is only visible here for a body
560		// delivered as received: a decoded body has had its Content-Length stripped, so the
561		// bytes written equal the wire bytes wherever this is Some (spec:BODY#tofile, ENC).
562		let content_length = self
563			.headers
564			.get(CONTENT_LENGTH)
565			.and_then(|value| value.to_str().ok())
566			.and_then(|value| value.trim().parse::<u64>().ok());
567
568		// The destination is opened before any of the body is read, so a failure to open it
569		// leaves the body unread and undisturbed.
570		let mut file = open_destination(path, options).await?;
571
572		// Commit the read now the destination is in hand. A concurrent read that slipped in
573		// since the load above wins, and this one finds the body already spent.
574		self.check_stream_disturbed()?;
575
576		let stream = claim.reader()?;
577
578		// Reporting is rate limited rather than per chunk, so a large body does not cross a
579		// surface boundary thousands of times.
580		// spec:BODY#tofile
581		let mut report = |written: u64| {
582			on_progress(FileProgress {
583				bytes_written: written,
584				content_length,
585			});
586		};
587
588		let mut written: u64 = 0;
589		let mut reported_at = Instant::now();
590		futures::pin_mut!(stream);
591		while let Some(result) = stream.next().await {
592			let chunk = result?;
593			if let Some(checker) = checker.as_mut() {
594				checker.input(&chunk);
595			}
596			file.write_all(&chunk)
597				.await
598				.map_err(|err| FaithError::new(FaithErrorKind::FileWrite, err.to_string()))?;
599			written += chunk.len() as u64;
600			// A server cannot send more than it promised: once the bytes off the wire exceed
601			// the advertised length, the write fails and the bytes so far stay on disk
602			// (spec:BODY#tofile).
603			if let Some(limit) = content_length {
604				if written > limit {
605					return Err(FaithErrorKind::ContentLengthOverrun.into());
606				}
607			}
608			if reported_at.elapsed() >= PROGRESS_INTERVAL {
609				reported_at = Instant::now();
610				report(written);
611			}
612		}
613
614		file.flush()
615			.await
616			.map_err(|err| FaithError::new(FaithErrorKind::FileWrite, err.to_string()))?;
617
618		// The last report always lands, whatever the rate limit allowed along the way, so a
619		// caller's final view of a completed write is the whole body rather than the last
620		// interval boundary. An empty body reports once, with nothing written.
621		report(written);
622
623		// The digest is only known once the last byte has been written, so the file that
624		// fails verification is on disk when the error arrives (spec:SRI).
625		if let Some(checker) = checker {
626			finish_integrity(checker)?;
627		}
628
629		Ok(FileWritten {
630			// A relative path resolves against the process's working directory; the caller
631			// is handed the absolute path the bytes landed at.
632			path: std::path::absolute(path).unwrap_or_else(|_| path.to_path_buf()),
633			bytes_written: written,
634		})
635	}
636}
637
638/// A response's body plumbing, for a surface that drives it directly. Permanently unstable.
639#[cfg(feature = "unstable-internals")]
640impl Response {
641	/// Whether the body has already been read or handed out, marking it disturbed either way.
642	pub fn check_disturbed(&self) -> Result<(), FaithError> {
643		self.check_stream_disturbed()
644	}
645
646	/// Stop the body's transfer because the request's signal was aborted after its headers
647	/// arrived. Every reader of the response and its clones errors with `Aborted`.
648	// spec:CANCEL#abortsignal
649	pub fn abort_body(&self) {
650		if let Some(claim) = &self.claim {
651			claim.body().abort();
652		}
653	}
654}
655
656/// A [`Response`]'s body as an [`http_body::Body`].
657///
658/// For handing to a tower service, a hyper client, or anything else that takes one.
659pub struct ResponseBody {
660	chunks: Pin<Box<dyn Stream<Item = Result<Bytes, FaithError>> + Send>>,
661}
662
663impl Debug for ResponseBody {
664	fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
665		f.debug_struct("ResponseBody").finish_non_exhaustive()
666	}
667}
668
669impl http_body::Body for ResponseBody {
670	type Data = Bytes;
671	type Error = FaithError;
672
673	fn poll_frame(
674		mut self: Pin<&mut Self>,
675		cx: &mut Context<'_>,
676	) -> Poll<Option<Result<http_body::Frame<Self::Data>, Self::Error>>> {
677		self.chunks
678			.as_mut()
679			.poll_next(cx)
680			.map(|chunk| chunk.map(|chunk| chunk.map(http_body::Frame::data)))
681	}
682}
683
684impl Response {
685	/// Convert into an [`http::Response`], for code written against the wider ecosystem.
686	///
687	/// Fails if the body is already being consumed elsewhere. A response that cannot carry a body
688	/// yields an empty one.
689	pub fn into_http(self) -> Result<http::Response<ResponseBody>, FaithError> {
690		let chunks: Pin<Box<dyn Stream<Item = Result<Bytes, FaithError>> + Send>> =
691			match self.body_stream()? {
692				Some(stream) => Box::pin(stream),
693				None => Box::pin(stream::empty()),
694			};
695
696		let mut response = http::Response::new(ResponseBody { chunks });
697		*response.status_mut() = self.status_code;
698		*response.version_mut() = self.version;
699		*response.headers_mut() = self.headers.clone();
700		Ok(response)
701	}
702}
703
704impl TryFrom<Response> for http::Response<ResponseBody> {
705	type Error = FaithError;
706
707	fn try_from(response: Response) -> Result<Self, Self::Error> {
708		response.into_http()
709	}
710}