1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
//! Warming an agent's caches ahead of the request that needs them.
// spec:WARM
use std::{future::Future, sync::atomic::Ordering};
use moka::sync::Cache as MokaCache;
use reqwest::Version;
use crate::{
agent::Agent,
error::{FaithError, FaithErrorKind},
warm_up::{extract_host, origin_key, reduce_to_origin},
};
impl Agent {
/// Warm the DNS cache for `host`, so a later request to it skips the lookup.
///
/// The argument is a bare host; a scheme, port, or path in a fuller string is ignored. The
/// returned future completes when the answer lands in the cache and never fails, whatever
/// happens on the network: the work is advisory. Under the system resolver there is no cache to
/// warm, so it completes having done nothing. A host with nothing to resolve, or a closed agent,
/// is refused here rather than by the future.
// spec:WARM
pub fn prefetch_dns(&self, host: &str) -> Result<impl Future<Output = ()> + use<>, FaithError> {
if self.is_closed() {
return Err(FaithErrorKind::Closed.into());
}
let Some(host) = extract_host(host) else {
return Err(FaithErrorKind::AddressParse.into());
};
#[cfg(feature = "dns")]
let resolver = self.dns_resolver();
Ok(async move {
// Nothing to warm without Faith's own resolver: the platform's cache is not ours to fill.
#[cfg(feature = "dns")]
if let Some(resolver) = resolver {
resolver.prefetch(&host).await;
}
#[cfg(not(feature = "dns"))]
let _ = host;
})
}
/// Open a pooled connection to `origin`, so the first request to it skips DNS, TCP, and TLS
/// setup.
///
/// The argument is an origin (`scheme://host[:port]`); a longer URL is reduced to one. The
/// warm-up sends a synthetic `HEAD` to the origin's root -- the origin sees it -- over the
/// transport the next foreground request would use: a confirmed HTTP/3 origin gets a warm QUIC
/// connection, every other origin a TCP one. The returned future completes when the attempt
/// finishes and never fails: every network failure is quiet. Something that cannot be connected
/// to, or a closed agent, is refused here rather than by the future.
// spec:WARM
pub fn preconnect(&self, origin: &str) -> Result<impl Future<Output = ()> + use<>, FaithError> {
let Some(raw_client) = self.raw_client() else {
return Err(FaithErrorKind::Closed.into());
};
let Some(url) = reduce_to_origin(origin) else {
return Err(FaithErrorKind::AddressParse.into());
};
let key = origin_key(&url);
// Already warm within the idle window, or a warm-up for this origin already in flight:
// either way there is no new work to do, so finish without opening a duplicate.
let redundant = self.warmed.contains_key(&key)
|| !self.warming.entry(key.clone()).or_insert(()).is_fresh();
// The transport the next foreground request would take, decided exactly as the Alt-Svc
// layer decides it: nothing upgrades with the machinery off; with a prober, only a
// confirmed origin routes to QUIC (an advertisement is evidence worth probing, not worth
// routing on); without one, the inline upgrade acts on advertisements too. Diverging here
// would warm the wrong transport.
// spec:WARM#preconnect
#[cfg(feature = "http3")]
let h3_port = self
.alt_svc_cache()
.filter(|_| self.h3_upgrade_enabled)
.and_then(|cache| {
if self.h3_prober().is_some() {
cache.confirmed_port(&url)
} else {
cache.should_use_h3(&url)
}
});
#[cfg(not(feature = "http3"))]
let h3_port: Option<u16> = None;
#[cfg(feature = "connection-tracking")]
let conn_tracker = self.conn_tracker.clone();
let warmed = self.warmed.clone();
let warming = self.warming.clone();
// Read before the warm-up starts, to compare against once it finishes.
let warm_generation = self.warm_generation.clone();
let generation = warm_generation.load(Ordering::Relaxed);
Ok(async move {
if redundant {
return;
}
// Release the single-flight claim whatever happens, so a later warm-up is not blocked
// by this one having finished.
struct ReleaseClaim {
warming: MokaCache<String, ()>,
key: String,
}
impl Drop for ReleaseClaim {
fn drop(&mut self) {
self.warming.invalidate(&self.key);
}
}
let _release = ReleaseClaim {
warming,
key: key.clone(),
};
let request = match h3_port {
Some(port) => {
let mut h3_url = url.clone();
// A port differing from the origin's only comes back with the
// follow-advertised-port option on; rewriting the URL is how reqwest is told to
// connect there, mirroring the foreground path.
if Some(port) != h3_url.port_or_known_default() {
let _ = h3_url.set_port(Some(port));
}
raw_client.head(h3_url).version(Version::HTTP_3)
}
None => raw_client.head(url.clone()),
};
let outcome = request.send().await;
// A TCP warm-up leaves a pooled connection to track; a QUIC one does not (QUIC
// connections are not tracked, and a confirmed origin has nothing left to probe).
#[cfg(feature = "connection-tracking")]
if h3_port.is_none()
&& let Ok(response) = &outcome
&& let Some(info) = response
.extensions()
.get::<hyper_util::client::legacy::connect::HttpInfo>()
{
conn_tracker.track_warmup(info.local_addr(), info.remote_addr());
}
// A network change while this was in flight leaves the origin unmarked: the connection
// landed in the pool that change dropped, so it is not warm however well the request
// went.
// spec:NETCHG#reach-across-the-subsystems
if outcome.is_ok() && warm_generation.load(Ordering::Relaxed) == generation {
warmed.insert(key, ());
}
})
}
}