1use std::fmt;
9use std::sync::Arc;
10use std::time::{SystemTime, UNIX_EPOCH};
11
12use async_trait::async_trait;
13use chrono::{DateTime, Utc};
14use serde::{Deserialize, Deserializer, Serialize};
15use tokio::sync::mpsc;
16
17use crate::capability::CodecInfo;
18use crate::error::{Result, RvoipError};
19use crate::stream::MediaFrame;
20
21#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
23#[serde(rename_all = "kebab-case")]
24pub enum BroadcastTransport {
25 UctpQuic,
27 Moqt,
29}
30
31#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
36pub struct BroadcastDescriptor {
37 pub transport: BroadcastTransport,
38 pub namespace: String,
39 pub audio_track: String,
40 pub catalog_track: Option<String>,
41 pub protocol_version: String,
42}
43
44impl fmt::Debug for BroadcastDescriptor {
45 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
46 formatter
47 .debug_struct("BroadcastDescriptor")
48 .field("transport", &self.transport)
49 .field("namespace_bytes", &self.namespace.len())
50 .field("audio_track_bytes", &self.audio_track.len())
51 .field("catalog_track_present", &self.catalog_track.is_some())
52 .field(
53 "catalog_track_bytes",
54 &self.catalog_track.as_ref().map_or(0, String::len),
55 )
56 .field("protocol_version_bytes", &self.protocol_version.len())
57 .finish()
58 }
59}
60
61pub const MAX_BROADCAST_EVENT_JSON_INTEGER: u64 = (1_u64 << 53) - 1;
63
64#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
70#[serde(rename_all = "kebab-case")]
71#[non_exhaustive]
72pub enum BroadcastSanitizedEventKind {
73 CallConnecting,
74 CallConnected,
75 CallHeld,
76 CallResumed,
77 TransferStarted,
78 TransferCompleted,
79 TransferFailed,
80 CallEnding,
81 CallEnded,
82}
83
84#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
86#[serde(rename_all = "camelCase")]
87pub struct BroadcastSanitizedEvent {
88 kind: BroadcastSanitizedEventKind,
89 occurred_at_unix_millis: u64,
90}
91
92#[derive(Deserialize)]
93#[serde(rename_all = "camelCase", deny_unknown_fields)]
94struct BroadcastSanitizedEventWire {
95 kind: BroadcastSanitizedEventKind,
96 occurred_at_unix_millis: u64,
97}
98
99impl<'de> Deserialize<'de> for BroadcastSanitizedEvent {
100 fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
101 where
102 D: Deserializer<'de>,
103 {
104 let wire = BroadcastSanitizedEventWire::deserialize(deserializer)?;
105 Self::at_unix_millis(wire.kind, wire.occurred_at_unix_millis)
106 .map_err(serde::de::Error::custom)
107 }
108}
109
110impl BroadcastSanitizedEvent {
111 pub fn at_unix_millis(
114 kind: BroadcastSanitizedEventKind,
115 occurred_at_unix_millis: u64,
116 ) -> std::result::Result<Self, BroadcastSanitizedEventError> {
117 if occurred_at_unix_millis > MAX_BROADCAST_EVENT_JSON_INTEGER {
118 return Err(BroadcastSanitizedEventError::TimestampOutOfRange {
119 maximum: MAX_BROADCAST_EVENT_JSON_INTEGER,
120 actual: occurred_at_unix_millis,
121 });
122 }
123 Ok(Self {
124 kind,
125 occurred_at_unix_millis,
126 })
127 }
128
129 pub fn now(
130 kind: BroadcastSanitizedEventKind,
131 ) -> std::result::Result<Self, BroadcastSanitizedEventError> {
132 let occurred_at_unix_millis = SystemTime::now()
133 .duration_since(UNIX_EPOCH)
134 .unwrap_or_default()
135 .as_millis()
136 .try_into()
137 .unwrap_or(u64::MAX);
138 Self::at_unix_millis(kind, occurred_at_unix_millis)
139 }
140
141 pub const fn kind(&self) -> BroadcastSanitizedEventKind {
142 self.kind
143 }
144
145 pub const fn occurred_at_unix_millis(&self) -> u64 {
146 self.occurred_at_unix_millis
147 }
148}
149
150#[derive(Clone, Copy, Debug, Eq, PartialEq, thiserror::Error)]
151#[non_exhaustive]
152pub enum BroadcastSanitizedEventError {
153 #[error("sanitized broadcast event timestamp {actual} exceeds JSON-safe maximum {maximum}")]
154 TimestampOutOfRange { maximum: u64, actual: u64 },
155}
156
157#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
159#[serde(rename_all = "camelCase")]
160pub struct BroadcastSanitizedEventCapability {
161 pub queue_capacity: u32,
162 pub history_capacity: u32,
163}
164
165#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
167#[serde(tag = "kind", rename_all = "kebab-case")]
168pub enum BroadcastResource {
169 Uctp {
171 session_id: String,
172 stream_id: String,
173 },
174 Moqt {
176 namespace: String,
177 audio_track: String,
178 catalog_track: Option<String>,
179 events_track: Option<String>,
180 },
181}
182
183impl fmt::Debug for BroadcastResource {
184 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
185 match self {
186 Self::Uctp {
187 session_id,
188 stream_id,
189 } => formatter
190 .debug_struct("Uctp")
191 .field("session_id_bytes", &session_id.len())
192 .field("stream_id_bytes", &stream_id.len())
193 .finish(),
194 Self::Moqt {
195 namespace,
196 audio_track,
197 catalog_track,
198 events_track,
199 } => formatter
200 .debug_struct("Moqt")
201 .field("namespace_bytes", &namespace.len())
202 .field("audio_track_bytes", &audio_track.len())
203 .field("catalog_track_present", &catalog_track.is_some())
204 .field("events_track_present", &events_track.is_some())
205 .finish(),
206 }
207 }
208}
209
210impl BroadcastResource {
211 pub fn transport(&self) -> BroadcastTransport {
213 match self {
214 Self::Uctp { .. } => BroadcastTransport::UctpQuic,
215 Self::Moqt { .. } => BroadcastTransport::Moqt,
216 }
217 }
218}
219
220#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
222#[serde(rename_all = "kebab-case")]
223#[non_exhaustive]
224pub enum BroadcastRelayRole {
225 Origin,
226 Relay,
227 Edge,
228}
229
230#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
235pub struct BroadcastRelayHop {
236 pub role: BroadcastRelayRole,
237 pub uri: String,
238}
239
240impl fmt::Debug for BroadcastRelayHop {
241 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
242 formatter
243 .debug_struct("BroadcastRelayHop")
244 .field("role", &self.role)
245 .field("uri_present", &!self.uri.is_empty())
246 .field("uri_bytes", &self.uri.len())
247 .finish()
248 }
249}
250
251#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
253pub struct BroadcastEndpoint {
254 pub uri: Option<String>,
256 pub resource: BroadcastResource,
257 pub relay_path: Vec<BroadcastRelayHop>,
259}
260
261impl fmt::Debug for BroadcastEndpoint {
262 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
263 formatter
264 .debug_struct("BroadcastEndpoint")
265 .field("uri_present", &self.uri.is_some())
266 .field("uri_bytes", &self.uri.as_ref().map_or(0, String::len))
267 .field("resource", &self.resource)
268 .field("relay_hop_count", &self.relay_path.len())
269 .finish()
270 }
271}
272
273impl BroadcastEndpoint {
274 pub fn transport(&self) -> BroadcastTransport {
276 self.resource.transport()
277 }
278}
279
280#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
282#[serde(rename_all = "kebab-case")]
283#[non_exhaustive]
284pub enum BroadcastProtocolFamily {
285 Uctp,
286 Moqt,
287}
288
289#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
291#[serde(rename_all = "kebab-case")]
292#[non_exhaustive]
293pub enum BroadcastSubstrate {
294 RawQuic,
295 WebTransport,
296 WebSocket,
297}
298
299#[derive(Clone, Eq, PartialEq, Serialize, Deserialize)]
308pub struct BroadcastProtocolDescriptor {
309 pub family: BroadcastProtocolFamily,
310 pub substrate: Option<BroadcastSubstrate>,
311 pub transport_version: String,
312 pub media_format_version: Option<String>,
313 pub object_format_version: Option<String>,
314 pub media_profile: Option<String>,
315}
316
317impl fmt::Debug for BroadcastProtocolDescriptor {
318 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
319 formatter
320 .debug_struct("BroadcastProtocolDescriptor")
321 .field("family", &self.family)
322 .field("substrate", &self.substrate)
323 .field("transport_version_bytes", &self.transport_version.len())
324 .field(
325 "media_format_version_present",
326 &self.media_format_version.is_some(),
327 )
328 .field(
329 "object_format_version_present",
330 &self.object_format_version.is_some(),
331 )
332 .field("media_profile_present", &self.media_profile.is_some())
333 .finish()
334 }
335}
336
337#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
339#[serde(rename_all = "kebab-case")]
340#[non_exhaustive]
341pub enum BroadcastLifecycleState {
342 Starting,
343 Ready,
344 Degraded,
345 Reconnecting,
346 Draining,
347 Closed,
348 Failed,
349}
350
351#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
353pub struct BroadcastLifecycleDescriptor {
354 pub state: BroadcastLifecycleState,
355 pub since: Option<DateTime<Utc>>,
357}
358
359#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
361#[serde(rename_all = "kebab-case")]
362#[non_exhaustive]
363pub enum BroadcastHealthStatus {
364 Healthy,
365 Degraded,
366 Unhealthy,
367 Closed,
368}
369
370#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
373#[serde(rename_all = "kebab-case")]
374#[non_exhaustive]
375pub enum BroadcastHealthIssue {
376 TransportUnavailable,
377 RelayUnavailable,
378 AuthenticationUnavailable,
379 VersionMismatch,
380 CapacityExhausted,
381 MediaStalled,
382 Reconnecting,
383 Draining,
384}
385
386#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
388pub struct BroadcastHealthDescriptor {
389 pub status: BroadcastHealthStatus,
390 pub issues: Vec<BroadcastHealthIssue>,
391 pub active_subscribers: Option<u32>,
392 pub subscriber_capacity: Option<u32>,
393 pub checked_at: DateTime<Utc>,
394}
395
396impl BroadcastHealthDescriptor {
397 pub fn healthy() -> Self {
399 Self {
400 status: BroadcastHealthStatus::Healthy,
401 issues: Vec::new(),
402 active_subscribers: None,
403 subscriber_capacity: None,
404 checked_at: Utc::now(),
405 }
406 }
407}
408
409#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
411#[serde(rename_all = "kebab-case")]
412#[non_exhaustive]
413pub enum BroadcastDrainReason {
414 OperatorRequest,
415 Shutdown,
416 Reconfigure,
417 Unhealthy,
418}
419
420#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
422pub struct BroadcastDrainRequest {
423 pub reason: BroadcastDrainReason,
424 pub deadline: DateTime<Utc>,
425}
426
427impl BroadcastDrainRequest {
428 pub fn immediate() -> Self {
430 Self {
431 reason: BroadcastDrainReason::OperatorRequest,
432 deadline: Utc::now(),
433 }
434 }
435}
436
437#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
439#[serde(rename_all = "kebab-case")]
440#[non_exhaustive]
441pub enum BroadcastDrainState {
442 Draining,
443 Drained,
444 DeadlineExceeded,
445}
446
447#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
449pub struct BroadcastDrainDescriptor {
450 pub state: BroadcastDrainState,
451 pub reason: BroadcastDrainReason,
452 pub started_at: DateTime<Utc>,
453 pub deadline: DateTime<Utc>,
454 pub completed_at: Option<DateTime<Utc>>,
455 pub remaining_subscribers: u32,
456}
457
458impl BroadcastDescriptor {
459 pub fn endpoint(&self) -> BroadcastEndpoint {
461 let resource = match self.transport {
462 BroadcastTransport::UctpQuic => BroadcastResource::Uctp {
463 session_id: self.namespace.clone(),
464 stream_id: self.audio_track.clone(),
465 },
466 BroadcastTransport::Moqt => BroadcastResource::Moqt {
467 namespace: self.namespace.clone(),
468 audio_track: self.audio_track.clone(),
469 catalog_track: self.catalog_track.clone(),
470 events_track: None,
471 },
472 };
473 BroadcastEndpoint {
474 uri: None,
475 resource,
476 relay_path: Vec::new(),
477 }
478 }
479
480 pub fn protocol(&self) -> BroadcastProtocolDescriptor {
484 BroadcastProtocolDescriptor {
485 family: match self.transport {
486 BroadcastTransport::UctpQuic => BroadcastProtocolFamily::Uctp,
487 BroadcastTransport::Moqt => BroadcastProtocolFamily::Moqt,
488 },
489 substrate: None,
490 transport_version: self.protocol_version.clone(),
491 media_format_version: None,
492 object_format_version: None,
493 media_profile: None,
494 }
495 }
496}
497
498#[async_trait]
504pub trait BroadcastPublisher: Send + Sync {
505 fn descriptor(&self) -> BroadcastDescriptor;
506 fn codec(&self) -> CodecInfo;
507 fn frames_out(&self) -> mpsc::Sender<MediaFrame>;
508
509 fn sanitized_event_capability(&self) -> Option<BroadcastSanitizedEventCapability> {
512 None
513 }
514
515 fn try_publish_sanitized_event(&self, _event: BroadcastSanitizedEvent) -> Result<()> {
517 Err(RvoipError::NotImplemented(
518 "sanitized broadcast event publication",
519 ))
520 }
521
522 fn endpoint(&self) -> BroadcastEndpoint {
524 self.descriptor().endpoint()
525 }
526
527 fn protocol(&self) -> BroadcastProtocolDescriptor {
530 self.descriptor().protocol()
531 }
532
533 fn lifecycle(&self) -> BroadcastLifecycleDescriptor {
535 BroadcastLifecycleDescriptor {
536 state: BroadcastLifecycleState::Ready,
537 since: None,
538 }
539 }
540
541 fn health(&self) -> BroadcastHealthDescriptor {
543 BroadcastHealthDescriptor::healthy()
544 }
545
546 async fn drain(
551 self: Arc<Self>,
552 request: BroadcastDrainRequest,
553 ) -> Result<BroadcastDrainDescriptor> {
554 let started_at = Utc::now();
555 let missed_deadline = started_at > request.deadline;
556 self.close().await?;
557 Ok(BroadcastDrainDescriptor {
558 state: if missed_deadline {
559 BroadcastDrainState::DeadlineExceeded
560 } else {
561 BroadcastDrainState::Drained
562 },
563 reason: request.reason,
564 started_at,
565 deadline: request.deadline,
566 completed_at: Some(Utc::now()),
567 remaining_subscribers: 0,
568 })
569 }
570
571 async fn close(self: Arc<Self>) -> Result<()>;
572}
573
574#[cfg(test)]
575mod tests {
576 use std::sync::atomic::{AtomicBool, Ordering};
577
578 use super::*;
579
580 struct LegacyPublisher {
581 closed: AtomicBool,
582 frame_tx: mpsc::Sender<MediaFrame>,
583 }
584
585 #[async_trait]
586 impl BroadcastPublisher for LegacyPublisher {
587 fn descriptor(&self) -> BroadcastDescriptor {
588 BroadcastDescriptor {
589 transport: BroadcastTransport::UctpQuic,
590 namespace: "session-1".into(),
591 audio_track: "stream-2".into(),
592 catalog_track: None,
593 protocol_version: "uctp/0.2; rtp-datagram/1".into(),
594 }
595 }
596
597 fn codec(&self) -> CodecInfo {
598 CodecInfo::from_name_with_defaults("opus")
599 }
600
601 fn frames_out(&self) -> mpsc::Sender<MediaFrame> {
602 self.frame_tx.clone()
603 }
604
605 async fn close(self: Arc<Self>) -> Result<()> {
606 self.closed.store(true, Ordering::Release);
607 Ok(())
608 }
609 }
610
611 #[tokio::test]
612 async fn legacy_implementor_gets_typed_defaults_and_object_safe_drain() {
613 let (frame_tx, _) = mpsc::channel(1);
614 let publisher: Arc<dyn BroadcastPublisher> = Arc::new(LegacyPublisher {
615 closed: AtomicBool::new(false),
616 frame_tx,
617 });
618
619 assert_eq!(
620 publisher.endpoint().resource,
621 BroadcastResource::Uctp {
622 session_id: "session-1".into(),
623 stream_id: "stream-2".into(),
624 }
625 );
626 assert_eq!(publisher.protocol().family, BroadcastProtocolFamily::Uctp);
627 assert_eq!(publisher.lifecycle().state, BroadcastLifecycleState::Ready);
628 assert_eq!(publisher.health().status, BroadcastHealthStatus::Healthy);
629 assert_eq!(publisher.sanitized_event_capability(), None);
630 assert!(matches!(
631 publisher.try_publish_sanitized_event(
632 BroadcastSanitizedEvent::at_unix_millis(
633 BroadcastSanitizedEventKind::CallConnected,
634 1_000,
635 )
636 .unwrap(),
637 ),
638 Err(RvoipError::NotImplemented(_))
639 ));
640
641 let drained = Arc::clone(&publisher)
642 .drain(BroadcastDrainRequest {
643 reason: BroadcastDrainReason::Shutdown,
644 deadline: Utc::now() + chrono::Duration::seconds(1),
645 })
646 .await
647 .unwrap();
648 assert_eq!(drained.state, BroadcastDrainState::Drained);
649 }
650
651 #[test]
652 fn moqt_legacy_descriptor_maps_to_typed_tracks() {
653 let endpoint = BroadcastDescriptor {
654 transport: BroadcastTransport::Moqt,
655 namespace: "tenant/broadcast".into(),
656 audio_track: "audio/main".into(),
657 catalog_track: Some("catalog".into()),
658 protocol_version: "draft-19".into(),
659 }
660 .endpoint();
661
662 assert_eq!(endpoint.transport(), BroadcastTransport::Moqt);
663 assert!(matches!(
664 endpoint.resource,
665 BroadcastResource::Moqt {
666 events_track: None,
667 ..
668 }
669 ));
670 }
671
672 #[test]
673 fn sanitized_event_model_is_fixed_and_json_safe() {
674 let event = BroadcastSanitizedEvent::at_unix_millis(
675 BroadcastSanitizedEventKind::CallConnected,
676 MAX_BROADCAST_EVENT_JSON_INTEGER,
677 )
678 .unwrap();
679 assert_eq!(
680 event.occurred_at_unix_millis(),
681 MAX_BROADCAST_EVENT_JSON_INTEGER
682 );
683 assert!(matches!(
684 BroadcastSanitizedEvent::at_unix_millis(
685 BroadcastSanitizedEventKind::CallConnected,
686 MAX_BROADCAST_EVENT_JSON_INTEGER + 1,
687 ),
688 Err(BroadcastSanitizedEventError::TimestampOutOfRange { .. })
689 ));
690 assert_eq!(
691 serde_json::to_value(event).unwrap(),
692 serde_json::json!({
693 "kind": "call-connected",
694 "occurredAtUnixMillis": MAX_BROADCAST_EVENT_JSON_INTEGER,
695 })
696 );
697 assert!(
698 serde_json::from_value::<BroadcastSanitizedEvent>(serde_json::json!({
699 "kind": "call-connected",
700 "occurredAtUnixMillis": MAX_BROADCAST_EVENT_JSON_INTEGER + 1,
701 }))
702 .is_err()
703 );
704 assert!(
705 serde_json::from_value::<BroadcastSanitizedEvent>(serde_json::json!({
706 "kind": "call-connected",
707 "occurredAtUnixMillis": 1_000,
708 "metadata": "forbidden",
709 }))
710 .is_err()
711 );
712 }
713
714 #[test]
715 fn broadcast_diagnostics_redact_resource_and_network_identifiers() {
716 const CANARY: &str = "broadcast-canary\r\nAuthorization: exposed";
717 let descriptor = BroadcastDescriptor {
718 transport: BroadcastTransport::Moqt,
719 namespace: CANARY.into(),
720 audio_track: CANARY.into(),
721 catalog_track: Some(CANARY.into()),
722 protocol_version: CANARY.into(),
723 };
724 let endpoint = BroadcastEndpoint {
725 uri: Some(CANARY.into()),
726 resource: BroadcastResource::Moqt {
727 namespace: CANARY.into(),
728 audio_track: CANARY.into(),
729 catalog_track: Some(CANARY.into()),
730 events_track: Some(CANARY.into()),
731 },
732 relay_path: vec![BroadcastRelayHop {
733 role: BroadcastRelayRole::Relay,
734 uri: CANARY.into(),
735 }],
736 };
737 let protocol = BroadcastProtocolDescriptor {
738 family: BroadcastProtocolFamily::Moqt,
739 substrate: Some(BroadcastSubstrate::RawQuic),
740 transport_version: CANARY.into(),
741 media_format_version: Some(CANARY.into()),
742 object_format_version: Some(CANARY.into()),
743 media_profile: Some(CANARY.into()),
744 };
745 for debug in [
746 format!("{descriptor:?}"),
747 format!("{endpoint:?}"),
748 format!("{protocol:?}"),
749 ] {
750 assert!(!debug.contains(CANARY), "broadcast value leaked: {debug}");
751 }
752 }
753}