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
1 change: 0 additions & 1 deletion quest/m1/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,6 @@ the transport line in m2 assumes a single stack.
## Quests

- [One QUIC backend](/quest/m1/quic-one-backend.md) - quinn and quiche are deleted; noq (and iroh on it) is the only QUIC stack, with the qmux fallbacks untouched
- [One auth path](/quest/m1/auth-one-path.md) - Server, Public, and Refuse become tasks on the `Admissions` queue; `admit()` stays send-plus-await; `Mode` is gone
- [Announce event](/quest/m1/api-net-announce.md) - publishers announce prefixes on every wire, consumers scoped by a pattern read the covered path already trimmed, with no `as_prefix().expect()` at 89 call sites
- [Origin scoping](/quest/m1/api-net-origin.md) - `scope(root, patterns)` is one fallible call, a fresh origin has a random hop, and the handles stop derefing to `Hop`
- [Bindings announce match](/quest/m1/api-origin-scopes.md) - every binding takes a pattern scope and reports the announce match with its captures
Expand Down
78 changes: 0 additions & 78 deletions quest/m1/auth-one-path.md

This file was deleted.

3 changes: 1 addition & 2 deletions quest/m2/auth-embedder.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,7 @@ Additive on `moq-auth` and `moq-relay`, so on main after the merge:
`producer.due().await -> Due::{Revalidate, Expired}`, `update(grant)`
reschedules, `failed()` applies the backoff. `moq_auth::Client` becomes a
ten-line loop over it, and an embedder answering `Admissions` writes the
same ten lines instead of a second driver. [One auth path](/quest/m1/auth-one-path.md)
should reference this so its fold does not ship a third.
same ten lines instead of a second driver.
- `Cluster::admit(&self, auth: &Auth, request: moq_auth::Request) ->
Result<Admitted, auth::Error>` with
`Admitted { lease: Lease, publisher: Option<origin::Producer>, subscriber: Option<origin::Consumer>, stats: stats::Session }`
Expand Down
147 changes: 94 additions & 53 deletions rs/moq-relay/src/auth.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,9 @@ use serde_with::{OneOrMany, serde_as};
use tokio::sync::{mpsc, oneshot};
use url::Url;

/// The longest an embedder may take to answer an admission, the bound
/// The longest a decider may take to answer an admission, the bound
/// `moq_auth::Client` puts on a server, so a stalled decider refuses rather
/// than parks the sessions behind it.
/// than parks the session behind it.
const ADMIT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);

