Skip to main content

web_faith/
timing.rs

1//! Per-request timing.
2
3// spec:RESP#request-timing
4
5use std::time::Instant;
6
7use reqwest::{Url, Version};
8use tokio::sync::watch;
9
10/// The timing of one request, filled in as it progresses.
11///
12/// Timings are best-effort: internal limitations mean they are not always perfectly accurate.
13#[derive(Clone, Debug, Default)]
14#[non_exhaustive]
15pub struct RequestTiming {
16	/// Milliseconds from the start of the request to the response head being read.
17	pub headers_ms: f64,
18	/// Milliseconds from the start of the request to the body finishing, once it has.
19	pub body_ms: Option<f64>,
20	/// Whether the request travelled on a connection that was already in the pool.
21	pub reused: bool,
22	/// The ALPN Protocol ID of the protocol the request travelled over.
23	pub next_hop_protocol: String,
24	/// The response's `Content-Encoding`, captured before a decoded body's header is stripped.
25	pub content_encoding: Option<String>,
26	/// Whether the response was served by the HTTP cache.
27	pub from_cache: bool,
28}
29
30/// Where the timing lands: written by whoever finishes the body, awaited by `timing()`.
31///
32/// A watch channel for the same reason the trailers slot is one: a body that is never read never
33/// finishes, so the wait is unbounded.
34#[derive(Debug)]
35pub struct TimingSlot {
36	tx: watch::Sender<RequestTiming>,
37	started: Instant,
38}
39
40impl TimingSlot {
41	pub fn new(started: Instant, timing: RequestTiming) -> Self {
42		Self {
43			tx: watch::channel(timing).0,
44			started,
45		}
46	}
47
48	/// Record that the body ended, if nothing got there first.
49	///
50	/// `send_if_modified` keeps the read and write one step, and wakes waiters only from the call
51	/// that settled it. Every route out of a body lands here: the stream ending, `discard()`, and
52	/// the collector draining an abandoned body.
53	pub fn ended(&self) {
54		let elapsed = self.started.elapsed().as_secs_f64() * 1000.0;
55		self.tx.send_if_modified(|timing| {
56			if timing.body_ms.is_none() {
57				timing.body_ms = Some(elapsed);
58				true
59			} else {
60				false
61			}
62		});
63	}
64
65	/// Wait until the body has finished.
66	pub async fn settled(&self) -> RequestTiming {
67		let mut rx = self.tx.subscribe();
68		// `wait_for` tests the current value before waiting, so a body that already finished
69		// returns without yielding. Its error case is the sender being gone, which means the
70		// response was dropped with the body unread: the phases reached before that are all
71		// there is to report, so report them rather than waiting for a moment that can no
72		// longer come.
73		let settled = match rx.wait_for(|timing| timing.body_ms.is_some()).await {
74			Ok(timing) => Some(timing.clone()),
75			Err(_) => None,
76		};
77		settled.unwrap_or_else(|| rx.borrow().clone())
78	}
79}
80
81/// The ALPN Protocol ID (RFC 7301) for the protocol a response travelled over.
82///
83/// Reported whether or not ALPN negotiated it, as a browser does: cleartext HTTP/2 is `h2c` and
84/// cleartext HTTP/1.1 is `http/1.1`.
85pub(crate) fn alpn_protocol_id(version: Version, url: &Url) -> String {
86	let secure = url.scheme() == "https";
87	match version {
88		Version::HTTP_3 => "h3",
89		Version::HTTP_2 => {
90			if secure {
91				"h2"
92			} else {
93				"h2c"
94			}
95		}
96		Version::HTTP_11 => "http/1.1",
97		Version::HTTP_10 => "http/1.0",
98		Version::HTTP_09 => "http/0.9",
99		_ => "",
100	}
101	.to_owned()
102}