diff --git a/Cargo.lock b/Cargo.lock
index fdd2033eb..93cce3f49 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -4945,6 +4945,7 @@ dependencies = [
"bytes",
"bytesize",
"futures",
+ "hex",
"http-body",
"http-cache-reqwest",
"jsonwebtoken",
@@ -4963,6 +4964,7 @@ dependencies = [
"serde",
"serde_json",
"serde_with",
+ "sha2",
"sysinfo",
"tempfile",
"thiserror 2.0.20",
diff --git a/doc/bin/relay/auth.md b/doc/bin/relay/auth.md
index 7720cc2ce..dda86fb2c 100644
--- a/doc/bin/relay/auth.md
+++ b/doc/bin/relay/auth.md
@@ -247,11 +247,11 @@ A grant that is only checked at connect time can only stop NEW connections. Revo
So `--auth-api` keeps asking. The relay re-issues each live session's admission request on the endpoint's own `Cache-Control: max-age` cadence, and closes the session once the reply no longer grants what the session holds. There is no flag: an endpoint that can refuse a connection can stop one, or the two would disagree.
-- **Still granted**: the reply still authorizes at least the scope the session already has - a `key` that verifies its credential, or the `public` prefixes it was admitted under. The next re-check waits out the new `max-age`.
+- **Still granted**: the reply still authorizes at least the scope the session already has - a `key` that verifies its credential, the `public` prefixes it was admitted under, or (in `proxy` mode) a `grant`. The next re-check waits out the new `max-age`.
- **Refused** (404, or a reply that no longer grants the session's scope): the session closes immediately.
- **Unavailable** (network error, 5xx, unparseable body, or a 401/403 rejecting the *relay's own* credential): evidence of nothing about this session. The session keeps serving and the re-check retries with jittered backoff until the outage window passes without a success, then closes. A brief auth outage does not mass-disconnect; a sustained one still fails closed.
-The re-check REPLAYS the admission request rather than asking a narrower question, which is what makes one mechanism correct for every credential: a key replaced under an existing `kid` no longer verifies the retained JWT, and a withdrawn `public` block revokes anonymous sessions. `exp` still applies as the outer bound wherever the credential has one. mTLS peers are never revalidated, so a customer-facing decision cannot tear down the relay mesh.
+The re-check REPLAYS the admission request rather than asking a narrower question, which is what makes one mechanism correct in both modes and for every credential: a key replaced under an existing `kid` no longer verifies the retained JWT, a withdrawn `public` block revokes anonymous sessions, and a `proxy` session simply stops being granted. `exp` still applies as the outer bound wherever the credential has one. mTLS peers are never revalidated, so a customer-facing decision cannot tear down the relay mesh.
**`max-age` is the opt-in.** Revalidation is switched on by the endpoint, not by relay config: a reply that names a `max-age` is telling the relay how long its answer is good for, and that is the cadence. A reply with no usable `max-age` - none at all, `no-store`, `no-cache`, or `max-age=0` - has not asked to be re-consulted, so the session is never re-checked and its credential's own `exp` remains the only bound, exactly as before revalidation existed. Nothing is invented on the endpoint's behalf, and an existing deployment that sends no `Cache-Control` is unaffected until it opts in.
@@ -271,6 +271,38 @@ Either stale directive grants the outage window, and `stale-if-error` wins when
Note the asymmetry when choosing a long `max-age`: the cadence is set by the reply the relay is already holding, so shortening `max-age` later cannot pull in a re-check that is already scheduled. Whatever TTL you hand a healthy connection is how long an unannounced revocation takes to reach it.
+### Letting the endpoint decide (`--auth-api-mode proxy`)
+
+By default (`--auth-api-mode token`) the relay is the verifier: the endpoint hands back a `key` and the relay checks the credential against it.
+
+With `--auth-api-mode proxy` the endpoint is the decider. The relay forwards the connection verbatim - host (the URL authority, or `Host` on HTTP/1.1), path, transport, and the credential as `Authorization: Bearer` - and enforces whatever comes back:
+
+```
+GET ?root=demo&host=live.example.com&transport=quic
+Authorization: Bearer
+```
+```json
+{
+ "alias": "x7k2qp",
+ "tier": "region/sjc",
+ "grant": { "subscribe": ["room"], "publish": ["room/alice"], "exp": 1893456000 }
+}
+```
+
+The relay verifies nothing and holds no keys. Every policy decision - signature checking, scoping, expiry, per-viewer rules, even subdomain routing - belongs to the endpoint, so changing one is a deploy of the endpoint rather than a roll of the fleet. The credential need not be a JWT: it is an opaque string the endpoint alone interprets.
+
+`root` on the grant defaults to the connection path. `exp` (unix seconds) is the outer bound, and one already in the past is refused rather than admitted; an endpoint that omits it is asking for a session that ends only when revalidation says so. Each re-check's `exp` replaces the one before it, so an endpoint can cut a session short or extend a renewed one; in `token` mode the JWT's own `exp` is a ceiling a reply may lower but never raise. A reply with no grant, or one that authorizes nothing, is a refusal - there is no second shape to fall back to, so a `key` or `public` in a proxy-mode reply means nothing.
+
+**Refusing a viewer**: return `404`, an empty grant, or - in `proxy` mode only - `401`/`403`. The relay reads a `401`/`403` as a definitive rejection exactly where it forwarded a credential to be rejected; a `token`-mode request carries none, and neither does an anonymous `proxy` connection, so there the same status can only mean the relay's own identity or a gateway in front of the endpoint and is treated as an outage. Reading it otherwise would disconnect an entire audience over a gateway blip.
+
+`--auth-api-mode proxy` cannot be combined with `--auth-domain`: both decide how a hostname becomes a broadcast root, and proxy mode gives that job to the endpoint. To do that job the `alias` may RESHAPE the path, not just rename its leading segment - a connection to `/` can be aliased to `x7k2qp`, and `/room` to `x7k2qp/room` - so the endpoint can anchor a forwarded host at a root the client never dialed. Grant prefixes stay relative to the connection path, so they follow wherever the alias anchors them. In `token` mode the alias remains a rename of the leading segment and must keep the path's depth: there it resolves a vanity name to a pid, and a reply that reshaped the path would silently relocate the broadcast. `--auth-api-mode` without `--auth-api` is likewise a startup error - a mode with no endpoint to consult decides nothing.
+
+These are separate modes rather than two shapes of one reply, deliberately. Letting one endpoint answer either way per connection puts both paths inside a single request - which cache key applies, whether the credential may be sent, what "still vouched for" means. Choosing once, per relay, keeps each path independently simple.
+
+**The mode is a cost decision.** In `token` mode the request depends only on (`kid`, `root`, `transport`), so an audience sharing a signing key resolves to ONE cached request per relay however many distinct tokens they hold: auth cost tracks broadcasts and keys, not viewers. In `proxy` mode the credential is part of the request, so responses cache per credential and cost tracks concurrent viewers. A tokenless connection still caches per path in either mode. Pick `proxy` for control and simplicity, `token` when audience size would otherwise multiply your auth traffic.
+
+The relay keys its cache on the credential (a SHA-256 of it, so the secret stays out of logs and metrics), so a missing `Vary: Authorization` on the endpoint cannot leak one viewer's grant to another. Send `Vary: Authorization` anyway if anything else caches in front of it. Because that key already separates credentials, the relay's cache is a *private* one in HTTP's sense and stores a credentialed reply on a plain `max-age`; a shared cache would refuse it (RFC 9111 §3.5) unless the endpoint also sent `public`.
+
### Authenticating the relay to the auth API
The outbound HTTP the relay makes for auth (`--auth-api` requests and JWK fetches) reuses the cluster dial TLS configuration. The same `--connect-tls-cert` / `--connect-tls-key` the relay presents when dialing cluster peers also identifies it to the auth API, and `--connect-tls-root` trusts a private CA on the endpoint (env `MOQ_CONNECT_TLS_*`, or `[connect.tls]` in TOML). So an auth API can require mTLS and recognize the relay by the same certificate it uses for clustering.
diff --git a/rs/moq-relay/Cargo.toml b/rs/moq-relay/Cargo.toml
index a62f5aa70..9db339886 100644
--- a/rs/moq-relay/Cargo.toml
+++ b/rs/moq-relay/Cargo.toml
@@ -47,6 +47,7 @@ axum-server = { version = "0.8", features = ["tls-rustls"] }
bytes = "1"
bytesize = "2.4.2"
futures = "0.3"
+hex = "0.4"
http-body = "1"
http-cache-reqwest = { version = "1.0.0-alpha.6", features = ["manager-moka", "reqwest-middleware", "url-standard"], default-features = false }
jsonwebtoken = "11"
@@ -62,6 +63,7 @@ rustls = { version = "0.23", features = ["aws-lc-rs"], default-features = false
serde = { version = "1", features = ["derive"] }
serde_json = "1"
serde_with = { version = "3", features = ["json", "base64"] }
+sha2 = "0.11"
sysinfo = { version = "0.39", default-features = false, features = ["system"] }
thiserror = "2"
tokio = { workspace = true, features = ["full"] }
diff --git a/rs/moq-relay/src/auth.rs b/rs/moq-relay/src/auth.rs
index 0b8d82dab..fe8858e08 100644
--- a/rs/moq-relay/src/auth.rs
+++ b/rs/moq-relay/src/auth.rs
@@ -25,6 +25,10 @@ use url::Url;
pub struct AuthParams {
/// The URL path identifying the broadcast root.
pub path: String,
+ /// The URL host, forwarded to the auth API in [`AuthApiMode::Proxy`] so the
+ /// endpoint can do its own subdomain routing. `None` outside a URL-dialed
+ /// connection (the gateways pass a path directly).
+ pub host: Option,
/// A JWT token, if provided via the `jwt` query parameter.
pub jwt: Option,
/// The connection's transport, forwarded to the auth API as `transport=` so it
@@ -59,6 +63,7 @@ impl AuthParams {
Some(slug) => format!("/{slug}{}", url.path()),
None => url.path().to_string(),
};
+ let host = url.host_str().map(str::to_ascii_lowercase);
let mut jwt = None;
@@ -73,6 +78,7 @@ impl AuthParams {
Self {
path,
+ host,
jwt,
..Default::default()
}
@@ -158,6 +164,9 @@ pub enum AuthError {
#[error("key not found")]
KeyNotFound,
+ #[error("the auth API refused the credential")]
+ Refused,
+
#[error("the auth API has no grant for this connection")]
NotFound,
@@ -197,6 +206,7 @@ impl AuthError {
| Self::IncorrectRoot
| Self::KeyNotFound
| Self::MissingKeyId
+ | Self::Refused
| Self::NotFound
| Self::InvalidKeyId(_)
)
@@ -460,6 +470,23 @@ pub struct AuthConfig {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub auth_api: Option,
+ /// How `--auth-api` decides a connection: `token` (default) or `proxy`.
+ ///
+ /// `token` keeps the relay as the verifier: the endpoint returns a `key` for
+ /// the JWT's `kid`, or `public` prefixes for a tokenless connection.
+ ///
+ /// `proxy` makes the endpoint the decider: the relay forwards the connection
+ /// verbatim - host, path, mTLS flag, transport, and the credential as
+ /// `Authorization: Bearer` - and enforces the `grant` it gets back, verifying
+ /// nothing and holding no keys. The credential is part of the request, so
+ /// responses cache per credential and auth cost tracks concurrent viewers
+ /// rather than broadcasts.
+ ///
+ /// `Option` so a TOML value survives the CLI re-parse.
+ #[usage(long = "auth-api-mode", env = "MOQ_AUTH_API_MODE")]
+ #[serde(default, skip_serializing_if = "Option::is_none")]
+ pub api_mode: Option,
+
/// Billing tier label for mTLS peers when the auth API doesn't return one
/// (or no `--auth-api` is configured). Defaults to the unprefixed tier.
#[usage(long = "auth-mtls-tier", env = "MOQ_AUTH_MTLS_TIER")]
@@ -589,6 +616,59 @@ struct PublicResponse {
publish: Vec,
}
+/// How `--auth-api` decides a connection.
+///
+/// These are separate MODES rather than two shapes of one response, deliberately.
+/// Letting a single endpoint answer either way per connection means both paths
+/// live inside one request: which cache key to use, whether the credential may be
+/// sent, what "still vouched for" means. Choosing once, per relay, keeps each
+/// path independently simple.
+///
+/// `#[non_exhaustive]` so a third strategy is not a breaking release.
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
+#[serde(rename_all = "lowercase")]
+#[non_exhaustive]
+pub enum AuthApiMode {
+ /// The endpoint returns a verifying `key` for the JWT's `kid`, or the `public`
+ /// prefixes for a tokenless connection, and the relay checks the credential
+ /// itself.
+ ///
+ /// The request depends only on (`kid`, root, transport), so a whole audience
+ /// sharing a signing key resolves to ONE cached request per relay. Auth cost
+ /// tracks broadcasts and keys, not viewers.
+ #[default]
+ Token,
+
+ /// The relay forwards the connection verbatim - host, path, transport, and
+ /// the credential - and enforces the `grant` it gets back. It verifies
+ /// nothing itself and holds no keys. mTLS peers are unaffected either way:
+ /// they resolve through [`Auth::verify_mtls`], which never consults a mode.
+ ///
+ /// The credential is part of the request, so responses cache per credential
+ /// and auth cost tracks concurrent viewers rather than broadcasts. In exchange
+ /// the endpoint owns every policy decision - key rotation, scoping, expiry,
+ /// per-viewer rules - and changing one is a deploy of the endpoint rather than
+ /// a roll of the fleet.
+ Proxy,
+}
+
+/// Parses the same spellings the CLI and TOML accept, case-insensitively. An
+/// unrecognized mode is an ERROR, not a silent fall back to the default: the
+/// value decides who authorizes every connection.
+impl std::str::FromStr for AuthApiMode {
+ type Err = String;
+
+ fn from_str(s: &str) -> std::result::Result {
+ match s.to_ascii_lowercase().as_str() {
+ "token" => Ok(Self::Token),
+ "proxy" => Ok(Self::Proxy),
+ other => Err(format!(
+ "unknown --auth-api-mode `{other}`, expected `token` or `proxy`"
+ )),
+ }
+ }
+}
+
/// The configured `--auth-api`, and the revalidation state that belongs to it.
///
/// One struct rather than two `Option`s that were always in lockstep:
@@ -599,6 +679,10 @@ struct PublicResponse {
#[derive(Clone)]
struct AuthApi {
base: url::Url,
+ /// How this endpoint decides a connection. Part of the endpoint rather than
+ /// of [`Auth`], so a live session's re-check asks the question its admission
+ /// asked - see [`Revalidate::api`].
+ mode: AuthApiMode,
client: ClientWithMiddleware,
revalidator: Arc,
}
@@ -627,7 +711,12 @@ impl AuthApi {
/// so there is no safe fallback. Also returns the response's `Cache-Control`
/// timings, which drive revalidation.
async fn fetch(&self, request: &AuthApiRequest) -> Result<(AuthApiResponse, CacheHints), AuthError> {
- let response = self.client.get(request.url(&self.base)).send().await?;
+ let mut get = self.client.get(request.url(&self.base));
+ let forwarded = request.credential.is_some();
+ if let Some(credential) = &request.credential {
+ get = get.header(http::header::AUTHORIZATION, format!("Bearer {credential}"));
+ }
+ let response = get.send().await?;
// `Warning: 111` means the cache served a STALE entry because it could not
// reach the origin (RFC 2616 14.46). Treating that as a success would let a
@@ -644,6 +733,19 @@ impl AuthApi {
if response.status() == http::StatusCode::NOT_FOUND {
return Err(AuthError::NotFound);
}
+ // A 401/403 is only a statement about a VIEWER where the relay actually
+ // forwarded that viewer's credential - i.e. proxy mode with one present.
+ // A token-mode request carries no credential at all, and an anonymous proxy
+ // connection none either, so there the status can only be about the relay's
+ // own identity or a gateway in front of the endpoint; reading it as a
+ // per-viewer refusal would mass-disconnect an audience on a gateway blip.
+ if forwarded
+ && matches!(
+ response.status(),
+ http::StatusCode::UNAUTHORIZED | http::StatusCode::FORBIDDEN
+ ) {
+ return Err(AuthError::Refused);
+ }
let response = response.error_for_status()?;
let hints = CacheHints::from_headers(response.headers());
let body = response.text().await?;
@@ -656,8 +758,14 @@ impl AuthApi {
struct AuthApiRequest {
/// The connection path.
path: String,
- /// The JWT `kid` to resolve a verifying key for.
+ /// The connection's URL host. Sent in [`AuthApiMode::Proxy`] only, where the
+ /// endpoint does its own subdomain routing.
+ host: Option,
+ /// The JWT `kid` to resolve a verifying key for ([`AuthApiMode::Token`]).
kid: Option,
+ /// The credential, forwarded as `Authorization: Bearer`
+ /// ([`AuthApiMode::Proxy`]).
+ credential: Option,
/// Set only after the relay has verified the peer's client certificate.
mtls: bool,
transport: Option<&'static str>,
@@ -666,12 +774,17 @@ struct AuthApiRequest {
impl AuthApiRequest {
/// The request URL. Everything the endpoint keys on is a query param on the
/// base URL - never a path segment - so client-controlled values are
- /// percent-encoded by `query_pairs_mut` and can't retarget the path/query.
+ /// percent-encoded by `query_pairs_mut` and can't retarget the path/query. The
+ /// credential is never a query param: it is a bearer secret and would land in
+ /// access logs.
fn url(&self, base: &url::Url) -> url::Url {
let mut url = base.clone();
{
let mut q = url.query_pairs_mut();
q.append_pair("root", self.path.trim_matches('/'));
+ if let Some(host) = &self.host {
+ q.append_pair("host", host);
+ }
if let Some(kid) = &self.kid {
q.append_pair("kid", kid);
}
@@ -687,13 +800,14 @@ impl AuthApiRequest {
/// What two sessions must share before they can share one re-check.
///
- /// Derived from the request that actually gets sent, never rebuilt alongside
- /// it: a key assembled from parts can describe a different request than the
- /// one issued, which silently costs a re-check per viewer instead of per
- /// broadcast.
+ /// Its URL plus the credential that does not appear in the URL, derived from
+ /// the request that actually gets sent rather than rebuilt alongside it: a key
+ /// assembled from parts can describe a different request than the one issued,
+ /// which silently costs a re-check per viewer instead of per broadcast.
fn identity(&self, base: &url::Url) -> FlightKey {
FlightKey {
url: self.url(base).into(),
+ credential: self.credential.clone(),
}
}
}
@@ -712,6 +826,11 @@ struct AuthApiResponse {
/// moq-token's serde); absent -> not found.
#[serde(default)]
key: Option,
+ /// A grant the endpoint resolved from the credential itself, for a credential
+ /// the relay cannot verify locally. Read only in [`AuthApiMode::Proxy`]; token
+ /// mode ignores it entirely. See [`GrantResponse`].
+ #[serde(default)]
+ grant: Option,
/// Billing tier label for this connection (e.g. `region/sjc`).
/// The relay sends `mtls=true` when the peer presented a verified client
/// cert and lets the API decide. Absent or empty selects the default
@@ -728,6 +847,66 @@ impl AuthApiResponse {
}
}
+/// A grant the auth API resolved itself, instead of handing back a key for the
+/// relay to verify a JWT against.
+///
+/// This is what lets a credential the relay cannot parse authorize a connection:
+/// the relay forwards it as `Authorization: Bearer` and the endpoint answers with
+/// the permissions directly. One endpoint and one response type covers both, so
+/// it can answer with a `key` for one connection and a `grant` for another; there
+/// is no second flag and an operator migrates per connection.
+#[derive(Debug, Default, Deserialize)]
+struct GrantResponse {
+ /// Root the permissions below are relative to; absent -> the connection path.
+ #[serde(default)]
+ root: Option,
+ #[serde(default)]
+ subscribe: Vec,
+ #[serde(default)]
+ publish: Vec,
+ /// Unix seconds after which the session closes. There is no JWT to read an
+ /// `exp` from, so this is the outer bound; an endpoint that omits it is asking
+ /// for a session that ends only when revalidation says so.
+ #[serde(default)]
+ exp: Option,
+}
+
+impl GrantResponse {
+ /// Claims equivalent to what a JWT carrying this grant would have decoded to.
+ ///
+ /// An `exp` already in the past is refused rather than admitted, matching what
+ /// `Key::verify` does with an expired JWT. Admitting it would hand back a
+ /// session that closes on its next tick, which looks like a flap rather than a
+ /// refusal.
+ fn to_claims(&self, path: &str) -> Result {
+ // `SystemTime + Duration` PANICS on overflow, so an endpoint answering with
+ // a huge `exp` would take down the connection task rather than be refused.
+ let expires = match self.exp {
+ Some(exp) => Some(
+ std::time::UNIX_EPOCH
+ .checked_add(Duration::from_secs(exp))
+ .ok_or(AuthError::Refused)?,
+ ),
+ None => None,
+ };
+ if expires.is_some_and(|expires| expires <= std::time::SystemTime::now()) {
+ return Err(AuthError::Refused);
+ }
+ let mut claims = moq_token::Claims::default()
+ .with_root(self.root.clone().unwrap_or_else(|| path.to_string()))
+ .with_subscribe(self.subscribe.clone())
+ .with_publish(self.publish.clone());
+ claims.expires = expires;
+ Ok(claims)
+ }
+
+ /// True when the endpoint returned a grant that authorizes nothing, which is
+ /// a refusal rather than an empty success.
+ fn is_empty(&self) -> bool {
+ self.subscribe.is_empty() && self.publish.is_empty()
+ }
+}
+
/// Resolved public access configuration.
#[derive(Clone, Default)]
struct PublicAccess {
@@ -803,13 +982,7 @@ pub struct AuthToken {
impl AuthToken {
/// Wait until the backing credential expires, or forever when it has no expiry.
pub(crate) async fn expired(&self) {
- match self.expires {
- Some(expires) => {
- let remaining = expires.duration_since(std::time::SystemTime::now()).unwrap_or_default();
- tokio::time::sleep(remaining).await
- }
- None => std::future::pending().await,
- }
+ elapsed(self.expires).await
}
/// Construct a token for a peer that was authenticated at the TLS layer
@@ -868,6 +1041,12 @@ pub(crate) struct Revalidate {
/// The schedule admission resolved. Its existence IS the opt-in: no `max-age`
/// on the admission reply means no `Revalidate` at all.
schedule: Schedule,
+ /// The outer bound admission granted, and in [`AuthApiMode::Token`] a CEILING
+ /// a later reply may lower but never raise: there the bound comes from a
+ /// signed JWT, and no endpoint reply gets to extend a signed credential's
+ /// life. Proxy mode has no signature to respect - the endpoint IS the
+ /// authority - so its latest word replaces this outright.
+ expires: Option,
}
/// The part of an [`AuthToken`] a re-check has to keep vouching for.
@@ -951,11 +1130,47 @@ enum Fetched {
Unavailable,
}
+/// Resolves once `bound` passes, or pends forever without one.
+///
+/// A bound already in the past resolves immediately rather than saturating, so a
+/// caller racing this never serves past it.
+async fn elapsed(bound: Option) {
+ match bound {
+ Some(bound) => {
+ let remaining = bound.duration_since(std::time::SystemTime::now()).unwrap_or_default();
+ tokio::time::sleep(remaining).await
+ }
+ None => std::future::pending().await,
+ }
+}
+
+/// What the auth API's `alias` may do to the connection path.
+///
+/// The permission prefixes are relative to the connection path either way, so
+/// this governs only where they get ANCHORED - which is what decides where the
+/// broadcast lands on the backbone.
+#[derive(Clone, Copy, PartialEq, Eq)]
+enum Alias {
+ /// The alias renames the leading segment and nothing else, so it must match
+ /// the connection path's depth. That IS the token-mode contract - a project
+ /// stays reachable by vanity name and pid alike - and a reply that changes the
+ /// shape is a mistake worth refusing rather than a relocation worth honoring.
+ Rename,
+ /// The endpoint owns the whole mapping and may add or drop segments, which is
+ /// what lets it resolve a forwarded host to a root the client never dialed.
+ /// Only [`AuthApiMode::Proxy`], where the endpoint decides everything anyway.
+ Rewrite,
+}
+
/// One session's conclusion, drawn from a [`Fetched`] against its own scope.
#[derive(Debug, Clone, Copy)]
enum Recheck {
- /// Still vouched for; check again after the new max-age.
- Valid { hints: CacheHints },
+ /// Still vouched for; check again after the new max-age. `expires` is the
+ /// bound THIS reply granted, which may differ from admission's.
+ Valid {
+ hints: CacheHints,
+ expires: Option,
+ },
/// The reply no longer grants what this session holds.
Revoked,
/// The API could not answer.
@@ -972,13 +1187,14 @@ struct FlightSlot {
/// One auth-API re-check request; sessions that would issue the identical
/// request share a flight. Built by [`AuthApiRequest::identity`].
///
-/// The credential is not part of it: the response depends only on (`kid`, root,
-/// transport), so an audience sharing one `kid` shares one re-check however many
-/// distinct tokens they hold, and auth cost tracks broadcasts rather than
-/// viewers.
+/// In [`AuthApiMode::Token`] the credential is absent, so an audience sharing one
+/// `kid` shares one re-check however many distinct tokens they hold. In
+/// [`AuthApiMode::Proxy`] the credential is the authorization, so nothing merges
+/// across viewers - that is the cost of the mode, not a defect in the key.
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct FlightKey {
url: String,
+ credential: Option,
}
/// Shared state for live-session revalidation.
@@ -1119,6 +1335,24 @@ impl Auth {
"--auth-api cannot be combined with --auth-key/--auth-key-dir/--auth-public/--auth-public-api"
);
+ // A mode with no endpoint to consult decides nothing: `verify` never reaches
+ // `api_mode` without an `auth_api`, so the relay would silently start with
+ // token or mTLS behavior after being told to proxy every decision.
+ anyhow::ensure!(
+ config.api_mode.is_none() || config.auth_api.is_some(),
+ "--auth-api-mode requires --auth-api"
+ );
+
+ // Both answer "who turns a hostname into a broadcast root", and proxy mode's
+ // premise is that the endpoint does. Applying both routes the subdomain
+ // twice: the relay rewrites `customer.example.com/foo` to `/customer/foo`
+ // and still forwards the host, so an endpoint doing its own routing
+ // prepends `customer` again.
+ anyhow::ensure!(
+ config.api_mode != Some(AuthApiMode::Proxy) || config.domains.is_empty(),
+ "--auth-api-mode proxy cannot be combined with --auth-domain: the endpoint owns subdomain routing"
+ );
+
// Outbound auth HTTP (JWK + auth/public-API fetches) reuses the cluster
// client's --client-tls-* identity. The deprecated --auth-tls-* flags
// still override it when set.
@@ -1239,6 +1473,7 @@ impl Auth {
let auth_api = match config.auth_api {
Some(url_str) => Some(AuthApi {
base: Url::parse(&url_str).context("invalid --auth-api URL")?,
+ mode: config.api_mode.unwrap_or_default(),
client: Self::build_client(&tls)?,
revalidator: Arc::default(),
}),
@@ -1303,7 +1538,9 @@ impl Auth {
let request = AuthApiRequest {
path: path.to_string(),
+ host: None,
kid: None,
+ credential: None,
mtls: true,
transport: transport.map(Transport::as_str),
};
@@ -1317,71 +1554,102 @@ impl Auth {
/// alias (root), the billing tier, and EITHER something to verify the
/// credential against (a `key`) or the answer itself (a `grant`).
async fn verify_via_api(&self, api: &AuthApi, params: &AuthParams) -> Result<(AuthToken, CacheHints), AuthError> {
- let request = self.api_request(params)?;
+ let request = Self::api_request(api, params)?;
let (resp, hints) = api.fetch(&request).await?;
- Ok((self.authorize(params, &resp)?, hints))
+ Ok((Self::authorize(api, params, &resp)?, hints))
}
/// Turn one auth-API reply into this connection's token.
///
/// Split out from the fetch because the two have different scopes: the reply
- /// depends only on (`kid`, root, transport) and is shared, while the
- /// authorization depends on the credential and is emphatically NOT. See
- /// [`Fetched`].
- fn authorize(&self, params: &AuthParams, resp: &AuthApiResponse) -> Result {
- let claims = match params.jwt.as_deref() {
- Some(token) => {
- let key = resp.key.as_ref().ok_or(AuthError::KeyNotFound)?;
- // claims.root is the token's own root (a vanity name OR a pid); it is
- // checked against the ORIGINAL connection path below, not the alias, so
- // a vanity token matches a vanity URL and a pid token matches a pid URL.
- key.verify(token).map_err(|_| AuthError::DecodeFailed)?
- }
- None => {
- let public = resp.public.as_ref();
- let subscribe = public.map(|p| p.subscribe.clone()).unwrap_or_default();
- let publish = public.map(|p| p.publish.clone()).unwrap_or_default();
- if subscribe.is_empty() && publish.is_empty() {
- return Err(AuthError::ExpectedToken);
+ /// depends only on the request and is shared, while the authorization depends
+ /// on the credential and is emphatically NOT. See [`Fetched`].
+ fn authorize(api: &AuthApi, params: &AuthParams, resp: &AuthApiResponse) -> Result {
+ let claims = match api.mode {
+ // The endpoint hands back something to check the credential AGAINST,
+ // and the relay does the checking.
+ AuthApiMode::Token => match params.jwt.as_deref() {
+ Some(token) => {
+ let key = resp.key.as_ref().ok_or(AuthError::KeyNotFound)?;
+ // claims.root is the token's own root (a vanity name OR a pid); it
+ // is checked against the ORIGINAL connection path below, not the
+ // alias, so a vanity token matches a vanity URL and a pid token
+ // matches a pid URL.
+ key.verify(token).map_err(|_| AuthError::DecodeFailed)?
}
- // Anonymous access: anchor the public claims at the connection path so
- // the overlap check below is a no-op; routing lands on the alias.
- moq_token::Claims::default()
- .with_root(params.path.clone())
- .with_subscribe(subscribe)
- .with_publish(publish)
+ None => {
+ let public = resp.public.as_ref();
+ let subscribe = public.map(|p| p.subscribe.clone()).unwrap_or_default();
+ let publish = public.map(|p| p.publish.clone()).unwrap_or_default();
+ if subscribe.is_empty() && publish.is_empty() {
+ return Err(AuthError::ExpectedToken);
+ }
+ // Anonymous access: anchor the public claims at the connection path
+ // so the overlap check below is a no-op; routing lands on the alias.
+ moq_token::Claims::default()
+ .with_root(params.path.clone())
+ .with_subscribe(subscribe)
+ .with_publish(publish)
+ }
+ },
+ // The endpoint already decided. A reply with no usable grant is a
+ // refusal; there is no second shape to fall back to.
+ AuthApiMode::Proxy => {
+ let grant = resp.grant.as_ref().ok_or(AuthError::Refused)?;
+ if grant.is_empty() {
+ return Err(AuthError::Refused);
+ }
+ grant.to_claims(¶ms.path)?
}
};
- Self::finalize_api(params, resp.alias.clone(), resp.tier(), claims)
+ Self::finalize_api(params, api.mode, resp.alias.clone(), resp.tier(), claims)
}
- /// The auth-API request for a connection.
+ /// The auth-API request for a connection, decided once by the mode.
///
/// Admission and every re-check build the request here, so the flight key can
- /// be taken from the request itself rather than reconstructed beside it. The
- /// credential is never sent: the response depends only on (kid, root,
- /// transport), which is what lets a whole audience sharing a signing key
- /// resolve to one cached request per relay.
- fn api_request(&self, params: &AuthParams) -> Result {
- Ok(AuthApiRequest {
- path: params.path.clone(),
- kid: match params.jwt.as_deref() {
- Some(token) => {
- jsonwebtoken::decode_header(token)
- .map_err(|_| AuthError::DecodeFailed)?
- .kid
- }
- None => None,
+ /// be taken from the request itself rather than reconstructed beside it.
+ fn api_request(api: &AuthApi, params: &AuthParams) -> Result {
+ let transport = params.transport.map(Transport::as_str);
+ Ok(match api.mode {
+ // The credential is NEVER sent: the response depends only on (kid, root,
+ // transport), which is what lets a whole audience sharing a signing key
+ // resolve to one cached request per relay.
+ AuthApiMode::Token => AuthApiRequest {
+ path: params.path.clone(),
+ host: None,
+ kid: match params.jwt.as_deref() {
+ Some(token) => {
+ jsonwebtoken::decode_header(token)
+ .map_err(|_| AuthError::DecodeFailed)?
+ .kid
+ }
+ None => None,
+ },
+ credential: None,
+ mtls: false,
+ transport,
+ },
+ // The connection goes over verbatim and the relay verifies nothing.
+ AuthApiMode::Proxy => AuthApiRequest {
+ path: params.path.clone(),
+ host: params.host.clone(),
+ kid: None,
+ credential: params.jwt.clone(),
+ mtls: false,
+ transport,
},
- mtls: false,
- transport: params.transport.map(Transport::as_str),
})
}
/// The flight key for a live session's re-check.
- fn flight_key(&self, grant: &Revalidate) -> Option {
- Some(self.api_request(&grant.params).ok()?.identity(&grant.api.base))
+ fn flight_key(grant: &Revalidate) -> Option {
+ Some(
+ Self::api_request(&grant.api, &grant.params)
+ .ok()?
+ .identity(&grant.api.base),
+ )
}
/// Anchor verified claims on the API's alias, shared by both modes.
@@ -1393,14 +1661,19 @@ impl Auth {
/// tier; the API may bucket specific ones under a named tier.
fn finalize_api(
params: &AuthParams,
+ mode: AuthApiMode,
alias: Option,
tier: Option,
claims: moq_token::Claims,
) -> Result {
- let alias = alias.unwrap_or_else(|| params.path.clone());
+ let route_root = alias.unwrap_or_else(|| params.path.clone());
// Check the token root against the ORIGINAL connection path (vanity or
// pid); anchor the resulting scope on the alias (canonical pid).
- let mut token = Self::finalize(¶ms.path, &alias, claims)?;
+ let alias = match mode {
+ AuthApiMode::Token => Alias::Rename,
+ AuthApiMode::Proxy => Alias::Rewrite,
+ };
+ let mut token = Self::finalize(¶ms.path, &route_root, alias, claims)?;
token.tier = tier.unwrap_or_default();
Ok(token)
}
@@ -1422,6 +1695,7 @@ impl Auth {
params: Arc::new(params.clone()),
scope: Scope::new(&token),
schedule,
+ expires: token.expires,
});
Ok(token)
}
@@ -1481,7 +1755,7 @@ impl Auth {
return Err(AuthError::ExpectedToken);
};
- Self::finalize(¶ms.path, ¶ms.path, claims)
+ Self::finalize(¶ms.path, ¶ms.path, Alias::Rename, claims)
}
/// Reduce verified `claims` into an [`AuthToken`].
@@ -1496,7 +1770,12 @@ impl Auth {
/// (same depth), so the rebased relative prefixes anchor unchanged. The standalone
/// path passes the same value for both (no alias). Shared by the standalone and
/// `--auth-api` paths.
- fn finalize(check_root: &str, route_root: &str, claims: moq_token::Claims) -> Result {
+ fn finalize(
+ check_root: &str,
+ route_root: &str,
+ alias: Alias,
+ claims: moq_token::Claims,
+ ) -> Result {
let root = Path::new(check_root);
let route_root = Path::new(route_root);
let depth = |path: &Path<'_>| {
@@ -1507,7 +1786,7 @@ impl Auth {
}
};
- if depth(&root) != depth(&route_root) {
+ if alias == Alias::Rename && depth(&root) != depth(&route_root) {
return Err(AuthError::IncorrectRoot);
}
@@ -1515,8 +1794,8 @@ impl Auth {
// another root, so both reduce to IncorrectRoot.
let permissions = claims.authorize(check_root).map_err(|_| AuthError::IncorrectRoot)?;
- // authorize() returns paths already normalized and relative to check_root,
- // which route_root matches in depth.
+ // authorize() returns paths already normalized and RELATIVE to check_root, so
+ // they anchor under route_root whatever its depth.
let rebase = |paths: Vec| -> PathPrefixes { paths.iter().map(|p| Path::new(p).to_owned()).collect() };
Ok(AuthToken {
@@ -1546,16 +1825,16 @@ impl Auth {
/// token itself, so a process holding several differently-configured `Auth`
/// instances still judges each token against the authority that issued it.
pub async fn expired(&self, token: &AuthToken) -> Expired {
- let revoked = async {
- match &token.revalidate {
- Some(grant) => self.revalidate(grant).await,
- None => std::future::pending().await,
+ match &token.revalidate {
+ // The loop starts from this same credential bound and each reply may move
+ // it, so it subsumes the timer below rather than racing it. Racing would
+ // pin the session to admission's bound, letting a re-check shorten one but
+ // never extend a renewed grant.
+ Some(grant) => self.revalidate(grant).await,
+ None => {
+ token.expired().await;
+ Expired::Credential
}
- };
-
- tokio::select! {
- _ = token.expired() => Expired::Credential,
- reason = revoked => reason,
}
}
@@ -1592,9 +1871,18 @@ impl Auth {
let mut next = Instant::now() + schedule.cadence;
let mut deadline = next + schedule.staleness;
let mut backoff = Revalidator::BACKOFF;
+ // The bound currently in force. Admission's to begin with; each re-check may
+ // move it (see `Revalidate::expires`).
+ let mut bound = grant.expires;
loop {
- tokio::time::sleep_until(next).await;
+ // Race the cadence against the bound: a re-check that shortened `exp`
+ // below the next cadence has to close the session at the new bound, not
+ // at whenever the endpoint next happens to be asked.
+ tokio::select! {
+ _ = tokio::time::sleep_until(next) => {}
+ _ = elapsed(bound) => return Expired::Credential,
+ }
// Bound the attempt by the deadline so a peer that accepts a request and
// then stalls cannot carry a revoked session past its window - but never
@@ -1602,12 +1890,30 @@ impl Auth {
// would otherwise cancel the very re-check that was about to RENEW the
// grant, closing every session without the endpoint ever being asked.
let budget = deadline.max(Instant::now() + crate::http_client::REQUEST_TIMEOUT);
- let outcome = match tokio::time::timeout_at(budget, self.recheck(grant)).await {
- Ok(outcome) => outcome,
- Err(_) => return Expired::Stale,
+ let outcome = tokio::select! {
+ outcome = tokio::time::timeout_at(budget, self.recheck(grant)) => match outcome {
+ Ok(outcome) => outcome,
+ Err(_) => return Expired::Stale,
+ },
+ // The bound passing mid-request ends the session there; a stalled
+ // endpoint must not carry it to the request timeout.
+ _ = elapsed(bound) => return Expired::Credential,
};
match outcome {
- Recheck::Valid { hints } => {
+ Recheck::Valid { hints, expires } => {
+ // The endpoint's latest word on the bound. Token mode clamps to
+ // admission's, which came off a signed JWT; proxy mode takes the reply
+ // outright, so an endpoint can extend a renewed session as well as cut
+ // one short. A reply that names no bound lifts it only where there was
+ // never a signature to respect.
+ bound = match grant.api.mode {
+ AuthApiMode::Token => match (grant.expires, expires) {
+ (Some(ceiling), Some(expires)) => Some(ceiling.min(expires)),
+ (ceiling, None) => ceiling,
+ (None, expires) => expires,
+ },
+ AuthApiMode::Proxy => expires,
+ };
// A reply that stops naming `max-age` keeps the schedule the session
// already opted into, rather than silently becoming unrevocable.
schedule = hints.schedule().unwrap_or(schedule);
@@ -1635,8 +1941,11 @@ impl Auth {
async fn recheck(&self, grant: &Revalidate) -> Recheck {
match self.fetch_shared(grant).await {
// The reply is shared; the verdict is this session's alone.
- Fetched::Ok { resp, hints } => match self.authorize(&grant.params, &resp) {
- Ok(token) if grant.scope.covered_by(&token) => Recheck::Valid { hints },
+ Fetched::Ok { resp, hints } => match Self::authorize(&grant.api, &grant.params, &resp) {
+ Ok(token) if grant.scope.covered_by(&token) => Recheck::Valid {
+ hints,
+ expires: token.expires,
+ },
Ok(_) => Recheck::Revoked,
Err(err) if err.is_refusal() => Recheck::Revoked,
Err(_) => Recheck::Unavailable,
@@ -1648,7 +1957,7 @@ impl Auth {
/// The shared auth-API fetch, joining an in-flight one for the same request.
async fn fetch_shared(&self, grant: &Revalidate) -> Fetched {
- let Some(key) = self.flight_key(grant) else {
+ let Some(key) = Self::flight_key(grant) else {
return Fetched::Unavailable;
};
let revalidator = &grant.api.revalidator;
@@ -1687,7 +1996,7 @@ impl Auth {
/// check ("does a key still exist for this kid?") cannot see a key REPLACED
/// under that kid, and cannot see an anonymous grant withdrawn.
async fn recheck_fetch(&self, grant: &Revalidate) -> Fetched {
- let request = match self.api_request(&grant.params) {
+ let request = match Self::api_request(&grant.api, &grant.params) {
Ok(request) => request,
// The credential parsed at admission, so this cannot be transient.
Err(_) => return Fetched::Refused,
@@ -2925,7 +3234,7 @@ api = "https://api.example.com/access"
// HTTP-based tests (URL key-dir + public API) using wiremock.
// ---------------------------------------------------------------------
- use wiremock::matchers::{method, path as path_matcher, query_param};
+ use wiremock::matchers::{header, method, path as path_matcher, query_param};
use wiremock::{Mock, MockServer, ResponseTemplate};
/// Serialize a key as JSON for serving from a mock URL endpoint.
@@ -3599,6 +3908,17 @@ api = "https://api.example.com/access"
.unwrap()
}
+ /// The same endpoint, with the relay forwarding instead of verifying.
+ async fn auth_with_api_proxy(server: &MockServer) -> Auth {
+ Auth::new(AuthConfig {
+ auth_api: Some(format!("{}/auth", server.uri())),
+ api_mode: Some(AuthApiMode::Proxy),
+ ..Default::default()
+ })
+ .await
+ .unwrap()
+ }
+
#[tokio::test]
async fn auth_api_jwt_scopes_to_alias() -> anyhow::Result<()> {
// JWT connection: the token root is the vanity path the client dialed
@@ -4034,6 +4354,7 @@ api = "https://api.example.com/access"
fn test_grant_schedule(auth: &Auth, jwt: Option, schedule: Schedule) -> Revalidate {
Revalidate {
api: test_api(auth),
+ expires: None,
params: Arc::new(AuthParams {
path: "demo".into(),
jwt,
@@ -4618,6 +4939,7 @@ api = "https://api.example.com/access"
let auth = auth_with_api(&server).await;
let scoped = |subscribe: &str| Revalidate {
api: test_api(&auth),
+ expires: None,
params: Arc::new(AuthParams {
path: "demo".into(),
..Default::default()
@@ -4723,6 +5045,29 @@ api = "https://api.example.com/access"
Ok(())
}
+ /// In proxy mode the credential IS the authorization, so two viewers holding
+ /// different ones must not share a flight even on the same path. This is the
+ /// cost of the mode, and the reason token mode does not send the credential.
+ #[tokio::test]
+ async fn revalidate_does_not_coalesce_across_credentials_in_proxy_mode() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .respond_with(ResponseTemplate::new(404).set_delay(Duration::from_millis(300)))
+ .expect(2)
+ .mount(&server)
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let a = test_grant(&auth, Some("credential-a".into()), Duration::from_millis(100));
+ let b = test_grant(&auth, Some("credential-b".into()), Duration::from_millis(100));
+
+ let (a, b) = tokio::join!(auth.revalidate(&a), auth.revalidate(&b));
+ assert_eq!(a, Expired::Revoked);
+ assert_eq!(b, Expired::Revoked);
+ Ok(())
+ }
+
#[tokio::test]
async fn revalidate_drops_an_abandoned_flight() -> anyhow::Result<()> {
let server = MockServer::start().await;
@@ -4816,6 +5161,643 @@ api = "https://api.example.com/access"
Ok(())
}
+ // --- Proxy mode: the endpoint decides, the relay enforces ---
+
+ /// The relay forwards the connection verbatim and enforces what comes back.
+ /// The credential need not be a JWT at all - in token mode this one is
+ /// rejected before any request is made.
+ #[tokio::test]
+ async fn proxy_forwards_the_connection_and_enforces_the_grant() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .and(query_param("root", "demo"))
+ .and(query_param("host", "live.example.com"))
+ .and(header("Authorization", "Bearer opaque-session-cookie"))
+ .respond_with(
+ ResponseTemplate::new(200)
+ .set_body_string(r#"{"alias":"x7k2qp","grant":{"subscribe":["room"],"publish":["room/alice"]}}"#),
+ )
+ .expect(1)
+ .mount(&server)
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let token = auth
+ .verify(&AuthParams {
+ path: "demo".into(),
+ host: Some("live.example.com".into()),
+ jwt: Some("opaque-session-cookie".into()),
+ ..Default::default()
+ })
+ .await?;
+
+ assert_eq!(token.root.as_str(), "x7k2qp");
+ assert_eq!(
+ token.subscribe.iter().map(|p| p.to_string()).collect::>(),
+ ["room"]
+ );
+ assert_eq!(
+ token.publish.iter().map(|p| p.to_string()).collect::>(),
+ ["room/alice"]
+ );
+ Ok(())
+ }
+
+ /// There is no JWT to read an `exp` from, so the grant's own `exp` is the
+ /// outer bound.
+ #[tokio::test]
+ async fn proxy_grant_exp_bounds_the_session() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ mount_auth(
+ &server,
+ "max-age=60",
+ r#"{"grant":{"subscribe":[""],"exp":1893456000}}"#.to_string(),
+ )
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let token = auth
+ .verify(&AuthParams {
+ path: "demo".into(),
+ jwt: Some("opaque".into()),
+ ..Default::default()
+ })
+ .await?;
+ assert_eq!(
+ token.expires,
+ Some(std::time::UNIX_EPOCH + Duration::from_secs(1893456000))
+ );
+ Ok(())
+ }
+
+ /// An `exp` that would overflow `SystemTime` is refused, not panicked on.
+ #[tokio::test]
+ async fn proxy_absurd_grant_exp_is_refused() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ mount_auth(
+ &server,
+ "max-age=60",
+ format!(r#"{{"grant":{{"subscribe":[""],"exp":{}}}}}"#, u64::MAX),
+ )
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let result = auth
+ .verify(&AuthParams {
+ path: "demo".into(),
+ jwt: Some("opaque".into()),
+ ..Default::default()
+ })
+ .await;
+ assert!(matches!(result, Err(AuthError::Refused)), "got {result:?}");
+ Ok(())
+ }
+
+ /// An `exp` already in the past is refused, not admitted into a session that
+ /// closes on its next tick.
+ #[tokio::test]
+ async fn proxy_expired_grant_is_refused() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ mount_auth(
+ &server,
+ "max-age=60",
+ r#"{"grant":{"subscribe":[""],"exp":1}}"#.to_string(),
+ )
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let result = auth
+ .verify(&AuthParams {
+ path: "demo".into(),
+ jwt: Some("opaque".into()),
+ ..Default::default()
+ })
+ .await;
+ assert!(matches!(result, Err(AuthError::Refused)), "got {result:?}");
+ Ok(())
+ }
+
+ /// A reply with no usable grant is a refusal. There is no second shape to fall
+ /// back to, so an empty grant, a missing one, and a `key` are all the same
+ /// answer: proxy mode holds no keys and verifies nothing.
+ #[tokio::test]
+ async fn proxy_requires_a_non_empty_grant() -> anyhow::Result<()> {
+ let key = create_test_key_with_kid("test-key");
+ for body in [
+ r#"{"grant":{}}"#.to_string(),
+ r#"{"alias":"x7k2qp"}"#.to_string(),
+ format!(r#"{{"key":{}}}"#, jwk_body(&key)),
+ r#"{"public":{"subscribe":[""]}}"#.to_string(),
+ ] {
+ let server = MockServer::start().await;
+ mount_auth(&server, "max-age=60", body.clone()).await;
+ let auth = auth_with_api_proxy(&server).await;
+ let result = auth
+ .verify(&AuthParams {
+ path: "demo".into(),
+ jwt: Some("opaque".into()),
+ ..Default::default()
+ })
+ .await;
+ assert!(matches!(result, Err(AuthError::Refused)), "{body} gave {result:?}");
+ }
+ Ok(())
+ }
+
+ /// Anonymous connections go through the same call - the endpoint knows there
+ /// was no credential and answers accordingly.
+ #[tokio::test]
+ async fn proxy_authorizes_an_anonymous_connection() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ mount_auth(&server, "max-age=60", r#"{"grant":{"subscribe":[""]}}"#.to_string()).await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let token = auth.verify(&AuthParams::new("demo")).await?;
+ assert!(token.revalidate.is_some(), "anonymous proxy sessions revalidate too");
+
+ let sent = server.received_requests().await.expect("recorded requests");
+ assert!(
+ sent[0].headers.get("Authorization").is_none(),
+ "there is no credential to forward"
+ );
+ Ok(())
+ }
+
+ // --- The modes are independent ---
+
+ /// Token mode NEVER sends the credential: the response depends only on (kid,
+ /// root, transport), which is what lets an audience share one cached request.
+ #[tokio::test]
+ async fn token_mode_never_sends_the_credential() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ let key = create_test_key_with_kid("test-key");
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .and(query_param("kid", "test-key"))
+ .respond_with(ResponseTemplate::new(200).set_body_string(format!(r#"{{"key":{}}}"#, jwk_body(&key))))
+ .mount(&server)
+ .await;
+
+ let auth = auth_with_api(&server).await;
+ let jwt = key.sign(&moq_token::Claims::default().with_root("demo").with_subscribe([""]))?;
+ auth.verify(&AuthParams {
+ path: "demo".into(),
+ host: Some("live.example.com".into()),
+ jwt: Some(jwt),
+ ..Default::default()
+ })
+ .await?;
+
+ let sent = server.received_requests().await.expect("recorded requests");
+ assert_eq!(sent.len(), 1);
+ assert!(
+ sent[0].headers.get("Authorization").is_none(),
+ "a kid lookup must not carry the credential"
+ );
+ assert!(
+ !sent[0].url.query_pairs().any(|(k, _)| k == "host"),
+ "the relay does its own subdomain routing in token mode"
+ );
+ Ok(())
+ }
+
+ /// A `grant` is inert in token mode. Honoring one would authorize a signature
+ /// the relay never checked - and that reply is cached per `kid` and shared
+ /// across the audience, so a forged token with a known kid would inherit it.
+ #[tokio::test]
+ async fn token_mode_ignores_a_grant() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .and(query_param("kid", "test-key"))
+ .respond_with(ResponseTemplate::new(200).set_body_string(r#"{"grant":{"subscribe":[""]}}"#))
+ .mount(&server)
+ .await;
+
+ // A token nobody signed, carrying a known kid: header {"alg":"HS256",
+ // "kid":"test-key"} and claims {"root":"demo","sub":[""]}.
+ let forged = concat!(
+ "eyJhbGciOiJIUzI1NiIsImtpZCI6InRlc3Qta2V5In0.",
+ "eyJyb290IjoiZGVtbyIsInN1YiI6WyIiXX0.",
+ "bm90LWEtc2lnbmF0dXJl"
+ );
+
+ let auth = auth_with_api(&server).await;
+ let result = auth
+ .verify(&AuthParams {
+ path: "demo".into(),
+ jwt: Some(forged.into()),
+ ..Default::default()
+ })
+ .await;
+ assert!(
+ matches!(result, Err(AuthError::KeyNotFound)),
+ "a kid lookup must require a key, got {result:?}"
+ );
+ Ok(())
+ }
+
+ /// A proxy session revalidates like any other, against the same grant reply.
+ #[tokio::test]
+ async fn proxy_session_survives_revalidation() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ mount_auth(&server, "max-age=1", r#"{"grant":{"subscribe":[""]}}"#.to_string()).await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let token = auth
+ .verify(&AuthParams {
+ path: "demo".into(),
+ jwt: Some("opaque".into()),
+ ..Default::default()
+ })
+ .await?;
+ let grant = token.revalidate.clone().expect("proxy sessions revalidate");
+
+ let pending = tokio::time::timeout(Duration::from_millis(2500), auth.revalidate(&grant)).await;
+ assert!(pending.is_err(), "a still-granted proxy session must keep serving");
+ Ok(())
+ }
+
+ /// Withdrawing the grant closes a proxy session, the same way a withdrawn key
+ /// or `public` block closes a token-mode one.
+ #[tokio::test]
+ async fn proxy_session_closes_when_the_grant_is_withdrawn() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ mount_auth(&server, "max-age=1", r#"{"alias":"demo"}"#.to_string()).await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let grant = test_grant(&auth, Some("opaque".into()), Duration::from_millis(200));
+ let reason = tokio::time::timeout(Duration::from_secs(5), auth.revalidate(&grant))
+ .await
+ .expect("a withdrawn grant must close the session");
+ assert_eq!(reason, Expired::Revoked);
+ Ok(())
+ }
+
+ /// The MODE rides on the grant alongside the endpoint, for the same reason
+ /// #3069 moved the endpoint there: a re-check must ask the question admission
+ /// asked. Judged by a token-mode `Auth` instead, this session's re-check would
+ /// send no credential, find no `key`, and close a session the endpoint is
+ /// still granting.
+ #[tokio::test]
+ async fn revalidate_keeps_the_granting_mode() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .and(header("Authorization", "Bearer opaque"))
+ .respond_with(
+ ResponseTemplate::new(200)
+ .insert_header("Cache-Control", "max-age=1")
+ .set_body_string(r#"{"grant":{"subscribe":[""]}}"#),
+ )
+ .mount(&server)
+ .await;
+
+ let issuer = auth_with_api_proxy(&server).await;
+ let token = issuer
+ .verify(&AuthParams {
+ path: "demo".into(),
+ jwt: Some("opaque".into()),
+ ..Default::default()
+ })
+ .await?;
+ let grant = token.revalidate.clone().expect("proxy sessions revalidate");
+
+ // The mock only answers a request carrying the credential, so a token-mode
+ // re-check would 404 and revoke.
+ let token_mode = auth_with_api(&server).await;
+ let pending = tokio::time::timeout(Duration::from_millis(2500), token_mode.revalidate(&grant)).await;
+ assert!(pending.is_err(), "a proxy grant must be re-checked in proxy mode");
+ Ok(())
+ }
+
+ /// A mode with no endpoint decides nothing, so saying so is a startup error
+ /// rather than a relay that quietly keeps token behavior.
+ #[tokio::test]
+ async fn proxy_mode_requires_an_auth_api() {
+ let err = Auth::new(AuthConfig {
+ api_mode: Some(AuthApiMode::Proxy),
+ key: Some("/dev/null".into()),
+ ..Default::default()
+ })
+ .await
+ .map(|_| ())
+ .expect_err("--auth-api-mode without --auth-api must fail");
+ assert!(err.to_string().contains("--auth-api-mode requires --auth-api"), "{err}");
+ }
+
+ /// `--auth-domain` and proxy mode both answer "who turns a hostname into a
+ /// root", and applying both routes the subdomain twice.
+ #[tokio::test]
+ async fn proxy_mode_rejects_auth_domain() {
+ let err = Auth::new(AuthConfig {
+ auth_api: Some("https://api.example.com/auth".into()),
+ api_mode: Some(AuthApiMode::Proxy),
+ domains: vec!["cdn.moq.dev".into()],
+ ..Default::default()
+ })
+ .await
+ .map(|_| ())
+ .expect_err("--auth-domain with proxy mode must fail");
+ assert!(
+ err.to_string().contains("cannot be combined with --auth-domain"),
+ "{err}"
+ );
+ }
+
+ /// Token mode is unaffected: the relay resolves the subdomain itself there.
+ #[tokio::test]
+ async fn token_mode_still_accepts_auth_domain() -> anyhow::Result<()> {
+ Auth::new(AuthConfig {
+ auth_api: Some("https://api.example.com/auth".into()),
+ domains: vec!["cdn.moq.dev".into()],
+ ..Default::default()
+ })
+ .await?;
+ Ok(())
+ }
+
+ /// An unrecognized mode must fail rather than silently resolve to the default,
+ /// since the value decides who authorizes every connection.
+ #[test]
+ fn unknown_api_mode_is_rejected() {
+ assert_eq!("proxy".parse::(), Ok(AuthApiMode::Proxy));
+ assert_eq!("TOKEN".parse::(), Ok(AuthApiMode::Token));
+ assert!("prxy".parse::().is_err());
+ }
+
+ /// A proxy endpoint rejecting the forwarded credential with 401/403 is saying
+ /// "no", not "I'm broken". Read as an outage it would keep a revoked viewer
+ /// serving for the whole staleness window.
+ #[tokio::test]
+ async fn proxy_credential_rejection_closes_the_session() -> anyhow::Result<()> {
+ for status in [401u16, 403] {
+ let server = MockServer::start().await;
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .respond_with(ResponseTemplate::new(status))
+ .mount(&server)
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let grant = test_grant(&auth, Some("opaque".into()), Duration::from_millis(100));
+ let reason = tokio::time::timeout(Duration::from_secs(5), auth.revalidate(&grant))
+ .await
+ .unwrap_or_else(|_| panic!("a {status} must close a proxy session immediately"));
+ assert_eq!(reason, Expired::Revoked);
+ }
+ Ok(())
+ }
+
+ /// The other half of the rule: a token-mode request carries NO credential, so
+ /// a 401/403 cannot be about a viewer - only the relay's own identity or a
+ /// gateway. Refusing there would mass-disconnect an audience on a blip.
+ #[tokio::test]
+ async fn token_mode_credential_rejection_is_an_outage() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .respond_with(ResponseTemplate::new(403))
+ .mount(&server)
+ .await;
+
+ let auth = auth_with_api(&server).await;
+ let grant = test_grant(&auth, None, Duration::from_millis(100));
+ // Revoked would be immediate; an outage rides out the staleness window.
+ let pending = tokio::time::timeout(Duration::from_millis(1500), auth.revalidate(&grant)).await;
+ assert!(pending.is_err(), "a token-mode 403 must not revoke immediately");
+ Ok(())
+ }
+
+ /// An anonymous proxy connection forwards no credential either, so it keeps
+ /// outage semantics for the same reason token mode does.
+ #[tokio::test]
+ async fn anonymous_proxy_rejection_is_an_outage() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .respond_with(ResponseTemplate::new(401))
+ .mount(&server)
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let grant = test_grant(&auth, None, Duration::from_millis(100));
+ let pending = tokio::time::timeout(Duration::from_millis(1500), auth.revalidate(&grant)).await;
+ assert!(pending.is_err(), "an anonymous proxy 401 must not revoke immediately");
+ Ok(())
+ }
+
+ /// A proxy endpoint shortening `exp` on a re-check has to bound the session
+ /// there. Dropped, the session would run to admission's bound - and with a
+ /// long `max-age`, well past what the endpoint now grants.
+ #[tokio::test]
+ async fn proxy_revalidation_applies_a_shortened_exp() -> anyhow::Result<()> {
+ // `exp` is whole unix seconds, so a sub-second bound truncates into the PAST
+ // and gets refused outright rather than exercising the bound.
+ let exp = std::time::SystemTime::now() + Duration::from_secs(2);
+ let exp = exp.duration_since(std::time::UNIX_EPOCH)?.as_secs();
+
+ let server = MockServer::start().await;
+ mount_auth(
+ &server,
+ "max-age=3600",
+ format!(r#"{{"grant":{{"subscribe":[""],"exp":{exp}}}}}"#),
+ )
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ // Admitted with no bound at all, on a cadence far under the hour max-age.
+ let grant = test_grant(&auth, Some("opaque".into()), Duration::from_millis(100));
+ let reason = tokio::time::timeout(Duration::from_secs(5), auth.revalidate(&grant))
+ .await
+ .expect("a revalidated exp must bound the session");
+ assert_eq!(reason, Expired::Credential);
+ Ok(())
+ }
+
+ /// Proxy mode delegates subdomain routing to the endpoint, so its alias has to
+ /// be able to PREPEND a root the client never dialed. The connection path is
+ /// `/` (depth 0) and the host-derived root is depth 1.
+ #[tokio::test]
+ async fn proxy_alias_may_add_a_host_derived_root() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .and(query_param("host", "customer.example.com"))
+ .respond_with(
+ ResponseTemplate::new(200).set_body_string(r#"{"alias":"x7k2qp","grant":{"subscribe":[""]}}"#),
+ )
+ .mount(&server)
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let token = auth
+ .verify(&AuthParams {
+ path: "".into(),
+ host: Some("customer.example.com".into()),
+ jwt: Some("opaque".into()),
+ ..Default::default()
+ })
+ .await?;
+ assert_eq!(token.root, "x7k2qp".as_path(), "the endpoint owns the whole mapping");
+ Ok(())
+ }
+
+ /// A deeper path keeps its tail under the host-derived root, so the prefixes
+ /// the grant named still anchor where the endpoint put them.
+ #[tokio::test]
+ async fn proxy_alias_may_nest_a_deeper_path() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ mount_auth(
+ &server,
+ "no-store",
+ r#"{"alias":"x7k2qp/room","grant":{"subscribe":["cam"]}}"#.to_string(),
+ )
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let token = auth
+ .verify(&AuthParams {
+ path: "room".into(),
+ host: Some("customer.example.com".into()),
+ jwt: Some("opaque".into()),
+ ..Default::default()
+ })
+ .await?;
+ assert_eq!(token.root, "x7k2qp/room".as_path());
+ assert!(token.subscribe.contains(&Path::new("cam").to_owned()));
+ Ok(())
+ }
+
+ /// Token mode keeps the depth rule: there the alias is a RENAME of the leading
+ /// segment, and a reply that changes the shape would silently relocate a
+ /// broadcast rather than resolve a vanity name.
+ #[tokio::test]
+ async fn token_alias_must_keep_the_path_depth() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ mount_auth(
+ &server,
+ "no-store",
+ r#"{"alias":"x7k2qp/extra","public":{"subscribe":[""]}}"#.to_string(),
+ )
+ .await;
+
+ let auth = auth_with_api(&server).await;
+ let result = auth
+ .verify(&AuthParams {
+ path: "demo".into(),
+ ..Default::default()
+ })
+ .await;
+ assert!(
+ matches!(result, Err(AuthError::IncorrectRoot)),
+ "a reshaping alias must be refused in token mode, got {result:?}"
+ );
+ Ok(())
+ }
+
+ /// The bound has to hold while a re-check is IN FLIGHT, not just between them.
+ /// A stalled endpoint would otherwise carry an expired session all the way to
+ /// the request timeout, well past what the credential granted.
+ #[tokio::test]
+ async fn a_stalled_recheck_still_honors_the_bound() -> anyhow::Result<()> {
+ let server = MockServer::start().await;
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ // Far longer than the bound below, and longer than the cadence.
+ .respond_with(ResponseTemplate::new(200).set_delay(Duration::from_secs(30)))
+ .mount(&server)
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let mut grant = test_grant(&auth, Some("opaque".into()), Duration::from_millis(100));
+ grant.expires = Some(std::time::SystemTime::now() + Duration::from_millis(600));
+
+ let start = std::time::Instant::now();
+ let reason = tokio::time::timeout(Duration::from_secs(5), auth.revalidate(&grant))
+ .await
+ .expect("the bound must resolve while the re-check hangs");
+ assert_eq!(reason, Expired::Credential);
+ assert!(
+ start.elapsed() < Duration::from_secs(3),
+ "closed on the request timeout, not the bound: {:?}",
+ start.elapsed()
+ );
+ Ok(())
+ }
+
+ /// The other direction, and the one an admission-time timer running beside the
+ /// loop would silently break: a proxy endpoint EXTENDING a renewed grant. The
+ /// session must outlive the bound it was admitted with.
+ #[tokio::test]
+ async fn proxy_revalidation_applies_an_extended_exp() -> anyhow::Result<()> {
+ let near = std::time::SystemTime::now() + Duration::from_secs(2);
+ let near = near.duration_since(std::time::UNIX_EPOCH)?.as_secs();
+ let far = std::time::SystemTime::now() + Duration::from_secs(3600);
+ let far = far.duration_since(std::time::UNIX_EPOCH)?.as_secs();
+
+ let server = MockServer::start().await;
+ // Admission sees the near bound; every re-check after it sees the extension.
+ Mock::given(method("GET"))
+ .and(path_matcher("/auth"))
+ .respond_with(
+ ResponseTemplate::new(200)
+ .insert_header("Cache-Control", "max-age=1")
+ .set_body_string(format!(r#"{{"grant":{{"subscribe":[""],"exp":{near}}}}}"#)),
+ )
+ .up_to_n_times(1)
+ .mount(&server)
+ .await;
+ mount_auth(
+ &server,
+ "max-age=1",
+ format!(r#"{{"grant":{{"subscribe":[""],"exp":{far}}}}}"#),
+ )
+ .await;
+
+ let auth = auth_with_api_proxy(&server).await;
+ let token = auth
+ .verify(&AuthParams {
+ path: "demo".into(),
+ jwt: Some("opaque".into()),
+ ..Default::default()
+ })
+ .await?;
+
+ // Well past the 2s the session was admitted with: the extension must hold.
+ let pending = tokio::time::timeout(Duration::from_millis(3500), auth.expired(&token)).await;
+ assert!(pending.is_err(), "an extended grant must outlive admission's exp");
+ Ok(())
+ }
+
+ /// Token mode's bound comes off a SIGNED credential, so a reply may lower it
+ /// but never raise it. Here the reply names an hour and admission named a
+ /// moment; the signature wins.
+ #[tokio::test]
+ async fn token_mode_reply_cannot_extend_a_signed_exp() -> anyhow::Result<()> {
+ let far = std::time::SystemTime::now() + Duration::from_secs(3600);
+ let far = far.duration_since(std::time::UNIX_EPOCH)?.as_secs();
+
+ let server = MockServer::start().await;
+ mount_auth(
+ &server,
+ "max-age=3600",
+ format!(r#"{{"public":{{"subscribe":[""]}},"grant":{{"subscribe":[""],"exp":{far}}}}}"#),
+ )
+ .await;
+
+ let auth = auth_with_api(&server).await;
+ let mut grant = test_grant(&auth, None, Duration::from_millis(100));
+ grant.expires = Some(std::time::SystemTime::now() + Duration::from_millis(600));
+
+ let reason = tokio::time::timeout(Duration::from_secs(5), auth.revalidate(&grant))
+ .await
+ .expect("the signed bound must still close the session");
+ assert_eq!(reason, Expired::Credential);
+ Ok(())
+ }
+
#[tokio::test(start_paused = true)]
async fn expired_resolves_at_credential_expiry() {
let auth = Auth::default();
diff --git a/rs/moq-relay/src/config.rs b/rs/moq-relay/src/config.rs
index bd12ada7c..e49fac299 100644
--- a/rs/moq-relay/src/config.rs
+++ b/rs/moq-relay/src/config.rs
@@ -715,6 +715,33 @@ auth_api = "https://api.moq.dev/cluster/auth"
);
}
+ /// Same clap+TOML clobber guard for `auth.api_mode`. It's an `Option` so an
+ /// absent `--auth-api-mode` must not wipe a TOML-configured value during the
+ /// `update_from` re-parse.
+ #[test]
+ fn cli_does_not_clobber_toml_auth_api_mode() {
+ let _env = EnvGuard::clear(&["MOQ_AUTH_API", "MOQ_AUTH_API_MODE"]);
+
+ let toml = r#"
+[auth]
+auth_api = "https://api.moq.dev/cluster/auth"
+api_mode = "proxy"
+"#;
+ let dir = std::env::temp_dir().join("moq-relay-config-test");
+ std::fs::create_dir_all(&dir).unwrap();
+ let path = dir.join("auth-api-mode-toml-wins.toml");
+ std::fs::write(&path, toml).unwrap();
+
+ let args = vec![std::ffi::OsString::from("moq-relay"), std::ffi::OsString::from(&path)];
+ let config = Config::parse_and_merge(args).expect("config load");
+
+ assert_eq!(
+ config.auth.api_mode,
+ Some(crate::AuthApiMode::Proxy),
+ "TOML's auth.api_mode must not be clobbered by the CLI re-parse"
+ );
+ }
+
/// The optional system-roots policy loaded from TOML survives when omitted on the CLI.
#[test]
fn cli_does_not_clobber_toml_system_roots() {
diff --git a/rs/moq-relay/src/http_client.rs b/rs/moq-relay/src/http_client.rs
index 1d78cc3b7..f00d4c804 100644
--- a/rs/moq-relay/src/http_client.rs
+++ b/rs/moq-relay/src/http_client.rs
@@ -1,4 +1,5 @@
use anyhow::Context;
+use axum::http;
use http_cache_reqwest::{Cache, CacheMode, HttpCache, HttpCacheOptions, MokaCache, MokaManager};
use reqwest_middleware::ClientWithMiddleware;
use std::time::Duration;
@@ -27,7 +28,47 @@ pub(crate) fn build(tls: &rustls::ClientConfig) -> anyhow::Result String {
+ let key = format!("{}:{}", parts.method, parts.uri);
+ match parts.headers.get(http::header::AUTHORIZATION) {
+ // Digested rather than interpolated: the key is a moka map key that can
+ // reach logs and metrics, and the raw value is a bearer secret. SHA-256
+ // rather than a `Hash` impl because the split is a security boundary and
+ // the credentials are attacker-chosen: two of them landing on one key
+ // would serve one viewer another's grant.
+ Some(auth) => {
+ use sha2::Digest;
+ let digest = sha2::Sha256::digest(auth.as_bytes());
+ format!("{key}:{}", hex::encode(digest))
+ }
+ None => key,
+ }
+}
diff --git a/rs/moq-relay/src/web.rs b/rs/moq-relay/src/web.rs
index 274deedac..6cb3c5aba 100644
--- a/rs/moq-relay/src/web.rs
+++ b/rs/moq-relay/src/web.rs
@@ -657,11 +657,35 @@ impl<'de> serde::Deserialize<'de> for FetchGroup {
}
}
+/// The host this request was addressed to, which `AuthApiMode::Proxy` forwards so
+/// the endpoint can do its own routing. These handlers build their params from a
+/// path rather than a URL, so it has to come off the request headers.
+fn request_host(uri: &http::Uri, headers: &http::HeaderMap) -> Option {
+ // HTTP/2 carries the host in `:authority`, which hyper surfaces on the URI and
+ // usually WITHOUT a `Host` header. The HTTPS listener advertises h2, so reading
+ // the header alone loses the tenant for ordinary clients.
+ uri.authority()
+ .map(|authority| authority.host().to_string())
+ .or_else(|| {
+ headers
+ .get(http::header::HOST)
+ .and_then(|host| host.to_str().ok())
+ // Parsed rather than split on the last colon, which would truncate a
+ // bracketed IPv6 literal (`[2001:db8::1]`) at its own separator.
+ .and_then(|host| host.parse::().ok())
+ .map(|authority| authority.host().to_string())
+ })
+ .map(|host| host.to_ascii_lowercase())
+ .filter(|host| !host.is_empty())
+}
+
/// Serve the announced broadcasts for a given prefix.
async fn serve_announced(
path: Option>,
Query(query): Query,
mtls: Option>,
+ uri: http::Uri,
+ headers: http::HeaderMap,
State(state): State>,
) -> axum::response::Result {
let prefix = match path {
@@ -671,6 +695,7 @@ async fn serve_announced(
let params = AuthParams {
path: prefix,
+ host: request_host(&uri, &headers),
jwt: query.jwt,
..Default::default()
};
@@ -709,6 +734,8 @@ async fn serve_fetch(
Path(path): Path,
Query(params): Query,
mtls: Option>,
+ uri: http::Uri,
+ headers: http::HeaderMap,
State(state): State>,
) -> axum::response::Result {
// The path containts a broadcast/track
@@ -722,6 +749,7 @@ async fn serve_fetch(
let auth = AuthParams {
path: path.join("/"),
+ host: request_host(&uri, &headers),
jwt: params.auth.jwt,
..Default::default()
};
@@ -904,6 +932,42 @@ mod tests {
);
}
+ /// HTTP/2 puts the host in `:authority` and usually sends no `Host` header, and
+ /// the HTTPS listener advertises h2 - so a header-only read loses the tenant.
+ #[test]
+ fn request_host_prefers_the_uri_authority() {
+ let uri: http::Uri = "https://customer.example.com/announced".parse().unwrap();
+ assert_eq!(
+ request_host(&uri, &http::HeaderMap::new()),
+ Some("customer.example.com".to_string())
+ );
+ }
+
+ /// HTTP/1.1 leaves no authority on the URI, so the header is the only source. A
+ /// port is stripped either way, so both paths agree on the host.
+ #[test]
+ fn request_host_falls_back_to_the_header() {
+ let uri: http::Uri = "/announced".parse().unwrap();
+ let mut headers = http::HeaderMap::new();
+ headers.insert(http::header::HOST, "Customer.Example.com:4443".parse().unwrap());
+ assert_eq!(request_host(&uri, &headers), Some("customer.example.com".to_string()));
+ assert_eq!(request_host(&uri, &http::HeaderMap::new()), None);
+ }
+
+ /// A bracketed IPv6 literal carries its own colons, so splitting on the last
+ /// one truncates the address instead of stripping a port.
+ #[test]
+ fn request_host_keeps_ipv6_literals_intact() {
+ let uri: http::Uri = "/announced".parse().unwrap();
+ let host = |value: &str| {
+ let mut headers = http::HeaderMap::new();
+ headers.insert(http::header::HOST, value.parse().unwrap());
+ request_host(&uri, &headers)
+ };
+ assert_eq!(host("[2001:db8::1]"), Some("[2001:db8::1]".to_string()));
+ assert_eq!(host("[2001:db8::1]:4443"), Some("[2001:db8::1]".to_string()));
+ }
+
#[test]
fn https_watch_paths_include_roots() {
let cert = PathBuf::from("cert.pem");