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: 4 additions & 1 deletion doc/bin/relay/auth.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,10 @@ ingest here). A 2xx with a grant admits. A 401 or 403 refuses.
Anything else at connect, a timeout, a 5xx, or an unparseable body, refuses and
logs an error; nothing is admitted because the server was down. A grant that
names nothing refuses, and one with `revalidate` but no `expires` is refused
as invalid. A few seconds of clock skew are tolerated on `expires`.
as invalid. A grant already less than five seconds past `expires` is accepted
for the remainder of that clock-skew window; future expiries are unchanged.
The client and relay snapshot each accepted grant's deadline on a monotonic
clock, so later polls, outages, and wall-clock adjustments do not restart it.

**Revalidate and outage.** On the cadence the relay POSTs `revalidate` with the
same request. A grant applies: a changed `root` or one that no longer covers
Expand Down
2 changes: 1 addition & 1 deletion doc/lib/rs/moq-auth.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ answers the relay, in a service that mints tokens for clients, or in your own
accept loop that decides in process.

- **Request and grant**: `Request` is the JSON a relay POSTs per session event (`connect`, `revalidate`, `end`) with everything it knows: id, node, transport, addresses, SNI and ALPN, the raw path and query, the declared role, and the verified certificate facts. `Grant` is the answer: `publish` and `subscribe` pattern unions, an optional `root` alias, `expires`, `revalidate`, `tier`, and `peer`, which marks the session as a cluster peer so the routes it announces report `Source::Peer`. `Grant::validate` refuses a grant that names nothing, asks to be revalidated without a bound or at no interval, or has already expired, with a few seconds of clock skew on `expires`.
- **Lease**: `lease::Producer` and `lease::Consumer` are the handle a session holds for its grant. The consumer reads the current grant, waits for a change, and learns why the lease ended; the producer updates and revokes. Either side's terminal call returns the reason the lease actually ended with, so whichever got there first is what both report. `Consumer::fixed` is a grant nobody drives. `Consumer::revalidate` nudges a re-check now and the producer observes it via `poll_revalidate` or `revalidate_requested`, which is how the relay's session push lands. Whoever runs the accept loop builds the producer, so an embedder decides in process with no trait and no HTTP. Enforcing `expires` is the holder's job; the `Client` driver also revokes at expiry so its `end` event goes out.
- **Lease**: `lease::Producer` and `lease::Consumer` are the handle a session holds for its grant. The consumer reads the current grant, waits for a change, and learns why the lease ended; the producer updates and revokes. Either side's terminal call returns the reason the lease actually ended with, so whichever got there first is what both report. `Consumer::fixed` is a grant nobody drives. `Consumer::revalidate` nudges a re-check now and the producer observes it via `poll_revalidate` or `revalidate_requested`, which is how the relay's session push lands. Whoever runs the accept loop builds the producer, so an embedder decides in process with no trait and no HTTP. Enforcing `expires` is the holder's job; the `Client` driver also revokes at expiry so its `end` event goes out. `Grant::deadline()` (feature `tokio`, also enabled by `client` and `serve`) snapshots expiry on Tokio's clock. Call it once per accepted grant and retain the deadline: future expiries are unchanged, and one already less than five seconds late gets the remainder of that skew window.
- **Client**: `Client::new(url, tls)` and `Client::connect(request)` drive a lease against an auth server over `https://`, `unix://`, or loopback `http://`: revalidate on cadence with jittered backoff through an outage until `expires`, revoke on a 401/403 or an invalid grant, and POST `end` with the reason, duration, and byte totals the session reported through `lease::Consumer::close` when it ended. Dropping the consumer reports zero bytes. `end.reason` is `dropped`, `expired`, `refused`, `invalid`, or the session's own classification.
- **Server**: `serve::Policy` and `serve::Server` (feature `serve`) are the reference auth server behind `moq auth serve`: a `jwt` in the query verified against a key file or a `{kid}.jwk` directory, an explicit grant for verified certificates, the anonymous permissions, a tier, the revalidation cadence, a default `expires`, and live session caps per token and per remote address. A token is authorized at the dialed path with `Claims::authorize`; residuals become the grant. `Server::router` is an axum `POST /` you can mount in your own service.
- **Keys**: generate HS256/384/512, RS256/384/512, PS256/384/512, ES256/384, or EdDSA keys as JWKs, with a `kid` for rotation and an optional immutable scope that caps every token the key signs.
Expand Down
1 change: 0 additions & 1 deletion quest/m1/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ transport, benchmark tooling); worktrees isolate commits, not semantics.
- [Interop flakes](/quest/m1/interop-flakes.md) - the interop harness passes with other runs sharing the machine
- [Signal.race cleanup](/quest/m1/signal-race.md) - `Signal.race` releases its signal listeners when its result loses a race
- [Origin narrowing](/quest/m1/origin-narrowing.md) - a live origin grant narrows in place and ends the subscriptions it no longer covers, the deafen boundary #2714 asked for
- [Auth expiry clock](/quest/m1/auth-expiry-clock.md) - moq-auth and the relay hold one fixed expiry deadline and honour the same skew allowance
- [Binding surface](/quest/m1/binding-surface.md) - moq-ffi, libmoq, and every wrapper expose the decode delay, route source, and connection timing
- [FFI shape](/quest/m1/ffi-shape/README.md) - the bindings mirror Rust's layers: net at the root, then media, json, audio, and video namespaces built from the handle below
- [Track demand](/quest/m1/track-demand.md) - Rust and JS watch a track's subscribers through `demand()` alone
Expand Down
24 changes: 0 additions & 24 deletions quest/m1/auth-expiry-clock.md

