diff --git a/doc/bin/relay/index.md b/doc/bin/relay/index.md index 0d449cb0bd..969ab37397 100644 --- a/doc/bin/relay/index.md +++ b/doc/bin/relay/index.md @@ -66,8 +66,9 @@ let web = relay.web().routes().route("/hello", get(|| async { "hello" })); relay.with_web(web).run().await?; ``` -`Relay::load` binds QUIC and web sockets. Read their actual addresses with -`quic_addr()` and `web_addrs()`, including ports assigned for `:0`. Clone +`Relay::load` binds every socket, so a taken port fails there. Read the actual +addresses with `quic_addr()`, `tcp_addr()`, `web_addrs()`, and +`internal().addr()`, including ports assigned for `:0`. Clone `ready()` before spawning `run`, then await `ready.wait()` when startup must finish before other workers begin. `config()` returns the resolved settings; `cluster().id()` returns the chosen origin ID. `with_listeners()` registers an diff --git a/rs/moq-cli/src/auth.rs b/rs/moq-cli/src/auth.rs index 6eeced4201..3061bbc1ea 100644 --- a/rs/moq-cli/src/auth.rs +++ b/rs/moq-cli/src/auth.rs @@ -823,30 +823,19 @@ mod tests { request.query = Some("jwt=secret".into()); let _reg = sessions.register(request); - let listen = std::net::TcpListener::bind("127.0.0.1:0") - .expect("probe bind") - .local_addr() - .expect("probe addr"); let mut internal_config = moq_relay::internal::Config::default(); - internal_config.listen = Some(listen); + internal_config.listen = Some("127.0.0.1:0".parse().unwrap()); let internal = moq_relay::internal::Internal::new(internal_config, moq_tokio::moq_net::stats::Registry::disabled()) - .with_sessions(sessions); + .with_sessions(sessions) + .bind() + .expect("bind internal listener"); + let listen = internal.addr().expect("internal listener is configured"); tokio::spawn(async move { let _ = internal.run().await; }); let url = format!("http://{listen}"); - let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); - while reqwest::Client::new() - .get(format!("{url}/health")) - .send() - .await - .is_err() - { - assert!(std::time::Instant::now() < deadline, "internal listener never came up"); - tokio::time::sleep(std::time::Duration::from_millis(25)).await; - } run(&["moq", "auth", "revalidate", "--internal-url", &url, "--id", "abc"]) .await diff --git a/rs/moq-relay/src/internal.rs b/rs/moq-relay/src/internal.rs index c3e576a8a8..fa97f3dfc8 100644 --- a/rs/moq-relay/src/internal.rs +++ b/rs/moq-relay/src/internal.rs @@ -89,6 +89,8 @@ pub struct Internal { health: moq_tokio::accept::Health, listeners: Vec, uring: Vec, + listener: Option, + addr: Option, } #[derive(Clone)] @@ -122,9 +124,31 @@ impl Internal { health, listeners, uring: Vec::new(), + listener: None, + addr: None, } } + /// Bind the configured listener now, so [`addr`](Self::addr) reports an ephemeral port before serving. + pub fn bind(mut self) -> anyhow::Result { + if let Some(listen) = self.config.listen + && self.listener.is_none() + { + let listener = moq_tokio::bind::tcp(listen).context("failed to bind internal listener")?; + let addr = listener + .local_addr() + .context("failed to resolve internal bind address")?; + self.addr = Some(addr); + self.listener = Some(listener); + } + Ok(self) + } + + /// The bound address after [`bind`](Self::bind) or [`crate::Relay::load`], if configured. + pub fn addr(&self) -> Option { + self.addr + } + /// Report other listeners' accept health at `/metrics`. /// /// Takes an iterator so the accessors feed it directly, however many listeners @@ -220,18 +244,7 @@ impl Internal { /// resolves), so it drops cleanly into a `select!` as a disabled no-op - /// mirroring how the relay treats other optional services. pub async fn serve(self, app: Router) -> anyhow::Result<()> { - let listener = self.bind()?; - self.serve_bound(app, listener).await - } - - pub(crate) fn bind(&self) -> anyhow::Result> { - self.config - .listen - .map(|listen| moq_tokio::bind::tcp(listen).context("failed to bind internal listener")) - .transpose() - } - - pub(crate) async fn serve_bound(self, app: Router, listener: Option) -> anyhow::Result<()> { + let Internal { listener, health, .. } = self.bind()?; let Some(listener) = listener else { std::future::pending::<()>().await; return Ok(()); @@ -241,7 +254,7 @@ impl Internal { // that single top-level layer, matching `Web::serve` / `Cluster::run`. // No accept-time work: the ops router never hands a connection to qmux, so // capturing a descriptor per health check would spend one for nothing. - crate::listener::server(listener, self.health, DefaultAcceptor::new())? + crate::listener::server(listener, health, DefaultAcceptor::new())? .serve(app.into_make_service()) .await?; Ok(()) @@ -803,32 +816,23 @@ mod tests { /// and forking both the socket options and the disabled-listener contract. #[tokio::test] async fn serve_hosts_merged_routes_alongside_the_defaults() { - // A throwaway bind picks a free port, released before `serve` claims it - // for real (`bind::tcp` sets SO_REUSEADDR, and nothing ever connected). - let listen = std::net::TcpListener::bind("127.0.0.1:0") - .expect("probe bind") - .local_addr() - .expect("probe addr"); - - let internal = Internal::new(Config { listen: Some(listen) }, moq_net::stats::Registry::disabled()); + let listen = Some("127.0.0.1:0".parse().unwrap()); + let internal = Internal::new(Config { listen }, moq_net::stats::Registry::disabled()) + .bind() + .expect("bind internal listener"); + let addr = internal.addr().expect("internal listener is configured"); let app = internal .routes() .merge(Router::new().route("/embedder", get(async || "embedded\n"))); let server = tokio::spawn(internal.serve(app)); - // `serve` binds inside the task, so poll rather than assume it is up the - // instant the spawn returns. let client = reqwest::Client::new(); - let url = format!("http://{listen}"); - let mut embedder = None; - for _ in 0..200 { - if let Ok(res) = client.get(format!("{url}/embedder")).send().await { - embedder = Some(res); - break; - } - tokio::time::sleep(std::time::Duration::from_millis(10)).await; - } - let embedder = embedder.expect("internal listener never accepted a connection"); + let url = format!("http://{addr}"); + let embedder = client + .get(format!("{url}/embedder")) + .send() + .await + .expect("embedder request"); assert_eq!(embedder.status(), reqwest::StatusCode::OK); assert_eq!(embedder.text().await.expect("embedder body"), "embedded\n"); diff --git a/rs/moq-relay/src/relay.rs b/rs/moq-relay/src/relay.rs index ba04e67570..082ab27978 100644 --- a/rs/moq-relay/src/relay.rs +++ b/rs/moq-relay/src/relay.rs @@ -71,7 +71,7 @@ impl Ready { pub struct Relay { ready: tokio::sync::watch::Sender, config: Config, - server: moq_tokio::Server, + server: moq_tokio::Listener, client: moq_tokio::Client, auth: auth::Auth, /// The sessions the embedder decides, until it takes them. `None` when the @@ -107,10 +107,10 @@ impl Relay { /// Assemble a relay from its configuration: bind the listeners, resolve /// auth, and build the cluster with its cache and stats attached. /// - /// This performs the side effects of starting up (binding sockets, reading - /// key material, spawning the cache governor), so a returned `Relay` is - /// ready to serve; nothing accepts a connection until [`Self::run`] drives - /// it. + /// This performs the side effects of starting up (binding every socket, + /// reading key material, spawning the cache governor), so a returned `Relay` + /// reports its ephemeral ports and a taken port fails here; no session is + /// admitted until [`Self::run`] drives it. pub async fn load(mut config: Config) -> anyhow::Result { config.resolve()?; let resolved_config = config.clone(); @@ -274,6 +274,10 @@ impl Relay { .with_versions(server_versions) .with_sessions(sessions.clone()) .bind()?; + // `bind`, not `listen`: the TCP/Unix accept loops handshake as soon as + // they run, and `load` is not yet willing to take a session. `run` + // starts them on its first accept. + let server = server.bind().await.context("failed to bind listeners")?; // Internal (ops) listener (plain HTTP, opt-in via `--internal-listen`) for // /metrics + /health + /nodes, separate from the customer-facing web server. No-op @@ -284,7 +288,8 @@ impl Relay { .with_cluster(&cluster) .with_sessions(sessions.clone()) .with_listeners(web.accept_health()) - .with_listeners(server.accept_health()); + .with_listeners(server.accept_health()) + .bind()?; // Bound but not yet serving: registering here (rather than after the // threads start) is what gives every worker a series from the first // scrape, including one that is about to fail setup. @@ -352,6 +357,11 @@ impl Relay { self.web.addrs() } + /// The bound plain TCP (qmux) address, or `None` when `listen.tcp.bind` is unset. + pub fn tcp_addr(&self) -> Option { + self.server.tcp_local_addr() + } + /// The client used to dial cluster peers. Already handed to [`Self::cluster`]; /// clone it for your own outbound dials so they share the connection config. pub fn client(&self) -> &moq_tokio::Client { @@ -505,10 +515,6 @@ impl Relay { .context("failed to start the io_uring QUIC workers")?; } - // Bind the shared TCP/Unix and optional internal sockets before reporting - // readiness, so an unavailable port cannot leave a falsely ready relay. - let server = server.listen().await.context("failed to bind listeners")?; - let internal_listener = internal.bind()?; ready.send_replace(true); #[cfg(unix)] @@ -616,7 +622,7 @@ impl Relay { let result = tokio::select! { Err(err) = started.run() => Err(err).context("cluster failed"), Err(err) = web.serve(web_routes) => Err(err).context("web server failed"), - Err(err) = internal.serve_bound(internal_routes, internal_listener) => Err(err).context("internal server failed"), + Err(err) = internal.serve(internal_routes) => Err(err).context("internal server failed"), Err(err) = serve_shared => Err(err).context("server failed"), Err(err) = quic_workers => Err(err).context("QUIC workers failed"), err = uring_failed => Err(err).context("io_uring QUIC workers failed"), diff --git a/rs/moq-relay/tests/auth_lifetime.rs b/rs/moq-relay/tests/auth_lifetime.rs index a4e421c8ad..63b2aa18a4 100644 --- a/rs/moq-relay/tests/auth_lifetime.rs +++ b/rs/moq-relay/tests/auth_lifetime.rs @@ -9,7 +9,6 @@ //! The last tests swap the server for an in-process decider answering //! `Admissions`, and prove the lease it drives reaches the session the same way. -use std::net::TcpListener; use std::sync::{Arc, Mutex}; use std::time::{Duration, SystemTime}; @@ -125,23 +124,6 @@ fn build_auth(url: url::Url) -> moq_relay::auth::Auth { .expect("auth init") } -/// Wait for a TCP listener to become dialable, or panic. -async fn wait_for_listener(port: u16) { - let deadline = std::time::Instant::now() + Duration::from_secs(5); - while tokio::net::TcpStream::connect(("127.0.0.1", port)).await.is_err() { - assert!( - std::time::Instant::now() < deadline, - "relay listener never became ready on port {port}" - ); - tokio::time::sleep(Duration::from_millis(25)).await; - } -} - -fn free_port() -> u16 { - let probe = TcpListener::bind("127.0.0.1:0").expect("bind probe"); - probe.local_addr().expect("local addr").port() -} - /// Stand up the relay's accept loop on a plain-TCP qmux listener and return the /// port plus an abort handle. async fn spawn_relay(auth: moq_relay::auth::Auth) -> (u16, tokio::task::JoinHandle<()>) { @@ -155,12 +137,12 @@ async fn spawn_relay_with( cluster: cluster::Cluster, ) -> (u16, tokio::task::JoinHandle<()>) { let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); - let port = free_port(); let mut config = moq_tokio::listen::Config::default(); - config.tcp.bind = Some(format!("127.0.0.1:{port}").parse().expect("parse addr")); + config.tcp.bind = Some("127.0.0.1:0".parse().expect("parse addr")); let server = config.init(Default::default()).expect("server init"); let mut server = server.listen().await.expect("listen"); + let port = server.tcp_local_addr().expect("TCP listener is configured").port(); let handle = tokio::spawn(async move { let mut id = 0; @@ -173,7 +155,6 @@ async fn spawn_relay_with( } }); - wait_for_listener(port).await; (port, handle) } @@ -181,7 +162,6 @@ async fn spawn_relay_with( /// port plus an abort handle. async fn spawn_ws_relay(auth: moq_relay::auth::Auth) -> (u16, tokio::task::JoinHandle<()>) { let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); - let port = free_port(); let cluster = cluster::Cluster::new(cluster::Options::default()).expect("cluster init"); // Stream listeners bind lazily, so this server never opens a socket; only @@ -196,14 +176,16 @@ async fn spawn_ws_relay(auth: moq_relay::auth::Auth) -> (u16, tokio::task::JoinH let mut web_config = web::Config::default(); web_config.ws = true; - web_config.http.listen = Some(format!("127.0.0.1:{port}").parse().expect("parse listen")); - let web = web::Web::new(auth, cluster, certificates, web_config); + web_config.http.listen = Some("127.0.0.1:0".parse().expect("parse listen")); + let web = web::Web::new(auth, cluster, certificates, web_config) + .bind() + .expect("bind web listener"); + let port = web.addrs().http.expect("HTTP listener is configured").port(); let handle = tokio::spawn(async move { let _ = web.run().await; }); - wait_for_listener(port).await; (port, handle) } @@ -1023,26 +1005,25 @@ async fn a_fixed_lease_still_expires() { #[tokio::test] async fn a_relay_without_an_auth_source_is_decided_by_the_embedder() { let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); - let config = |port: u16| { + let config = || { let mut config = Config::default(); - config.listen.tcp.bind = Some(format!("127.0.0.1:{port}").parse().expect("parse addr")); + config.listen.tcp.bind = Some("127.0.0.1:0".parse().expect("parse addr")); // The sessions are gone by the time the trigger fires; no need to wait out the default window. config.drain_timeout = Duration::from_millis(100); config }; - let untaken = Relay::load(config(free_port())).await.expect("load relay"); + let untaken = Relay::load(config()).await.expect("load relay"); let err = untaken.run().await.expect_err("nobody can authenticate"); assert!(err.to_string().contains("nobody can authenticate"), "{err}"); - let port = free_port(); - let mut relay = Relay::load(config(port)).await.expect("load relay"); + let mut relay = Relay::load(config()).await.expect("load relay"); + let port = relay.tcp_addr().expect("TCP listener bound").port(); let admissions = relay.admissions().expect("an empty [auth] hands over the admissions"); assert!(relay.admissions().is_none(), "taken once"); let decider = Decider::spawn(admissions, grant(Duration::from_secs(3600))); let trigger = relay.shutdown_trigger().clone(); let running = tokio::spawn(relay.run()); - wait_for_listener(port).await; let (pub_session, sub_session) = connect_and_round_trip(&room_url("tcp", port)).await; assert_eq!(decider.seen.lock().unwrap().len(), 2, "one admission per session"); diff --git a/rs/moq-relay/tests/cluster_unknown.rs b/rs/moq-relay/tests/cluster_unknown.rs index 9ad71cdde3..5409b85818 100644 --- a/rs/moq-relay/tests/cluster_unknown.rs +++ b/rs/moq-relay/tests/cluster_unknown.rs @@ -2,7 +2,7 @@ //! The relay records that publisher as `Hop::UNKNOWN`; reflected cluster paths //! must not replace it while gossiping around a redundant mesh. -use std::{net::TcpListener, time::Duration}; +use std::time::Duration; use moq_relay::{Config, Relay}; use url::Url; @@ -10,24 +10,15 @@ use url::Url; const TIMEOUT: Duration = Duration::from_secs(10); const PATH: &str = "opalin/cell-clumsy-octopus/cameras/left.hang"; -fn free_tcp_port() -> u16 { - TcpListener::bind("127.0.0.1:0") - .expect("bind probe") - .local_addr() - .expect("local addr") - .port() -} - async fn spawn_relay( id: u64, connect: Vec, cluster_version: Option, ) -> (u16, tokio::task::JoinHandle<()>) { let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); - let port = free_tcp_port(); let mut config = Config::default(); - config.listen.tcp.bind = Some(format!("127.0.0.1:{port}").parse().expect("parse bind")); + config.listen.tcp.bind = Some("127.0.0.1:0".parse().expect("parse bind")); config.connect.bind = Some("127.0.0.1:0".parse().expect("parse client bind")); config.connect.tls.insecure = Some(true); config.connect.version.extend(cluster_version); @@ -38,20 +29,13 @@ async fn spawn_relay( config.cluster.id = Some(id); config.cluster.connect = connect.into_iter().map(moq_relay::cluster::Peer::new).collect(); + // `load` binds the TCP listener, so the port is ours before anyone dials it. let relay = Relay::load(config).await.expect("relay load"); + let port = relay.tcp_addr().expect("TCP listener is configured").port(); let handle = tokio::spawn(async move { let _ = relay.run().await; }); - let deadline = std::time::Instant::now() + Duration::from_secs(5); - loop { - if tokio::net::TcpStream::connect(("127.0.0.1", port)).await.is_ok() { - break; - } - assert!(std::time::Instant::now() < deadline, "relay {id} never became ready"); - tokio::time::sleep(Duration::from_millis(25)).await; - } - (port, handle) } diff --git a/rs/moq-relay/tests/embed.rs b/rs/moq-relay/tests/embed.rs index e65b5577e9..bd8b4bfbb8 100644 --- a/rs/moq-relay/tests/embed.rs +++ b/rs/moq-relay/tests/embed.rs @@ -9,8 +9,6 @@ #![cfg(feature = "_quic")] -#[cfg(target_os = "linux")] -use std::net::UdpSocket; use std::net::{SocketAddr, TcpListener}; use std::time::Duration; @@ -19,23 +17,6 @@ use moq_tokio::moq_net; const TIMEOUT: Duration = Duration::from_secs(10); -fn free_tcp_port() -> u16 { - let probe = TcpListener::bind("127.0.0.1:0").expect("bind probe"); - let port = probe.local_addr().expect("local addr").port(); - drop(probe); - port -} - -/// Only used by the Linux-only worker/uring tests below; without the gate the -/// macOS test build fails `-D warnings` on dead code. -#[cfg(target_os = "linux")] -fn free_udp_port() -> u16 { - let probe = UdpSocket::bind("127.0.0.1:0").expect("bind probe"); - let port = probe.local_addr().expect("local addr").port(); - drop(probe); - port -} - fn certificate(dir: &std::path::Path) -> (std::path::PathBuf, std::path::PathBuf) { let key = rcgen::KeyPair::generate().expect("keypair"); let params = rcgen::CertificateParams::new(vec!["localhost".to_string()]).expect("cert params"); @@ -59,19 +40,6 @@ fn client() -> moq_tokio::Client { config.init(Default::default()).expect("client init") } -async fn wait_for_http(port: u16) { - let deadline = std::time::Instant::now() + Duration::from_secs(5); - loop { - if tokio::net::TcpStream::connect(("127.0.0.1", port)).await.is_ok() { - return; - } - if std::time::Instant::now() >= deadline { - panic!("relay http listener never became ready on port {port}"); - } - tokio::time::sleep(Duration::from_millis(25)).await; - } -} - async fn assert_owner_stopped(quic: SocketAddr, http: SocketAddr) { assert!( tokio::net::TcpStream::connect(http).await.is_err(), @@ -109,17 +77,17 @@ async fn embed_and_stop(mut config: Config) { // No drain window: the sessions are already gone by the time the owner // stops, and the test should not wait out the default. config.drain_timeout = Duration::ZERO; - let http = config.web.http.listen.expect("http listener configured"); let relay = Relay::load(config.clone()).await.expect("load relay"); let quic = relay.quic_addr().expect("quic listener bound"); - assert_eq!(relay.web_addrs().http, Some(http)); + let http = relay.web_addrs().http.expect("http listener bound"); assert_eq!( relay.config().quic.max_streams, Some(moq_tokio::quic::DEFAULT_MAX_STREAMS) ); assert_eq!(relay.cluster().id(), relay.cluster().origin.hop().id()); - // Pin the replacement to the same ports, including a `:0` first bind. + // Pin the replacement to the same ports the `:0` first binds got. config.listen.bind = Some(moq_tokio::listen::Bind::Addr(quic)); + config.web.http.listen = Some(http); // The application handles: in-process workers publish into the origin the // QUIC sessions see, and the trigger stops the owner from any task. Both @@ -138,8 +106,6 @@ async fn embed_and_stop(mut config: Config) { .route("/plain-post", axum::routing::post(|| async { "plain" })); let running = tokio::spawn(relay.with_web(web).run()); ready.wait().await.expect("relay ready"); - - wait_for_http(http.port()).await; assert!(!running.is_finished(), "the relay stopped while serving"); let body = reqwest::get(format!("http://127.0.0.1:{}/embedded", http.port())) @@ -287,7 +253,6 @@ async fn embed_and_stop(mut config: Config) { assert_eq!(replacement.addr(), Some(quic), "replacement bound a different address"); let trigger = replacement.shutdown_trigger().clone(); let replacing = tokio::spawn(replacement.run()); - wait_for_http(http.port()).await; let health = reqwest::get(format!("http://127.0.0.1:{}/health", http.port())) .await .expect("replacement health") @@ -304,55 +269,39 @@ fn http_and_quic(cert: &std::path::Path, key: &std::path::Path, quic_bind: Strin config.listen.bind = Some(quic_bind.parse().unwrap()); config.listen.tls.cert = vec![cert.to_path_buf()]; config.listen.tls.key = vec![key.to_path_buf()]; - config.web.http.listen = Some(format!("127.0.0.1:{}", free_tcp_port()).parse().expect("parse http")); + config.web.http.listen = Some("127.0.0.1:0".parse().expect("parse http")); config.web.ws = false; public_auth(&mut config); config } -/// A late TCP bind failure must close readiness without reporting success. +/// An occupied TCP port fails `load`, before anything could report readiness. #[tokio::test] -async fn tcp_bind_failure_does_not_report_ready() { +async fn tcp_bind_failure_fails_load() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); let occupied = TcpListener::bind("127.0.0.1:0").expect("reserve TCP port"); let mut config = http_and_quic(&cert, &key, "127.0.0.1:0".into()); config.listen.tcp.bind = Some(occupied.local_addr().expect("reserved address")); - let relay = Relay::load(config).await.expect("load relay before TCP bind"); - let ready = relay.ready(); - let running = tokio::spawn(relay.run()); - - let result = tokio::time::timeout(TIMEOUT, ready.wait()) + let error = Relay::load(config) .await - .expect("readiness never resolved"); - assert!(result.is_err(), "failed TCP bind reported readiness"); - let error = running - .await - .expect("run panicked") - .expect_err("run accepted an occupied TCP port"); + .err() + .expect("load accepted an occupied TCP port"); assert!(error.to_string().contains("failed to bind listeners"), "{error:#}"); } -/// An occupied internal port must fail before the relay reports readiness. +/// An occupied internal port fails `load`, before anything could report readiness. #[tokio::test] -async fn internal_bind_failure_does_not_report_ready() { +async fn internal_bind_failure_fails_load() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); let occupied = TcpListener::bind("127.0.0.1:0").expect("reserve internal port"); let mut config = http_and_quic(&cert, &key, "127.0.0.1:0".into()); config.internal.listen = Some(occupied.local_addr().expect("reserved address")); - let relay = Relay::load(config).await.expect("load relay before internal bind"); - let ready = relay.ready(); - let running = tokio::spawn(relay.run()); - - let result = tokio::time::timeout(TIMEOUT, ready.wait()) - .await - .expect("readiness never resolved"); - assert!(result.is_err(), "failed internal bind reported readiness"); - let error = running + let error = Relay::load(config) .await - .expect("run panicked") - .expect_err("run accepted an occupied internal port"); + .err() + .expect("load accepted an occupied internal port"); assert!( error.to_string().contains("failed to bind internal listener"), "{error:#}" @@ -373,7 +322,7 @@ async fn shared_tokio_custom_route_and_quic() { async fn worker_tokio_custom_route_and_quic() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); - let mut config = http_and_quic(&cert, &key, format!("127.0.0.1:{}", free_udp_port())); + let mut config = http_and_quic(&cert, &key, "127.0.0.1:0".into()); config.runtime.workers = Some(2); config.runtime.pin = false; embed_and_stop(config).await; @@ -395,7 +344,7 @@ async fn uring_custom_route_and_quic() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); - let mut config = http_and_quic(&cert, &key, format!("127.0.0.1:{}", free_udp_port())); + let mut config = http_and_quic(&cert, &key, "127.0.0.1:0".into()); config.runtime.workers = Some(2); config.runtime.pin = false; config.runtime.io_uring = true; @@ -408,10 +357,10 @@ async fn embedded_listener_health_reaches_metrics() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); let mut config = http_and_quic(&cert, &key, "127.0.0.1:0".into()); - config.internal.listen = Some(format!("127.0.0.1:{}", free_tcp_port()).parse().unwrap()); + config.internal.listen = Some("127.0.0.1:0".parse().unwrap()); config.drain_timeout = Duration::ZERO; - let internal = config.internal.listen.unwrap(); let relay = Relay::load(config).await.expect("load relay"); + let internal = relay.internal().addr().expect("internal listener bound"); let ready = relay.ready(); let trigger = relay.shutdown_trigger().clone(); let health = moq_tokio::accept::Health::new("embedded"); diff --git a/rs/moq-relay/tests/hidden_cluster.rs b/rs/moq-relay/tests/hidden_cluster.rs index 2132607289..e49437910f 100644 --- a/rs/moq-relay/tests/hidden_cluster.rs +++ b/rs/moq-relay/tests/hidden_cluster.rs @@ -1,7 +1,6 @@ //! A cluster peer that predates the hidden opt-in still discovers the relay's //! `.`-named broadcasts, so a mixed-version mesh keeps `.internal/origins`. -use std::net::TcpListener; use std::time::Duration; use moq_relay::cluster::{self, Peer}; @@ -28,21 +27,16 @@ where .expect("test thread panicked"); } -/// A stream-only moq server on a free loopback TCP port, speaking only `version`. -fn bind_free_tcp_server(version: moq_net::Version) -> (u16, moq_tokio::Server) { - for _ in 0..20 { - let probe = TcpListener::bind("127.0.0.1:0").expect("bind probe"); - let port = probe.local_addr().expect("local addr").port(); - drop(probe); - - let mut config = moq_tokio::listen::Config::default(); - config.tcp.bind = Some(format!("127.0.0.1:{port}").parse().expect("parse addr")); - config.version = vec![version]; - if let Ok(server) = config.init(Default::default()) { - return (port, server); - } - } - panic!("could not bind a free TCP port after 20 attempts"); +/// A stream-only moq server listening on an ephemeral loopback TCP port, +/// speaking only `version`, and that port. +async fn listen_tcp(version: moq_net::Version) -> (u16, moq_tokio::Listener) { + let mut config = moq_tokio::listen::Config::default(); + config.tcp.bind = Some("127.0.0.1:0".parse().expect("parse addr")); + config.version = vec![version]; + let server = config.init(Default::default()).expect("server init"); + let listener = server.listen().await.expect("listen"); + let port = listener.tcp_local_addr().expect("TCP listener is configured").port(); + (port, listener) } #[test] @@ -58,11 +52,10 @@ async fn old_peer_keeps_hidden_paths(version: moq_net::Version) { let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); tokio::time::timeout(TEST_TIMEOUT, async { - let (port, server) = bind_free_tcp_server(version); + let (port, mut server) = listen_tcp(version).await; let peer_origin = moq_tokio::origin::spawn(); let mut discovered = peer_origin.consume().with_hidden(true).announced(); let accept = tokio::spawn(async move { - let mut server = server.listen().await.expect("listen"); let request = server.accept().await.expect("cluster dial"); let session = request.with_subscriber(peer_origin).ok().await.expect("accept"); session.closed().await; diff --git a/rs/moq-relay/tests/runtime_uring.rs b/rs/moq-relay/tests/runtime_uring.rs index 6fbebee5c8..0d8be349f9 100644 --- a/rs/moq-relay/tests/runtime_uring.rs +++ b/rs/moq-relay/tests/runtime_uring.rs @@ -7,7 +7,7 @@ //! floor (GitHub-hosted CI), where it skips loudly. #![cfg(all(target_os = "linux", feature = "_uring"))] -use std::net::{SocketAddr, UdpSocket}; +use std::net::SocketAddr; use std::time::Duration; use moq_relay::{Config, Relay}; @@ -28,15 +28,6 @@ fn supported() -> bool { } } -/// A UDP port nothing is bound to. Every worker binds the same port, so this -/// cannot be `:0`. -fn free_udp_port() -> u16 { - let probe = UdpSocket::bind("127.0.0.1:0").expect("bind probe"); - let port = probe.local_addr().expect("local addr").port(); - drop(probe); - port -} - /// A CA on disk plus a certificate it signed, for the mTLS test. Returns the /// root, the client's certificate, and its key. fn signed_client(dir: &std::path::Path) -> (std::path::PathBuf, std::path::PathBuf, std::path::PathBuf) { @@ -78,9 +69,9 @@ fn certificate(dir: &std::path::Path) -> (std::path::PathBuf, std::path::PathBuf /// A relay config serving QUIC from io_uring workers. Pinning is off because a /// CI container may restrict which cores it may run on. -fn uring_config(cert: &std::path::Path, key: &std::path::Path, port: u16) -> Config { +fn uring_config(cert: &std::path::Path, key: &std::path::Path) -> Config { let mut config = Config::default(); - config.listen.bind = Some(format!("127.0.0.1:{port}").parse().unwrap()); + config.listen.bind = Some("127.0.0.1:0".parse().unwrap()); config.listen.tls.cert = vec![cert.to_path_buf()]; config.listen.tls.key = vec![key.to_path_buf()]; config.runtime.workers = Some(WORKERS); @@ -118,11 +109,8 @@ async fn uring_workers_serve_webtransport_and_raw_quic() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); - let port = free_udp_port(); - - let relay = Relay::load(uring_config(&cert, &key, port)).await.expect("load relay"); - let expected: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap(); - assert_eq!(relay.addr(), Some(expected), "workers bound a different address"); + let relay = Relay::load(uring_config(&cert, &key)).await.expect("load relay"); + let port = relay.addr().expect("workers bound an address").port(); // The stock loop serves everything: the uring workers own QUIC, the shared // runtime owns auth and supervision. @@ -209,10 +197,9 @@ async fn uring_workers_report_link_facts() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); - let port = free_udp_port(); - - let relay = Relay::load(uring_config(&cert, &key, port)).await.expect("load relay"); + let relay = Relay::load(uring_config(&cert, &key)).await.expect("load relay"); let local = relay.addr().expect("bound address"); + let port = local.port(); let sessions = relay.sessions().clone(); let running = tokio::spawn(relay.run()); @@ -278,9 +265,7 @@ async fn uring_workers_publish_their_certificate_fingerprint() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); - let relay = Relay::load(uring_config(&cert, &key, free_udp_port())) - .await - .expect("load relay"); + let relay = Relay::load(uring_config(&cert, &key)).await.expect("load relay"); // The relay's own web listener needs TLS; its router does not, and the // handler reads the same certificate handle either way. @@ -319,14 +304,13 @@ async fn an_mtls_client_authenticates_without_a_token() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); let (root, client_cert, client_key) = signed_client(dir.path()); - let port = free_udp_port(); - - let mut config = uring_config(&cert, &key, port); + let mut config = uring_config(&cert, &key); config.auth.public = Vec::new(); config.auth.url = Some(spawn_auth_server(mtls_only()).await); config.listen.tls.root = vec![root]; let relay = Relay::load(config).await.expect("load relay"); + let port = relay.addr().expect("workers bound an address").port(); let running = tokio::spawn(relay.run()); let client = || { @@ -414,11 +398,10 @@ async fn uring_workers_write_qlog_traces() { let (cert, key) = certificate(dir.path()); let traces = dir.path().join("qlog"); std::fs::create_dir(&traces).expect("create qlog dir"); - let port = free_udp_port(); - - let mut config = uring_config(&cert, &key, port); + let mut config = uring_config(&cert, &key); config.quic.qlog = Some(traces.clone()); let relay = Relay::load(config).await.expect("load relay"); + let port = relay.addr().expect("workers bound an address").port(); let running = tokio::spawn(relay.run()); // A real session, so a trace covers a handshake and application data diff --git a/rs/moq-relay/tests/runtime_workers.rs b/rs/moq-relay/tests/runtime_workers.rs index 7f6040ec1d..97cce1fd36 100644 --- a/rs/moq-relay/tests/runtime_workers.rs +++ b/rs/moq-relay/tests/runtime_workers.rs @@ -6,7 +6,6 @@ //! a build without one refuses `runtime.workers` at load. #![cfg(all(target_os = "linux", feature = "_quic"))] -use std::net::{SocketAddr, UdpSocket}; use std::time::Duration; use moq_relay::{Config, Relay}; @@ -15,17 +14,6 @@ use moq_tokio::moq_net; const TIMEOUT: Duration = Duration::from_secs(10); const WORKERS: u16 = 4; -/// A UDP port nothing is bound to. -/// -/// Every worker binds the same port, so this cannot be `:0`: each would pick an -/// ephemeral port of its own and they would not form a group. -fn free_udp_port() -> u16 { - let probe = UdpSocket::bind("127.0.0.1:0").expect("bind probe"); - let port = probe.local_addr().expect("local addr").port(); - drop(probe); - port -} - /// A self-signed certificate on disk. Workers refuse `listen.tls.generate`, /// since each would generate one of its own and serve a different identity. fn certificate(dir: &std::path::Path) -> (std::path::PathBuf, std::path::PathBuf) { @@ -43,9 +31,9 @@ fn certificate(dir: &std::path::Path) -> (std::path::PathBuf, std::path::PathBuf /// /// Pinning is off because a CI container may restrict which cores it may run on, /// and none of these tests are about placement. -fn worker_config(cert: &std::path::Path, key: &std::path::Path, port: u16, workers: u16) -> Config { +fn worker_config(cert: &std::path::Path, key: &std::path::Path, workers: u16) -> Config { let mut config = Config::default(); - config.listen.bind = Some(format!("127.0.0.1:{port}").parse().unwrap()); + config.listen.bind = Some("127.0.0.1:0".parse().unwrap()); config.listen.tls.cert = vec![cert.to_path_buf()]; config.listen.tls.key = vec![key.to_path_buf()]; config.runtime.workers = Some(workers); @@ -82,13 +70,10 @@ async fn workers_serve_quic_and_share_one_origin() { let dir = tempfile::tempdir().expect("tempdir"); let (cert, key) = certificate(dir.path()); - let port = free_udp_port(); - - let config = worker_config(&cert, &key, port, WORKERS); + let config = worker_config(&cert, &key, WORKERS); let relay = Relay::load(config).await.expect("load relay"); - let expected: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap(); - assert_eq!(relay.addr(), Some(expected), "workers bound a different address"); + let port = relay.addr().expect("workers bound an address").port(); // The owner keeps the worker group; aborting this task is what joins them. let running = tokio::spawn(relay.run()); diff --git a/rs/moq-relay/tests/session_revalidate.rs b/rs/moq-relay/tests/session_revalidate.rs index 5d777ae0ba..77225c7802 100644 --- a/rs/moq-relay/tests/session_revalidate.rs +++ b/rs/moq-relay/tests/session_revalidate.rs @@ -5,7 +5,6 @@ //! an authority: the server's next word is what kicks, reties, or keeps the //! session. -use std::net::TcpListener; use std::sync::{Arc, Mutex}; use std::time::{Duration, SystemTime}; @@ -113,22 +112,6 @@ fn build_auth(url: url::Url) -> moq_relay::auth::Auth { .expect("auth init") } -async fn wait_for_listener(port: u16) { - let deadline = std::time::Instant::now() + Duration::from_secs(5); - while tokio::net::TcpStream::connect(("127.0.0.1", port)).await.is_err() { - assert!( - std::time::Instant::now() < deadline, - "listener never became ready on port {port}" - ); - tokio::time::sleep(Duration::from_millis(25)).await; - } -} - -fn free_port() -> u16 { - let probe = TcpListener::bind("127.0.0.1:0").expect("bind probe"); - probe.local_addr().expect("local addr").port() -} - struct Fixture { url: url::Url, internal: url::Url, @@ -148,19 +131,18 @@ impl Fixture { async fn spawn(auth: moq_relay::auth::Auth, websocket: bool) -> Self { let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); let sessions = moq_relay::session::Registry::new(); - let internal_port = free_port(); - let internal_addr: std::net::SocketAddr = format!("127.0.0.1:{internal_port}").parse().unwrap(); let mut internal_config = internal::Config::default(); - internal_config.listen = Some(internal_addr); + internal_config.listen = Some("127.0.0.1:0".parse().unwrap()); let internal = internal::Internal::new(internal_config, moq_net::stats::Registry::disabled()) - .with_sessions(sessions.clone()); + .with_sessions(sessions.clone()) + .bind() + .expect("bind internal listener"); + let internal_addr = internal.addr().expect("internal listener is configured"); tokio::spawn(async move { let _ = internal.run().await; }); - wait_for_listener(internal_port).await; let (url, handle) = if websocket { - let port = free_port(); let cluster = cluster::Cluster::new(cluster::Options::default()).expect("cluster init"); let mut server_config = moq_tokio::listen::Config::default(); server_config.bind = Some("[::]:0".parse().unwrap()); @@ -171,19 +153,22 @@ impl Fixture { .certificates(); let mut web_config = web::Config::default(); web_config.ws = true; - web_config.http.listen = Some(format!("127.0.0.1:{port}").parse().expect("parse listen")); - let web = web::Web::new(auth, cluster, certificates, web_config).with_sessions(sessions.clone()); + web_config.http.listen = Some("127.0.0.1:0".parse().expect("parse listen")); + let web = web::Web::new(auth, cluster, certificates, web_config) + .with_sessions(sessions.clone()) + .bind() + .expect("bind web listener"); + let port = web.addrs().http.expect("HTTP listener is configured").port(); let handle = tokio::spawn(async move { let _ = web.run().await; }); - wait_for_listener(port).await; (format!("ws://127.0.0.1:{port}").parse().unwrap(), handle) } else { - let port = free_port(); let mut config = moq_tokio::listen::Config::default(); - config.tcp.bind = Some(format!("127.0.0.1:{port}").parse().expect("parse addr")); + config.tcp.bind = Some("127.0.0.1:0".parse().expect("parse addr")); let server = config.init(Default::default()).expect("server init"); let mut server = server.listen().await.expect("listen"); + let port = server.tcp_local_addr().expect("TCP listener is configured").port(); let cluster = cluster::Cluster::new(cluster::Options::default()).expect("cluster init"); let sessions = sessions.clone(); let handle = tokio::spawn(async move { @@ -198,13 +183,12 @@ impl Fixture { }); } }); - wait_for_listener(port).await; (format!("tcp://127.0.0.1:{port}").parse().unwrap(), handle) }; Self { url, - internal: format!("http://127.0.0.1:{internal_port}").parse().unwrap(), + internal: format!("http://{internal_addr}").parse().unwrap(), sessions, _relay: handle, } diff --git a/rs/moq-relay/tests/shutdown_signal.rs b/rs/moq-relay/tests/shutdown_signal.rs index 23cf66e264..3f45135fd0 100644 --- a/rs/moq-relay/tests/shutdown_signal.rs +++ b/rs/moq-relay/tests/shutdown_signal.rs @@ -14,7 +14,7 @@ #![cfg(unix)] -use std::{net::TcpListener, time::Duration}; +use std::time::Duration; use moq_relay::{Config, Relay, auth}; @@ -51,10 +51,9 @@ async fn sigint_drains_sessions_before_exiting_inner() { let _interrupt = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt()).expect("register SIGINT"); - let (port, config) = relay_config(); - let relay = Relay::load(config).await.expect("load relay"); + let relay = Relay::load(relay_config()).await.expect("load relay"); + let port = relay.tcp_addr().expect("TCP listener bound").port(); let run = tokio::spawn(relay.run()); - wait_listening(port).await; let mut client_config = moq_tokio::connect::Config::default(); client_config.tls.insecure = Some(true); @@ -106,37 +105,16 @@ async fn sigint_drains_sessions_before_exiting_inner() { ); } -/// A stream-only relay on a free loopback TCP port, fully public, with a short -/// drain window. Returns the port and the config to hand [`Relay::load`]. -fn relay_config() -> (u16, Config) { - // The listener is bound by `Relay::run`, not here, so this leaves the usual - // probe/bind gap; on loopback it is not worth retrying around. - let probe = TcpListener::bind("127.0.0.1:0").expect("bind probe"); - let port = probe.local_addr().expect("local addr").port(); - drop(probe); - +/// A stream-only relay on an ephemeral loopback TCP port, fully public, with a +/// short drain window, to hand [`Relay::load`]. +fn relay_config() -> Config { // Fully public auth: any no-JWT stream client gets the whole root. let mut auth = auth::Config::default(); auth.public = vec![moq_auth::Pattern::all()]; let mut config = Config::default(); - config.listen.tcp.bind = Some(format!("127.0.0.1:{port}").parse().expect("parse addr")); + config.listen.tcp.bind = Some("127.0.0.1:0".parse().expect("parse addr")); config.auth = auth; config.drain_timeout = DRAIN_TIMEOUT; - - (port, config) -} - -async fn wait_listening(port: u16) { - let deadline = std::time::Instant::now() + Duration::from_secs(5); - loop { - if tokio::net::TcpStream::connect(("127.0.0.1", port)).await.is_ok() { - break; - } - assert!( - std::time::Instant::now() < deadline, - "relay never became ready on port {port}" - ); - tokio::time::sleep(Duration::from_millis(25)).await; - } + config } diff --git a/rs/moq-tokio/src/server.rs b/rs/moq-tokio/src/server.rs index dacd68c671..c06beaf474 100644 --- a/rs/moq-tokio/src/server.rs +++ b/rs/moq-tokio/src/server.rs @@ -484,32 +484,63 @@ impl Server { health } + /// Bind whatever is still unbound and hand back the [`Listener`], without + /// accepting. + /// + /// Same terminal bind as [`listen`](Self::listen), including ephemeral ports + /// via [`Listener::tcp_local_addr`]. Stream accept loops stay stopped until + /// [`Listener::accept`], so nothing is read off the socket before the caller + /// is ready to take sessions. [`listen`](Self::listen) is this plus starting + /// those loops immediately. + /// + /// A bind failure is the error, not a silent `None` from a later accept. It + /// leaves nothing bound: the partially built `Listener` drops here, closing + /// whatever it opened. Build a fresh `Server` from the (cloneable) config to try + /// again. + // `mut` is only needed to bind the stream listeners, which a QUIC-only build has none of. + #[allow(unused_mut)] + pub async fn bind(mut self) -> crate::Result { + #[cfg(any(feature = "tcp", all(feature = "uds", unix)))] + self.streams.bind().await?; + Ok(Listener { server: self }) + } + /// Start serving: bind whatever is still unbound and hand back the /// [`Listener`] to accept sessions from. /// /// Terminal, and that is the point: it consumes the `Server`, so the builders /// above cannot run afterwards and every session is served the configuration /// this call captured. The QUIC socket is bound by [`crate::listen::Config::init`], but - /// the stream (`tcp`/`unix`) listeners need a runtime, so they bind here. + /// the stream (`tcp`/`unix`) listeners need a runtime, so they bind here, and + /// their accept loops start here too. Use [`bind`](Self::bind) to learn the + /// port without accepting yet. /// /// A bind failure is the error, not a silent `None` from a later accept. It /// leaves nothing bound: the partially built `Listener` drops here, closing /// whatever it opened. Build a fresh `Server` from the (cloneable) config to try /// again. - // `mut` is only needed to bind the stream listeners, which a QUIC-only build has none of. - #[allow(unused_mut)] - pub async fn listen(mut self) -> crate::Result { + pub async fn listen(self) -> crate::Result { + #[cfg(not(any(feature = "tcp", all(feature = "uds", unix))))] + { + self.bind().await + } #[cfg(any(feature = "tcp", all(feature = "uds", unix)))] { - // The stream listeners offer a wider version set than the server's own - // (see `stream_versions`), against the same configuration. - let server = self.moq.clone().with_versions(self.streams.versions.clone()); - self.streams.start(&server).await?; + let mut listener = self.bind().await?; + // Accept immediately, so a dial can finish its handshake before the + // caller polls `accept`. `bind` leaves that until the first accept. + let server = listener + .server + .moq + .clone() + .with_versions(listener.server.streams.versions.clone()); + listener.server.streams.serve(&server); + Ok(listener) } - Ok(Listener { server: self }) } - /// The body of [`Listener::accept`]; the listeners are already running. + /// The body of [`Listener::accept`]. Stream accept loops start here if + /// [`Server::bind`] left them stopped; [`Server::listen`] already started them. #[cfg(any( feature = "noq", feature = "iroh", @@ -518,6 +549,14 @@ impl Server { all(feature = "uds", unix) ))] async fn accept_next(&mut self) -> Option { + // A `bind` (rather than `listen`) leaves the stream loops stopped so an + // embedder can read the port before it is willing to take sessions. + // Starting them here, on the first accept, is that moment. + #[cfg(any(feature = "tcp", all(feature = "uds", unix)))] + { + let server = self.moq.clone().with_versions(self.streams.versions.clone()); + self.streams.serve(&server); + } loop { // The QUIC endpoint address, reported as a QUIC session's local side. #[cfg(feature = "noq")] @@ -825,10 +864,11 @@ impl StreamBind { /// The stream (`tcp`/`unix`) listeners owned by a [`Server`]. /// -/// Bound by [`Server::listen`] (they need a runtime), after which each runs an -/// accept loop in its own task and feeds completed [`Request`]s back over a channel. -/// The tasks own their listeners and are stopped when the [`Listener`] closes or -/// drops, so bound sockets don't linger. +/// Bound by [`Server::bind`] (they need a runtime). [`Server::listen`] then starts +/// each accept loop; a bind-only listener starts them on the first +/// [`Listener::accept`] instead. Each loop feeds completed [`Request`]s back over +/// a channel. The tasks own their listeners and are stopped when the [`Listener`] +/// closes or drops, so bound sockets don't linger. #[cfg(any(feature = "tcp", all(feature = "uds", unix)))] struct StreamListeners { binds: Vec, @@ -839,7 +879,10 @@ struct StreamListeners { versions: moq_net::Versions, #[cfg(all(feature = "uds", unix))] unix_allow: Option, - /// The address the TCP listener bound, once [`Self::start`] has run. + /// Bound sockets whose accept loops have not started. Empty once [`Self::serve`] + /// runs, and when nothing was configured. + pending: Vec, + /// The address the TCP listener bound, once [`Self::bind`] has run. #[cfg(feature = "tcp")] tcp_local_addr: Option, rx: Option>, @@ -863,6 +906,7 @@ impl StreamListeners { versions, #[cfg(all(feature = "uds", unix))] unix_allow, + pending: Vec::new(), #[cfg(feature = "tcp")] tcp_local_addr: None, rx: None, @@ -870,21 +914,17 @@ impl StreamListeners { } } - /// Bind every configured listener and spawn its accept loop. - /// - /// Called once, from [`Server::listen`]. Everything binds before anything is - /// spawned, so a failure part-way drops the listeners already opened and frees - /// their sockets there and then, rather than leaving accept loops to be aborted - /// at some later point. + /// Bind every configured listener without accepting. /// - /// `server` is the configuration each accepted session handshakes against, - /// already fixed by the time this runs. - async fn start(&mut self, server: &moq_net::Server) -> crate::Result<()> { + /// Everything binds before [`Self::serve`] starts a loop, so a failure + /// part-way drops the listeners already opened and frees their sockets there + /// and then, rather than leaving accept loops to be aborted later. + async fn bind(&mut self) -> crate::Result<()> { if self.binds.is_empty() { return Ok(()); } - let mut bound = Vec::with_capacity(self.binds.len()); + let mut pending = Vec::with_capacity(self.binds.len()); for (bind, health) in self.binds.drain(..).zip(self.health.iter().cloned()) { let alpns = self.versions.alpns(); match bind { @@ -900,7 +940,7 @@ impl StreamListeners { let local = listener.local_addr()?; tracing::info!(addr = %local, "listening (tcp)"); self.tcp_local_addr = Some(local); - bound.push(BoundListener::Tcp(listener)); + pending.push(BoundListener::Tcp(listener)); } #[cfg(all(feature = "uds", unix))] StreamBind::Unix(path) => { @@ -912,13 +952,26 @@ impl StreamListeners { // directory or uid/gid/pid allowlist is the access gate. listener.set_mode(0o666)?; tracing::info!(path = %path.display(), allow = ?self.unix_allow, "listening (unix)"); - bound.push(BoundListener::Unix(listener)); + pending.push(BoundListener::Unix(listener)); } } } + self.pending = pending; + Ok(()) + } + + /// Spawn an accept loop for every listener [`Self::bind`] opened. + /// + /// No-op when nothing is waiting: either nothing was configured, or the loops + /// are already running. `server` is the configuration each accepted session + /// handshakes against, fixed by the time this runs. + fn serve(&mut self, server: &moq_net::Server) { + if self.pending.is_empty() { + return; + } let (tx, rx) = tokio::sync::mpsc::channel(16); - for listener in bound { + for listener in self.pending.drain(..) { let task = match listener { #[cfg(feature = "tcp")] BoundListener::Tcp(listener) => spawn_tcp_loop(listener, server.clone(), tx.clone()), @@ -927,9 +980,7 @@ impl StreamListeners { }; self.tasks.push(task); } - self.rx = Some(rx); - Ok(()) } /// Yield the next stream [`Request`], or `None` if no listener is running: @@ -944,6 +995,9 @@ impl StreamListeners { /// Stop every accept loop and wait until its listener has released the socket. async fn shutdown(&mut self) { self.binds.clear(); + // Drop sockets that were bound but never accepted, so the address frees + // even when `run` never started the loops. + self.pending.clear(); self.rx = None; for task in self.tasks.drain(..) { task.abort();