use std::fmt;
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Deserializer, Serialize};
use tokio::sync::mpsc;
use crate::capability::CodecInfo;
use crate::error::{Result, RvoipError};
use crate::stream::MediaFrame;
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum BroadcastTransport {
UctpQuic,
Moqt,
}
#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
pub struct BroadcastDescriptor {
pub transport: BroadcastTransport,
pub namespace: String,
pub audio_track: String,
pub catalog_track: Option<String>,
pub protocol_version: String,
}
impl fmt::Debug for BroadcastDescriptor {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("BroadcastDescriptor")
.field("transport", &self.transport)
.field("namespace_bytes", &self.namespace.len())
.field("audio_track_bytes", &self.audio_track.len())
.field("catalog_track_present", &self.catalog_track.is_some())
.field(
"catalog_track_bytes",
&self.catalog_track.as_ref().map_or(0, String::len),
)
.field("protocol_version_bytes", &self.protocol_version.len())
.finish()
}
}
pub const MAX_BROADCAST_EVENT_JSON_INTEGER: u64 = (1_u64 << 53) - 1;
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum BroadcastSanitizedEventKind {
CallConnecting,
CallConnected,
CallHeld,
CallResumed,
TransferStarted,
TransferCompleted,
TransferFailed,
CallEnding,
CallEnded,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastSanitizedEvent {
kind: BroadcastSanitizedEventKind,
occurred_at_unix_millis: u64,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct BroadcastSanitizedEventWire {
kind: BroadcastSanitizedEventKind,
occurred_at_unix_millis: u64,
}
impl<'de> Deserialize<'de> for BroadcastSanitizedEvent {
fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let wire = BroadcastSanitizedEventWire::deserialize(deserializer)?;
Self::at_unix_millis(wire.kind, wire.occurred_at_unix_millis)
.map_err(serde::de::Error::custom)
}
}
impl BroadcastSanitizedEvent {
pub fn at_unix_millis(
kind: BroadcastSanitizedEventKind,
occurred_at_unix_millis: u64,
) -> std::result::Result<Self, BroadcastSanitizedEventError> {
if occurred_at_unix_millis > MAX_BROADCAST_EVENT_JSON_INTEGER {
return Err(BroadcastSanitizedEventError::TimestampOutOfRange {
maximum: MAX_BROADCAST_EVENT_JSON_INTEGER,
actual: occurred_at_unix_millis,
});
}
Ok(Self {
kind,
occurred_at_unix_millis,
})
}
pub fn now(
kind: BroadcastSanitizedEventKind,
) -> std::result::Result<Self, BroadcastSanitizedEventError> {
let occurred_at_unix_millis = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
.try_into()
.unwrap_or(u64::MAX);
Self::at_unix_millis(kind, occurred_at_unix_millis)
}
pub const fn kind(&self) -> BroadcastSanitizedEventKind {
self.kind
}
pub const fn occurred_at_unix_millis(&self) -> u64 {
self.occurred_at_unix_millis
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, thiserror::Error)]
#[non_exhaustive]
pub enum BroadcastSanitizedEventError {
#[error("sanitized broadcast event timestamp {actual} exceeds JSON-safe maximum {maximum}")]
TimestampOutOfRange { maximum: u64, actual: u64 },
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct BroadcastSanitizedEventCapability {
pub queue_capacity: u32,
pub history_capacity: u32,
}
#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "kebab-case")]
pub enum BroadcastResource {
Uctp {
session_id: String,
stream_id: String,
},
Moqt {
namespace: String,
audio_track: String,
catalog_track: Option<String>,
events_track: Option<String>,
},
}
impl fmt::Debug for BroadcastResource {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Uctp {
session_id,
stream_id,
} => formatter
.debug_struct("Uctp")
.field("session_id_bytes", &session_id.len())
.field("stream_id_bytes", &stream_id.len())
.finish(),
Self::Moqt {
namespace,
audio_track,
catalog_track,
events_track,
} => formatter
.debug_struct("Moqt")
.field("namespace_bytes", &namespace.len())
.field("audio_track_bytes", &audio_track.len())
.field("catalog_track_present", &catalog_track.is_some())
.field("events_track_present", &events_track.is_some())
.finish(),
}
}
}
impl BroadcastResource {
pub fn transport(&self) -> BroadcastTransport {
match self {
Self::Uctp { .. } => BroadcastTransport::UctpQuic,
Self::Moqt { .. } => BroadcastTransport::Moqt,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum BroadcastRelayRole {
Origin,
Relay,
Edge,
}
#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
pub struct BroadcastRelayHop {
pub role: BroadcastRelayRole,
pub uri: String,
}
impl fmt::Debug for BroadcastRelayHop {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("BroadcastRelayHop")
.field("role", &self.role)
.field("uri_present", &!self.uri.is_empty())
.field("uri_bytes", &self.uri.len())
.finish()
}
}
#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
pub struct BroadcastEndpoint {
pub uri: Option<String>,
pub resource: BroadcastResource,
pub relay_path: Vec<BroadcastRelayHop>,
}
impl fmt::Debug for BroadcastEndpoint {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("BroadcastEndpoint")
.field("uri_present", &self.uri.is_some())
.field("uri_bytes", &self.uri.as_ref().map_or(0, String::len))
.field("resource", &self.resource)
.field("relay_hop_count", &self.relay_path.len())
.finish()
}
}
impl BroadcastEndpoint {
pub fn transport(&self) -> BroadcastTransport {
self.resource.transport()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum BroadcastProtocolFamily {
Uctp,
Moqt,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum BroadcastSubstrate {
RawQuic,
WebTransport,
WebSocket,
}
#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
pub struct BroadcastProtocolDescriptor {
pub family: BroadcastProtocolFamily,
pub substrate: Option<BroadcastSubstrate>,
pub transport_version: String,
pub media_format_version: Option<String>,
pub object_format_version: Option<String>,
pub media_profile: Option<String>,
}
impl fmt::Debug for BroadcastProtocolDescriptor {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("BroadcastProtocolDescriptor")
.field("family", &self.family)
.field("substrate", &self.substrate)
.field("transport_version_bytes", &self.transport_version.len())
.field(
"media_format_version_present",
&self.media_format_version.is_some(),
)
.field(
"object_format_version_present",
&self.object_format_version.is_some(),
)
.field("media_profile_present", &self.media_profile.is_some())
.finish()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum BroadcastLifecycleState {
Starting,
Ready,
Degraded,
Reconnecting,
Draining,
Closed,
Failed,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct BroadcastLifecycleDescriptor {
pub state: BroadcastLifecycleState,
pub since: Option<DateTime<Utc>>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum BroadcastHealthStatus {
Healthy,
Degraded,
Unhealthy,
Closed,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum BroadcastHealthIssue {
TransportUnavailable,
RelayUnavailable,
AuthenticationUnavailable,
VersionMismatch,
CapacityExhausted,
MediaStalled,
Reconnecting,
Draining,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct BroadcastHealthDescriptor {
pub status: BroadcastHealthStatus,
pub issues: Vec<BroadcastHealthIssue>,
pub active_subscribers: Option<u32>,
pub subscriber_capacity: Option<u32>,
pub checked_at: DateTime<Utc>,
}
impl BroadcastHealthDescriptor {
pub fn healthy() -> Self {
Self {
status: BroadcastHealthStatus::Healthy,
issues: Vec::new(),
active_subscribers: None,
subscriber_capacity: None,
checked_at: Utc::now(),
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum BroadcastDrainReason {
OperatorRequest,
Shutdown,
Reconfigure,
Unhealthy,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct BroadcastDrainRequest {
pub reason: BroadcastDrainReason,
pub deadline: DateTime<Utc>,
}
impl BroadcastDrainRequest {
pub fn immediate() -> Self {
Self {
reason: BroadcastDrainReason::OperatorRequest,
deadline: Utc::now(),
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
#[non_exhaustive]
pub enum BroadcastDrainState {
Draining,
Drained,
DeadlineExceeded,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct BroadcastDrainDescriptor {
pub state: BroadcastDrainState,
pub reason: BroadcastDrainReason,
pub started_at: DateTime<Utc>,
pub deadline: DateTime<Utc>,
pub completed_at: Option<DateTime<Utc>>,
pub remaining_subscribers: u32,
}
impl BroadcastDescriptor {
pub fn endpoint(&self) -> BroadcastEndpoint {
let resource = match self.transport {
BroadcastTransport::UctpQuic => BroadcastResource::Uctp {
session_id: self.namespace.clone(),
stream_id: self.audio_track.clone(),
},
BroadcastTransport::Moqt => BroadcastResource::Moqt {
namespace: self.namespace.clone(),
audio_track: self.audio_track.clone(),
catalog_track: self.catalog_track.clone(),
events_track: None,
},
};
BroadcastEndpoint {
uri: None,
resource,
relay_path: Vec::new(),
}
}
pub fn protocol(&self) -> BroadcastProtocolDescriptor {
BroadcastProtocolDescriptor {
family: match self.transport {
BroadcastTransport::UctpQuic => BroadcastProtocolFamily::Uctp,
BroadcastTransport::Moqt => BroadcastProtocolFamily::Moqt,
},
substrate: None,
transport_version: self.protocol_version.clone(),
media_format_version: None,
object_format_version: None,
media_profile: None,
}
}
}
#[async_trait]
pub trait BroadcastPublisher: Send + Sync {
fn descriptor(&self) -> BroadcastDescriptor;
fn codec(&self) -> CodecInfo;
fn frames_out(&self) -> mpsc::Sender<MediaFrame>;
fn sanitized_event_capability(&self) -> Option<BroadcastSanitizedEventCapability> {
None
}
fn try_publish_sanitized_event(&self, _event: BroadcastSanitizedEvent) -> Result<()> {
Err(RvoipError::NotImplemented(
"sanitized broadcast event publication",
))
}
fn endpoint(&self) -> BroadcastEndpoint {
self.descriptor().endpoint()
}
fn protocol(&self) -> BroadcastProtocolDescriptor {
self.descriptor().protocol()
}
fn lifecycle(&self) -> BroadcastLifecycleDescriptor {
BroadcastLifecycleDescriptor {
state: BroadcastLifecycleState::Ready,
since: None,
}
}
fn health(&self) -> BroadcastHealthDescriptor {
BroadcastHealthDescriptor::healthy()
}
async fn drain(
self: Arc<Self>,
request: BroadcastDrainRequest,
) -> Result<BroadcastDrainDescriptor> {
let started_at = Utc::now();
let missed_deadline = started_at > request.deadline;
self.close().await?;
Ok(BroadcastDrainDescriptor {
state: if missed_deadline {
BroadcastDrainState::DeadlineExceeded
} else {
BroadcastDrainState::Drained
},
reason: request.reason,
started_at,
deadline: request.deadline,
completed_at: Some(Utc::now()),
remaining_subscribers: 0,
})
}
async fn close(self: Arc<Self>) -> Result<()>;
}
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicBool, Ordering};
use super::*;
struct LegacyPublisher {
closed: AtomicBool,
frame_tx: mpsc::Sender<MediaFrame>,
}
#[async_trait]
impl BroadcastPublisher for LegacyPublisher {
fn descriptor(&self) -> BroadcastDescriptor {
BroadcastDescriptor {
transport: BroadcastTransport::UctpQuic,
namespace: "session-1".into(),
audio_track: "stream-2".into(),
catalog_track: None,
protocol_version: "uctp/0.2; rtp-datagram/1".into(),
}
}
fn codec(&self) -> CodecInfo {
CodecInfo::from_name_with_defaults("opus")
}
fn frames_out(&self) -> mpsc::Sender<MediaFrame> {
self.frame_tx.clone()
}
async fn close(self: Arc<Self>) -> Result<()> {
self.closed.store(true, Ordering::Release);
Ok(())
}
}
#[tokio::test]
async fn legacy_implementor_gets_typed_defaults_and_object_safe_drain() {
let (frame_tx, _) = mpsc::channel(1);
let publisher: Arc<dyn BroadcastPublisher> = Arc::new(LegacyPublisher {
closed: AtomicBool::new(false),
frame_tx,
});
assert_eq!(
publisher.endpoint().resource,
BroadcastResource::Uctp {
session_id: "session-1".into(),
stream_id: "stream-2".into(),
}
);
assert_eq!(publisher.protocol().family, BroadcastProtocolFamily::Uctp);
assert_eq!(publisher.lifecycle().state, BroadcastLifecycleState::Ready);
assert_eq!(publisher.health().status, BroadcastHealthStatus::Healthy);
assert_eq!(publisher.sanitized_event_capability(), None);
assert!(matches!(
publisher.try_publish_sanitized_event(
BroadcastSanitizedEvent::at_unix_millis(
BroadcastSanitizedEventKind::CallConnected,
1_000,
)
.unwrap(),
),
Err(RvoipError::NotImplemented(_))
));
let drained = Arc::clone(&publisher)
.drain(BroadcastDrainRequest {
reason: BroadcastDrainReason::Shutdown,
deadline: Utc::now() + chrono::Duration::seconds(1),
})
.await
.unwrap();
assert_eq!(drained.state, BroadcastDrainState::Drained);
}
#[test]
fn moqt_legacy_descriptor_maps_to_typed_tracks() {
let endpoint = BroadcastDescriptor {
transport: BroadcastTransport::Moqt,
namespace: "tenant/broadcast".into(),
audio_track: "audio/main".into(),
catalog_track: Some("catalog".into()),
protocol_version: "draft-19".into(),
}
.endpoint();
assert_eq!(endpoint.transport(), BroadcastTransport::Moqt);
assert!(matches!(
endpoint.resource,
BroadcastResource::Moqt {
events_track: None,
..
}
));
}
#[test]
fn sanitized_event_model_is_fixed_and_json_safe() {
let event = BroadcastSanitizedEvent::at_unix_millis(
BroadcastSanitizedEventKind::CallConnected,
MAX_BROADCAST_EVENT_JSON_INTEGER,
)
.unwrap();
assert_eq!(
event.occurred_at_unix_millis(),
MAX_BROADCAST_EVENT_JSON_INTEGER
);
assert!(matches!(
BroadcastSanitizedEvent::at_unix_millis(
BroadcastSanitizedEventKind::CallConnected,
MAX_BROADCAST_EVENT_JSON_INTEGER + 1,
),
Err(BroadcastSanitizedEventError::TimestampOutOfRange { .. })
));
assert_eq!(
serde_json::to_value(event).unwrap(),
serde_json::json!({
"kind": "call-connected",
"occurredAtUnixMillis": MAX_BROADCAST_EVENT_JSON_INTEGER,
})
);
assert!(
serde_json::from_value::<BroadcastSanitizedEvent>(serde_json::json!({
"kind": "call-connected",
"occurredAtUnixMillis": MAX_BROADCAST_EVENT_JSON_INTEGER + 1,
}))
.is_err()
);
assert!(
serde_json::from_value::<BroadcastSanitizedEvent>(serde_json::json!({
"kind": "call-connected",
"occurredAtUnixMillis": 1_000,
"metadata": "forbidden",
}))
.is_err()
);
}
#[test]
fn broadcast_diagnostics_redact_resource_and_network_identifiers() {
const CANARY: &str = "broadcast-canary\r\nAuthorization: exposed";
let descriptor = BroadcastDescriptor {
transport: BroadcastTransport::Moqt,
namespace: CANARY.into(),
audio_track: CANARY.into(),
catalog_track: Some(CANARY.into()),
protocol_version: CANARY.into(),
};
let endpoint = BroadcastEndpoint {
uri: Some(CANARY.into()),
resource: BroadcastResource::Moqt {
namespace: CANARY.into(),
audio_track: CANARY.into(),
catalog_track: Some(CANARY.into()),
events_track: Some(CANARY.into()),
},
relay_path: vec![BroadcastRelayHop {
role: BroadcastRelayRole::Relay,
uri: CANARY.into(),
}],
};
let protocol = BroadcastProtocolDescriptor {
family: BroadcastProtocolFamily::Moqt,
substrate: Some(BroadcastSubstrate::RawQuic),
transport_version: CANARY.into(),
media_format_version: Some(CANARY.into()),
object_format_version: Some(CANARY.into()),
media_profile: Some(CANARY.into()),
};
for debug in [
format!("{descriptor:?}"),
format!("{endpoint:?}"),
format!("{protocol:?}"),
] {
assert!(!debug.contains(CANARY), "broadcast value leaked: {debug}");
}
}
}