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}