acme_proxy/listener.rs
1//! The sockets, and replacing one while it is serving.
2//!
3//! This server binds up to three sockets — ACME, the web admin, the metrics
4//! endpoint — and until this module existed each was handed straight to
5//! `axum::serve`, which **consumes** its listener. That made the socket the one
6//! thing a configuration generation could not replace: `server.bind_address`,
7//! `admin.enabled` and both `tls.enabled` flips were refused by name and needed
8//! a restart, while everything else about a reload was a rebuild and a swap.
9//!
10//! What replaces it is one `axum::serve` per role, for the life of the process,
11//! over a [`RoleListener`] that owns the accept loop itself. Two things are
12//! swappable underneath it, and both are read **per connection**:
13//!
14//! - the TCP socket, replaced through a channel — a rebind, with connections
15//! already established untouched, because hyper owns those and only the
16//! socket beneath them changes;
17//! - `Option<TlsSettings>`, so turning TLS on or off is the same kind of change
18//! as renewing a certificate rather than a different listener type. **A
19//! `tls.enabled` flip on an unchanged address therefore rebinds nothing at
20//! all**, which is what removes the one case a bind-then-drain scheme cannot
21//! serve: two listeners cannot hold one port while the old one drains.
22//!
23//! A role with no socket **parks**. That is how `admin.enabled = false` and
24//! `metrics.enabled = false` are expressed, at startup as on a reload, and it is
25//! the same answer this loop already gave when its accept task ended: there is
26//! no error in [`Listener::accept`]'s signature, and returning a connection is
27//! impossible, so parking is what is left.
28//!
29//! Provisioning the TLS material — reading or generating the files, building the
30//! rustls acceptor — stays in [`crate::tls`]. The line between the two modules is
31//! the one that file's own documentation already drew: resolving configuration
32//! at startup on one side, accepting connections on the other.
33
34use std::io;
35use std::net::SocketAddr;
36use std::pin::Pin;
37use std::task::{Context, Poll};
38
39use axum::serve::{Listener, ListenerExt, TapIo};
40use tokio::io::{AsyncRead, AsyncWrite, ReadBuf};
41use tokio::net::{TcpListener, TcpStream};
42use tokio::sync::{mpsc, watch};
43use tokio_rustls::server::TlsStream;
44use tracing::{debug, warn};
45
46use crate::tls::TlsSettings;
47
48/// How many connections may sit between the TCP accept and `axum` — handshakes
49/// in flight plus finished streams not yet picked up.
50///
51/// This is the accept loop's backpressure: it reserves a slot *before* accepting,
52/// so a flood of half-open TLS connections cannot spawn tasks without bound. It
53/// now bounds the cleartext path too, which used to be `axum::serve`'s own
54/// unbounded accept loop — backpressure it did not have, ahead of the admission
55/// control that bounds what happens next.
56const MAX_PENDING_CONNECTIONS: usize = 256;
57
58/// Pause after a failed `accept()`, so a listener that is refusing connections
59/// (out of file descriptors, say) does not spin a core. Mirrors what
60/// `axum::serve` does for the same case.
61const ACCEPT_ERROR_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
62
63/// A connection, however it was accepted.
64///
65/// The `Io` type all three roles share, and the reason one `axum::serve` can
66/// outlive a `tls.enabled` flip: the listener answers with a different variant
67/// from one connection to the next without its own type changing.
68///
69/// The TLS half is boxed because `TlsStream` carries a whole rustls connection
70/// state — roughly two orders of magnitude larger than a `TcpStream` — and this
71/// enum is as big as its largest variant on every cleartext connection too.
72pub enum MaybeTls {
73 Plain(TcpStream),
74 Tls(Box<TlsStream<TcpStream>>),
75}
76
77impl AsyncRead for MaybeTls {
78 fn poll_read(
79 self: Pin<&mut Self>,
80 context: &mut Context<'_>,
81 buffer: &mut ReadBuf<'_>,
82 ) -> Poll<io::Result<()>> {
83 match self.get_mut() {
84 Self::Plain(stream) => Pin::new(stream).poll_read(context, buffer),
85 Self::Tls(stream) => Pin::new(stream.as_mut()).poll_read(context, buffer),
86 }
87 }
88}
89
90impl AsyncWrite for MaybeTls {
91 fn poll_write(
92 self: Pin<&mut Self>,
93 context: &mut Context<'_>,
94 buffer: &[u8],
95 ) -> Poll<io::Result<usize>> {
96 match self.get_mut() {
97 Self::Plain(stream) => Pin::new(stream).poll_write(context, buffer),
98 Self::Tls(stream) => Pin::new(stream.as_mut()).poll_write(context, buffer),
99 }
100 }
101
102 fn poll_flush(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<io::Result<()>> {
103 match self.get_mut() {
104 Self::Plain(stream) => Pin::new(stream).poll_flush(context),
105 Self::Tls(stream) => Pin::new(stream.as_mut()).poll_flush(context),
106 }
107 }
108
109 fn poll_shutdown(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<io::Result<()>> {
110 match self.get_mut() {
111 Self::Plain(stream) => Pin::new(stream).poll_shutdown(context),
112 Self::Tls(stream) => Pin::new(stream.as_mut()).poll_shutdown(context),
113 }
114 }
115
116 /// Delegated rather than left to the default, which would write one buffer
117 /// per call: hyper writes a response's head and body as separate slices, so
118 /// the vectored path is the ordinary one rather than an optimisation.
119 fn poll_write_vectored(
120 self: Pin<&mut Self>,
121 context: &mut Context<'_>,
122 buffers: &[io::IoSlice<'_>],
123 ) -> Poll<io::Result<usize>> {
124 match self.get_mut() {
125 Self::Plain(stream) => Pin::new(stream).poll_write_vectored(context, buffers),
126 Self::Tls(stream) => Pin::new(stream.as_mut()).poll_write_vectored(context, buffers),
127 }
128 }
129
130 fn is_write_vectored(&self) -> bool {
131 match self {
132 Self::Plain(stream) => stream.is_write_vectored(),
133 Self::Tls(stream) => stream.is_write_vectored(),
134 }
135 }
136}
137
138/// What a reload does to one role's socket.
139///
140/// The `Serve` variant carries an **already bound** listener, which is the whole
141/// ordering rule: binding happens where a failure can still refuse the reload,
142/// not here, where nothing could be done about it.
143pub enum SocketCommand {
144 /// Serve this socket from now on, dropping whatever was being served.
145 Serve(TcpListener),
146 /// Stop accepting and release the socket — the role was switched off.
147 Close,
148}
149
150/// The write side of one role's listener: its socket, and its TLS mode.
151///
152/// Held by the reload supervisor and by nothing else. **Both sends are
153/// synchronous** (`watch::Sender::send_replace`, and an unbounded
154/// `mpsc::Sender::send`), which is what lets a socket change sit in the same
155/// uninterruptible publishing run as the routers and the job registry — see
156/// `cli::apply_reload`.
157///
158/// It deliberately remembers **nothing** about what it is serving. Whether a
159/// role should rebind is decided by comparing the applied configuration against
160/// the proposed one — never against the address actually bound, since a caller
161/// supplying its own socket (`serve_on_with`, and every test that binds
162/// `127.0.0.1:0`) is entitled to one that does not match what the file says.
163pub struct ListenerHandle {
164 sockets: mpsc::UnboundedSender<SocketCommand>,
165 tls: watch::Sender<Option<TlsSettings>>,
166}
167
168impl ListenerHandle {
169 /// Replaces the socket this role serves.
170 pub fn serve(&self, listener: TcpListener) {
171 let _ = self.sockets.send(SocketCommand::Serve(listener));
172 }
173
174 /// Stops serving this role, releasing its socket.
175 pub fn close(&self) {
176 let _ = self.sockets.send(SocketCommand::Close);
177 }
178
179 /// Publishes the TLS mode the *next* connection is accepted under. `None`
180 /// is cleartext.
181 pub fn set_tls(&self, settings: Option<TlsSettings>) {
182 self.tls.send_replace(settings);
183 }
184}
185
186/// Binds `address` without awaiting.
187///
188/// `std::net::TcpListener::bind` rather than tokio's, so a rebind can happen
189/// inside `cli::apply_reload` — which is deliberately not `async`, so that its
190/// publishing run has no await point another task could interleave with. The
191/// blocking part is name resolution, on the reload supervisor's own task, where
192/// building a generation already reads and writes files.
193///
194/// # Errors
195///
196/// Whatever the bind failed with — a port in use, an address that does not
197/// resolve, a privileged port. The caller turns it into a refused reload with
198/// the running socket untouched.
199pub fn bind_blocking(address: &str) -> io::Result<TcpListener> {
200 let listener = std::net::TcpListener::bind(address)?;
201 listener.set_nonblocking(true)?;
202 TcpListener::from_std(listener)
203}
204
205/// A listener whose socket and TLS mode can both be replaced while it serves.
206///
207/// Implements `axum::serve::Listener`, so the server keeps being run by
208/// `axum::serve` — including `into_make_service_with_connect_info`, which is what
209/// the IP filters depend on (see [`RoleSocket`]).
210///
211/// **Handshakes happen off the accept path.** A background task accepts TCP
212/// connections and spawns each handshake, so one slow client cannot hold up
213/// every other connection for the length of the timeout; `accept()` only picks
214/// up the results. A handshake that fails or runs out of budget is logged and
215/// dropped — the trait's `accept()` returns no `Result`, which suits this
216/// exactly: a bad connection simply never becomes one.
217pub struct RoleListener {
218 /// Connections ready to be served, oldest first.
219 incoming: mpsc::Receiver<(MaybeTls, SocketAddr)>,
220 /// The last address this role was bound to, for `local_addr` — which axum
221 /// calls per connection, including after the socket it named has gone.
222 local_addr: SocketAddr,
223 /// Republished by the accept task on every rebind, so `local_addr` follows
224 /// the socket rather than answering the address of a listener that is no
225 /// longer there.
226 bound: watch::Receiver<SocketAddr>,
227}
228
229/// The listener type [`spawn`] hands to `axum::serve`.
230///
231/// The `tap_io` wrapper is **functional, not decorative**.
232/// `into_make_service_with_connect_info::<SocketAddr>()` requires
233/// `SocketAddr: Connected<IncomingStream<'_, L>>`, and axum implements that for
234/// exactly two listeners: the concrete `TcpListener`, and — blanket — any
235/// `TapIo<L, F>` whose `L::Addr` is `Clone + Sync + 'static`. We cannot write the
236/// missing impl ourselves: foreign trait, foreign `Self` type, coherence refuses
237/// it. So the wrapper is what makes the peer address reach the request
238/// extensions, and without it `add_filter_middleware` sees no client address and
239/// the IP filters fail closed. Returning the wrapped type from `spawn` is what
240/// keeps a caller from forgetting.
241pub type RoleSocket = TapIo<RoleListener, fn(&mut MaybeTls)>;
242
243/// Starts a role's accept loop.
244///
245/// `initial` is the socket to serve at once, or `None` for a role that is
246/// switched off — which parks until a reload hands it one. `tls` is the mode
247/// every connection is accepted under, read fresh each time so a renewed
248/// certificate, a changed `handshake_timeout_ms` and a `tls.enabled` flip all
249/// land on the next client without disturbing anyone already connected.
250///
251/// The address reported by `local_addr` before anything is bound is
252/// `0.0.0.0:0`; nothing consults it until a connection arrives, and by then the
253/// accept task has published the real one.
254#[must_use]
255pub fn spawn(
256 role: &'static str,
257 initial: Option<TcpListener>,
258 tls: Option<TlsSettings>,
259) -> (RoleSocket, ListenerHandle) {
260 let (sender, incoming) = mpsc::channel(MAX_PENDING_CONNECTIONS);
261 let (sockets_tx, sockets_rx) = mpsc::unbounded_channel();
262 let (tls_tx, tls_rx) = watch::channel(tls);
263
264 let first_addr = initial
265 .as_ref()
266 .and_then(|listener| listener.local_addr().ok())
267 .unwrap_or_else(|| SocketAddr::from(([0, 0, 0, 0], 0)));
268 let (bound_tx, bound_rx) = watch::channel(first_addr);
269
270 tokio::spawn(async move {
271 let mut current = initial;
272 // Taken away once the last [`ListenerHandle`] is dropped, which is not
273 // an ending: `reload::Reloads::none()` drops the whole set the moment
274 // the supervisor sees it will never fire, and every caller that serves
275 // no reloads goes through it. It means only that this socket is now
276 // whatever it is for good.
277 let mut commands = Some(sockets_rx);
278 loop {
279 if current.is_none() && commands.is_none() {
280 // Nothing to accept, and nothing that could ever hand this role
281 // a socket again. Waiting on the listener rather than returning
282 // outright: ending here would close `incoming`, and `accept`
283 // reads a closed `incoming` as the accept task having *died*.
284 sender.closed().await;
285 debug!(
286 event = "server_accept_loop_ended",
287 outcome = "success",
288 listener = role
289 );
290 return;
291 }
292
293 // Reserving before accepting is the backpressure: at most
294 // `MAX_PENDING_CONNECTIONS` connections are in flight, and a
295 // dropped listener closes the channel, which ends this task.
296 let Ok(permit) = sender.clone().reserve_owned().await else {
297 debug!(
298 event = "server_accept_loop_ended",
299 outcome = "success",
300 listener = role
301 );
302 return;
303 };
304
305 // With no socket there is only the command channel to wait on,
306 // which is exactly the parking this module's documentation
307 // describes: a role that is switched off costs one idle task.
308 //
309 // The result is carried out of the `match` rather than acted on
310 // inside it, because the accept arm borrows `current` and the
311 // command arm has to replace it.
312 let next = match (current.as_ref(), commands.as_mut()) {
313 (None, None) => unreachable!("checked at the top of the loop"),
314 (None, Some(commands)) => Next::Command(commands.recv().await),
315 (Some(listener), None) => accept_one(listener, role).await,
316 (Some(listener), Some(commands)) => tokio::select! {
317 // Biased so a pending rebind is taken before another
318 // connection is accepted on the socket being replaced.
319 // Both arms are cancellation-safe.
320 biased;
321 command = commands.recv() => Next::Command(command),
322 accepted = accept_one(listener, role) => accepted,
323 },
324 };
325
326 let (stream, peer) = match next {
327 Next::Connection(connection) => connection,
328 Next::Command(Some(command)) => {
329 apply(&mut current, command, &bound_tx, role);
330 continue;
331 }
332 Next::Command(None) => {
333 commands = None;
334 continue;
335 }
336 Next::Backoff => {
337 tokio::time::sleep(ACCEPT_ERROR_BACKOFF).await;
338 continue;
339 }
340 };
341
342 // Read here rather than captured above: this is the point a
343 // reloaded certificate — or a `tls.enabled` flip — takes effect. The
344 // `Ref` guard is dropped before the spawn, so nothing holds the lock
345 // across an await.
346 let settings = tls_rx.borrow().clone();
347 match settings {
348 // No handshake to run, so the connection is ready as it stands.
349 // `send` hands back the cloned sender the permit came from,
350 // which has nothing left to do.
351 None => {
352 permit.send((MaybeTls::Plain(stream), peer));
353 }
354 Some(TlsSettings {
355 acceptor,
356 handshake_timeout,
357 }) => {
358 tokio::spawn(async move {
359 match tokio::time::timeout(handshake_timeout, acceptor.accept(stream)).await
360 {
361 Ok(Ok(tls)) => {
362 permit.send((MaybeTls::Tls(Box::new(tls)), peer));
363 }
364 // Neither is the server's problem: a port scan, a
365 // cleartext client, a stalled handshake. Never fatal.
366 Ok(Err(error)) => {
367 debug!(event = "tls_handshake_failed", outcome = "failure", peer = %peer, error = %error);
368 }
369 Err(_) => {
370 debug!(event = "tls_handshake_timeout", outcome = "failure", peer = %peer)
371 }
372 }
373 });
374 }
375 }
376 }
377 });
378
379 let listener = RoleListener {
380 incoming,
381 local_addr: first_addr,
382 bound: bound_rx,
383 }
384 // See `RoleSocket`: this is what carries the peer address into the request
385 // extensions. The closure itself has nothing to do.
386 .tap_io(noop_tap as fn(&mut MaybeTls));
387
388 (
389 listener,
390 ListenerHandle {
391 sockets: sockets_tx,
392 tls: tls_tx,
393 },
394 )
395}
396
397/// What one pass of the accept loop produced.
398///
399/// A value rather than three branches acting in place, because the arm that
400/// accepts borrows the current socket and the arm that takes a command has to
401/// replace it — which the borrow checker will not allow inside one `select!`.
402enum Next {
403 Connection((TcpStream, SocketAddr)),
404 /// `None` means every [`ListenerHandle`] has been dropped.
405 Command(Option<SocketCommand>),
406 /// `accept()` failed; pause before trying again so a socket that is
407 /// refusing connections does not spin a core.
408 Backoff,
409}
410
411/// One `accept()`, with a failure turned into a pause rather than an end.
412///
413/// A function so the accept loop can await it both inside a `select!` and on its
414/// own, the two differing only in whether a rebind can interrupt it.
415/// Cancellation-safe, because `TcpListener::accept` is and nothing here holds
416/// state across it.
417async fn accept_one(listener: &TcpListener, role: &'static str) -> Next {
418 match listener.accept().await {
419 Ok(connection) => Next::Connection(connection),
420 Err(error) => {
421 warn!(
422 event = "server_accept_failed",
423 outcome = "failure",
424 listener = role,
425 error = %error
426 );
427 Next::Backoff
428 }
429 }
430}
431
432/// Applies one [`SocketCommand`] to the accept loop's current socket.
433///
434/// Dropping the old listener is what stops new connections reaching the old
435/// address; everything already accepted is unaffected, and everything already
436/// established belongs to hyper.
437fn apply(
438 current: &mut Option<TcpListener>,
439 command: SocketCommand,
440 bound: &watch::Sender<SocketAddr>,
441 role: &'static str,
442) {
443 match command {
444 SocketCommand::Serve(listener) => {
445 if let Ok(address) = listener.local_addr() {
446 bound.send_replace(address);
447 }
448 *current = Some(listener);
449 }
450 SocketCommand::Close => {
451 *current = None;
452 debug!(
453 event = "server_socket_closed",
454 outcome = "success",
455 listener = role
456 );
457 }
458 }
459}
460
461/// The `tap_io` callback. Named rather than a closure so [`RoleSocket`] can
462/// spell its type out.
463fn noop_tap(_stream: &mut MaybeTls) {}
464
465impl Listener for RoleListener {
466 type Io = MaybeTls;
467 type Addr = SocketAddr;
468
469 async fn accept(&mut self) -> (Self::Io, Self::Addr) {
470 match self.incoming.recv().await {
471 Some(connection) => {
472 // Cheap, and the only place it can be refreshed: the accept task
473 // owns the socket, so this is how a rebind reaches `local_addr`.
474 if self.bound.has_changed().unwrap_or(false) {
475 self.local_addr = *self.bound.borrow_and_update();
476 }
477 connection
478 }
479 // The accept task only ends when this listener is dropped, so the
480 // channel closing means it died. There is no error to return in this
481 // signature, and returning a connection is impossible; parking is
482 // what is left.
483 None => {
484 tracing::error!(
485 event = "server_acceptor_stopped",
486 outcome = "failure",
487 "the accept task ended: no further connection will be served"
488 );
489 std::future::pending().await
490 }
491 }
492 }
493
494 fn local_addr(&self) -> io::Result<Self::Addr> {
495 Ok(self.local_addr)
496 }
497}
498
499#[cfg(test)]
500mod tests {
501 use super::*;
502 use crate::config::{ServerConfig, TlsConfig};
503 use crate::testutil::TempDir;
504 use axum::extract::ConnectInfo;
505 use axum::routing::get;
506 use axum::{Router, serve};
507 use rustls::pki_types::ServerName;
508 use std::time::Duration;
509 use tokio::io::{AsyncReadExt, AsyncWriteExt};
510 use tokio_rustls::TlsConnector;
511
512 /// A `ServerConfig` with TLS enabled, its material inside `dir`.
513 fn tls_config(dir: &TempDir, base_url: &str) -> ServerConfig {
514 ServerConfig {
515 bind_address: "127.0.0.1:0".to_string(),
516 base_url: base_url.to_string(),
517 tls: TlsConfig {
518 enabled: true,
519 cert_path: dir.join("server.pem").display().to_string(),
520 key_path: dir.join("server.key").display().to_string(),
521 handshake_timeout_ms: 5_000,
522 },
523 ..ServerConfig::default()
524 }
525 }
526
527 /// One freshly provisioned acceptor, under a directory of its own so two
528 /// calls give two provably distinct certificates.
529 fn settings(name: &str, timeout: Duration) -> TlsSettings {
530 let dir = TempDir::new(name);
531 let acceptor = crate::tls::from_config(&tls_config(&dir, "https://localhost"))
532 .unwrap()
533 .unwrap();
534 // The directory may go: the acceptor holds the parsed material, and
535 // nothing reads the files again.
536 drop(dir);
537 TlsSettings::new(acceptor, timeout)
538 }
539
540 /// Serves `/peer` — which answers with the peer address axum saw, the value
541 /// every IP filter runs on — on a role listener the caller can then poke.
542 ///
543 /// Returns the port it started on and the handle a reload would use.
544 async fn serve_peer(tls: Option<TlsSettings>) -> (u16, ListenerHandle) {
545 let tcp = TcpListener::bind("127.0.0.1:0").await.unwrap();
546 let port = tcp.local_addr().unwrap().port();
547 let (socket, handle) = spawn("test", Some(tcp), tls);
548
549 let app = Router::new().route(
550 "/peer",
551 get(|ConnectInfo(peer): ConnectInfo<SocketAddr>| async move { peer.to_string() }),
552 );
553 tokio::spawn(async move {
554 serve(
555 socket,
556 app.into_make_service_with_connect_info::<SocketAddr>(),
557 )
558 .await
559 .unwrap();
560 });
561 (port, handle)
562 }
563
564 /// One `GET /peer` over TLS, returning the raw response and the address the
565 /// client used.
566 async fn get_peer(port: u16) -> (String, SocketAddr) {
567 let config =
568 crate::challenge::tls_alpn_01::accept_any_client_config(&[b"http/1.1"]).unwrap();
569 let stream = TcpStream::connect(("127.0.0.1", port)).await.unwrap();
570 let client_addr = stream.local_addr().unwrap();
571
572 let mut tls = TlsConnector::from(config)
573 .connect(ServerName::try_from("localhost").unwrap(), stream)
574 .await
575 .unwrap();
576 tls.write_all(b"GET /peer HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n")
577 .await
578 .unwrap();
579
580 let mut response = String::new();
581 tls.read_to_string(&mut response).await.unwrap();
582 (response, client_addr)
583 }
584
585 /// The same request in cleartext, for the arm that speaks no TLS.
586 async fn get_plain(port: u16) -> String {
587 let mut stream = TcpStream::connect(("127.0.0.1", port)).await.unwrap();
588 stream
589 .write_all(b"GET /peer HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n")
590 .await
591 .unwrap();
592 let mut response = String::new();
593 stream.read_to_string(&mut response).await.unwrap();
594 response
595 }
596
597 /// The certificate the server presented on one fresh connection.
598 async fn peer_certificate(port: u16) -> Vec<u8> {
599 let config =
600 crate::challenge::tls_alpn_01::accept_any_client_config(&[b"http/1.1"]).unwrap();
601 let stream = TcpStream::connect(("127.0.0.1", port)).await.unwrap();
602 let tls = TlsConnector::from(config)
603 .connect(ServerName::try_from("localhost").unwrap(), stream)
604 .await
605 .unwrap();
606 tls.get_ref()
607 .1
608 .peer_certificates()
609 .expect("the server presented a certificate")[0]
610 .to_vec()
611 }
612
613 /// Whether anything is listening on `port` at all.
614 async fn refused(port: u16) -> bool {
615 matches!(
616 tokio::time::timeout(
617 Duration::from_secs(2),
618 TcpStream::connect(("127.0.0.1", port)),
619 )
620 .await,
621 Ok(Err(_))
622 )
623 }
624
625 /// The headline case: the socket moves and the `axum::serve` above it does
626 /// not. One request answered on the first port, one on the second, and the
627 /// first refusing afterwards — which together are the whole feature.
628 #[tokio::test]
629 async fn a_replaced_socket_serves_the_new_port_and_releases_the_old() {
630 let (first, handle) = serve_peer(None).await;
631 assert!(get_plain(first).await.starts_with("HTTP/1.1 200 OK"));
632
633 let replacement = TcpListener::bind("127.0.0.1:0").await.unwrap();
634 let second = replacement.local_addr().unwrap().port();
635 handle.serve(replacement);
636
637 // The same server, the same router, a different socket.
638 let response = get_plain(second).await;
639 assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}");
640 assert!(
641 refused(first).await,
642 "the old socket must be released, not merely ignored"
643 );
644 }
645
646 /// `tls.enabled` flipped with the address unchanged, which is the case a
647 /// bind-first-then-drain scheme could not serve at all: two listeners
648 /// cannot hold one port. Here nothing is rebound — the mode is read per
649 /// connection, so the very next client speaks the new protocol.
650 #[tokio::test]
651 async fn tls_can_be_switched_on_without_the_socket_moving() {
652 let (port, handle) = serve_peer(None).await;
653 assert!(get_plain(port).await.starts_with("HTTP/1.1 200 OK"));
654
655 handle.set_tls(Some(settings("listener-flip", Duration::from_secs(5))));
656
657 let (response, _) = get_peer(port).await;
658 assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}");
659
660 // And back again, since an operator who turns TLS on can turn it off.
661 handle.set_tls(None);
662 assert!(get_plain(port).await.starts_with("HTTP/1.1 200 OK"));
663 }
664
665 /// Closing a role releases its socket and leaves the listener parked rather
666 /// than dead: a later reload hands it a new one and it serves again. That
667 /// is `admin.enabled` and `metrics.enabled` going off and on.
668 #[tokio::test]
669 async fn a_closed_role_refuses_connections_and_can_be_reopened() {
670 let (port, handle) = serve_peer(None).await;
671 assert!(get_plain(port).await.starts_with("HTTP/1.1 200 OK"));
672
673 handle.close();
674 // The command is applied by the accept task, so give it a turn.
675 tokio::task::yield_now().await;
676 assert!(refused(port).await, "a closed role must not be listening");
677
678 let reopened = TcpListener::bind("127.0.0.1:0").await.unwrap();
679 let again = reopened.local_addr().unwrap().port();
680 handle.serve(reopened);
681 assert!(get_plain(again).await.starts_with("HTTP/1.1 200 OK"));
682 }
683
684 /// A renewed certificate reaches the next client without the socket
685 /// moving — the reason the accept loop reads its settings per connection
686 /// instead of capturing them.
687 ///
688 /// Two whole certificates rather than a renewed one: provisioning generates
689 /// a fresh key each time, so two directories give two provably distinct
690 /// DERs, which is all the assertion needs.
691 #[tokio::test]
692 async fn a_swapped_certificate_is_served_to_the_next_connection() {
693 let (port, handle) =
694 serve_peer(Some(settings("listener-first", Duration::from_secs(5)))).await;
695
696 let before = peer_certificate(port).await;
697 handle.set_tls(Some(settings("listener-second", Duration::from_secs(5))));
698 let after = peer_certificate(port).await;
699
700 // Both connections went to the same `port`, which is the half that
701 // matters: the certificate changed without the socket being rebound.
702 assert_ne!(
703 before, after,
704 "the connection after the swap must see the new certificate"
705 );
706 }
707
708 /// The end-to-end proof: a real handshake, a real request, and — the point
709 /// of the test — `ConnectInfo` surviving the TLS wrapper. Without it every
710 /// IP filter would fail closed.
711 #[tokio::test]
712 async fn a_request_is_served_with_the_peer_address_intact() {
713 let (port, _handle) =
714 serve_peer(Some(settings("listener-peer", Duration::from_secs(5)))).await;
715 let (response, client_addr) = get_peer(port).await;
716
717 assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}");
718 assert!(
719 response.ends_with(&client_addr.to_string()),
720 "expected the body to be {client_addr}, got {response}"
721 );
722 }
723
724 /// The same guarantee on the cleartext arm, which used to be
725 /// `axum::serve`'s own accept loop and is now this one: the peer address
726 /// has to survive `MaybeTls::Plain` as well.
727 #[tokio::test]
728 async fn a_cleartext_request_keeps_its_peer_address_too() {
729 let (port, _handle) = serve_peer(None).await;
730 let mut stream = TcpStream::connect(("127.0.0.1", port)).await.unwrap();
731 let client_addr = stream.local_addr().unwrap();
732 stream
733 .write_all(b"GET /peer HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n")
734 .await
735 .unwrap();
736 let mut response = String::new();
737 stream.read_to_string(&mut response).await.unwrap();
738
739 assert!(response.ends_with(&client_addr.to_string()), "{response}");
740 }
741
742 /// A cleartext client (or a port scan) fails the handshake without taking
743 /// the listener down with it.
744 #[tokio::test]
745 async fn a_failed_handshake_does_not_stop_the_listener() {
746 let (port, _handle) =
747 serve_peer(Some(settings("listener-scan", Duration::from_secs(5)))).await;
748
749 let mut plain = TcpStream::connect(("127.0.0.1", port)).await.unwrap();
750 plain.write_all(b"GET / HTTP/1.1\r\n\r\n").await.unwrap();
751 // The server answers a TLS alert, not HTTP.
752 let mut ignored = Vec::new();
753 let _ = plain.read_to_end(&mut ignored).await;
754
755 let (response, _) = get_peer(port).await;
756 assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}");
757 }
758
759 /// A client that connects and then says nothing is dropped when its budget
760 /// runs out — and, crucially, does not hold up anyone else while it stalls.
761 /// That is the whole reason handshakes are spawned rather than run inside
762 /// `accept()`.
763 #[tokio::test]
764 async fn a_stalled_handshake_times_out_without_blocking_others() {
765 let (port, _handle) =
766 serve_peer(Some(settings("listener-stall", Duration::from_millis(300)))).await;
767
768 let mut stalled = TcpStream::connect(("127.0.0.1", port)).await.unwrap();
769
770 // Served while the first connection is still stuck mid-handshake.
771 let (response, _) = tokio::time::timeout(Duration::from_secs(5), get_peer(port))
772 .await
773 .expect("a stalled handshake must not block the accept loop");
774 assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}");
775
776 // And the stalled one is eventually dropped, not held forever.
777 let mut buffer = [0u8; 1];
778 let read = tokio::time::timeout(Duration::from_secs(5), stalled.read(&mut buffer))
779 .await
780 .expect("the handshake timeout must close the connection");
781 assert!(
782 matches!(read, Ok(0) | Err(_)),
783 "expected EOF after the handshake timeout, got {read:?}"
784 );
785 }
786
787 /// `bind_blocking` is what a reload binds through, so both of its answers
788 /// matter: a usable socket, and an error the refusal can quote rather than
789 /// a panic.
790 #[tokio::test]
791 async fn bind_blocking_binds_or_says_why_not() {
792 let listener = bind_blocking("127.0.0.1:0").expect("an ephemeral port must bind");
793 let port = listener.local_addr().unwrap().port();
794
795 let error = bind_blocking(&format!("127.0.0.1:{port}"))
796 .expect_err("the port is taken, and saying so is what refuses a reload");
797 assert_eq!(error.kind(), std::io::ErrorKind::AddrInUse);
798
799 assert!(
800 bind_blocking("not-an-address").is_err(),
801 "an unparseable address must not panic on the reload path"
802 );
803 }
804}