use std::{
collections::{HashMap, HashSet, VecDeque},
sync::{Arc, Mutex},
};
use base64::{engine::general_purpose::URL_SAFE_NO_PAD, Engine as _};
use ed25519_dalek::{Signature, Signer as _, SigningKey, Verifier as _, VerifyingKey};
use crate::media::{
EncodedMediaChunk, EncodedMediaSample, MediaControlFrame, MediaPublicationConfig,
PortableMediaSession, PublicationId,
};
pub const BROADCAST_PROTOCOL: &str = "openrtc-broadcast/1";
pub const MAX_BROADCAST_ACTIONS: usize = 128;
pub const MAX_BROADCAST_MEDIA_BYTES: u64 = 8 * 1024 * 1024;
const MAX_GRANT_BYTES: usize = 16 * 1024;
const MAX_PENDING_BINDING_CHALLENGES: usize = 128;
const BINDING_CHALLENGE_TTL_MS: u64 = 15_000;
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum BroadcastRole {
Owner,
Publisher,
Subscriber,
Processor,
}
impl BroadcastRole {
fn can_publish(self) -> bool {
matches!(self, Self::Owner | Self::Publisher | Self::Processor)
}
fn can_subscribe(self) -> bool {
matches!(self, Self::Owner | Self::Subscriber | Self::Processor)
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastBudget {
pub max_publishers: u16,
pub max_subscribers: u32,
pub max_bitrate_bps: u64,
pub max_egress_bytes: u64,
pub expires_at_ms: u64,
}
impl BroadcastBudget {
fn validate(&self, now_ms: u64) -> Result<(), BroadcastError> {
if self.expires_at_ms <= now_ms {
return Err(BroadcastError::GrantExpired);
}
if self.max_publishers == 0
|| self.max_subscribers == 0
|| self.max_bitrate_bps == 0
|| self.max_egress_bytes == 0
{
return Err(BroadcastError::InvalidGrant("budget must be non-zero"));
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastSourceAuthorization {
pub source_slot: String,
pub grant_generation: u64,
pub verifying_key: [u8; 32],
}
#[derive(Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastGrantClaims {
pub protocol: String,
pub grant_id: String,
pub broadcast_id: String,
pub broadcast_generation: u64,
pub grant_generation: u64,
pub role: BroadcastRole,
pub source_slots: Vec<String>,
pub source_authorizations: Vec<BroadcastSourceAuthorization>,
pub publication_verifying_key: Option<[u8; 32]>,
pub binding_key: [u8; 32],
pub not_before_ms: u64,
pub budget: BroadcastBudget,
pub relay_allocation_id: String,
}
impl std::fmt::Debug for BroadcastGrantClaims {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("BroadcastGrantClaims")
.field("grant_id", &self.grant_id)
.field("broadcast_id", &self.broadcast_id)
.field("broadcast_generation", &self.broadcast_generation)
.field("grant_generation", &self.grant_generation)
.field("role", &self.role)
.field("source_slots", &self.source_slots)
.field("budget", &self.budget)
.finish_non_exhaustive()
}
}
impl BroadcastGrantClaims {
fn validate(&self, now_ms: u64) -> Result<(), BroadcastError> {
if self.protocol != BROADCAST_PROTOCOL {
return Err(BroadcastError::InvalidGrant("unsupported protocol"));
}
if self.grant_id.is_empty() || self.broadcast_id.is_empty() {
return Err(BroadcastError::InvalidGrant("missing grant identity"));
}
if self.broadcast_generation == 0 || self.grant_generation == 0 {
return Err(BroadcastError::InvalidGrant("generation must be non-zero"));
}
if now_ms < self.not_before_ms {
return Err(BroadcastError::GrantNotYetValid);
}
self.budget.validate(now_ms)?;
if self.role.can_publish()
&& (self.source_slots.is_empty() || self.publication_verifying_key.is_none())
{
return Err(BroadcastError::InvalidGrant(
"publisher grant needs a source slot and signing key",
));
}
if self.role.can_subscribe() && self.source_authorizations.is_empty() {
return Err(BroadcastError::InvalidGrant(
"subscriber grant needs source authorization",
));
}
if self.relay_allocation_id.is_empty() || self.relay_allocation_id.len() > 160 {
return Err(BroadcastError::InvalidGrant("invalid relay allocation"));
}
let mut unique = HashSet::new();
if self.source_slots.iter().any(|slot| {
slot.is_empty()
|| slot.len() > 96
|| !unique.insert(slot)
|| !slot
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || b"_.-".contains(&byte))
}) {
return Err(BroadcastError::InvalidGrant("invalid source slot"));
}
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct VerifiedBroadcastGrant {
claims: BroadcastGrantClaims,
}
impl VerifiedBroadcastGrant {
pub fn grant_generation(&self) -> u64 {
self.claims.grant_generation
}
}
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastBindingChallenge {
pub handle: String,
pub signing_bytes: Vec<u8>,
pub expires_at_ms: u64,
}
pub fn sign_broadcast_grant(
claims: &BroadcastGrantClaims,
issuer: &SigningKey,
) -> Result<String, BroadcastError> {
let payload = serde_json::to_vec(claims).map_err(|_| BroadcastError::InvalidGrant("encode"))?;
if payload.len() > MAX_GRANT_BYTES {
return Err(BroadcastError::InvalidGrant("grant is too large"));
}
let payload = URL_SAFE_NO_PAD.encode(payload);
let signed = format!("orb1.{payload}");
let signature = issuer.sign(signed.as_bytes());
Ok(format!(
"{signed}.{}",
URL_SAFE_NO_PAD.encode(signature.to_bytes())
))
}
fn verify_broadcast_grant_envelope(
token: &str,
issuer: &VerifyingKey,
now_ms: u64,
) -> Result<BroadcastGrantClaims, BroadcastError> {
if token.len() > MAX_GRANT_BYTES * 2 {
return Err(BroadcastError::InvalidGrant("grant is too large"));
}
let mut parts = token.split('.');
let (Some(prefix), Some(payload), Some(signature), None) =
(parts.next(), parts.next(), parts.next(), parts.next())
else {
return Err(BroadcastError::InvalidGrant("malformed grant"));
};
if prefix != "orb1" {
return Err(BroadcastError::InvalidGrant("unsupported grant envelope"));
}
let signed = format!("{prefix}.{payload}");
let signature = URL_SAFE_NO_PAD
.decode(signature)
.ok()
.and_then(|bytes| Signature::from_slice(&bytes).ok())
.ok_or(BroadcastError::InvalidGrant("invalid issuer signature"))?;
issuer
.verify(signed.as_bytes(), &signature)
.map_err(|_| BroadcastError::InvalidGrant("invalid issuer signature"))?;
let claims: BroadcastGrantClaims = URL_SAFE_NO_PAD
.decode(payload)
.ok()
.and_then(|bytes| serde_json::from_slice(&bytes).ok())
.ok_or(BroadcastError::InvalidGrant("invalid claims"))?;
claims.validate(now_ms)?;
VerifyingKey::from_bytes(&claims.binding_key)
.map_err(|_| BroadcastError::InvalidGrant("invalid binding key"))?;
Ok(claims)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum BroadcastState {
Opening,
Live,
Retrying,
Closed,
Failed,
}
#[derive(Clone, serde::Serialize, serde::Deserialize)]
#[serde(
tag = "type",
rename_all = "kebab-case",
rename_all_fields = "camelCase"
)]
pub enum BroadcastAdapterAction {
Open {
broadcast_generation: u64,
grant_generation: u64,
role: BroadcastRole,
relay_allocation_id: String,
},
Publish {
broadcast_generation: u64,
grant_generation: u64,
source_slot: String,
object: Vec<u8>,
discardable: bool,
control: bool,
},
Subscribe {
broadcast_generation: u64,
grant_generation: u64,
source_slots: Vec<String>,
},
Close {
broadcast_generation: u64,
grant_generation: u64,
reason: String,
},
}
impl std::fmt::Debug for BroadcastAdapterAction {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Open {
broadcast_generation,
grant_generation,
role,
..
} => formatter
.debug_struct("Open")
.field("broadcast_generation", broadcast_generation)
.field("grant_generation", grant_generation)
.field("role", role)
.finish(),
Self::Publish {
broadcast_generation,
grant_generation,
source_slot,
object,
discardable,
control,
} => formatter
.debug_struct("Publish")
.field("broadcast_generation", broadcast_generation)
.field("grant_generation", grant_generation)
.field("source_slot", source_slot)
.field("bytes", &object.len())
.field("discardable", discardable)
.field("control", control)
.finish(),
Self::Subscribe {
broadcast_generation,
grant_generation,
source_slots,
} => formatter
.debug_struct("Subscribe")
.field("broadcast_generation", broadcast_generation)
.field("grant_generation", grant_generation)
.field("source_slots", source_slots)
.finish(),
Self::Close {
broadcast_generation,
grant_generation,
reason,
} => formatter
.debug_struct("Close")
.field("broadcast_generation", broadcast_generation)
.field("grant_generation", grant_generation)
.field("reason", reason)
.finish(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(
tag = "type",
rename_all = "kebab-case",
rename_all_fields = "camelCase"
)]
pub enum BroadcastAdapterObservation {
Opened {
broadcast_generation: u64,
grant_generation: u64,
},
Usage {
broadcast_generation: u64,
grant_generation: u64,
delivered_bytes: u64,
},
Failed {
broadcast_generation: u64,
grant_generation: u64,
retryable: bool,
code: String,
},
Closed {
broadcast_generation: u64,
grant_generation: u64,
code: String,
},
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastStats {
pub queued_actions: u64,
pub queued_media_bytes: u64,
pub dropped_media_objects: u64,
pub published_bytes: u64,
pub delivered_bytes: u64,
pub retries: u8,
}
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct AcceptedBroadcastMedia {
pub source_slot: String,
pub publication_id: String,
pub media_generation: u32,
#[serde(skip_serializing_if = "Option::is_none")]
pub control_type: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub control: Option<Vec<u8>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub media: Option<Vec<u8>>,
}
#[derive(Debug, PartialEq, Eq)]
pub enum BroadcastError {
InvalidGrant(&'static str),
GrantNotYetValid,
GrantExpired,
InvalidBindingProof,
InvalidBindingChallenge,
MissingSigner,
InvalidSigner,
UnauthorizedRole,
UnauthorizedSource,
StaleGeneration,
Revoked,
BudgetExceeded,
QueueFull,
Media(String),
InvalidMediaSignature,
Closed,
}
impl std::fmt::Display for BroadcastError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::InvalidGrant(reason) => write!(formatter, "invalid broadcast grant: {reason}"),
Self::GrantNotYetValid => formatter.write_str("broadcast grant is not valid yet"),
Self::GrantExpired => formatter.write_str("broadcast grant has expired"),
Self::InvalidBindingProof => {
formatter.write_str("broadcast installation binding proof is invalid")
}
Self::InvalidBindingChallenge => {
formatter.write_str("broadcast installation binding challenge is invalid or spent")
}
Self::MissingSigner => formatter.write_str("broadcast publisher signer is required"),
Self::InvalidSigner => formatter.write_str("broadcast publisher signer is invalid"),
Self::UnauthorizedRole => formatter.write_str("broadcast role is not authorized"),
Self::UnauthorizedSource => formatter.write_str("broadcast source is not authorized"),
Self::StaleGeneration => formatter.write_str("broadcast generation is stale"),
Self::Revoked => formatter.write_str("broadcast grant has been revoked"),
Self::BudgetExceeded => formatter.write_str("broadcast hard budget was exceeded"),
Self::QueueFull => formatter.write_str("broadcast queue is full"),
Self::Media(reason) => write!(formatter, "broadcast media error: {reason}"),
Self::InvalidMediaSignature => {
formatter.write_str("broadcast media signature is invalid")
}
Self::Closed => formatter.write_str("broadcast session is closed"),
}
}
}
impl std::error::Error for BroadcastError {}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
struct BroadcastMediaObject {
protocol: String,
broadcast_id: String,
broadcast_generation: u64,
grant_generation: u64,
source_slot: String,
control: Option<Vec<u8>>,
media: Option<Vec<u8>>,
signature: Vec<u8>,
}
impl BroadcastMediaObject {
fn signing_bytes(&self) -> Result<Vec<u8>, BroadcastError> {
serde_json::to_vec(&(
&self.protocol,
&self.broadcast_id,
self.broadcast_generation,
self.grant_generation,
&self.source_slot,
&self.control,
&self.media,
))
.map_err(|error| BroadcastError::Media(error.to_string()))
}
}
struct PendingBroadcastBindingChallenge {
token_hash: [u8; 32],
binding_key: [u8; 32],
signing_bytes: Vec<u8>,
expires_at_ms: u64,
}
#[derive(Default)]
struct BroadcastCatalog {
active: Mutex<HashSet<(String, u64, u64)>>,
pending_binding_challenges: Mutex<HashMap<String, PendingBroadcastBindingChallenge>>,
}
pub trait BroadcastSigner: Send + Sync {
fn verifying_key(&self) -> [u8; 32];
fn sign(&self, canonical_media_group: &[u8]) -> Result<Vec<u8>, BroadcastError>;
}
#[derive(Clone)]
pub struct BroadcastPublisherSigner {
signing_key: Arc<SigningKey>,
}
impl BroadcastPublisherSigner {
pub fn generate() -> Result<Self, BroadcastError> {
let mut seed = [0_u8; 32];
getrandom::getrandom(&mut seed)
.map_err(|_| BroadcastError::Media("publisher signer entropy failed".to_string()))?;
let signing_key = Arc::new(SigningKey::from_bytes(&seed));
seed.fill(0);
Ok(Self { signing_key })
}
pub fn verifying_key(&self) -> [u8; 32] {
self.signing_key.verifying_key().to_bytes()
}
}
impl BroadcastSigner for BroadcastPublisherSigner {
fn verifying_key(&self) -> [u8; 32] {
BroadcastPublisherSigner::verifying_key(self)
}
fn sign(&self, canonical_media_group: &[u8]) -> Result<Vec<u8>, BroadcastError> {
Ok(self
.signing_key
.sign(canonical_media_group)
.to_bytes()
.to_vec())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastAdmissionLimits {
pub max_publishers: u16,
pub max_subscribers: u32,
pub max_total_bitrate_bps: u64,
pub max_duration_ms: u64,
pub max_egress_bytes: u64,
}
impl BroadcastAdmissionLimits {
fn validate(self) -> Result<(), BroadcastError> {
if self.max_publishers == 0
|| self.max_subscribers == 0
|| self.max_total_bitrate_bps == 0
|| self.max_duration_ms == 0
|| self.max_egress_bytes == 0
{
return Err(BroadcastError::InvalidGrant(
"managed allocation limits must be non-zero",
));
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastAdmissionRequest {
pub role: BroadcastRole,
pub bitrate_bps: u64,
pub duration_ms: u64,
pub reserved_egress_bytes: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastAdmissionReservation {
pub reservation_id: String,
pub generation: u64,
pub role: BroadcastRole,
pub bitrate_bps: u64,
pub expires_at_ms: u64,
pub reserved_egress_bytes: u64,
}
#[derive(Default)]
struct ManagedBroadcastAllocationState {
closed: bool,
next_generation: u64,
reservations: HashMap<String, BroadcastAdmissionReservation>,
}
#[derive(Clone)]
pub struct ManagedBroadcastAllocation {
relay_allocation_id: String,
opened_at_ms: u64,
limits: BroadcastAdmissionLimits,
state: Arc<Mutex<ManagedBroadcastAllocationState>>,
}
impl ManagedBroadcastAllocation {
pub fn create(
relay_allocation_id: String,
opened_at_ms: u64,
limits: BroadcastAdmissionLimits,
) -> Result<Self, BroadcastError> {
limits.validate()?;
if relay_allocation_id.is_empty() || relay_allocation_id.len() > 160 {
return Err(BroadcastError::InvalidGrant("invalid relay allocation"));
}
Ok(Self {
relay_allocation_id,
opened_at_ms,
limits,
state: Arc::new(Mutex::new(ManagedBroadcastAllocationState::default())),
})
}
pub fn relay_allocation_id(&self) -> &str {
&self.relay_allocation_id
}
pub fn reserve(
&self,
reservation_id: String,
request: BroadcastAdmissionRequest,
now_ms: u64,
) -> Result<BroadcastAdmissionReservation, BroadcastError> {
self.update_reservation(reservation_id, request, now_ms, false)
}
pub fn renew(
&self,
reservation_id: String,
request: BroadcastAdmissionRequest,
now_ms: u64,
) -> Result<BroadcastAdmissionReservation, BroadcastError> {
self.update_reservation(reservation_id, request, now_ms, true)
}
pub fn revoke(&self, reservation_id: &str) -> bool {
self.state
.lock()
.map(|mut state| state.reservations.remove(reservation_id).is_some())
.unwrap_or(false)
}
pub fn close(&self) {
if let Ok(mut state) = self.state.lock() {
state.closed = true;
state.reservations.clear();
}
}
pub fn reservations(&self) -> Vec<BroadcastAdmissionReservation> {
self.state
.lock()
.map(|state| state.reservations.values().cloned().collect())
.unwrap_or_default()
}
fn update_reservation(
&self,
reservation_id: String,
request: BroadcastAdmissionRequest,
now_ms: u64,
renewing: bool,
) -> Result<BroadcastAdmissionReservation, BroadcastError> {
if reservation_id.is_empty() || reservation_id.len() > 160 {
return Err(BroadcastError::InvalidGrant("invalid reservation identity"));
}
if request.duration_ms == 0
|| request.bitrate_bps == 0
|| request.reserved_egress_bytes == 0
{
return Err(BroadcastError::InvalidGrant(
"reservation quantities must be non-zero",
));
}
let allocation_expires = self
.opened_at_ms
.checked_add(self.limits.max_duration_ms)
.ok_or(BroadcastError::BudgetExceeded)?;
let expires_at_ms = now_ms
.checked_add(request.duration_ms)
.ok_or(BroadcastError::BudgetExceeded)?;
if now_ms >= allocation_expires || expires_at_ms > allocation_expires {
return Err(BroadcastError::BudgetExceeded);
}
let mut state = self.state.lock().map_err(|_| BroadcastError::Closed)?;
if state.closed {
return Err(BroadcastError::Closed);
}
state
.reservations
.retain(|_, reservation| reservation.expires_at_ms > now_ms);
let previous = state.reservations.remove(&reservation_id);
if renewing != previous.is_some() {
if let Some(previous) = previous {
state.reservations.insert(reservation_id, previous);
}
return Err(if renewing {
BroadcastError::Revoked
} else {
BroadcastError::InvalidGrant("reservation is already active")
});
}
let publisher_count = state
.reservations
.values()
.filter(|reservation| reservation.role.can_publish())
.count() as u16
+ u16::from(request.role.can_publish());
let subscriber_count = state
.reservations
.values()
.filter(|reservation| reservation.role.can_subscribe())
.count() as u32
+ u32::from(request.role.can_subscribe());
let bitrate = state
.reservations
.values()
.fold(request.bitrate_bps, |total, reservation| {
total.saturating_add(reservation.bitrate_bps)
});
let egress = state
.reservations
.values()
.fold(request.reserved_egress_bytes, |total, reservation| {
total.saturating_add(reservation.reserved_egress_bytes)
});
if publisher_count > self.limits.max_publishers
|| subscriber_count > self.limits.max_subscribers
|| bitrate > self.limits.max_total_bitrate_bps
|| egress > self.limits.max_egress_bytes
{
if let Some(previous) = previous {
state.reservations.insert(reservation_id, previous);
}
return Err(BroadcastError::BudgetExceeded);
}
state.next_generation = state.next_generation.saturating_add(1).max(1);
let reservation = BroadcastAdmissionReservation {
reservation_id: reservation_id.clone(),
generation: state.next_generation,
role: request.role,
bitrate_bps: request.bitrate_bps,
expires_at_ms,
reserved_egress_bytes: request.reserved_egress_bytes,
};
state
.reservations
.insert(reservation_id, reservation.clone());
Ok(reservation)
}
}
#[derive(Clone, Default)]
pub struct Broadcasts {
catalog: Arc<BroadcastCatalog>,
}
impl Broadcasts {
pub fn prepare_grant_verification(
&self,
token: &str,
issuer: &VerifyingKey,
now_ms: u64,
) -> Result<BroadcastBindingChallenge, BroadcastError> {
let claims = verify_broadcast_grant_envelope(token, issuer, now_ms)?;
let expires_at_ms = now_ms
.checked_add(BINDING_CHALLENGE_TTL_MS)
.ok_or(BroadcastError::InvalidBindingChallenge)?;
let token_hash = *blake3::hash(token.as_bytes()).as_bytes();
let mut nonce = [0_u8; 32];
getrandom::getrandom(&mut nonce).map_err(|_| BroadcastError::InvalidBindingChallenge)?;
let handle = URL_SAFE_NO_PAD.encode(nonce);
let signing_bytes = serde_json::to_vec(&(
BROADCAST_PROTOCOL,
"installation-binding",
claims.grant_id.as_str(),
claims.broadcast_id.as_str(),
claims.broadcast_generation,
claims.grant_generation,
URL_SAFE_NO_PAD.encode(token_hash),
handle.as_str(),
expires_at_ms,
))
.map_err(|_| BroadcastError::InvalidBindingChallenge)?;
let pending = PendingBroadcastBindingChallenge {
token_hash,
binding_key: claims.binding_key,
signing_bytes: signing_bytes.clone(),
expires_at_ms,
};
let mut challenges = self
.catalog
.pending_binding_challenges
.lock()
.map_err(|_| BroadcastError::Closed)?;
challenges.retain(|_, challenge| challenge.expires_at_ms > now_ms);
if challenges.len() >= MAX_PENDING_BINDING_CHALLENGES {
return Err(BroadcastError::QueueFull);
}
challenges.insert(handle.clone(), pending);
Ok(BroadcastBindingChallenge {
handle,
signing_bytes,
expires_at_ms,
})
}
pub fn complete_grant_verification(
&self,
token: &str,
issuer: &VerifyingKey,
challenge_handle: &str,
binding_signature: &[u8],
now_ms: u64,
) -> Result<VerifiedBroadcastGrant, BroadcastError> {
let pending = self
.catalog
.pending_binding_challenges
.lock()
.map_err(|_| BroadcastError::Closed)?
.remove(challenge_handle)
.ok_or(BroadcastError::InvalidBindingChallenge)?;
if pending.expires_at_ms <= now_ms
|| pending.token_hash != *blake3::hash(token.as_bytes()).as_bytes()
{
return Err(BroadcastError::InvalidBindingChallenge);
}
let claims = verify_broadcast_grant_envelope(token, issuer, now_ms)?;
if claims.binding_key != pending.binding_key {
return Err(BroadcastError::InvalidBindingChallenge);
}
let binding_key = VerifyingKey::from_bytes(&pending.binding_key)
.map_err(|_| BroadcastError::InvalidGrant("invalid binding key"))?;
let binding_signature = Signature::from_slice(binding_signature)
.map_err(|_| BroadcastError::InvalidBindingProof)?;
binding_key
.verify(&pending.signing_bytes, &binding_signature)
.map_err(|_| BroadcastError::InvalidBindingProof)?;
Ok(VerifiedBroadcastGrant { claims })
}
pub fn open(
&self,
grant: VerifiedBroadcastGrant,
now_ms: u64,
) -> Result<BroadcastSession, BroadcastError> {
self.open_inner(grant, None, now_ms)
}
pub fn open_with_signer(
&self,
grant: VerifiedBroadcastGrant,
signer: Arc<dyn BroadcastSigner>,
now_ms: u64,
) -> Result<BroadcastSession, BroadcastError> {
self.open_inner(grant, Some(signer), now_ms)
}
pub fn open_publisher(
&self,
grant: VerifiedBroadcastGrant,
signer: &BroadcastPublisherSigner,
now_ms: u64,
) -> Result<BroadcastSession, BroadcastError> {
self.open_inner(grant, Some(Arc::new(signer.clone())), now_ms)
}
fn open_inner(
&self,
grant: VerifiedBroadcastGrant,
signer: Option<Arc<dyn BroadcastSigner>>,
now_ms: u64,
) -> Result<BroadcastSession, BroadcastError> {
grant.claims.validate(now_ms)?;
if grant.claims.role.can_publish() {
let signer = signer.as_ref().ok_or(BroadcastError::MissingSigner)?;
if Some(signer.verifying_key()) != grant.claims.publication_verifying_key {
return Err(BroadcastError::InvalidSigner);
}
}
let key = (
grant.claims.broadcast_id.clone(),
grant.claims.broadcast_generation,
grant.claims.grant_generation,
);
let mut active = self
.catalog
.active
.lock()
.map_err(|_| BroadcastError::Closed)?;
if !active.insert(key.clone()) {
return Err(BroadcastError::InvalidGrant(
"grant generation is already active",
));
}
drop(active);
let mut queue = VecDeque::new();
queue.push_back(BroadcastAdapterAction::Open {
broadcast_generation: grant.claims.broadcast_generation,
grant_generation: grant.claims.grant_generation,
role: grant.claims.role,
relay_allocation_id: grant.claims.relay_allocation_id.clone(),
});
if grant.claims.role.can_subscribe() {
queue.push_back(BroadcastAdapterAction::Subscribe {
broadcast_generation: grant.claims.broadcast_generation,
grant_generation: grant.claims.grant_generation,
source_slots: grant
.claims
.source_authorizations
.iter()
.map(|authorization| authorization.source_slot.clone())
.collect(),
});
}
let publish_tokens_bytes = grant.claims.budget.max_bitrate_bps.saturating_add(7) / 8;
Ok(BroadcastSession {
inner: Arc::new(Mutex::new(BroadcastSessionInner {
grant: grant.claims,
state: BroadcastState::Opening,
actions: queue,
media: PortableMediaSession::default(),
signer,
opened_at_ms: now_ms,
opened_at: web_time::Instant::now(),
publish_tokens_bytes,
publish_tokens_refilled_at: web_time::Instant::now(),
publication_slots: HashMap::new(),
publication_generations: HashMap::new(),
stats: BroadcastStats::default(),
revoked: false,
})),
catalog: self.catalog.clone(),
key,
})
}
}
struct BroadcastSessionInner {
grant: BroadcastGrantClaims,
state: BroadcastState,
actions: VecDeque<BroadcastAdapterAction>,
media: PortableMediaSession,
signer: Option<Arc<dyn BroadcastSigner>>,
opened_at_ms: u64,
opened_at: web_time::Instant,
publish_tokens_bytes: u64,
publish_tokens_refilled_at: web_time::Instant,
publication_slots: HashMap<PublicationId, String>,
publication_generations: HashMap<PublicationId, u32>,
stats: BroadcastStats,
revoked: bool,
}
#[derive(Clone)]
pub struct BroadcastSession {
inner: Arc<Mutex<BroadcastSessionInner>>,
catalog: Arc<BroadcastCatalog>,
key: (String, u64, u64),
}
impl BroadcastSession {
pub fn id(&self) -> String {
self.inner
.lock()
.expect("broadcast session lock")
.grant
.broadcast_id
.clone()
}
pub fn role(&self) -> BroadcastRole {
self.inner
.lock()
.expect("broadcast session lock")
.grant
.role
}
pub fn source_slots(&self) -> Vec<String> {
self.inner
.lock()
.map(|inner| inner.grant.source_slots.clone())
.unwrap_or_default()
}
pub fn state(&self) -> BroadcastState {
self.inner.lock().expect("broadcast session lock").state
}
pub fn budget(&self) -> BroadcastBudget {
self.inner
.lock()
.expect("broadcast session lock")
.grant
.budget
.clone()
}
pub fn stats(&self) -> BroadcastStats {
self.inner.lock().expect("broadcast session lock").stats
}
pub fn begin_publication(
&self,
source_slot: &str,
publication: MediaPublicationConfig,
) -> Result<MediaPublicationConfig, BroadcastError> {
let mut inner = self.inner.lock().map_err(|_| BroadcastError::Closed)?;
ensure_open(&mut inner)?;
if !inner.grant.role.can_publish() {
return Err(BroadcastError::UnauthorizedRole);
}
if !inner
.grant
.source_slots
.iter()
.any(|slot| slot == source_slot)
{
return Err(BroadcastError::UnauthorizedSource);
}
let (publication, control) = inner
.media
.begin_publication(publication)
.map_err(|error| BroadcastError::Media(error.to_string()))?;
inner
.publication_slots
.insert(publication.publication_id, source_slot.to_string());
inner
.publication_generations
.insert(publication.publication_id, publication.media_generation);
let object = sign_media_object(&inner, source_slot, Some(control), None)?;
push_publish(&mut inner, source_slot, object, false, true)?;
Ok(publication)
}
pub fn publish(
&self,
publication_id: PublicationId,
sample: EncodedMediaSample,
) -> Result<(), BroadcastError> {
let mut inner = self.inner.lock().map_err(|_| BroadcastError::Closed)?;
ensure_open(&mut inner)?;
let source_slot = inner
.publication_slots
.get(&publication_id)
.cloned()
.ok_or(BroadcastError::UnauthorizedSource)?;
if sample.payload.len() as u64 > inner.grant.budget.max_egress_bytes {
return Err(BroadcastError::BudgetExceeded);
}
let discardable = sample.discardable;
let media = inner
.media
.encode_sample(publication_id, sample)
.map_err(|error| BroadcastError::Media(error.to_string()))?;
let object = sign_media_object(&inner, &source_slot, None, Some(media))?;
push_publish(&mut inner, &source_slot, object, discardable, false)
}
pub fn retire_publication(&self, publication_id: PublicationId) -> Result<(), BroadcastError> {
let mut inner = self.inner.lock().map_err(|_| BroadcastError::Closed)?;
ensure_open(&mut inner)?;
let source_slot = inner
.publication_slots
.get(&publication_id)
.cloned()
.ok_or(BroadcastError::UnauthorizedSource)?;
let media_generation = inner
.publication_generations
.get(&publication_id)
.copied()
.ok_or(BroadcastError::StaleGeneration)?;
inner.media.retire_publication(publication_id);
let control = MediaControlFrame::Stop {
publication_id,
media_generation,
reason: "publisher-retired".to_string(),
}
.encode()
.map_err(|error| BroadcastError::Media(error.to_string()))?;
let object = sign_media_object(&inner, &source_slot, Some(control), None)?;
push_publish(&mut inner, &source_slot, object, false, true)?;
inner.publication_slots.remove(&publication_id);
inner.publication_generations.remove(&publication_id);
Ok(())
}
pub fn pause_publication(&self, publication_id: PublicationId) -> Result<(), BroadcastError> {
let mut inner = self.inner.lock().map_err(|_| BroadcastError::Closed)?;
ensure_open(&mut inner)?;
inner
.media
.pause_publication(publication_id)
.map_err(|error| BroadcastError::Media(error.to_string()))
}
pub fn retire_receiver(
&self,
publication_id: PublicationId,
media_generation: u32,
) -> Result<(), BroadcastError> {
let mut inner = self.inner.lock().map_err(|_| BroadcastError::Closed)?;
inner
.media
.retire_receiver(publication_id, media_generation);
Ok(())
}
pub fn subscribe(&self, source_slots: Vec<String>) -> Result<(), BroadcastError> {
let mut inner = self.inner.lock().map_err(|_| BroadcastError::Closed)?;
ensure_open(&mut inner)?;
if !inner.grant.role.can_subscribe() {
return Err(BroadcastError::UnauthorizedRole);
}
if source_slots.is_empty()
|| source_slots.iter().any(|slot| {
!inner
.grant
.source_authorizations
.iter()
.any(|authorized| authorized.source_slot == *slot)
})
{
return Err(BroadcastError::UnauthorizedSource);
}
push_control(
&mut inner,
BroadcastAdapterAction::Subscribe {
broadcast_generation: self.key.1,
grant_generation: self.key.2,
source_slots,
},
)
}
pub fn accept_object(
&self,
encoded: &[u8],
) -> Result<Option<EncodedMediaChunk>, BroadcastError> {
match self.accept_media(encoded)? {
Some(value) if value.media.is_some() => value
.media
.as_deref()
.map(EncodedMediaChunk::decode)
.transpose()
.map_err(|error| BroadcastError::Media(error.to_string())),
_ => Ok(None),
}
}
pub fn accept_media(
&self,
encoded: &[u8],
) -> Result<Option<AcceptedBroadcastMedia>, BroadcastError> {
self.accept_media_from(None, encoded)
}
pub fn accept_media_for_source(
&self,
source_slot: &str,
encoded: &[u8],
) -> Result<Option<AcceptedBroadcastMedia>, BroadcastError> {
self.accept_media_from(Some(source_slot), encoded)
}
fn accept_media_from(
&self,
expected_source_slot: Option<&str>,
encoded: &[u8],
) -> Result<Option<AcceptedBroadcastMedia>, BroadcastError> {
let mut inner = self.inner.lock().map_err(|_| BroadcastError::Closed)?;
ensure_open(&mut inner)?;
if !inner.grant.role.can_subscribe() {
return Err(BroadcastError::UnauthorizedRole);
}
let object: BroadcastMediaObject = serde_json::from_slice(encoded)
.map_err(|error| BroadcastError::Media(error.to_string()))?;
if expected_source_slot.is_some_and(|expected| expected != object.source_slot) {
return Err(BroadcastError::UnauthorizedSource);
}
if object.protocol != BROADCAST_PROTOCOL
|| object.broadcast_id != inner.grant.broadcast_id
|| object.broadcast_generation != inner.grant.broadcast_generation
|| object.grant_generation > inner.grant.grant_generation
{
return Err(BroadcastError::StaleGeneration);
}
let authorization = inner
.grant
.source_authorizations
.iter()
.find(|authorization| {
authorization.source_slot == object.source_slot
&& authorization.grant_generation == object.grant_generation
})
.ok_or(BroadcastError::UnauthorizedSource)?;
let verifying_key = VerifyingKey::from_bytes(&authorization.verifying_key)
.map_err(|_| BroadcastError::InvalidMediaSignature)?;
let signature = Signature::from_slice(&object.signature)
.map_err(|_| BroadcastError::InvalidMediaSignature)?;
verifying_key
.verify(&object.signing_bytes()?, &signature)
.map_err(|_| BroadcastError::InvalidMediaSignature)?;
match (object.control, object.media) {
(Some(control), None) => {
let frame = inner
.media
.decode_control(&control)
.map_err(|error| BroadcastError::Media(error.to_string()))?;
let (publication_id, media_generation, control_type) = match frame {
MediaControlFrame::Publish {
publication_id,
media_generation,
..
} => (publication_id, media_generation, "publish"),
MediaControlFrame::SetEnabled {
publication_id,
media_generation,
..
} => (publication_id, media_generation, "set-enabled"),
MediaControlFrame::RequestKeyframe {
publication_id,
media_generation,
} => (publication_id, media_generation, "request-keyframe"),
MediaControlFrame::Stop {
publication_id,
media_generation,
..
} => (publication_id, media_generation, "stop"),
};
Ok(Some(AcceptedBroadcastMedia {
source_slot: object.source_slot,
publication_id: publication_id.to_string(),
media_generation,
control_type: Some(control_type.to_string()),
control: Some(control),
media: None,
}))
}
(None, Some(media)) => {
let chunk = inner
.media
.decode_chunk(&media)
.map_err(|error| BroadcastError::Media(error.to_string()))?;
Ok(Some(AcceptedBroadcastMedia {
source_slot: object.source_slot,
publication_id: chunk.publication_id.to_string(),
media_generation: chunk.media_generation,
control_type: None,
control: None,
media: Some(media),
}))
}
_ => Err(BroadcastError::Media(
"broadcast media object must contain exactly one payload".to_string(),
)),
}
}
pub fn observe(&self, observation: BroadcastAdapterObservation) -> Result<(), BroadcastError> {
let mut inner = self.inner.lock().map_err(|_| BroadcastError::Closed)?;
ensure_open(&mut inner)?;
let generations = match &observation {
BroadcastAdapterObservation::Opened {
broadcast_generation,
grant_generation,
}
| BroadcastAdapterObservation::Usage {
broadcast_generation,
grant_generation,
..
}
| BroadcastAdapterObservation::Failed {
broadcast_generation,
grant_generation,
..
}
| BroadcastAdapterObservation::Closed {
broadcast_generation,
grant_generation,
..
} => (*broadcast_generation, *grant_generation),
};
if generations
!= (
inner.grant.broadcast_generation,
inner.grant.grant_generation,
)
{
return Err(BroadcastError::StaleGeneration);
}
match observation {
BroadcastAdapterObservation::Opened { .. } => inner.state = BroadcastState::Live,
BroadcastAdapterObservation::Usage {
delivered_bytes, ..
} => {
inner.stats.delivered_bytes =
inner.stats.delivered_bytes.saturating_add(delivered_bytes);
if inner.stats.delivered_bytes > inner.grant.budget.max_egress_bytes {
close_inner(&mut inner, "budget-exceeded");
return Err(BroadcastError::BudgetExceeded);
}
}
BroadcastAdapterObservation::Failed {
retryable: true, ..
} if inner.stats.retries == 0 => {
inner.stats.retries = 1;
inner.state = BroadcastState::Retrying;
let open = BroadcastAdapterAction::Open {
broadcast_generation: inner.grant.broadcast_generation,
grant_generation: inner.grant.grant_generation,
role: inner.grant.role,
relay_allocation_id: inner.grant.relay_allocation_id.clone(),
};
push_control(&mut inner, open)?;
}
BroadcastAdapterObservation::Failed { .. } => {
close_inner(&mut inner, "adapter-failed");
inner.state = BroadcastState::Failed;
}
BroadcastAdapterObservation::Closed { .. } => inner.state = BroadcastState::Closed,
}
Ok(())
}
pub fn take_actions(&self, max: usize) -> Vec<BroadcastAdapterAction> {
let mut inner = self.inner.lock().expect("broadcast session lock");
let count = max.min(MAX_BROADCAST_ACTIONS).min(inner.actions.len());
let actions: Vec<_> = inner.actions.drain(..count).collect();
let removed = actions.iter().map(action_media_bytes).sum::<u64>();
inner.stats.queued_media_bytes = inner.stats.queued_media_bytes.saturating_sub(removed);
inner.stats.queued_actions = inner.actions.len() as u64;
actions
}
pub fn revoke(&self, grant_generation: u64) -> Result<(), BroadcastError> {
let mut inner = self.inner.lock().map_err(|_| BroadcastError::Closed)?;
if grant_generation < inner.grant.grant_generation {
return Err(BroadcastError::StaleGeneration);
}
inner.revoked = true;
close_inner(&mut inner, "grant-revoked");
self.release();
Ok(())
}
pub fn close(&self) {
if let Ok(mut inner) = self.inner.lock() {
close_inner(&mut inner, "consumer-close");
}
self.release();
}
fn release(&self) {
if let Ok(mut active) = self.catalog.active.lock() {
active.remove(&self.key);
}
}
}
fn ensure_open(inner: &mut BroadcastSessionInner) -> Result<(), BroadcastError> {
if inner.revoked {
return Err(BroadcastError::Revoked);
}
if matches!(inner.state, BroadcastState::Closed | BroadcastState::Failed) {
return Err(BroadcastError::Closed);
}
let elapsed_ms = inner.opened_at.elapsed().as_millis().min(u64::MAX as u128) as u64;
if inner.opened_at_ms.saturating_add(elapsed_ms) >= inner.grant.budget.expires_at_ms {
close_inner(inner, "grant-expired");
return Err(BroadcastError::GrantExpired);
}
Ok(())
}
fn sign_media_object(
inner: &BroadcastSessionInner,
source_slot: &str,
control: Option<Vec<u8>>,
media: Option<Vec<u8>>,
) -> Result<Vec<u8>, BroadcastError> {
let signer = inner.signer.as_ref().ok_or(BroadcastError::MissingSigner)?;
let mut object = BroadcastMediaObject {
protocol: BROADCAST_PROTOCOL.to_string(),
broadcast_id: inner.grant.broadcast_id.clone(),
broadcast_generation: inner.grant.broadcast_generation,
grant_generation: inner.grant.grant_generation,
source_slot: source_slot.to_string(),
control,
media,
signature: Vec::new(),
};
let canonical = object.signing_bytes()?;
let signature = signer.sign(&canonical)?;
let signature = Signature::from_slice(&signature).map_err(|_| BroadcastError::InvalidSigner)?;
let verifying_key = inner
.grant
.publication_verifying_key
.and_then(|key| VerifyingKey::from_bytes(&key).ok())
.ok_or(BroadcastError::InvalidSigner)?;
verifying_key
.verify(&canonical, &signature)
.map_err(|_| BroadcastError::InvalidSigner)?;
object.signature = signature.to_bytes().to_vec();
serde_json::to_vec(&object).map_err(|error| BroadcastError::Media(error.to_string()))
}
fn push_publish(
inner: &mut BroadcastSessionInner,
source_slot: &str,
object: Vec<u8>,
discardable: bool,
control: bool,
) -> Result<(), BroadcastError> {
let bytes = object.len() as u64;
if bytes > MAX_BROADCAST_MEDIA_BYTES {
return Err(BroadcastError::QueueFull);
}
prepare_publish_budget(inner, bytes)?;
while inner.actions.len() >= MAX_BROADCAST_ACTIONS
|| inner.stats.queued_media_bytes.saturating_add(bytes) > MAX_BROADCAST_MEDIA_BYTES
{
let position = inner
.actions
.iter()
.position(|action| {
matches!(
action,
BroadcastAdapterAction::Publish {
control: false,
discardable: true,
..
}
)
})
.or_else(|| {
inner.actions.iter().position(|action| {
matches!(
action,
BroadcastAdapterAction::Publish { control: false, .. }
)
})
});
let Some(position) = position else {
return Err(BroadcastError::QueueFull);
};
let removed = inner
.actions
.remove(position)
.expect("queued action exists");
inner.stats.queued_media_bytes = inner
.stats
.queued_media_bytes
.saturating_sub(action_media_bytes(&removed));
inner.stats.dropped_media_objects = inner.stats.dropped_media_objects.saturating_add(1);
}
inner.publish_tokens_bytes -= bytes;
inner.actions.push_back(BroadcastAdapterAction::Publish {
broadcast_generation: inner.grant.broadcast_generation,
grant_generation: inner.grant.grant_generation,
source_slot: source_slot.to_string(),
object,
discardable,
control,
});
inner.stats.queued_media_bytes = inner.stats.queued_media_bytes.saturating_add(bytes);
inner.stats.published_bytes = inner.stats.published_bytes.saturating_add(bytes);
inner.stats.queued_actions = inner.actions.len() as u64;
Ok(())
}
fn prepare_publish_budget(
inner: &mut BroadcastSessionInner,
bytes: u64,
) -> Result<(), BroadcastError> {
if inner.stats.published_bytes.saturating_add(bytes) > inner.grant.budget.max_egress_bytes {
return Err(BroadcastError::BudgetExceeded);
}
let capacity = inner.grant.budget.max_bitrate_bps.saturating_add(7) / 8;
let elapsed_ms = inner
.publish_tokens_refilled_at
.elapsed()
.as_millis()
.min(u64::MAX as u128) as u64;
if elapsed_ms > 0 {
let refill =
(inner.grant.budget.max_bitrate_bps as u128).saturating_mul(elapsed_ms as u128) / 8_000;
inner.publish_tokens_bytes = inner
.publish_tokens_bytes
.saturating_add(refill.min(u64::MAX as u128) as u64)
.min(capacity);
inner.publish_tokens_refilled_at = web_time::Instant::now();
}
if bytes > inner.publish_tokens_bytes {
return Err(BroadcastError::BudgetExceeded);
}
Ok(())
}
fn push_control(
inner: &mut BroadcastSessionInner,
action: BroadcastAdapterAction,
) -> Result<(), BroadcastError> {
if inner.actions.len() >= MAX_BROADCAST_ACTIONS {
if let Some(position) = inner
.actions
.iter()
.position(|queued| {
matches!(
queued,
BroadcastAdapterAction::Publish {
control: false,
discardable: true,
..
}
)
})
.or_else(|| {
inner.actions.iter().position(|queued| {
matches!(
queued,
BroadcastAdapterAction::Publish { control: false, .. }
)
})
})
{
let removed = inner
.actions
.remove(position)
.expect("queued action exists");
inner.stats.queued_media_bytes = inner
.stats
.queued_media_bytes
.saturating_sub(action_media_bytes(&removed));
inner.stats.dropped_media_objects = inner.stats.dropped_media_objects.saturating_add(1);
} else {
return Err(BroadcastError::QueueFull);
}
}
inner.actions.push_back(action);
inner.stats.queued_actions = inner.actions.len() as u64;
Ok(())
}
fn close_inner(inner: &mut BroadcastSessionInner, reason: &str) {
if matches!(inner.state, BroadcastState::Closed | BroadcastState::Failed) {
return;
}
inner.state = BroadcastState::Closed;
inner
.actions
.retain(|action| !matches!(action, BroadcastAdapterAction::Publish { .. }));
inner.stats.queued_media_bytes = 0;
let action = BroadcastAdapterAction::Close {
broadcast_generation: inner.grant.broadcast_generation,
grant_generation: inner.grant.grant_generation,
reason: reason.to_string(),
};
let _ = push_control(inner, action);
}
fn action_media_bytes(action: &BroadcastAdapterAction) -> u64 {
match action {
BroadcastAdapterAction::Publish { object, .. } => object.len() as u64,
_ => 0,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::media::{EncodedMediaSample, MediaCodec, MediaKind};
#[derive(Clone)]
struct TestSigner(SigningKey);
impl BroadcastSigner for TestSigner {
fn verifying_key(&self) -> [u8; 32] {
self.0.verifying_key().to_bytes()
}
fn sign(&self, canonical_media_group: &[u8]) -> Result<Vec<u8>, BroadcastError> {
Ok(self.0.sign(canonical_media_group).to_bytes().to_vec())
}
}
#[derive(Clone)]
struct InvalidSignatureSigner {
verifying_key: [u8; 32],
}
impl BroadcastSigner for InvalidSignatureSigner {
fn verifying_key(&self) -> [u8; 32] {
self.verifying_key
}
fn sign(&self, _canonical_media_group: &[u8]) -> Result<Vec<u8>, BroadcastError> {
Ok(vec![0_u8; 64])
}
}
fn claims(role: BroadcastRole, source_key: VerifyingKey) -> BroadcastGrantClaims {
BroadcastGrantClaims {
protocol: BROADCAST_PROTOCOL.to_string(),
grant_id: format!("grant-{role:?}"),
broadcast_id: "broadcast-1".to_string(),
broadcast_generation: 2,
grant_generation: 3,
role,
source_slots: if role.can_publish() {
vec!["camera-a".to_string()]
} else {
vec![]
},
source_authorizations: if role.can_subscribe() {
vec![BroadcastSourceAuthorization {
source_slot: "camera-a".to_string(),
grant_generation: 3,
verifying_key: source_key.to_bytes(),
}]
} else {
vec![]
},
publication_verifying_key: role.can_publish().then_some(source_key.to_bytes()),
binding_key: SigningKey::from_bytes(&[7_u8; 32])
.verifying_key()
.to_bytes(),
not_before_ms: 100,
budget: BroadcastBudget {
max_publishers: 4,
max_subscribers: 1_000,
max_bitrate_bps: 4_000_000,
max_egress_bytes: 32 * 1024 * 1024,
expires_at_ms: 10_000,
},
relay_allocation_id: "allocation-1".to_string(),
}
}
fn verified(claims: BroadcastGrantClaims) -> VerifiedBroadcastGrant {
let issuer = SigningKey::from_bytes(&[5_u8; 32]);
let binding = SigningKey::from_bytes(&[7_u8; 32]);
let token = sign_broadcast_grant(&claims, &issuer).unwrap();
let broadcasts = Broadcasts::default();
let challenge = broadcasts
.prepare_grant_verification(&token, &issuer.verifying_key(), 200)
.unwrap();
let proof = binding.sign(&challenge.signing_bytes);
broadcasts
.complete_grant_verification(
&token,
&issuer.verifying_key(),
&challenge.handle,
&proof.to_bytes(),
200,
)
.unwrap()
}
fn publisher_session(role: BroadcastRole, source: &SigningKey) -> BroadcastSession {
Broadcasts::default()
.open_with_signer(
verified(claims(role, source.verifying_key())),
Arc::new(TestSigner(source.clone())),
200,
)
.unwrap()
}
fn assert_json_has_no_private_key_material(value: &serde_json::Value) {
match value {
serde_json::Value::Object(fields) => {
for (name, value) in fields {
let normalized = name.to_ascii_lowercase();
assert!(!normalized.contains("private"), "private field: {name}");
assert!(!normalized.contains("secret"), "secret field: {name}");
assert!(!normalized.contains("seed"), "seed field: {name}");
assert_json_has_no_private_key_material(value);
}
}
serde_json::Value::Array(values) => {
for value in values {
assert_json_has_no_private_key_material(value);
}
}
_ => {}
}
}
#[test]
fn copied_grant_and_replayed_binding_proof_are_rejected_and_debug_is_redacted() {
let issuer = SigningKey::from_bytes(&[5_u8; 32]);
let source = SigningKey::from_bytes(&[9_u8; 32]);
let claims = claims(BroadcastRole::Publisher, source.verifying_key());
let token = sign_broadcast_grant(&claims, &issuer).unwrap();
let payload = token.split('.').nth(1).unwrap();
let decoded: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(payload).expect("grant payload"))
.expect("grant JSON");
assert_json_has_no_private_key_material(&decoded);
assert!(decoded.get("publicationSigningSeed").is_none());
assert!(decoded.get("publicationPrivateKey").is_none());
assert!(decoded.get("publicationVerifyingKey").is_some());
assert_eq!(decoded["relayAllocationId"], "allocation-1");
let broadcasts = Broadcasts::default();
let challenge = broadcasts
.prepare_grant_verification(&token, &issuer.verifying_key(), 200)
.unwrap();
let attacker = SigningKey::from_bytes(&[8_u8; 32]).sign(&challenge.signing_bytes);
assert!(matches!(
broadcasts.complete_grant_verification(
&token,
&issuer.verifying_key(),
&challenge.handle,
&attacker.to_bytes(),
200,
),
Err(BroadcastError::InvalidBindingProof)
));
let correct_but_late = SigningKey::from_bytes(&[7_u8; 32]).sign(&challenge.signing_bytes);
assert!(matches!(
Broadcasts::default().complete_grant_verification(
&token,
&issuer.verifying_key(),
&challenge.handle,
&correct_but_late.to_bytes(),
200,
),
Err(BroadcastError::InvalidBindingChallenge)
));
assert!(matches!(
broadcasts.complete_grant_verification(
&token,
&issuer.verifying_key(),
&challenge.handle,
&correct_but_late.to_bytes(),
200,
),
Err(BroadcastError::InvalidBindingChallenge)
));
let debug = format!("{:?}", claims);
assert!(!debug.contains("publication_signing_seed"));
assert!(!debug.contains("binding_key"));
let fresh = broadcasts
.prepare_grant_verification(&token, &issuer.verifying_key(), 200)
.unwrap();
let legitimate_proof = SigningKey::from_bytes(&[7_u8; 32]).sign(&fresh.signing_bytes);
let copied = broadcasts
.complete_grant_verification(
&token,
&issuer.verifying_key(),
&fresh.handle,
&legitimate_proof.to_bytes(),
200,
)
.unwrap();
assert!(matches!(
broadcasts.complete_grant_verification(
&token,
&issuer.verifying_key(),
&fresh.handle,
&legitimate_proof.to_bytes(),
200,
),
Err(BroadcastError::InvalidBindingChallenge)
));
let attacker_signer = Arc::new(TestSigner(SigningKey::from_bytes(&[8_u8; 32])));
assert!(matches!(
Broadcasts::default().open_with_signer(copied, attacker_signer, 200),
Err(BroadcastError::InvalidSigner)
));
}
#[test]
fn opaque_runtime_signer_opens_a_bound_publisher_and_core_verifies_its_output() {
let signer = BroadcastPublisherSigner::generate().unwrap();
let grant = verified(claims(
BroadcastRole::Publisher,
VerifyingKey::from_bytes(&signer.verifying_key()).unwrap(),
));
let session = Broadcasts::default()
.open_publisher(grant, &signer, 200)
.unwrap();
session.take_actions(8);
session
.begin_publication(
"camera-a",
MediaPublicationConfig {
publication_id: PublicationId([4_u8; 16]),
media_generation: 0,
kind: MediaKind::Audio,
codec: MediaCodec::Opus,
clock_rate: 48_000,
coded_width: None,
coded_height: None,
channels: Some(2),
},
)
.unwrap();
session
.publish(
PublicationId([4_u8; 16]),
EncodedMediaSample {
timestamp_us: 1,
duration_us: 20_000,
keyframe: true,
discardable: false,
payload: vec![1, 2, 3],
},
)
.unwrap();
assert!(session
.take_actions(8)
.into_iter()
.any(|action| { matches!(action, BroadcastAdapterAction::Publish { .. }) }));
let invalid = Arc::new(InvalidSignatureSigner {
verifying_key: signer.verifying_key(),
});
let invalid_session = Broadcasts::default()
.open_with_signer(
verified(claims(
BroadcastRole::Publisher,
VerifyingKey::from_bytes(&signer.verifying_key()).unwrap(),
)),
invalid,
200,
)
.unwrap();
invalid_session.take_actions(8);
assert_eq!(
invalid_session.begin_publication(
"camera-a",
MediaPublicationConfig {
publication_id: PublicationId([5_u8; 16]),
media_generation: 0,
kind: MediaKind::Audio,
codec: MediaCodec::Opus,
clock_rate: 48_000,
coded_width: None,
coded_height: None,
channels: Some(2),
},
),
Err(BroadcastError::InvalidSigner)
);
}
#[test]
fn cohost_publication_is_signed_and_viewer_rejects_replay() {
let source = SigningKey::from_bytes(&[9_u8; 32]);
let publisher = publisher_session(BroadcastRole::Publisher, &source);
publisher.take_actions(8);
let config = MediaPublicationConfig {
publication_id: PublicationId([1_u8; 16]),
media_generation: 0,
kind: MediaKind::Video,
codec: MediaCodec::Vp9,
clock_rate: 90_000,
coded_width: Some(1280),
coded_height: Some(720),
channels: None,
};
publisher.begin_publication("camera-a", config).unwrap();
publisher
.publish(
PublicationId([1_u8; 16]),
EncodedMediaSample {
timestamp_us: 1,
duration_us: 33_333,
keyframe: true,
discardable: false,
payload: vec![1, 2, 3],
},
)
.unwrap();
let actions = publisher.take_actions(8);
let objects: Vec<_> = actions
.into_iter()
.filter_map(|action| match action {
BroadcastAdapterAction::Publish { object, .. } => Some(object),
_ => None,
})
.collect();
let viewer = Broadcasts::default()
.open(
verified(claims(BroadcastRole::Subscriber, source.verifying_key())),
200,
)
.unwrap();
assert!(matches!(
viewer.accept_media_for_source("camera-b", &objects[0]),
Err(BroadcastError::UnauthorizedSource)
));
assert_eq!(viewer.accept_object(&objects[0]).unwrap(), None);
assert_eq!(
viewer.accept_object(&objects[1]).unwrap().unwrap().sequence,
0
);
assert!(matches!(
viewer.accept_object(&objects[1]),
Err(BroadcastError::Media(_))
));
publisher
.retire_publication(PublicationId([1_u8; 16]))
.unwrap();
let stop = publisher
.take_actions(8)
.into_iter()
.find_map(|action| match action {
BroadcastAdapterAction::Publish { object, .. } => Some(object),
_ => None,
})
.expect("signed stop object");
let stopped = viewer
.accept_media_for_source("camera-a", &stop)
.unwrap()
.expect("accepted stop");
assert_eq!(stopped.control_type.as_deref(), Some("stop"));
}
#[test]
fn stale_observation_is_fenced_and_one_retry_is_owned_by_rust() {
let source = SigningKey::from_bytes(&[9_u8; 32]);
let session = publisher_session(BroadcastRole::Publisher, &source);
assert_eq!(
session.observe(BroadcastAdapterObservation::Opened {
broadcast_generation: 1,
grant_generation: 3,
}),
Err(BroadcastError::StaleGeneration)
);
session
.observe(BroadcastAdapterObservation::Failed {
broadcast_generation: 2,
grant_generation: 3,
retryable: true,
code: "network-change".to_string(),
})
.unwrap();
assert_eq!(session.state(), BroadcastState::Retrying);
assert_eq!(session.stats().retries, 1);
session
.observe(BroadcastAdapterObservation::Failed {
broadcast_generation: 2,
grant_generation: 3,
retryable: true,
code: "network-change".to_string(),
})
.unwrap();
assert_eq!(session.state(), BroadcastState::Failed);
}
#[test]
fn grant_revocation_closes_and_prevents_further_publication() {
let source = SigningKey::from_bytes(&[9_u8; 32]);
let session = publisher_session(BroadcastRole::Publisher, &source);
session.revoke(3).unwrap();
assert_eq!(session.state(), BroadcastState::Closed);
assert_eq!(
session.subscribe(vec!["camera-a".to_string()]),
Err(BroadcastError::Revoked)
);
}
#[test]
fn role_source_and_egress_budgets_fail_closed() {
let source = SigningKey::from_bytes(&[9_u8; 32]);
let viewer = Broadcasts::default()
.open(
verified(claims(BroadcastRole::Subscriber, source.verifying_key())),
200,
)
.unwrap();
let config = MediaPublicationConfig {
publication_id: PublicationId([2_u8; 16]),
media_generation: 0,
kind: MediaKind::Audio,
codec: MediaCodec::Opus,
clock_rate: 48_000,
coded_width: None,
coded_height: None,
channels: Some(2),
};
assert!(matches!(
viewer.begin_publication("camera-a", config.clone()),
Err(BroadcastError::UnauthorizedRole)
));
let publisher = publisher_session(BroadcastRole::Publisher, &source);
assert!(matches!(
publisher.begin_publication("camera-b", config),
Err(BroadcastError::UnauthorizedSource)
));
let limit = publisher.budget().max_egress_bytes;
assert_eq!(
publisher.observe(BroadcastAdapterObservation::Usage {
broadcast_generation: 2,
grant_generation: 3,
delivered_bytes: limit + 1,
}),
Err(BroadcastError::BudgetExceeded)
);
assert_eq!(publisher.state(), BroadcastState::Closed);
}
#[test]
fn publisher_cumulative_bytes_are_bounded_before_enqueue() {
let source = SigningKey::from_bytes(&[9_u8; 32]);
let publisher = publisher_session(BroadcastRole::Publisher, &source);
publisher.take_actions(8);
let publication_id = PublicationId([8_u8; 16]);
publisher
.begin_publication(
"camera-a",
MediaPublicationConfig {
publication_id,
media_generation: 0,
kind: MediaKind::Audio,
codec: MediaCodec::Opus,
clock_rate: 48_000,
coded_width: None,
coded_height: None,
channels: Some(2),
},
)
.unwrap();
publisher.take_actions(8);
{
let mut inner = publisher.inner.lock().unwrap();
inner.grant.budget.max_egress_bytes = inner.stats.published_bytes + 1;
}
assert_eq!(
publisher.publish(
publication_id,
EncodedMediaSample {
timestamp_us: 1,
duration_us: 20_000,
keyframe: true,
discardable: false,
payload: vec![1, 2, 3],
},
),
Err(BroadcastError::BudgetExceeded)
);
assert!(publisher.take_actions(8).is_empty());
}
#[test]
fn publisher_bitrate_bucket_is_bounded_before_enqueue() {
let source = SigningKey::from_bytes(&[9_u8; 32]);
let publisher = publisher_session(BroadcastRole::Publisher, &source);
publisher.take_actions(8);
let publication_id = PublicationId([9_u8; 16]);
publisher
.begin_publication(
"camera-a",
MediaPublicationConfig {
publication_id,
media_generation: 0,
kind: MediaKind::Audio,
codec: MediaCodec::Opus,
clock_rate: 48_000,
coded_width: None,
coded_height: None,
channels: Some(2),
},
)
.unwrap();
publisher.take_actions(8);
{
let mut inner = publisher.inner.lock().unwrap();
inner.grant.budget.max_bitrate_bps = 1;
inner.publish_tokens_bytes = 0;
inner.publish_tokens_refilled_at = web_time::Instant::now();
}
assert_eq!(
publisher.publish(
publication_id,
EncodedMediaSample {
timestamp_us: 1,
duration_us: 20_000,
keyframe: true,
discardable: false,
payload: vec![1, 2, 3],
},
),
Err(BroadcastError::BudgetExceeded)
);
assert!(publisher.take_actions(8).is_empty());
}
#[test]
fn managed_allocation_reserves_before_provider_work_and_fences_renewal() {
let allocation = ManagedBroadcastAllocation::create(
"allocation-public-label".to_string(),
1_000,
BroadcastAdmissionLimits {
max_publishers: 2,
max_subscribers: 2,
max_total_bitrate_bps: 3_000,
max_duration_ms: 10_000,
max_egress_bytes: 10_000,
},
)
.unwrap();
assert_eq!(allocation.relay_allocation_id(), "allocation-public-label");
let publisher = allocation
.reserve(
"host-a".to_string(),
BroadcastAdmissionRequest {
role: BroadcastRole::Publisher,
bitrate_bps: 2_000,
duration_ms: 5_000,
reserved_egress_bytes: 4_000,
},
1_100,
)
.unwrap();
assert_eq!(publisher.generation, 1);
assert_eq!(
allocation.reserve(
"host-b".to_string(),
BroadcastAdmissionRequest {
role: BroadcastRole::Publisher,
bitrate_bps: 2_000,
duration_ms: 5_000,
reserved_egress_bytes: 4_000,
},
1_100,
),
Err(BroadcastError::BudgetExceeded)
);
let renewed = allocation
.renew(
"host-a".to_string(),
BroadcastAdmissionRequest {
role: BroadcastRole::Publisher,
bitrate_bps: 1_500,
duration_ms: 5_000,
reserved_egress_bytes: 3_000,
},
1_200,
)
.unwrap();
assert!(renewed.generation > publisher.generation);
assert!(allocation.revoke("host-a"));
assert_eq!(
allocation.renew(
"host-a".to_string(),
BroadcastAdmissionRequest {
role: BroadcastRole::Publisher,
bitrate_bps: 1_000,
duration_ms: 1_000,
reserved_egress_bytes: 1_000,
},
1_300,
),
Err(BroadcastError::Revoked)
);
}
#[test]
fn managed_allocation_rejects_expiry_subscriber_and_egress_overcommit() {
let allocation = ManagedBroadcastAllocation::create(
"allocation-public-label".to_string(),
100,
BroadcastAdmissionLimits {
max_publishers: 1,
max_subscribers: 1,
max_total_bitrate_bps: 2_000,
max_duration_ms: 1_000,
max_egress_bytes: 2_000,
},
)
.unwrap();
let subscriber = BroadcastAdmissionRequest {
role: BroadcastRole::Subscriber,
bitrate_bps: 1_000,
duration_ms: 500,
reserved_egress_bytes: 1_500,
};
allocation
.reserve("viewer-a".to_string(), subscriber, 200)
.unwrap();
assert_eq!(
allocation.reserve("viewer-b".to_string(), subscriber, 200),
Err(BroadcastError::BudgetExceeded)
);
allocation
.reserve(
"viewer-b".to_string(),
BroadcastAdmissionRequest {
duration_ms: 300,
..subscriber
},
701,
)
.expect("expired viewer reservation releases its hard capacity");
assert_eq!(
allocation.renew(
"viewer-a".to_string(),
BroadcastAdmissionRequest {
duration_ms: 300,
..subscriber
},
701,
),
Err(BroadcastError::Revoked)
);
assert_eq!(
allocation.reserve(
"late".to_string(),
BroadcastAdmissionRequest {
duration_ms: 200,
..subscriber
},
1_050,
),
Err(BroadcastError::BudgetExceeded)
);
allocation.close();
assert_eq!(
allocation.reserve("closed".to_string(), subscriber, 200),
Err(BroadcastError::Closed)
);
}
}