/// Where every session's grant comes from. Exactly one of `url` and the public
Expand Down Expand Up @@ -101,21 +101,21 @@ impl Config {

/// Build the [`Auth`] this configuration describes. `tls` is the client
/// identity an `https://` server is dialed with; `node` names this relay in
/// every request.
/// every request. Must be called within a Tokio runtime, which drives the
/// admission decider.
pub fn init(&self, node: impl Into<String>, tls: &moq_tokio::tls::Connect) -> anyhow::Result<Auth> {
self.validate()?;
let mode = match (&self.url, self.public_grant()) {
let decider = match (&self.url, self.public_grant()) {
(Some(url), _) => {
let tls = tls.build()?;
Mode::Server(moq_auth::Client::new(url.clone(), Some(tls))?)
Decider::Server(moq_auth::Client::new(url.clone(), Some(tls))?)
}
(None, Some(grant)) => Mode::Public(grant),
(None, Some(grant)) => Decider::Public(grant),
(None, None) => unreachable!("validated above"),
};
Ok(Auth {
mode: Arc::new(mode),
node: Arc::from(node.into()),
})
let (auth, admissions) = Auth::embedded(node);
decider.spawn(admissions);
Ok(auth)
}
}

Expand All @@ -140,7 +140,7 @@ pub enum Error {
impl From<moq_auth::Error> for Error {
fn from(err: moq_auth::Error) -> Self {
match err {
moq_auth::Error::Refused => Self::Refused,
moq_auth::Error::Refused | moq_auth::Error::UselessGrant => Self::Refused,
other => Self::Unavailable(other.to_string()),
}
}
Expand Down Expand Up @@ -301,23 +301,38 @@ impl Lease {
}
}

enum Mode {
enum Decider {
Server(moq_auth::Client),
Public(Grant),
/// The embedding process decides: every session is queued for whoever holds
/// the [`Admissions`], and admits nothing once they are gone.
Embedded(mpsc::UnboundedSender<Admission>),
/// Nothing admits an ordinary session; only a locally decided grant (the LAN
/// mesh credential) gets through. What a `--cluster-lan` process with no
/// listener of its own runs.
Refuse,
}

/// Admits sessions: asks the server, hands out the static public grant, or queues
/// the session for the embedder.
impl Decider {
fn spawn(self, mut admissions: Admissions) {
tokio::spawn(async move {
while let Some(admission) = admissions.next().await {
match &self {
Self::Server(client) => {
let client = client.clone();
tokio::spawn(async move {
match client.connect(admission.request.clone()).await {
Ok(lease) => admission.grant(lease),
Err(err) => admission.refuse(err.into()),
}
});
}
Self::Public(grant) => admission.grant(lease::Consumer::fixed(grant.clone())),
Self::Refuse => admission.refuse(Error::Refused),
}
}
});
}
}

/// Admits sessions by queueing every request for one admission decider.
#[derive(Clone)]
pub struct Auth {
mode: Arc<Mode>,
admissions: mpsc::UnboundedSender<Admission>,
node: Arc<str>,
}

Expand All @@ -328,19 +343,19 @@ impl Auth {
pub fn embedded(node: impl Into<String>) -> (Self, Admissions) {
let (sender, receiver) = mpsc::unbounded_channel();
let auth = Self {
mode: Arc::new(Mode::Embedded(sender)),
admissions: sender,
node: Arc::from(node.into()),
};
(auth, Admissions(receiver))
}

/// An `Auth` that refuses every session a server or a public grant would have
/// decided, admitting only what the relay decides for itself (a LAN peer).
/// Must be called within a Tokio runtime, which drives the refusal decider.
pub fn refuse(node: impl Into<String>) -> Self {
Self {
mode: Arc::new(Mode::Refuse),
node: Arc::from(node.into()),
}
let (auth, admissions) = Self::embedded(node);
Decider::Refuse.spawn(admissions);
auth
}

/// The name this relay puts in every request.
Expand All @@ -356,27 +371,17 @@ impl Auth {
/// Admit a session: the lease it holds, carrying the scope the origin applies.
pub async fn admit(&self, request: Request) -> Result<Lease, Error> {
let path = request.path.clone();
let consumer = match self.mode.as_ref() {
Mode::Server(client) => client.connect(request).await?,
// A certificate is a fact for a server to weigh; with no server it admits
// nothing on its own, so the peer gets what any anonymous session gets.
Mode::Public(grant) => lease::Consumer::fixed(grant.clone()),
Mode::Embedded(admissions) => {
let (reply, answer) = oneshot::channel();
admissions
.send(Admission { request, reply })
.map_err(|_| Error::Unavailable("nobody is answering admissions".into()))?;
let consumer = tokio::time::timeout(ADMIT_TIMEOUT, answer)
.await
.map_err(|_| Error::Unavailable("the admission timed out".into()))?
.map_err(|_| Error::Unavailable("the admission went unanswered".into()))??;
// Held to what a server's answer is held to: a grant that admits nothing
// or asks for a re-check without a bound is the decider's bug, not a refusal.
consumer.grant().validate()?;
consumer
}
Mode::Refuse => return Err(Error::Refused),
};
let (reply, answer) = oneshot::channel();
self.admissions
.send(Admission { request, reply })
.map_err(|_| Error::Unavailable("nobody is answering admissions".into()))?;
let consumer = tokio::time::timeout(ADMIT_TIMEOUT, answer)
.await
.map_err(|_| Error::Unavailable("the admission timed out".into()))?
.map_err(|_| Error::Unavailable("the admission went unanswered".into()))??;
// Every decider is held to the same answer contract: a grant that admits
// nothing or asks for a re-check without a bound is a bug, not a refusal.
consumer.grant().validate()?;
Ok(Lease::new(&path, consumer))
}

Expand Down Expand Up @@ -506,18 +511,34 @@ mod tests {
assert!(grant.publish.is_empty());
}

#[test]
fn a_public_config_admits_anonymous_and_certificate_alike() {
#[tokio::test]
async fn a_public_config_admits_anonymous_and_certificate_alike() {
let auth = config(None, &["anon/**"])
.init("relay-1", &moq_tokio::tls::Connect::default())
.unwrap();
let request = auth.request(moq_auth::Transport::Quic, "/anon/room");
let lease = futures::executor::block_on(auth.admit(request)).unwrap();
let lease = auth.admit(request).await.unwrap();
assert_eq!(lease.token().root, Path::new("anon/room").to_owned());
assert_eq!(lease.token().subscribe, patterns(&["anon/**"]));
assert_eq!(lease.token().tier, Tier::default());
}

#[test]
fn config_init_requires_a_runtime() {
let result = std::panic::catch_unwind(|| {
let _ = config(None, &["anon/**"])
.init("relay-1", &moq_tokio::tls::Connect::default())
.unwrap();
});
assert!(result.is_err());
}

#[test]
fn refuse_requires_a_runtime() {
let result = std::panic::catch_unwind(|| Auth::refuse("relay-1"));
assert!(result.is_err());
}

/// The embedder's answer is the session's verdict; an answer that never comes,
/// or a decider that is gone, is an outage rather than a refusal.
#[tokio::test]
Expand All @@ -542,16 +563,36 @@ mod tests {
let lease = auth.admit(request()).await.expect("granted");
assert_eq!(lease.token().root, Path::new("anon/room").to_owned());
assert!(matches!(auth.admit(request()).await, Err(Error::Refused)));
for _ in 0..2 {
assert!(matches!(auth.admit(request()).await, Err(Error::Unavailable(_))));
}
assert!(matches!(auth.admit(request()).await, Err(Error::Refused)));
assert!(matches!(auth.admit(request()).await, Err(Error::Unavailable(_))));
};
let (admissions, ()) = tokio::join!(decide, admit);

drop(admissions);
assert!(matches!(auth.admit(request()).await, Err(Error::Unavailable(_))));
}

#[tokio::test]
async fn dropping_the_last_auth_ends_admissions() {
let (auth, mut admissions) = Auth::embedded("relay-1");
let clone = auth.clone();
drop(auth);
assert!(
tokio::time::timeout(std::time::Duration::ZERO, admissions.next())
.await
.is_err()
);
drop(clone);
assert!(admissions.next().await.is_none());
}

#[tokio::test]
async fn refuse_answers_the_admission_queue() {
let auth = Auth::refuse("relay-1");
let request = auth.request(moq_auth::Transport::Quic, "/anon/room");
assert!(matches!(auth.admit(request).await, Err(Error::Refused)));
}

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