Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions doc/bin/relay/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
21 changes: 5 additions & 16 deletions rs/moq-cli/src/auth.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
70 changes: 37 additions & 33 deletions rs/moq-relay/src/internal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,8 @@ pub struct Internal {
health: moq_tokio::accept::Health,
listeners: Vec<moq_tokio::accept::Health>,
uring: Vec<UringWorker>,
listener: Option<net::TcpListener>,
addr: Option<net::SocketAddr>,
}

#[derive(Clone)]
Expand Down Expand Up @@ -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<Self> {
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<net::SocketAddr> {
self.addr
}

/// Report other listeners' accept health at `/metrics`.
///
/// Takes an iterator so the accessors feed it directly, however many listeners
Expand Down Expand Up @@ -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<Option<net::TcpListener>> {
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<net::TcpListener>) -> anyhow::Result<()> {
let Internal { listener, health, .. } = self.bind()?;
let Some(listener) = listener else {
std::future::pending::<()>().await;
return Ok(());
Expand All @@ -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(())
Expand Down Expand Up @@ -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");
Expand Down
28 changes: 17 additions & 11 deletions rs/moq-relay/src/relay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,7 @@ impl Ready {
pub struct Relay {
ready: tokio::sync::watch::Sender<bool>,
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
Expand Down Expand Up @@ -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<Self> {
config.resolve()?;
let resolved_config = config.clone();
Expand Down Expand Up @@ -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
Expand All @@ -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.
Expand Down Expand Up @@ -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<std::net::SocketAddr> {
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 {
Expand Down Expand Up @@ -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)]
Expand Down Expand Up @@ -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"),
Expand Down
43 changes: 12 additions & 31 deletions rs/moq-relay/tests/auth_lifetime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};

Expand Down Expand Up @@ -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<()>) {
Expand All @@ -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;
Expand All @@ -173,15 +155,13 @@ async fn spawn_relay_with(
}
});

wait_for_listener(port).await;
(port, handle)
}

/// Stand up the relay's axum web stack with WebSocket enabled and return the
/// 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
Expand All @@ -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)
}

Expand Down Expand Up @@ -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");
Expand Down
Loading
Loading