This file was deleted.

4 changes: 2 additions & 2 deletions rs/moq-auth/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,9 @@ rust-version.workspace = true
[features]
default = ["client"]
# The HTTP client that drives a lease against an auth server.
client = ["dep:reqwest", "dep:rustls", "dep:tokio", "dep:tracing", "dep:rand"]
client = ["dep:reqwest", "dep:rustls", "tokio", "dep:tracing", "dep:rand"]
# The reference auth server behind `moq auth serve`: the policy a relay held, on axum.
serve = ["dep:axum", "dep:tokio", "dep:tracing", "tokio/net"]
serve = ["dep:axum", "tokio", "dep:tracing", "tokio/net"]
tokio = ["dep:tokio"]

[dependencies]
Expand Down
87 changes: 59 additions & 28 deletions rs/moq-auth/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,7 @@ impl Client {
client: self.clone(),
request,
producer: Some(producer),
expires: grant.deadline(),
started: Instant::now(),
};
tokio::spawn(driver.run(grant));
Expand Down Expand Up @@ -115,6 +116,7 @@ struct Driver {
request: Request,
producer: Option<lease::Producer>,
started: Instant,
expires: Option<tokio::time::Instant>,
}

impl Driver {
Expand All @@ -140,7 +142,7 @@ impl Driver {
.take()
.expect("the driver owns the producer until it ends");
let mut failures = 0u32;
let mut next = grant.revalidate.map(|cadence| Instant::now() + cadence);
let mut next = grant.revalidate.map(|cadence| tokio::time::Instant::now() + cadence);
// The re-check in flight, kept out of the select so expiry and the session's
// close are still polled while a stalled server holds the reply.
let mut inflight: Option<Pin<Box<dyn Future<Output = crate::Result<Grant>> + Send>>> = None;
Expand All @@ -149,16 +151,15 @@ impl Driver {
let mut pending = false;

loop {
let expires = grant.expires.map(crate::grant::until);
let revalidate = async {
match next {
Some(at) => tokio::time::sleep_until(at.into()).await,
Some(at) => tokio::time::sleep_until(at).await,
None => std::future::pending().await,
}
};
let expire = async {
match expires {
Some(after) => tokio::time::sleep(after).await,
match self.expires {
Some(at) => tokio::time::sleep_until(at).await,
None => std::future::pending().await,
}
};
Expand All @@ -177,7 +178,8 @@ impl Driver {
match result {
Ok(fresh) => {
failures = 0;
next = fresh.revalidate.map(|cadence| Instant::now() + cadence);
self.expires = fresh.deadline();
next = fresh.revalidate.map(|cadence| tokio::time::Instant::now() + cadence);
producer.update(fresh.clone());
grant = fresh;
}
Expand All @@ -190,7 +192,7 @@ impl Driver {
failures += 1;
let delay = backoff(failures, grant.revalidate.unwrap_or(BACKOFF_MAX));
tracing::warn!(id = %self.request.id, %err, ?delay, "auth revalidation failed; retrying");
next = Some(Instant::now() + delay);
next = Some(tokio::time::Instant::now() + delay);
}
}
if pending {
Expand Down Expand Up @@ -314,6 +316,37 @@ mod tests {
server
}

/// Keep HTTP handling on the paused test clock instead of wiremock's separate runtime.
async fn clock_server(log: Log, grant: Grant, stall: bool) -> Client {
use axum::{Json, Router, extract::State, http::StatusCode, response::IntoResponse, routing::post};
let router = Router::new()
.route(
"/",
post(
|State((log, grant, stall)): State<(Log, Grant, bool)>, Json(request): Json<Request>| async move {
let event = request.event.clone();
log.0.lock().push(request);
match event {
Event::Connect => Json(grant).into_response(),
Event::Revalidate if stall => std::future::pending().await,
Event::Revalidate => StatusCode::SERVICE_UNAVAILABLE.into_response(),
Event::End { .. } => StatusCode::OK.into_response(),
}
},
),
)
.with_state((log, grant, stall));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let url = format!("http://{}/", listener.local_addr().unwrap());
tokio::spawn(async move { axum::serve(listener, router).await.unwrap() });
// Only the lease clock is under test; an HTTP timeout can auto-advance
// Tokio's paused clock before loopback I/O gets its first reactor turn.
Client {
http: reqwest::Client::builder().no_proxy().build().unwrap(),
url: url.parse().unwrap(),
}
}

fn client(server: &MockServer) -> Client {
Client::new(server.uri().parse().unwrap(), None).unwrap()
}
Expand Down Expand Up @@ -461,14 +494,12 @@ mod tests {

#[tokio::test]
async fn a_grant_within_clock_skew_stays_live() {
let server = server(Log::default(), |_| {
let mut grant = Grant::new(patterns(&["**"]), Patterns::new());
grant.expires = Some(SystemTime::now() - Duration::from_secs(1));
ResponseTemplate::new(200).set_body_json(grant)
})
.await;
tokio::time::pause();
let mut grant = Grant::new(patterns(&["**"]), Patterns::new());
grant.expires = Some(SystemTime::now() - Duration::from_secs(1));
let client = clock_server(Log::default(), grant, false).await;
let consumer = client.connect(request()).await.unwrap();

let consumer = client(&server).connect(request()).await.unwrap();
tokio::time::sleep(Duration::from_millis(500)).await;
assert!(
tokio::time::timeout(Duration::from_millis(100), consumer.closed())
Expand All @@ -485,15 +516,16 @@ mod tests {

#[tokio::test]
async fn an_outage_keeps_the_grant_until_expires() {
tokio::time::pause();
let log = Log::default();
let server = server(log.clone(), |request| match request.event {
Event::Connect => ResponseTemplate::new(200)
.set_body_json(grant(Some(Duration::from_secs(3)), Some(Duration::from_secs(1)))),
_ => ResponseTemplate::new(503),
})
let client = clock_server(
log.clone(),
grant(Some(Duration::from_secs(3)), Some(Duration::from_secs(1))),
false,
)
.await;
let consumer = client.connect(request()).await.unwrap();

let consumer = client(&server).connect(request()).await.unwrap();
tokio::time::sleep(Duration::from_millis(1500)).await;
assert!(log.revalidates() >= 1, "re-checks happened");
assert_eq!(
Expand All @@ -517,16 +549,15 @@ mod tests {

#[tokio::test]
async fn expiry_fires_while_a_recheck_is_stalled() {
let server = server(Log::default(), |request| match request.event {
Event::Connect => ResponseTemplate::new(200)
.set_body_json(grant(Some(Duration::from_secs(3)), Some(Duration::from_secs(1)))),
_ => ResponseTemplate::new(200)
.set_body_json(grant(Some(Duration::from_secs(3600)), Some(Duration::from_secs(60))))
.set_delay(Duration::from_secs(30)),
})
tokio::time::pause();
let client = clock_server(
Log::default(),
grant(Some(Duration::from_secs(3)), Some(Duration::from_secs(1))),
true,
)
.await;
let consumer = client.connect(request()).await.unwrap();

let consumer = client(&server).connect(request()).await.unwrap();
let reason = tokio::time::timeout(Duration::from_secs(5), consumer.closed())
.await
.expect("expired while the re-check was in flight");
Expand Down
8 changes: 7 additions & 1 deletion rs/moq-auth/src/grant.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ pub(crate) const CLOCK_SKEW: Duration = Duration::from_secs(5);
/// How long until `at`. A deadline up to [`CLOCK_SKEW`] in the past still has the
/// remaining window; anything older is zero. Future deadlines are unchanged, so a
/// grant that expires in ten seconds still expires in ten seconds.
pub(crate) fn until(at: SystemTime) -> Duration {
fn until(at: SystemTime) -> Duration {
match at.duration_since(SystemTime::now()) {
Ok(remaining) => remaining,
Err(late) => CLOCK_SKEW.saturating_sub(late.duration()),
Expand Down Expand Up @@ -67,6 +67,12 @@ impl Grant {
}
}

/// Snapshot the expiry on Tokio's clock, allowing five seconds of past clock skew.
#[cfg(feature = "tokio")]
pub fn deadline(&self) -> Option<tokio::time::Instant> {
self.expires.map(|at| tokio::time::Instant::now() + until(at))
}

/// Refuse a grant that admits nothing, asks to be revalidated without a bound or
/// at no interval, or has already expired. A few seconds of clock skew are
/// tolerated so an auth server whose clock runs behind still admits.
Expand Down
35 changes: 26 additions & 9 deletions rs/moq-relay/src/auth.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@
//! certificate is reported to whoever decides as a fact.

use std::sync::Arc;
use std::time::SystemTime;

use axum::http;
use moq_auth::{Bytes, Grant, Request, lease};
Expand Down Expand Up @@ -248,7 +247,7 @@ impl Lease {
let grant = consumer.grant();
Self {
token: Token::new(path, &grant),
expires: deadline(&grant),
expires: grant.deadline(),
consumer,
stats: Default::default(),
}
Expand Down Expand Up @@ -306,7 +305,7 @@ impl Lease {
self.stats.set_tier(fresh.tier.clone());
self.token.tier = fresh.tier;
}
self.expires = deadline(&grant);
self.expires = grant.deadline();
},
Err(reason) => return reason,
},
Expand All @@ -323,12 +322,6 @@ impl Lease {
}
}

/// The grant's `expires` as a deadline on tokio's clock, counted from now.
fn deadline(grant: &Grant) -> Option<tokio::time::Instant> {
let at = grant.expires?;
Some(tokio::time::Instant::now() + at.duration_since(SystemTime::now()).unwrap_or_default())
}

enum Decider {
Server(moq_auth::Client),
Public(Grant),
Expand Down Expand Up @@ -509,6 +502,7 @@ pub(crate) fn peer(identity: &moq_tokio::tls::PeerIdentity) -> Option<moq_auth::
#[cfg(test)]
mod tests {
use super::*;
use std::time::SystemTime;

fn patterns(texts: &[&str]) -> Patterns {
texts.iter().map(|text| text.parse().unwrap()).collect()
Expand Down Expand Up @@ -643,6 +637,29 @@ mod tests {
assert_eq!(reason, lease::Reason::Expired);
}

#[tokio::test]
async fn a_grant_within_clock_skew_stays_live() {
tokio::time::pause();
use std::time::Duration;

let mut grant = Grant::new(patterns(&["**"]), patterns(&["**"]));
grant.expires = Some(SystemTime::now() - Duration::from_secs(1));
grant.validate().expect("accepted inside the skew window");
let mut lease = Lease::new("/room", lease::Consumer::fixed(grant));
assert!(
tokio::time::timeout(Duration::from_secs(1), lease.ended())
.await
.is_err(),
"still live inside the skew window"
);
assert_eq!(
tokio::time::timeout(Duration::from_secs(4), lease.ended())
.await
.expect("expired once the skew window ended"),
lease::Reason::Expired
);
}

#[test]
fn token_keeps_every_grant_pattern() {
let mut grant = Grant::new(patterns(&["alice/**"]), patterns(&["**"]));
Expand Down
Loading