Skip to main content

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}