use std::collections::{HashMap, VecDeque};
use std::fmt;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;
use base64::Engine as _;
use bytes::Bytes;
use moq_transport_draft16::coding::TrackNamespace;
use moq_transport_draft16::serve::{self, TrackReaderMode};
use openrtc::broadcast::{BroadcastAdapterAction, BroadcastAdapterObservation, BroadcastSession};
use serde::{Deserialize, Serialize};
use sha2::{Digest as _, Sha256};
use tokio::sync::Mutex;
use url::Url;
use crate::{BroadcastAdapterFuture, BroadcastObjectsFuture, NativeBroadcastAdapter};
const BROADCAST_TRACK_NAME: &str = "openrtc-broadcast-1";
const MAX_RELAY_TOKEN_BYTES: usize = 16 * 1024;
const MAX_RELAY_NAMESPACE_BYTES: usize = 512;
const MAX_INBOUND_OBJECTS: usize = 128;
const MAX_INBOUND_BYTES: usize = openrtc::broadcast::MAX_BROADCAST_MEDIA_BYTES as usize;
const OBJECT_WRITE_TIMEOUT: Duration = Duration::from_secs(2);
const CONTROL_PLANE_TIMEOUT: Duration = Duration::from_secs(15);
#[derive(Clone, PartialEq, Eq, Hash)]
struct SessionKey {
broadcast_id: String,
broadcast_generation: u64,
grant_generation: u64,
}
impl SessionKey {
fn new(session: &BroadcastSession, broadcast_generation: u64, grant_generation: u64) -> Self {
Self {
broadcast_id: session.id(),
broadcast_generation,
grant_generation,
}
}
}
#[derive(Clone, PartialEq, Eq, Hash)]
struct PendingAccessKey {
broadcast_id: String,
grant_generation: u64,
}
#[derive(Clone)]
pub(super) struct NativeBroadcastRelayAccess {
relay_url: Url,
token: String,
namespace: String,
}
pub(super) struct ResolvedNativeBroadcastAccess {
pub issuer_public_key: [u8; 32],
pub relay: NativeBroadcastRelayAccess,
}
impl fmt::Debug for NativeBroadcastRelayAccess {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("NativeBroadcastRelayAccess")
.field("relay_origin", &relay_diagnostic_label(&self.relay_url))
.field("namespace", &self.namespace)
.field("token", &"[redacted]")
.finish()
}
}
impl NativeBroadcastRelayAccess {
fn parse(relay_url: String, token: String, namespace: String) -> Result<Self, String> {
let relay_url = Url::parse(relay_url.trim())
.map_err(|_| "broadcast relay URL is invalid".to_string())?;
if relay_url.scheme() != "https"
|| relay_url.username() != ""
|| relay_url.password().is_some()
|| relay_url.query().is_some()
|| relay_url.fragment().is_some()
{
return Err("broadcast relay URL is invalid".to_string());
}
if !(16..=MAX_RELAY_TOKEN_BYTES).contains(&token.len())
|| !token
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || b"._-".contains(&byte))
{
return Err("broadcast relay allocation is invalid".to_string());
}
if namespace.is_empty()
|| namespace.len() > MAX_RELAY_NAMESPACE_BYTES
|| !namespace
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || b"._/-".contains(&byte))
{
return Err("broadcast relay allocation is invalid".to_string());
}
Ok(Self {
relay_url,
token,
namespace,
})
}
fn connection_url(&self) -> Result<Url, String> {
let mut url = self.relay_url.clone();
url.path_segments_mut()
.map_err(|_| "broadcast relay URL is invalid".to_string())?
.pop_if_empty()
.push(&self.token);
Ok(url)
}
fn source_namespace(&self, source_slot: &str) -> TrackNamespace {
TrackNamespace::from_utf8_path(&format!(
"{}/source/{source_slot}",
self.namespace.trim_end_matches('/')
))
}
#[cfg(test)]
fn redact(&self, value: &str) -> String {
value.replace(&self.token, "[redacted]")
}
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub(super) struct NativeBroadcastDeviceProof {
pub public_key_jwk: serde_json::Value,
pub signature: String,
pub nonce: String,
pub issued_at: u64,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct BroadcastAccessRequest<'a> {
api_key: &'a str,
grant: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
publication_verifying_key: Option<&'a str>,
device_proof: NativeBroadcastDeviceProof,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct BroadcastAccessResponse {
grant: String,
issuer_public_key: String,
relay: BroadcastRelayResponse,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct BroadcastRelayResponse {
relay_url: String,
token: String,
namespace: String,
}
pub(super) fn broadcast_access_challenge(
api_key: &str,
grant: &str,
publication_verifying_key: Option<&[u8; 32]>,
nonce: &str,
issued_at: u64,
) -> (String, Option<String>) {
let grant_digest =
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(Sha256::digest(grant.as_bytes()));
let publication_key = publication_verifying_key
.map(|key| base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(key));
let challenge = [
"openrtc:v2:broadcast-access".to_string(),
api_key.to_string(),
grant_digest,
publication_key.clone().unwrap_or_default(),
nonce.to_string(),
issued_at.to_string(),
]
.join(":");
(challenge, publication_key)
}
pub(super) async fn resolve_broadcast_access(
endpoint: &str,
api_key: &str,
grant: &str,
publication_verifying_key: Option<&str>,
device_proof: NativeBroadcastDeviceProof,
) -> Result<ResolvedNativeBroadcastAccess, String> {
let endpoint = endpoint.trim_end_matches('/');
let url = format!("{endpoint}/v2/broadcasts/access");
let response = reqwest::Client::new()
.post(url)
.timeout(CONTROL_PLANE_TIMEOUT)
.json(&BroadcastAccessRequest {
api_key,
grant,
publication_verifying_key,
device_proof,
})
.send()
.await
.map_err(|_| "broadcast access request failed".to_string())?;
if !response.status().is_success() {
return Err(format!(
"broadcast access request was rejected ({})",
response.status().as_u16()
));
}
let response = response
.json::<BroadcastAccessResponse>()
.await
.map_err(|_| "broadcast access response is invalid".to_string())?;
if response.grant != grant {
return Err("broadcast access grant mismatch".to_string());
}
let issuer_public_key: [u8; 32] = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(response.issuer_public_key)
.map_err(|_| "broadcast issuer key is invalid".to_string())?
.try_into()
.map_err(|_| "broadcast issuer key is invalid".to_string())?;
Ok(ResolvedNativeBroadcastAccess {
issuer_public_key,
relay: NativeBroadcastRelayAccess::parse(
response.relay.relay_url,
response.relay.token,
response.relay.namespace,
)?,
})
}
#[derive(Default)]
struct InboundObjects {
objects: VecDeque<NativeBroadcastInboundObject>,
bytes: usize,
}
pub(super) struct NativeBroadcastInboundObject {
pub source_slot: String,
pub object: Vec<u8>,
}
impl InboundObjects {
fn push(&mut self, source_slot: &str, object: Bytes) -> Result<(), String> {
if object.len() > MAX_INBOUND_BYTES
|| self.objects.len() >= MAX_INBOUND_OBJECTS
|| self.bytes.saturating_add(object.len()) > MAX_INBOUND_BYTES
{
return Err("broadcast inbound object queue exceeded its bound".to_string());
}
self.bytes = self.bytes.saturating_add(object.len());
self.objects.push_back(NativeBroadcastInboundObject {
source_slot: source_slot.to_string(),
object: object.to_vec(),
});
Ok(())
}
fn take(&mut self, max: usize) -> Vec<NativeBroadcastInboundObject> {
let count = max.min(MAX_INBOUND_OBJECTS).min(self.objects.len());
let objects = self.objects.drain(..count).collect::<Vec<_>>();
let removed = objects.iter().map(|item| item.object.len()).sum::<usize>();
self.bytes = self.bytes.saturating_sub(removed);
objects
}
}
struct Draft16TrackWriter {
writer: serve::SubgroupWriter,
retained_sizes: VecDeque<usize>,
retained_bytes: usize,
}
impl Draft16TrackWriter {
fn new(writer: serve::SubgroupWriter) -> Self {
Self {
writer,
retained_sizes: VecDeque::new(),
retained_bytes: 0,
}
}
fn refresh_retention(&mut self) {
let retained_objects = self.writer.len();
while self.retained_sizes.len() > retained_objects {
if let Some(size) = self.retained_sizes.pop_front() {
self.retained_bytes = self.retained_bytes.saturating_sub(size);
}
}
}
async fn write(&mut self, object: Vec<u8>) -> Result<(), String> {
if object.len() > MAX_INBOUND_BYTES {
return Err("broadcast relay object exceeds its bound".to_string());
}
let object = Bytes::from(object);
let object_len = object.len();
tokio::time::timeout(OBJECT_WRITE_TIMEOUT, async {
loop {
self.refresh_retention();
if self.retained_sizes.len() < MAX_INBOUND_OBJECTS
&& self.retained_bytes.saturating_add(object_len) <= MAX_INBOUND_BYTES
{
self.writer
.write(object.clone())
.map_err(|_| "broadcast relay object write failed".to_string())?;
self.retained_sizes.push_back(object_len);
self.retained_bytes = self.retained_bytes.saturating_add(object_len);
return Ok(());
}
tokio::time::sleep(Duration::from_millis(1)).await;
}
})
.await
.map_err(|_| "broadcast relay object backpressure timed out".to_string())?
}
}
struct ActiveSession {
access: NativeBroadcastRelayAccess,
_endpoint_client: moq_native_ietf_draft16::quic::Client,
publisher: moq_transport_draft16::session::Publisher,
subscriber: moq_transport_draft16::session::Subscriber,
writers: HashMap<String, Draft16TrackWriter>,
subscriptions: HashMap<String, tokio::task::JoinHandle<()>>,
tasks: Vec<tokio::task::JoinHandle<()>>,
inbound: Arc<Mutex<InboundObjects>>,
retired: Arc<AtomicBool>,
failure_reported: Arc<AtomicBool>,
}
impl ActiveSession {
async fn shutdown(&mut self) {
self.retired.store(true, Ordering::Release);
for task in self.subscriptions.drain().map(|(_, task)| task) {
task.abort();
}
for task in self.tasks.drain(..) {
task.abort();
}
self.writers.clear();
}
}
pub(super) struct NativeBroadcastMoqAdapter {
pending_access: Mutex<HashMap<PendingAccessKey, NativeBroadcastRelayAccess>>,
authorized_access: Mutex<HashMap<SessionKey, NativeBroadcastRelayAccess>>,
active: Mutex<HashMap<SessionKey, Arc<Mutex<ActiveSession>>>>,
}
impl NativeBroadcastMoqAdapter {
pub(super) fn new() -> Self {
Self {
pending_access: Mutex::new(HashMap::new()),
authorized_access: Mutex::new(HashMap::new()),
active: Mutex::new(HashMap::new()),
}
}
pub(super) async fn authorize(
&self,
broadcast_id: String,
grant_generation: u64,
access: NativeBroadcastRelayAccess,
) {
let key = PendingAccessKey {
broadcast_id: broadcast_id.clone(),
grant_generation,
};
let mut pending = self.pending_access.lock().await;
pending.retain(|candidate, _| candidate.broadcast_id != broadcast_id);
pending.insert(key, access);
}
pub(super) async fn forget_authorization(&self, broadcast_id: &str, grant_generation: u64) {
self.pending_access.lock().await.remove(&PendingAccessKey {
broadcast_id: broadcast_id.to_string(),
grant_generation,
});
self.authorized_access.lock().await.retain(|key, _| {
key.broadcast_id != broadcast_id || key.grant_generation != grant_generation
});
}
async fn access_for_open(&self, key: &SessionKey) -> Option<NativeBroadcastRelayAccess> {
if let Some(access) = self.authorized_access.lock().await.get(key).cloned() {
return Some(access);
}
let pending_key = PendingAccessKey {
broadcast_id: key.broadcast_id.clone(),
grant_generation: key.grant_generation,
};
let access = self.pending_access.lock().await.remove(&pending_key)?;
self.authorized_access
.lock()
.await
.insert(key.clone(), access.clone());
Some(access)
}
async fn open(&self, session: BroadcastSession, key: SessionKey) -> Result<(), String> {
let access = self
.access_for_open(&key)
.await
.ok_or_else(|| "broadcast relay access is unavailable".to_string())?;
let connection_url = access.connection_url()?;
let tls = moq_native_ietf_draft16::tls::Args::default()
.load()
.map_err(|_| "broadcast relay TLS configuration failed".to_string())?;
let quic = moq_native_ietf_draft16::quic::Args::default()
.load()
.or_else(|_| {
moq_native_ietf_draft16::quic::Config::new(
"0.0.0.0:0".parse().expect("valid fallback bind address"),
None,
tls,
)
})
.map_err(|_| "broadcast relay QUIC configuration failed".to_string())?;
let endpoint = moq_native_ietf_draft16::quic::Endpoint::new(quic)
.map_err(|_| "broadcast relay endpoint creation failed".to_string())?;
let endpoint_client = endpoint.client;
let (webtransport, _, transport) = endpoint_client
.connect(&connection_url, None)
.await
.map_err(|_| "broadcast relay connection failed".to_string())?;
let (wire_session, publisher, subscriber) =
moq_transport_draft16::session::Session::connect(webtransport, None, transport)
.await
.map_err(|_| "broadcast relay Draft 16 setup failed".to_string())?;
let retired = Arc::new(AtomicBool::new(false));
let failure_reported = Arc::new(AtomicBool::new(false));
let run_retired = retired.clone();
let run_failure_reported = failure_reported.clone();
let run_session = session.clone();
let run_key = key.clone();
let run_task = tokio::spawn(async move {
let _ = wire_session.run().await;
report_background_failure(
&run_session,
&run_key,
&run_retired,
&run_failure_reported,
"relay-session-closed",
);
});
let active = Arc::new(Mutex::new(ActiveSession {
access,
_endpoint_client: endpoint_client,
publisher,
subscriber,
writers: HashMap::new(),
subscriptions: HashMap::new(),
tasks: vec![run_task],
inbound: Arc::new(Mutex::new(InboundObjects::default())),
retired,
failure_reported,
}));
let previous = {
let mut sessions = self.active.lock().await;
let stale = sessions
.keys()
.filter(|candidate| candidate.broadcast_id == key.broadcast_id)
.cloned()
.collect::<Vec<_>>();
let mut previous = stale
.into_iter()
.filter_map(|candidate| sessions.remove(&candidate))
.collect::<Vec<_>>();
if let Some(replaced) = sessions.insert(key, active) {
previous.push(replaced);
}
previous
};
for previous in previous {
previous.lock().await.shutdown().await;
}
Ok(())
}
async fn publish(
&self,
session: BroadcastSession,
key: &SessionKey,
source_slot: String,
object: Vec<u8>,
) -> Result<(), String> {
let active = self
.active
.lock()
.await
.get(key)
.cloned()
.ok_or_else(|| "broadcast relay session is unavailable".to_string())?;
let mut active = active.lock().await;
if !active.writers.contains_key(&source_slot) {
let namespace = active.access.source_namespace(&source_slot);
let (track_writer, track_reader) =
serve::Track::new(namespace, BROADCAST_TRACK_NAME).produce();
let writer = track_writer
.subgroups()
.and_then(|mut groups| groups.append(0))
.map_err(|_| "broadcast relay track creation failed".to_string())?;
let mut publisher = active.publisher.clone();
let published = publisher
.publish(track_reader, Default::default())
.await
.map_err(|_| "broadcast relay publication failed".to_string())?;
let retired = active.retired.clone();
let failure_reported = active.failure_reported.clone();
let published_session = session.clone();
let published_key = key.clone();
active.tasks.push(tokio::spawn(async move {
let _ = published.serve().await;
report_background_failure(
&published_session,
&published_key,
&retired,
&failure_reported,
"relay-publication-closed",
);
}));
active
.writers
.insert(source_slot.clone(), Draft16TrackWriter::new(writer));
}
active
.writers
.get_mut(&source_slot)
.expect("writer inserted above")
.write(object)
.await
}
async fn subscribe(
&self,
session: BroadcastSession,
key: &SessionKey,
source_slots: Vec<String>,
) -> Result<(), String> {
let active = self
.active
.lock()
.await
.get(key)
.cloned()
.ok_or_else(|| "broadcast relay session is unavailable".to_string())?;
let mut active = active.lock().await;
for source_slot in source_slots {
if active.subscriptions.contains_key(&source_slot) {
continue;
}
let namespace = active.access.source_namespace(&source_slot);
let subscriber = active.subscriber.clone();
let inbound = active.inbound.clone();
let retired = active.retired.clone();
let failure_reported = active.failure_reported.clone();
let subscription_session = session.clone();
let subscription_key = key.clone();
let subscription_slot = source_slot.clone();
let task = tokio::spawn(async move {
let (track_writer, track_reader) =
serve::Track::new(namespace, BROADCAST_TRACK_NAME).produce();
let mut subscriber = subscriber;
let result = async {
let subscription = subscriber
.subscribe_open(track_writer)
.await
.map_err(|_| "broadcast relay subscription failed".to_string())?;
let result = receive_track(&subscription_slot, track_reader, inbound).await;
drop(subscription);
result
}
.await;
if result.is_err() {
report_background_failure(
&subscription_session,
&subscription_key,
&retired,
&failure_reported,
"relay-subscription-closed",
);
}
});
active.subscriptions.insert(source_slot, task);
}
Ok(())
}
async fn close(&self, key: &SessionKey) {
if let Some(active) = self.active.lock().await.remove(key) {
active.lock().await.shutdown().await;
}
self.authorized_access.lock().await.remove(key);
self.pending_access.lock().await.remove(&PendingAccessKey {
broadcast_id: key.broadcast_id.clone(),
grant_generation: key.grant_generation,
});
}
async fn take(
&self,
key: &SessionKey,
max: usize,
) -> Result<Vec<NativeBroadcastInboundObject>, String> {
let active = self
.active
.lock()
.await
.get(key)
.cloned()
.ok_or_else(|| "broadcast relay session is unavailable".to_string())?;
let inbound = active.lock().await.inbound.clone();
let objects = inbound.lock().await.take(max);
Ok(objects)
}
pub(super) async fn take_inbound_objects_with_source(
&self,
session: &BroadcastSession,
max: usize,
) -> Result<Vec<NativeBroadcastInboundObject>, String> {
let key = self
.active_key_for_session(session)
.await
.ok_or_else(|| "broadcast relay generation is unavailable".to_string())?;
self.take(&key, max).await
}
pub(super) async fn record_delivered_bytes(
&self,
session: &BroadcastSession,
delivered_bytes: u64,
) -> Result<(), String> {
if delivered_bytes == 0 {
return Ok(());
}
let key = self
.active_key_for_session(session)
.await
.ok_or_else(|| "broadcast relay generation is unavailable".to_string())?;
session
.observe(BroadcastAdapterObservation::Usage {
broadcast_generation: key.broadcast_generation,
grant_generation: key.grant_generation,
delivered_bytes,
})
.map_err(|error| error.to_string())
}
async fn active_key_for_session(&self, session: &BroadcastSession) -> Option<SessionKey> {
let broadcast_id = session.id();
let sessions = self.active.lock().await;
let mut matches = sessions
.keys()
.filter(|key| key.broadcast_id == broadcast_id)
.cloned();
let key = matches.next()?;
matches.next().is_none().then_some(key)
}
}
impl NativeBroadcastAdapter for NativeBroadcastMoqAdapter {
fn apply(
&self,
session: BroadcastSession,
action: BroadcastAdapterAction,
) -> BroadcastAdapterFuture<'_> {
Box::pin(async move {
match action {
BroadcastAdapterAction::Open {
broadcast_generation,
grant_generation,
..
} => {
let key = SessionKey::new(&session, broadcast_generation, grant_generation);
Ok(Some(match self.open(session, key).await {
Ok(()) => BroadcastAdapterObservation::Opened {
broadcast_generation,
grant_generation,
},
Err(_) => BroadcastAdapterObservation::Failed {
broadcast_generation,
grant_generation,
retryable: true,
code: "relay-open-failed".to_string(),
},
}))
}
BroadcastAdapterAction::Publish {
broadcast_generation,
grant_generation,
source_slot,
object,
..
} => {
let key = SessionKey::new(&session, broadcast_generation, grant_generation);
Ok(self
.publish(session, &key, source_slot, object)
.await
.err()
.map(|_| BroadcastAdapterObservation::Failed {
broadcast_generation,
grant_generation,
retryable: true,
code: "relay-publish-failed".to_string(),
}))
}
BroadcastAdapterAction::Subscribe {
broadcast_generation,
grant_generation,
source_slots,
} => {
let key = SessionKey::new(&session, broadcast_generation, grant_generation);
Ok(self
.subscribe(session, &key, source_slots)
.await
.err()
.map(|_| BroadcastAdapterObservation::Failed {
broadcast_generation,
grant_generation,
retryable: true,
code: "relay-subscribe-failed".to_string(),
}))
}
BroadcastAdapterAction::Close {
broadcast_generation,
grant_generation,
reason,
} => {
let key = SessionKey::new(&session, broadcast_generation, grant_generation);
self.close(&key).await;
Ok(Some(BroadcastAdapterObservation::Closed {
broadcast_generation,
grant_generation,
code: reason,
}))
}
}
})
}
fn take_inbound_objects(
&self,
session: BroadcastSession,
max: usize,
) -> BroadcastObjectsFuture<'_> {
Box::pin(async move {
Ok(self
.take_inbound_objects_with_source(&session, max)
.await?
.into_iter()
.map(|item| item.object)
.collect())
})
}
}
fn report_background_failure(
session: &BroadcastSession,
key: &SessionKey,
retired: &AtomicBool,
failure_reported: &AtomicBool,
code: &str,
) {
if retired.load(Ordering::Acquire) || failure_reported.swap(true, Ordering::AcqRel) {
return;
}
let _ = session.observe(BroadcastAdapterObservation::Failed {
broadcast_generation: key.broadcast_generation,
grant_generation: key.grant_generation,
retryable: true,
code: code.to_string(),
});
}
async fn receive_track(
source_slot: &str,
track: serve::TrackReader,
inbound: Arc<Mutex<InboundObjects>>,
) -> Result<(), String> {
match track
.mode()
.await
.map_err(|_| "broadcast relay track failed".to_string())?
{
TrackReaderMode::Subgroups(mut groups) => {
while let Some(mut group) = groups
.next()
.await
.map_err(|_| "broadcast relay track failed".to_string())?
{
while let Some(object) = group
.read_next()
.await
.map_err(|_| "broadcast relay object failed".to_string())?
{
inbound.lock().await.push(source_slot, object)?;
}
}
Ok(())
}
_ => Err("broadcast relay selected an unsupported track mode".to_string()),
}
}
fn relay_diagnostic_label(url: &Url) -> String {
let host = url.host_str().unwrap_or("invalid");
match url.port() {
Some(port) => format!("{}://{host}:{port}", url.scheme()),
None => format!("{}://{host}", url.scheme()),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn access(token: &str) -> NativeBroadcastRelayAccess {
NativeBroadcastRelayAccess::parse(
"https://relay.cloudflare.com/moq".to_string(),
token.to_string(),
"openrtc/broadcast-1".to_string(),
)
.unwrap()
}
#[test]
fn access_challenge_matches_browser_contract() {
let key = [7_u8; 32];
let (challenge, encoded_key) = broadcast_access_challenge(
"pk_test_example",
"orb1.payload.signature",
Some(&key),
"nonce-1",
123,
);
let expected_digest = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(Sha256::digest(b"orb1.payload.signature"));
let expected_key = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(key);
assert_eq!(encoded_key.as_deref(), Some(expected_key.as_str()));
assert_eq!(
challenge,
format!(
"openrtc:v2:broadcast-access:pk_test_example:{expected_digest}:{expected_key}:nonce-1:123"
)
);
}
#[test]
fn credentials_are_redacted_from_debug_and_errors() {
let token = "secret-token-value-123";
let access = access(token);
let debug = format!("{access:?}");
assert!(!debug.contains(token));
assert!(debug.contains("[redacted]"));
let connection = access.connection_url().unwrap().to_string();
assert!(connection.contains(token));
assert!(!access.redact(&connection).contains(token));
}
#[test]
fn rejects_credentials_embedded_in_relay_url() {
assert!(NativeBroadcastRelayAccess::parse(
"https://relay.example/moq?jwt=secret".to_string(),
"valid-token-value-123".to_string(),
"openrtc/broadcast-1".to_string(),
)
.is_err());
}
#[test]
fn inbound_queue_is_bounded_and_accounted() {
let mut inbound = InboundObjects::default();
inbound.push("camera", Bytes::from_static(b"one")).unwrap();
inbound.push("camera", Bytes::from_static(b"two")).unwrap();
assert_eq!(inbound.bytes, 6);
let first = inbound.take(1).pop().unwrap();
assert_eq!(first.source_slot, "camera");
assert_eq!(first.object, b"one".to_vec());
assert_eq!(inbound.bytes, 3);
}
#[tokio::test]
async fn pending_access_is_scoped_to_one_grant_generation() {
let adapter = NativeBroadcastMoqAdapter::new();
adapter
.authorize("broadcast-1".to_string(), 4, access("token-value-123456"))
.await;
let stale = SessionKey {
broadcast_id: "broadcast-1".to_string(),
broadcast_generation: 7,
grant_generation: 3,
};
let current = SessionKey {
broadcast_id: "broadcast-1".to_string(),
broadcast_generation: 7,
grant_generation: 4,
};
assert!(adapter.access_for_open(&stale).await.is_none());
assert!(adapter.access_for_open(¤t).await.is_some());
assert!(adapter.access_for_open(¤t).await.is_some());
}
}