1use std::collections::{HashMap, HashSet, VecDeque};
9use std::fmt;
10use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
11use std::sync::{Arc, Mutex, RwLock as StdRwLock};
12use std::time::{Duration, Instant};
13
14use chrono::{DateTime, Utc};
15use rvoip_media_core::codec::audio::{
16 payload_type::PCM_S16LE, AudioCodec, OpusApplication, OpusCodec, OpusConfig, PcmS16LeCodec,
17};
18use rvoip_media_core::codec::factory::CodecFactory;
19use rvoip_media_core::error::CodecError;
20use rvoip_media_core::processing::format::{ConversionParams, FormatConverter};
21use rvoip_media_core::types::SampleRate;
22use serde::{Deserialize, Serialize};
23use tokio::sync::{mpsc, oneshot, watch, Notify};
24use tokio::task::AbortHandle;
25use tracing::{debug, warn};
26use uuid::Uuid;
27
28use crate::bridge::{codec_to_pt, frame_pump::DEFAULT_TELEPHONE_EVENT_PT};
29use crate::capability::CodecInfo;
30use crate::error::{Result, RvoipError};
31use crate::ids::MediaRouteId;
32use crate::stream::MediaFrame;
33
34const SNAPSHOT_TIMEOUT: Duration = Duration::from_secs(1);
35const SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);
36const SNAPSHOT_PUBLISH_INTERVAL: Duration = Duration::from_secs(1);
37pub const MEDIA_GRAPH_ACTIVITY_OBSERVATION_INTERVAL: Duration = Duration::from_secs(1);
39const RECENT_EVICTION_LIMIT: usize = 64;
40const CONTROL_QUEUE_CAPACITY: usize = 256;
41const SINK_EVENT_QUEUE_CAPACITY: usize = 256;
42
43pub const DEFAULT_MEDIA_GRAPH_MAX_SINKS: usize = 1_024;
47
48#[derive(Clone, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize, Deserialize)]
53pub struct MediaGraphId(String);
54
55impl MediaGraphId {
56 pub fn new() -> Self {
57 Self(format!("graph_{}", Uuid::new_v4().simple()))
58 }
59
60 pub fn from_string(value: impl Into<String>) -> Self {
61 Self(value.into())
62 }
63
64 pub fn as_str(&self) -> &str {
65 &self.0
66 }
67}
68
69impl Default for MediaGraphId {
70 fn default() -> Self {
71 Self::new()
72 }
73}
74
75impl fmt::Display for MediaGraphId {
76 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
77 f.write_str(&self.0)
78 }
79}
80
81impl fmt::Debug for MediaGraphId {
82 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
83 f.write_str("MediaGraphId([redacted])")
84 }
85}
86
87#[derive(Clone, Debug)]
88pub struct MediaGraphPolicy {
89 pub max_sinks: usize,
93 pub sink_queue_frames: usize,
94 pub pre_sink_buffer_frames: usize,
97 pub eviction_window: Duration,
98 pub eviction_drop_ratio: f32,
99 pub minimum_eviction_samples: usize,
100}
101
102impl Default for MediaGraphPolicy {
103 fn default() -> Self {
104 Self {
105 max_sinks: DEFAULT_MEDIA_GRAPH_MAX_SINKS,
106 sink_queue_frames: 10,
107 pre_sink_buffer_frames: 10,
108 eviction_window: Duration::from_secs(10),
109 eviction_drop_ratio: 0.25,
110 minimum_eviction_samples: 50,
111 }
112 }
113}
114
115#[derive(Clone, Debug, Eq, PartialEq)]
122pub struct MediaGraphActivityObservation {
123 pub source_frames: u64,
124 pub observed_at: DateTime<Utc>,
125}
126
127#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
129#[serde(rename_all = "snake_case")]
130pub enum MediaGraphSourceState {
131 Open,
132 Closed,
133 Shutdown,
134 Aborted,
135}
136
137#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
138#[serde(rename_all = "snake_case")]
139pub enum MediaGraphEvictionReason {
140 SlowConsumer,
141}
142
143#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
145#[serde(rename_all = "snake_case")]
146pub enum MediaGraphRouteTerminalReason {
147 OwnerRemoved,
148 TargetClosed,
149 SlowConsumerEvicted,
150 GraphShutdown,
151 SourceClosed,
152 GraphAborted,
153}
154
155#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
160#[serde(rename_all = "snake_case")]
161pub enum MediaGraphRouteState {
162 Pending,
163 Active,
164 Terminal(MediaGraphRouteTerminalReason),
165}
166
167#[derive(Clone)]
172pub struct MediaGraphRouteStatus {
173 route_id: MediaRouteId,
174 state: watch::Receiver<MediaGraphRouteState>,
175}
176
177impl MediaGraphRouteStatus {
178 pub fn id(&self) -> &MediaRouteId {
179 &self.route_id
180 }
181
182 pub fn state(&self) -> MediaGraphRouteState {
183 let state = *self.state.borrow();
184 if matches!(state, MediaGraphRouteState::Terminal(_)) {
185 return state;
186 }
187 let receiver = self.state.clone();
188 if receiver.has_changed().is_err() {
189 MediaGraphRouteState::Terminal(MediaGraphRouteTerminalReason::GraphAborted)
190 } else {
191 state
192 }
193 }
194
195 pub async fn wait_active(&self) -> std::result::Result<(), MediaGraphRouteTerminalReason> {
198 let mut receiver = self.state.clone();
199 loop {
200 match *receiver.borrow_and_update() {
201 MediaGraphRouteState::Active => return Ok(()),
202 MediaGraphRouteState::Terminal(reason) => return Err(reason),
203 MediaGraphRouteState::Pending => {}
204 }
205 if receiver.changed().await.is_err() {
206 return Err(MediaGraphRouteTerminalReason::GraphAborted);
207 }
208 }
209 }
210
211 pub async fn wait_terminal(&self) -> MediaGraphRouteTerminalReason {
212 let mut receiver = self.state.clone();
213 loop {
214 if let MediaGraphRouteState::Terminal(reason) = *receiver.borrow_and_update() {
215 return reason;
216 }
217 if receiver.changed().await.is_err() {
218 return MediaGraphRouteTerminalReason::GraphAborted;
219 }
220 }
221 }
222}
223
224#[must_use = "dropping the managed route immediately removes its sink"]
231pub struct ManagedMediaRoute {
232 status: MediaGraphRouteStatus,
233 commands: mpsc::Sender<Command>,
234 owner_liveness: Arc<RouteOwnerLiveness>,
235 remove_on_drop: bool,
236}
237
238impl ManagedMediaRoute {
239 pub fn id(&self) -> &MediaRouteId {
240 self.status.id()
241 }
242
243 pub fn status(&self) -> MediaGraphRouteStatus {
244 self.status.clone()
245 }
246
247 pub fn state(&self) -> MediaGraphRouteState {
248 self.status.state()
249 }
250
251 pub async fn wait_active(&self) -> std::result::Result<(), MediaGraphRouteTerminalReason> {
252 self.status.wait_active().await
253 }
254
255 pub async fn wait_terminal(&self) -> MediaGraphRouteTerminalReason {
256 self.status.wait_terminal().await
257 }
258
259 pub async fn remove(mut self) -> Result<bool> {
263 let (ack, done) = oneshot::channel();
264 tokio::time::timeout(
265 SNAPSHOT_TIMEOUT,
266 self.commands.send(Command::Remove {
267 route_id: self.status.route_id.clone(),
268 ack: Some(ack),
269 }),
270 )
271 .await
272 .map_err(|_| RvoipError::InvalidState("media graph command queue is full"))?
273 .map_err(|_| RvoipError::InvalidState("media graph is closed"))?;
274 self.remove_on_drop = false;
275 tokio::time::timeout(SNAPSHOT_TIMEOUT, done)
276 .await
277 .map_err(|_| RvoipError::InvalidState("media graph removal timed out"))?
278 .map_err(|_| RvoipError::InvalidState("media graph removal was cancelled"))
279 }
280
281 fn into_unmanaged_route_id(mut self) -> MediaRouteId {
282 self.remove_on_drop = false;
283 self.status.route_id.clone()
284 }
285}
286
287impl Drop for ManagedMediaRoute {
288 fn drop(&mut self) {
289 if self.remove_on_drop {
290 self.owner_liveness.cancel();
294 let _ = self.commands.try_send(Command::Remove {
295 route_id: self.status.route_id.clone(),
296 ack: None,
297 });
298 }
299 }
300}
301
302#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
304pub struct MediaGraphSinkSnapshot {
305 pub route_id: MediaRouteId,
306 pub target_codec: CodecInfo,
307 pub target_payload_type: u8,
308 pub queue_depth: usize,
309 pub queue_capacity: usize,
310 pub offered_frames: u64,
311 pub dropped_frames: u64,
312 pub rolling_samples: usize,
313 pub rolling_drops: usize,
314 pub rolling_drop_ratio: f32,
315}
316
317#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
320pub struct MediaGraphCodecGroupSnapshot {
321 pub target_codec: CodecInfo,
322 pub target_payload_type: u8,
323 pub sink_routes: Vec<MediaRouteId>,
324 pub transcoding: bool,
325 pub source_frames_routed: u64,
326 pub transcode_operations: u64,
327}
328
329#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
330pub struct MediaGraphEvictionSnapshot {
331 pub route_id: MediaRouteId,
332 pub reason: MediaGraphEvictionReason,
333 pub offered_frames: u64,
334 pub dropped_frames: u64,
335}
336
337#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
344pub struct MediaGraphSnapshot {
345 pub graph_id: MediaGraphId,
346 pub source_state: MediaGraphSourceState,
347 pub source_codec: CodecInfo,
348 pub source_payload_type: u8,
349 pub source_frames: u64,
350 pub sink_offers: u64,
351 pub dropped_frames: u64,
352 pub evictions: u64,
353 pub transcode_operations: u64,
354 pub transcode_errors: u64,
355 pub sinks: Vec<MediaGraphSinkSnapshot>,
356 pub codec_groups: Vec<MediaGraphCodecGroupSnapshot>,
357 pub recent_evictions: Vec<MediaGraphEvictionSnapshot>,
358}
359
360enum Command {
361 Add {
362 route_id: MediaRouteId,
363 codec: CodecInfo,
364 target: mpsc::Sender<MediaFrame>,
365 owner_liveness: Arc<RouteOwnerLiveness>,
366 admission: SinkAdmissionPermit,
367 },
368 Remove {
369 route_id: MediaRouteId,
370 ack: Option<oneshot::Sender<bool>>,
371 },
372 UpdateSourceCodec {
373 codec: CodecInfo,
374 source_pt: u8,
375 ack: oneshot::Sender<Result<()>>,
376 },
377 UpdateSinkCodec {
378 route_id: MediaRouteId,
379 codec: CodecInfo,
380 target_pt: u8,
381 ack: oneshot::Sender<Result<()>>,
382 },
383 UpdateRoute {
385 route_id: MediaRouteId,
386 source_codec: CodecInfo,
387 source_pt: u8,
388 target_codec: CodecInfo,
389 target_pt: u8,
390 ack: oneshot::Sender<Result<()>>,
391 },
392 Snapshot(oneshot::Sender<Arc<MediaGraphSnapshot>>),
393 Shutdown,
394}
395
396type RetainedSnapshot = Arc<StdRwLock<Arc<MediaGraphSnapshot>>>;
397type SinkTaskRegistry = Arc<Mutex<Vec<AbortHandle>>>;
398type RouteStatusRegistry = Arc<Mutex<HashMap<MediaRouteId, watch::Sender<MediaGraphRouteState>>>>;
399
400#[derive(Default)]
401struct RouteOwnerLiveness {
402 cancelled: AtomicBool,
403}
404
405impl RouteOwnerLiveness {
406 fn cancel(&self) {
407 self.cancelled.store(true, Ordering::Release);
408 }
409
410 fn is_cancelled(&self) -> bool {
411 self.cancelled.load(Ordering::Acquire)
412 }
413}
414
415struct SinkAdmissionState {
416 in_use: AtomicUsize,
417 maximum: usize,
418}
419
420impl SinkAdmissionState {
421 fn new(maximum: usize) -> Self {
422 Self {
423 in_use: AtomicUsize::new(0),
424 maximum,
425 }
426 }
427
428 fn try_acquire(self: &Arc<Self>) -> Option<SinkAdmissionPermit> {
429 let mut in_use = self.in_use.load(Ordering::Acquire);
430 loop {
431 if in_use >= self.maximum {
432 return None;
433 }
434 match self.in_use.compare_exchange_weak(
435 in_use,
436 in_use + 1,
437 Ordering::AcqRel,
438 Ordering::Acquire,
439 ) {
440 Ok(_) => {
441 return Some(SinkAdmissionPermit {
442 state: Arc::clone(self),
443 });
444 }
445 Err(observed) => in_use = observed,
446 }
447 }
448 }
449}
450
451struct SinkAdmissionPermit {
452 state: Arc<SinkAdmissionState>,
453}
454
455impl Drop for SinkAdmissionPermit {
456 fn drop(&mut self) {
457 let previous = self.state.in_use.fetch_sub(1, Ordering::AcqRel);
458 debug_assert!(previous > 0, "media graph sink admission underflow");
459 }
460}
461
462#[derive(Clone)]
463pub struct MediaGraphHandle {
464 graph_id: MediaGraphId,
465 commands: mpsc::Sender<Command>,
466 abort: AbortHandle,
467 latest_snapshot: RetainedSnapshot,
468 snapshot_in_flight: Arc<AtomicBool>,
469 completion: watch::Receiver<Option<MediaGraphSourceState>>,
470 activity: watch::Receiver<Option<MediaGraphActivityObservation>>,
471 route_statuses: RouteStatusRegistry,
472 sink_admission: Arc<SinkAdmissionState>,
473}
474
475impl MediaGraphHandle {
476 pub fn id(&self) -> &MediaGraphId {
477 &self.graph_id
478 }
479
480 pub fn subscribe_activity(&self) -> watch::Receiver<Option<MediaGraphActivityObservation>> {
486 self.activity.clone()
487 }
488
489 pub fn add_sink(
490 &self,
491 codec: CodecInfo,
492 target: mpsc::Sender<MediaFrame>,
493 ) -> Result<MediaRouteId> {
494 self.add_managed_sink(codec, target)
495 .map(ManagedMediaRoute::into_unmanaged_route_id)
496 }
497
498 pub fn add_managed_sink(
501 &self,
502 codec: CodecInfo,
503 target: mpsc::Sender<MediaFrame>,
504 ) -> Result<ManagedMediaRoute> {
505 payload_type_for_codec(&codec)
506 .ok_or_else(|| RvoipError::UnsupportedCodec(codec.name.clone()))?;
507 let Some(admission) = self.sink_admission.try_acquire() else {
508 metrics::counter!(
509 "rvoip_media_graph_sink_admission_rejections_total",
510 "reason" => "max-sinks"
511 )
512 .increment(1);
513 return Err(RvoipError::AdmissionRejected(
514 "media graph maximum sink count reached",
515 ));
516 };
517 let route_id = MediaRouteId::new();
518 let (status_tx, status_rx) = watch::channel(MediaGraphRouteState::Pending);
519 let owner_liveness = Arc::new(RouteOwnerLiveness::default());
520 self.route_statuses
521 .lock()
522 .unwrap_or_else(|poisoned| poisoned.into_inner())
523 .insert(route_id.clone(), status_tx);
524 if let Err(error) = self.commands.try_send(Command::Add {
525 route_id: route_id.clone(),
526 codec,
527 target,
528 owner_liveness: Arc::clone(&owner_liveness),
529 admission,
530 }) {
531 self.route_statuses
532 .lock()
533 .unwrap_or_else(|poisoned| poisoned.into_inner())
534 .remove(&route_id);
535 return Err(map_try_send_error(error));
536 }
537 Ok(ManagedMediaRoute {
538 status: MediaGraphRouteStatus {
539 route_id,
540 state: status_rx,
541 },
542 commands: self.commands.clone(),
543 owner_liveness,
544 remove_on_drop: true,
545 })
546 }
547
548 pub fn remove_sink(&self, route_id: MediaRouteId) -> bool {
552 self.commands
553 .try_send(Command::Remove {
554 route_id,
555 ack: None,
556 })
557 .is_ok()
558 }
559
560 pub async fn remove_sink_and_wait(&self, route_id: MediaRouteId) -> Result<bool> {
561 let (ack, done) = oneshot::channel();
562 self.send_control(Command::Remove {
563 route_id,
564 ack: Some(ack),
565 })
566 .await?;
567 tokio::time::timeout(SNAPSHOT_TIMEOUT, done)
568 .await
569 .map_err(|_| RvoipError::InvalidState("media graph removal timed out"))?
570 .map_err(|_| RvoipError::InvalidState("media graph removal was cancelled"))
571 }
572
573 pub async fn update_source_codec(&self, codec: CodecInfo) -> Result<()> {
575 let source_pt = payload_type_for_codec(&codec)
576 .ok_or_else(|| RvoipError::UnsupportedCodec(codec.name.clone()))?;
577 let (ack, done) = oneshot::channel();
578 self.send_control(Command::UpdateSourceCodec {
579 codec,
580 source_pt,
581 ack,
582 })
583 .await?;
584 await_update(done).await
585 }
586
587 pub async fn update_sink_codec(&self, route_id: MediaRouteId, codec: CodecInfo) -> Result<()> {
589 let target_pt = payload_type_for_codec(&codec)
590 .ok_or_else(|| RvoipError::UnsupportedCodec(codec.name.clone()))?;
591 let (ack, done) = oneshot::channel();
592 self.send_control(Command::UpdateSinkCodec {
593 route_id,
594 codec,
595 target_pt,
596 ack,
597 })
598 .await?;
599 await_update(done).await
600 }
601
602 pub async fn update_route(
606 &self,
607 route_id: MediaRouteId,
608 source_pt: u8,
609 target_pt: u8,
610 ) -> Result<()> {
611 let source_codec = codec_for_payload_type(source_pt)
612 .ok_or_else(|| RvoipError::UnsupportedCodec(format!("rtp-payload-type-{source_pt}")))?;
613 let target_codec = codec_for_payload_type(target_pt)
614 .ok_or_else(|| RvoipError::UnsupportedCodec(format!("rtp-payload-type-{target_pt}")))?;
615 let (ack, done) = oneshot::channel();
616 self.send_control(Command::UpdateRoute {
617 route_id,
618 source_codec,
619 source_pt,
620 target_codec,
621 target_pt,
622 ack,
623 })
624 .await?;
625 await_update(done).await
626 }
627
628 pub async fn snapshot(&self) -> MediaGraphSnapshot {
632 (*self.snapshot_arc().await).clone()
633 }
634
635 pub async fn snapshot_arc(&self) -> Arc<MediaGraphSnapshot> {
639 if self
640 .snapshot_in_flight
641 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
642 .is_err()
643 {
644 return self.latest_snapshot_arc();
645 }
646 let _request_guard = SnapshotRequestGuard(Arc::clone(&self.snapshot_in_flight));
647 let (send, receive) = oneshot::channel();
648 if self.commands.try_send(Command::Snapshot(send)).is_ok() {
649 if let Ok(Ok(snapshot)) = tokio::time::timeout(SNAPSHOT_TIMEOUT, receive).await {
650 return snapshot;
651 }
652 }
653 self.latest_snapshot_arc()
654 }
655
656 pub fn latest_snapshot(&self) -> MediaGraphSnapshot {
659 (*self.latest_snapshot_arc()).clone()
660 }
661
662 pub fn latest_snapshot_arc(&self) -> Arc<MediaGraphSnapshot> {
665 read_snapshot(&self.latest_snapshot)
666 }
667
668 pub fn shutdown(&self) {
669 if matches!(
670 self.commands.try_send(Command::Shutdown),
671 Err(mpsc::error::TrySendError::Full(_))
672 ) {
673 self.abort.abort();
675 }
676 }
677
678 pub async fn shutdown_and_wait(&self) -> Result<MediaGraphSourceState> {
681 self.shutdown();
682 self.wait_closed().await
683 }
684
685 pub async fn wait_closed(&self) -> Result<MediaGraphSourceState> {
687 let mut completion = self.completion.clone();
688 let actor = self.abort.clone();
689 tokio::time::timeout(SHUTDOWN_TIMEOUT, async move {
690 let state = loop {
691 if let Some(state) = *completion.borrow() {
692 break state;
693 }
694 completion
695 .changed()
696 .await
697 .map_err(|_| RvoipError::InvalidState("media graph completion was dropped"))?;
698 };
699 while !actor.is_finished() {
700 tokio::task::yield_now().await;
701 }
702 Ok(state)
703 })
704 .await
705 .map_err(|_| RvoipError::InvalidState("media graph shutdown timed out"))?
706 }
707
708 pub fn abort_handle(&self) -> AbortHandle {
709 self.abort.clone()
710 }
711
712 async fn send_control(&self, command: Command) -> Result<()> {
713 tokio::time::timeout(SNAPSHOT_TIMEOUT, self.commands.send(command))
714 .await
715 .map_err(|_| RvoipError::InvalidState("media graph command queue is full"))?
716 .map_err(|_| RvoipError::InvalidState("media graph is closed"))
717 }
718}
719
720struct SnapshotRequestGuard(Arc<AtomicBool>);
721
722impl Drop for SnapshotRequestGuard {
723 fn drop(&mut self) {
724 self.0.store(false, Ordering::Release);
725 }
726}
727
728fn map_try_send_error(error: mpsc::error::TrySendError<Command>) -> RvoipError {
729 match error {
730 mpsc::error::TrySendError::Full(_) => {
731 RvoipError::InvalidState("media graph command queue is full")
732 }
733 mpsc::error::TrySendError::Closed(_) => RvoipError::InvalidState("media graph is closed"),
734 }
735}
736
737fn payload_type_for_codec(codec: &CodecInfo) -> Option<u8> {
738 codec_to_pt(codec.name.trim())
739}
740
741pub fn validate_media_graph_codec(codec: &CodecInfo) -> Result<()> {
744 payload_type_for_codec(codec)
745 .map(|_| ())
746 .ok_or_else(|| RvoipError::UnsupportedCodec(codec.name.clone()))
747}
748
749async fn await_update(done: oneshot::Receiver<Result<()>>) -> Result<()> {
750 tokio::time::timeout(SNAPSHOT_TIMEOUT, done)
751 .await
752 .map_err(|_| RvoipError::InvalidState("media graph update timed out"))?
753 .map_err(|_| RvoipError::InvalidState("media graph update was cancelled"))?
754}
755
756struct SinkQueueState {
757 frames: VecDeque<MediaFrame>,
758 closed: bool,
759}
760
761struct SinkQueue {
762 capacity: usize,
763 state: Mutex<SinkQueueState>,
764 notify: Notify,
765}
766
767#[derive(Clone, Copy, Debug, Eq, PartialEq)]
768enum OfferResult {
769 Enqueued,
770 DroppedOldest,
771 Closed,
772}
773
774impl SinkQueue {
775 fn new(capacity: usize) -> Self {
776 let capacity = capacity.max(1);
777 Self {
778 capacity,
779 state: Mutex::new(SinkQueueState {
780 frames: VecDeque::with_capacity(capacity),
781 closed: false,
782 }),
783 notify: Notify::new(),
784 }
785 }
786
787 fn offer(&self, frame: MediaFrame) -> OfferResult {
790 let result = {
791 let mut state = self.state.lock().expect("media sink queue poisoned");
792 if state.closed {
793 return OfferResult::Closed;
794 }
795 let result = if state.frames.len() >= self.capacity {
796 state.frames.pop_front();
797 OfferResult::DroppedOldest
798 } else {
799 OfferResult::Enqueued
800 };
801 state.frames.push_back(frame);
802 result
803 };
804 self.notify.notify_one();
805 result
806 }
807
808 async fn receive(&self) -> Option<MediaFrame> {
809 loop {
810 let closed = {
811 let mut state = self.state.lock().expect("media sink queue poisoned");
812 if let Some(frame) = state.frames.pop_front() {
813 return Some(frame);
814 }
815 state.closed
816 };
817 if closed {
818 return None;
819 }
820 self.notify.notified().await;
821 }
822 }
823
824 fn depth(&self) -> usize {
825 self.state
826 .lock()
827 .expect("media sink queue poisoned")
828 .frames
829 .len()
830 }
831
832 fn close(&self) {
833 {
834 let mut state = self.state.lock().expect("media sink queue poisoned");
835 state.closed = true;
836 state.frames.clear();
837 }
838 self.notify.notify_waiters();
839 }
840}
841
842struct SinkRuntime {
843 target_codec: CodecInfo,
844 target_pt: u8,
845 group_key: CodecGroupKey,
846 owner_liveness: Arc<RouteOwnerLiveness>,
847 _admission: SinkAdmissionPermit,
850 clock: RtpClockTranslator,
854 queue: Arc<SinkQueue>,
855 task: AbortHandle,
856 history: VecDeque<(Instant, bool)>,
857 rolling_drops: usize,
858 offered_frames: u64,
859 dropped_frames: u64,
860}
861
862impl SinkRuntime {
863 fn record_offer(&mut self, now: Instant, dropped: bool, policy: &MediaGraphPolicy) -> bool {
864 self.offered_frames = self.offered_frames.saturating_add(1);
865 if dropped {
866 self.dropped_frames = self.dropped_frames.saturating_add(1);
867 self.rolling_drops = self.rolling_drops.saturating_add(1);
868 }
869 self.history.push_back((now, dropped));
870 while self
871 .history
872 .front()
873 .is_some_and(|(at, _)| now.saturating_duration_since(*at) > policy.eviction_window)
874 {
875 if self.history.pop_front().is_some_and(|(_, dropped)| dropped) {
876 self.rolling_drops = self.rolling_drops.saturating_sub(1);
877 }
878 }
879 if self.history.len() < policy.minimum_eviction_samples {
880 return false;
881 }
882 self.rolling_drops as f32 / self.history.len() as f32 > policy.eviction_drop_ratio
883 }
884
885 fn rolling_drop_counts(&self) -> (usize, usize) {
886 (self.history.len(), self.rolling_drops)
887 }
888}
889
890impl Drop for SinkRuntime {
891 fn drop(&mut self) {
892 self.queue.close();
893 self.task.abort();
894 }
895}
896
897#[derive(Clone, Eq, Hash, Ord, PartialEq, PartialOrd)]
898struct CodecGroupKey {
899 payload_type: u8,
900 name: String,
901 clock_rate_hz: u32,
902 channels: u8,
903 fmtp: Option<String>,
904}
905
906impl fmt::Debug for CodecGroupKey {
907 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
908 formatter
909 .debug_struct("CodecGroupKey")
910 .field("payload_type", &self.payload_type)
911 .field("name_present", &!self.name.is_empty())
912 .field("name_bytes", &self.name.len())
913 .field("clock_rate_hz", &self.clock_rate_hz)
914 .field("channels", &self.channels)
915 .field("fmtp_present", &self.fmtp.is_some())
916 .finish()
917 }
918}
919
920impl CodecGroupKey {
921 fn new(codec: &CodecInfo, payload_type: u8) -> Self {
922 Self {
923 payload_type,
924 name: canonical_codec_name(codec, payload_type),
925 clock_rate_hz: codec.clock_rate_hz,
926 channels: codec.channels,
927 fmtp: normalize_fmtp(codec.fmtp.as_deref()),
928 }
929 }
930}
931
932fn canonical_codec_name(codec: &CodecInfo, payload_type: u8) -> String {
933 match payload_type {
934 0 => "pcmu".into(),
935 8 => "pcma".into(),
936 18 => "g729".into(),
937 111 => "opus".into(),
938 PCM_S16LE => "pcm_s16le".into(),
939 _ => codec.name.trim().to_ascii_lowercase(),
940 }
941}
942
943fn normalize_fmtp(fmtp: Option<&str>) -> Option<String> {
944 let mut parameters: Vec<_> = fmtp
945 .unwrap_or_default()
946 .split(';')
947 .filter_map(|parameter| {
948 let parameter = parameter.trim();
949 if parameter.is_empty() {
950 return None;
951 }
952 let normalized = match parameter.split_once('=') {
953 Some((name, value)) => {
954 format!("{}={}", name.trim().to_ascii_lowercase(), value.trim())
955 }
956 None => parameter.to_ascii_lowercase(),
957 };
958 Some(normalized)
959 })
960 .collect();
961 parameters.sort();
962 (!parameters.is_empty()).then(|| parameters.join(";"))
963}
964
965struct RtpClockTranslator {
966 source_rate: u32,
967 target_rate: u32,
968 last_source: Option<u32>,
969 last_target: u32,
970 remainder: u64,
971}
972
973impl RtpClockTranslator {
974 fn new(source_rate: u32, target_rate: u32) -> Self {
975 Self {
976 source_rate: source_rate.max(1),
977 target_rate: target_rate.max(1),
978 last_source: None,
979 last_target: 0,
980 remainder: 0,
981 }
982 }
983
984 fn translate(&mut self, source_timestamp: u32) -> u32 {
985 let Some(last_source) = self.last_source.replace(source_timestamp) else {
986 self.last_target = source_timestamp;
987 return source_timestamp;
988 };
989 let source_delta = source_timestamp.wrapping_sub(last_source) as u64;
990 let numerator = source_delta
991 .saturating_mul(self.target_rate as u64)
992 .saturating_add(self.remainder);
993 let target_delta = numerator / self.source_rate as u64;
994 self.remainder = numerator % self.source_rate as u64;
995 self.last_target = self.last_target.wrapping_add(target_delta as u32);
996 self.last_target
997 }
998
999 fn reconfigure(&mut self, source_rate: u32, target_rate: u32) {
1000 self.source_rate = source_rate.max(1);
1001 self.target_rate = target_rate.max(1);
1002 self.remainder = 0;
1003 }
1004}
1005
1006struct ConfiguredTranscodingSession {
1007 source_codec: Box<dyn AudioCodec>,
1008 target_codec: Box<dyn AudioCodec>,
1009 format_converter: FormatConverter,
1010}
1011
1012impl ConfiguredTranscodingSession {
1013 fn new(
1014 source: &CodecInfo,
1015 source_pt: u8,
1016 target: &CodecInfo,
1017 target_pt: u8,
1018 ) -> rvoip_media_core::Result<Self> {
1019 Ok(Self {
1020 source_codec: create_configured_codec(source, source_pt)?,
1021 target_codec: create_configured_codec(target, target_pt)?,
1022 format_converter: FormatConverter::new(),
1023 })
1024 }
1025
1026 fn transcode(&mut self, encoded_data: &[u8]) -> rvoip_media_core::Result<Vec<u8>> {
1027 let source_frame = self.source_codec.decode(encoded_data)?;
1028 let target_info = self.target_codec.get_info();
1029 let converted = if source_frame.sample_rate != target_info.sample_rate
1030 || source_frame.channels != target_info.channels
1031 {
1032 let target_rate = SampleRate::from_hz(target_info.sample_rate).ok_or_else(|| {
1033 CodecError::InvalidParameters {
1034 details: format!("unsupported target clock rate {}", target_info.sample_rate),
1035 }
1036 })?;
1037 self.format_converter
1038 .convert_frame(
1039 &source_frame,
1040 &ConversionParams::new(target_rate, target_info.channels),
1041 )?
1042 .frame
1043 } else {
1044 source_frame
1045 };
1046 self.target_codec.encode(&converted)
1047 }
1048}
1049
1050struct ConfiguredTranscoder {
1051 source_codec: CodecInfo,
1052 source_pt: u8,
1053 target_codec: CodecInfo,
1054 target_pt: u8,
1055 session: Option<ConfiguredTranscodingSession>,
1056}
1057
1058impl ConfiguredTranscoder {
1059 fn new(source_codec: CodecInfo, source_pt: u8, target_codec: CodecInfo, target_pt: u8) -> Self {
1060 Self {
1061 source_codec,
1062 source_pt,
1063 target_codec,
1064 target_pt,
1065 session: None,
1066 }
1067 }
1068
1069 fn transcode(&mut self, payload: &[u8]) -> rvoip_media_core::Result<Vec<u8>> {
1070 if self.session.is_none() {
1071 self.session = Some(ConfiguredTranscodingSession::new(
1072 &self.source_codec,
1073 self.source_pt,
1074 &self.target_codec,
1075 self.target_pt,
1076 )?);
1077 }
1078 self.session
1079 .as_mut()
1080 .expect("configured transcoder session initialized")
1081 .transcode(payload)
1082 }
1083}
1084
1085fn create_configured_codec(
1086 codec: &CodecInfo,
1087 payload_type: u8,
1088) -> rvoip_media_core::Result<Box<dyn AudioCodec>> {
1089 match payload_type {
1090 0 | 8 | 18 => CodecFactory::create_codec(
1091 payload_type,
1092 Some(codec.clock_rate_hz),
1093 Some(codec.channels.into()),
1094 ),
1095 111 => {
1096 let sample_rate = SampleRate::from_hz(codec.clock_rate_hz).ok_or_else(|| {
1097 CodecError::InvalidParameters {
1098 details: format!("unsupported Opus clock rate {}", codec.clock_rate_hz),
1099 }
1100 })?;
1101 let mut config = OpusConfig {
1102 application: OpusApplication::Voip,
1103 ..OpusConfig::default()
1104 };
1105 for parameter in normalize_fmtp(codec.fmtp.as_deref())
1106 .as_deref()
1107 .unwrap_or_default()
1108 .split(';')
1109 {
1110 if let Some(("maxaveragebitrate", value)) = parameter.split_once('=') {
1111 if let Ok(bitrate) = value.parse::<u32>() {
1112 if (6_000..=510_000).contains(&bitrate) {
1113 config.bitrate = bitrate;
1114 }
1115 }
1116 }
1117 if parameter == "cbr=1" {
1118 config.vbr = false;
1119 }
1120 }
1121 Ok(Box::new(OpusCodec::new(
1122 sample_rate,
1123 codec.channels,
1124 config,
1125 )?))
1126 }
1127 PCM_S16LE => Ok(Box::new(PcmS16LeCodec::new(
1128 codec.clock_rate_hz,
1129 codec.channels,
1130 )?)),
1131 _ => Err(CodecError::UnsupportedPayloadType { payload_type }.into()),
1132 }
1133}
1134
1135struct CodecGroup {
1136 target_codec: CodecInfo,
1137 target_pt: u8,
1138 transcoder: Option<ConfiguredTranscoder>,
1139 sinks: HashSet<MediaRouteId>,
1140 source_frames_routed: u64,
1141 transcode_operations: u64,
1142}
1143
1144impl CodecGroup {
1145 fn new(
1146 source_codec: &CodecInfo,
1147 source_pt: u8,
1148 target_codec: CodecInfo,
1149 target_pt: u8,
1150 ) -> Self {
1151 Self {
1152 transcoder: make_transcoder(source_codec, source_pt, &target_codec, target_pt),
1153 target_codec,
1154 target_pt,
1155 sinks: HashSet::new(),
1156 source_frames_routed: 0,
1157 transcode_operations: 0,
1158 }
1159 }
1160}
1161
1162fn make_transcoder(
1163 source_codec: &CodecInfo,
1164 source_pt: u8,
1165 target_codec: &CodecInfo,
1166 target_pt: u8,
1167) -> Option<ConfiguredTranscoder> {
1168 (CodecGroupKey::new(source_codec, source_pt) != CodecGroupKey::new(target_codec, target_pt))
1169 .then(|| {
1170 ConfiguredTranscoder::new(
1171 source_codec.clone(),
1172 source_pt,
1173 target_codec.clone(),
1174 target_pt,
1175 )
1176 })
1177}
1178
1179#[derive(Default)]
1180struct GraphStats {
1181 source_frames: u64,
1182 sink_offers: u64,
1183 dropped_frames: u64,
1184 evictions: u64,
1185 transcode_operations: u64,
1186 transcode_errors: u64,
1187 recent_evictions: VecDeque<MediaGraphEvictionSnapshot>,
1188}
1189
1190impl GraphStats {
1191 fn record_eviction(&mut self, route_id: MediaRouteId, sink: &SinkRuntime) {
1192 self.evictions = self.evictions.saturating_add(1);
1193 if self.recent_evictions.len() == RECENT_EVICTION_LIMIT {
1194 self.recent_evictions.pop_front();
1195 }
1196 self.recent_evictions.push_back(MediaGraphEvictionSnapshot {
1197 route_id,
1198 reason: MediaGraphEvictionReason::SlowConsumer,
1199 offered_frames: sink.offered_frames,
1200 dropped_frames: sink.dropped_frames,
1201 });
1202 }
1203}
1204
1205struct AggregateMetricsGuard {
1208 sinks: usize,
1209 codec_groups: usize,
1210}
1211
1212impl AggregateMetricsGuard {
1213 fn new() -> Self {
1214 metrics::gauge!("rvoip_media_graphs_active").increment(1.0);
1215 Self {
1216 sinks: 0,
1217 codec_groups: 0,
1218 }
1219 }
1220
1221 fn set_sink_count(&mut self, count: usize) {
1222 adjust_gauge("rvoip_media_graph_sinks", self.sinks, count);
1223 self.sinks = count;
1224 }
1225
1226 fn set_codec_group_count(&mut self, count: usize) {
1227 adjust_gauge("rvoip_media_graph_codec_groups", self.codec_groups, count);
1228 self.codec_groups = count;
1229 }
1230}
1231
1232impl Drop for AggregateMetricsGuard {
1233 fn drop(&mut self) {
1234 adjust_gauge("rvoip_media_graph_sinks", self.sinks, 0);
1235 adjust_gauge("rvoip_media_graph_codec_groups", self.codec_groups, 0);
1236 metrics::gauge!("rvoip_media_graphs_active").decrement(1.0);
1237 }
1238}
1239
1240fn adjust_gauge(name: &'static str, old: usize, new: usize) {
1241 if new > old {
1242 metrics::gauge!(name).increment((new - old) as f64);
1243 } else if old > new {
1244 metrics::gauge!(name).decrement((old - new) as f64);
1245 }
1246}
1247
1248struct SnapshotTerminalGuard {
1250 snapshot: RetainedSnapshot,
1251 sink_tasks: SinkTaskRegistry,
1252 completion: watch::Sender<Option<MediaGraphSourceState>>,
1253 route_statuses: RouteStatusRegistry,
1254}
1255
1256impl Drop for SnapshotTerminalGuard {
1257 fn drop(&mut self) {
1258 let mut snapshot = (*read_snapshot(&self.snapshot)).clone();
1259 if snapshot.source_state == MediaGraphSourceState::Open {
1260 snapshot.source_state = MediaGraphSourceState::Aborted;
1261 }
1262 snapshot.sinks.clear();
1263 snapshot.codec_groups.clear();
1264 let terminal_state = snapshot.source_state;
1265 publish_snapshot(&self.snapshot, Arc::new(snapshot));
1266
1267 let reason = match terminal_state {
1268 MediaGraphSourceState::Closed => MediaGraphRouteTerminalReason::SourceClosed,
1269 MediaGraphSourceState::Shutdown => MediaGraphRouteTerminalReason::GraphShutdown,
1270 MediaGraphSourceState::Open | MediaGraphSourceState::Aborted => {
1271 MediaGraphRouteTerminalReason::GraphAborted
1272 }
1273 };
1274 terminate_all_routes(&self.route_statuses, reason);
1275
1276 let mut sink_tasks = self
1277 .sink_tasks
1278 .lock()
1279 .unwrap_or_else(|poisoned| poisoned.into_inner())
1280 .drain(..)
1281 .collect::<Vec<_>>();
1282 for task in &sink_tasks {
1283 task.abort();
1284 }
1285 let completion = self.completion.clone();
1286 if sink_tasks.iter().all(AbortHandle::is_finished) {
1287 let _ = completion.send(Some(terminal_state));
1288 } else if let Ok(runtime) = tokio::runtime::Handle::try_current() {
1289 runtime.spawn(async move {
1290 while sink_tasks.iter().any(|task| !task.is_finished()) {
1291 tokio::task::yield_now().await;
1292 }
1293 sink_tasks.clear();
1294 let _ = completion.send(Some(terminal_state));
1295 });
1296 } else {
1297 let _ = completion.send(Some(terminal_state));
1301 }
1302 }
1303}
1304
1305fn activate_route(statuses: &RouteStatusRegistry, route_id: &MediaRouteId) {
1306 if let Some(status) = statuses
1307 .lock()
1308 .unwrap_or_else(|poisoned| poisoned.into_inner())
1309 .get(route_id)
1310 {
1311 status.send_replace(MediaGraphRouteState::Active);
1312 }
1313}
1314
1315fn terminate_route(
1316 statuses: &RouteStatusRegistry,
1317 route_id: &MediaRouteId,
1318 reason: MediaGraphRouteTerminalReason,
1319) {
1320 if let Some(status) = statuses
1321 .lock()
1322 .unwrap_or_else(|poisoned| poisoned.into_inner())
1323 .remove(route_id)
1324 {
1325 status.send_replace(MediaGraphRouteState::Terminal(reason));
1326 }
1327}
1328
1329fn terminate_all_routes(statuses: &RouteStatusRegistry, reason: MediaGraphRouteTerminalReason) {
1330 let statuses = statuses
1331 .lock()
1332 .unwrap_or_else(|poisoned| poisoned.into_inner())
1333 .drain()
1334 .map(|(_, status)| status)
1335 .collect::<Vec<_>>();
1336 for status in statuses {
1337 status.send_replace(MediaGraphRouteState::Terminal(reason));
1338 }
1339}
1340
1341fn prune_sink_tasks(registry: &SinkTaskRegistry) {
1342 registry
1343 .lock()
1344 .unwrap_or_else(|poisoned| poisoned.into_inner())
1345 .retain(|task| !task.is_finished());
1346}
1347
1348pub fn start_media_graph(
1350 source: mpsc::Receiver<MediaFrame>,
1351 source_codec: CodecInfo,
1352 policy: MediaGraphPolicy,
1353) -> Result<MediaGraphHandle> {
1354 start_media_graph_with_activity_interval(
1355 source,
1356 source_codec,
1357 policy,
1358 MEDIA_GRAPH_ACTIVITY_OBSERVATION_INTERVAL,
1359 )
1360}
1361
1362fn start_media_graph_with_activity_interval(
1363 mut source: mpsc::Receiver<MediaFrame>,
1364 source_codec: CodecInfo,
1365 policy: MediaGraphPolicy,
1366 activity_observation_interval: Duration,
1367) -> Result<MediaGraphHandle> {
1368 let graph_id = MediaGraphId::new();
1369 validate_media_graph_codec(&source_codec)?;
1370 if activity_observation_interval.is_zero() {
1371 return Err(RvoipError::InvalidState(
1372 "media graph activity observation interval is invalid",
1373 ));
1374 }
1375 let initial_source_pt = payload_type_for_codec(&source_codec)
1376 .expect("validated media graph codec has an RTP payload type");
1377 let initial_snapshot = MediaGraphSnapshot {
1378 graph_id: graph_id.clone(),
1379 source_state: MediaGraphSourceState::Open,
1380 source_codec: source_codec.clone(),
1381 source_payload_type: initial_source_pt,
1382 source_frames: 0,
1383 sink_offers: 0,
1384 dropped_frames: 0,
1385 evictions: 0,
1386 transcode_operations: 0,
1387 transcode_errors: 0,
1388 sinks: Vec::new(),
1389 codec_groups: Vec::new(),
1390 recent_evictions: Vec::new(),
1391 };
1392 let latest_snapshot = Arc::new(StdRwLock::new(Arc::new(initial_snapshot)));
1393 let snapshot_for_task = Arc::clone(&latest_snapshot);
1394 let snapshot_in_flight = Arc::new(AtomicBool::new(false));
1395 let sink_tasks = Arc::new(Mutex::new(Vec::new()));
1396 let sink_tasks_for_actor = Arc::clone(&sink_tasks);
1397 let (completion_tx, completion_rx) = watch::channel(None);
1398 let (activity_tx, activity_rx) = watch::channel(None);
1399 let route_statuses: RouteStatusRegistry = Arc::new(Mutex::new(HashMap::new()));
1400 let route_statuses_for_actor = Arc::clone(&route_statuses);
1401 let sink_admission = Arc::new(SinkAdmissionState::new(policy.max_sinks));
1402 let graph_id_for_task = graph_id.clone();
1403 let (command_tx, mut command_rx) = mpsc::channel(CONTROL_QUEUE_CAPACITY);
1404 let (sink_event_tx, mut sink_event_rx) =
1405 mpsc::channel::<MediaRouteId>(SINK_EVENT_QUEUE_CAPACITY);
1406
1407 let terminal_guard = SnapshotTerminalGuard {
1410 snapshot: Arc::clone(&snapshot_for_task),
1411 sink_tasks: Arc::clone(&sink_tasks_for_actor),
1412 completion: completion_tx,
1413 route_statuses: Arc::clone(&route_statuses_for_actor),
1414 };
1415 let task = tokio::spawn(async move {
1416 let _terminal_guard = terminal_guard;
1417 let mut aggregate_metrics = AggregateMetricsGuard::new();
1418 let mut source_codec = source_codec;
1419 let mut source_pt = initial_source_pt;
1420 let mut source_state = MediaGraphSourceState::Open;
1421 let mut sinks: HashMap<MediaRouteId, SinkRuntime> = HashMap::new();
1422 let mut groups: HashMap<CodecGroupKey, CodecGroup> = HashMap::new();
1423 let mut stats = GraphStats::default();
1424 let mut pre_sink_buffer = VecDeque::with_capacity(policy.pre_sink_buffer_frames);
1425 let mut source_routing_started = false;
1428 let mut snapshot_dirty = false;
1429 let mut last_activity_at = None;
1430 let mut activity_pending = false;
1431 let mut snapshot_tick = tokio::time::interval(SNAPSHOT_PUBLISH_INTERVAL);
1432 snapshot_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1433 let mut activity_tick = tokio::time::interval(activity_observation_interval);
1434 activity_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1435 snapshot_tick.tick().await;
1438 activity_tick.tick().await;
1439
1440 loop {
1441 tokio::select! {
1442 command = command_rx.recv() => {
1443 let Some(command) = command else {
1444 source_state = MediaGraphSourceState::Shutdown;
1445 break;
1446 };
1447 match command {
1448 Command::Add {
1449 route_id,
1450 codec,
1451 target,
1452 owner_liveness,
1453 admission,
1454 } => {
1455 if owner_liveness.is_cancelled() {
1456 terminate_route(
1457 &route_statuses_for_actor,
1458 &route_id,
1459 MediaGraphRouteTerminalReason::OwnerRemoved,
1460 );
1461 metrics::counter!(
1462 "rvoip_media_graph_owner_prunes_total",
1463 "phase" => "before-install"
1464 )
1465 .increment(1);
1466 continue;
1467 }
1468 let Some(target_pt) = payload_type_for_codec(&codec) else {
1469 terminate_route(
1470 &route_statuses_for_actor,
1471 &route_id,
1472 MediaGraphRouteTerminalReason::GraphAborted,
1473 );
1474 continue;
1475 };
1476 let first_sink = !source_routing_started;
1477 let target_clock_rate_hz = codec.clock_rate_hz;
1478 let status_route_id = route_id.clone();
1479 let queue = Arc::new(SinkQueue::new(policy.sink_queue_frames));
1480 let queue_for_task = Arc::clone(&queue);
1481 let route_for_task = route_id.clone();
1482 let event_tx = sink_event_tx.clone();
1483 let task = tokio::spawn(async move {
1484 while let Some(frame) = queue_for_task.receive().await {
1485 if target.send(frame).await.is_err() {
1486 let _ = event_tx.send(route_for_task.clone()).await;
1487 return;
1488 }
1489 }
1490 });
1491 prune_sink_tasks(&sink_tasks_for_actor);
1492 sink_tasks_for_actor
1493 .lock()
1494 .unwrap_or_else(|poisoned| poisoned.into_inner())
1495 .push(task.abort_handle());
1496 let group_key = CodecGroupKey::new(&codec, target_pt);
1497 groups.entry(group_key.clone())
1498 .or_insert_with(|| CodecGroup::new(
1499 &source_codec,
1500 source_pt,
1501 codec.clone(),
1502 target_pt,
1503 ))
1504 .sinks.insert(route_id.clone());
1505 sinks.insert(route_id, SinkRuntime {
1506 target_codec: codec,
1507 target_pt,
1508 group_key,
1509 owner_liveness,
1510 _admission: admission,
1511 clock: RtpClockTranslator::new(
1512 source_codec.clock_rate_hz,
1513 target_clock_rate_hz,
1514 ),
1515 queue,
1516 task: task.abort_handle(),
1517 history: VecDeque::new(),
1518 rolling_drops: 0,
1519 offered_frames: 0,
1520 dropped_frames: 0,
1521 });
1522 source_routing_started = true;
1523 aggregate_metrics.set_sink_count(sinks.len());
1524 aggregate_metrics.set_codec_group_count(groups.len());
1525 publish_actor_snapshot(
1526 &snapshot_for_task,
1527 &graph_id_for_task,
1528 source_state,
1529 &source_codec,
1530 source_pt,
1531 &sinks,
1532 &groups,
1533 &stats,
1534 );
1535 snapshot_dirty = false;
1536 activate_route(&route_statuses_for_actor, &status_route_id);
1537
1538 if first_sink {
1539 let mut terminal = Vec::new();
1540 while let Some(frame) = pre_sink_buffer.pop_front() {
1541 terminal.extend(route_source_frame(
1542 frame,
1543 source_pt,
1544 &policy,
1545 &mut sinks,
1546 &mut groups,
1547 &mut stats,
1548 ));
1549 }
1550 snapshot_dirty = true;
1551 if !terminal.is_empty() {
1552 aggregate_metrics.set_sink_count(sinks.len());
1553 aggregate_metrics.set_codec_group_count(groups.len());
1554 publish_actor_snapshot(
1555 &snapshot_for_task,
1556 &graph_id_for_task,
1557 source_state,
1558 &source_codec,
1559 source_pt,
1560 &sinks,
1561 &groups,
1562 &stats,
1563 );
1564 snapshot_dirty = false;
1565 for (route_id, reason) in terminal {
1566 terminate_route(
1567 &route_statuses_for_actor,
1568 &route_id,
1569 reason,
1570 );
1571 }
1572 }
1573 }
1574 }
1575 Command::Remove { route_id, ack } => {
1576 let removed = remove_sink(&route_id, &mut sinks, &mut groups);
1577 prune_sink_tasks(&sink_tasks_for_actor);
1578 aggregate_metrics.set_sink_count(sinks.len());
1579 aggregate_metrics.set_codec_group_count(groups.len());
1580 if removed {
1581 publish_actor_snapshot(
1582 &snapshot_for_task,
1583 &graph_id_for_task,
1584 source_state,
1585 &source_codec,
1586 source_pt,
1587 &sinks,
1588 &groups,
1589 &stats,
1590 );
1591 snapshot_dirty = false;
1592 terminate_route(
1593 &route_statuses_for_actor,
1594 &route_id,
1595 MediaGraphRouteTerminalReason::OwnerRemoved,
1596 );
1597 }
1598 if let Some(ack) = ack {
1599 let _ = ack.send(removed);
1600 }
1601 }
1602 Command::UpdateSourceCodec { codec, source_pt: new_source_pt, ack } => {
1603 source_codec = codec;
1604 source_pt = new_source_pt;
1605 rebuild_transcoders(
1606 &source_codec,
1607 source_pt,
1608 &mut sinks,
1609 &mut groups,
1610 );
1611 publish_actor_snapshot(
1612 &snapshot_for_task,
1613 &graph_id_for_task,
1614 source_state,
1615 &source_codec,
1616 source_pt,
1617 &sinks,
1618 &groups,
1619 &stats,
1620 );
1621 snapshot_dirty = false;
1622 let _ = ack.send(Ok(()));
1623 }
1624 Command::UpdateSinkCodec { route_id, codec, target_pt, ack } => {
1625 let result = update_sink_group(
1626 &route_id,
1627 codec,
1628 target_pt,
1629 &source_codec,
1630 source_pt,
1631 &mut sinks,
1632 &mut groups,
1633 );
1634 aggregate_metrics.set_codec_group_count(groups.len());
1635 if result.is_ok() {
1636 publish_actor_snapshot(
1637 &snapshot_for_task,
1638 &graph_id_for_task,
1639 source_state,
1640 &source_codec,
1641 source_pt,
1642 &sinks,
1643 &groups,
1644 &stats,
1645 );
1646 snapshot_dirty = false;
1647 }
1648 let _ = ack.send(result);
1649 }
1650 Command::UpdateRoute {
1651 route_id,
1652 source_codec: new_source_codec,
1653 source_pt: new_source_pt,
1654 target_codec,
1655 target_pt,
1656 ack,
1657 } => {
1658 source_codec = new_source_codec;
1659 source_pt = new_source_pt;
1660 rebuild_transcoders(
1661 &source_codec,
1662 source_pt,
1663 &mut sinks,
1664 &mut groups,
1665 );
1666 if sinks.contains_key(&route_id) {
1668 let _ = update_sink_group(
1669 &route_id,
1670 target_codec,
1671 target_pt,
1672 &source_codec,
1673 source_pt,
1674 &mut sinks,
1675 &mut groups,
1676 );
1677 }
1678 aggregate_metrics.set_codec_group_count(groups.len());
1679 publish_actor_snapshot(
1680 &snapshot_for_task,
1681 &graph_id_for_task,
1682 source_state,
1683 &source_codec,
1684 source_pt,
1685 &sinks,
1686 &groups,
1687 &stats,
1688 );
1689 snapshot_dirty = false;
1690 let _ = ack.send(Ok(()));
1691 }
1692 Command::Snapshot(reply) => {
1693 if reply.is_closed() {
1696 continue;
1697 }
1698 let snapshot = Arc::new(build_snapshot(
1699 &graph_id_for_task,
1700 source_state,
1701 &source_codec,
1702 source_pt,
1703 &sinks,
1704 &groups,
1705 &stats,
1706 ));
1707 publish_snapshot(&snapshot_for_task, Arc::clone(&snapshot));
1708 snapshot_dirty = false;
1709 let _ = reply.send(snapshot);
1710 }
1711 Command::Shutdown => {
1712 source_state = MediaGraphSourceState::Shutdown;
1713 break;
1714 }
1715 }
1716 }
1717 closed_route = sink_event_rx.recv() => {
1718 if let Some(route_id) = closed_route {
1719 let removed = remove_sink(&route_id, &mut sinks, &mut groups);
1720 aggregate_metrics.set_sink_count(sinks.len());
1721 aggregate_metrics.set_codec_group_count(groups.len());
1722 if removed {
1723 publish_actor_snapshot(
1724 &snapshot_for_task,
1725 &graph_id_for_task,
1726 source_state,
1727 &source_codec,
1728 source_pt,
1729 &sinks,
1730 &groups,
1731 &stats,
1732 );
1733 snapshot_dirty = false;
1734 terminate_route(
1735 &route_statuses_for_actor,
1736 &route_id,
1737 MediaGraphRouteTerminalReason::TargetClosed,
1738 );
1739 }
1740 }
1741 }
1742 frame = source.recv() => {
1743 let owner_removed =
1747 prune_cancelled_owner_sinks(&mut sinks, &mut groups);
1748 if !owner_removed.is_empty() {
1749 aggregate_metrics.set_sink_count(sinks.len());
1750 aggregate_metrics.set_codec_group_count(groups.len());
1751 publish_actor_snapshot(
1752 &snapshot_for_task,
1753 &graph_id_for_task,
1754 source_state,
1755 &source_codec,
1756 source_pt,
1757 &sinks,
1758 &groups,
1759 &stats,
1760 );
1761 for route_id in owner_removed {
1762 terminate_route(
1763 &route_statuses_for_actor,
1764 &route_id,
1765 MediaGraphRouteTerminalReason::OwnerRemoved,
1766 );
1767 }
1768 }
1769 let Some(frame) = frame else {
1770 if !pre_sink_buffer.is_empty() {
1771 let dropped = pre_sink_buffer.len() as u64;
1772 pre_sink_buffer.clear();
1773 stats.dropped_frames =
1774 stats.dropped_frames.saturating_add(dropped);
1775 metrics::counter!(
1776 "rvoip_media_graph_drops_total",
1777 "reason" => "source-closed-before-sink"
1778 )
1779 .increment(dropped);
1780 }
1781 source_state = MediaGraphSourceState::Closed;
1782 break;
1783 };
1784 stats.source_frames = stats.source_frames.saturating_add(1);
1785 record_source_activity(
1786 &mut last_activity_at,
1787 &mut activity_pending,
1788 Utc::now(),
1789 );
1790 snapshot_dirty = true;
1791 metrics::counter!("rvoip_media_graph_source_frames_total").increment(1);
1792 if !source_routing_started {
1793 if policy.pre_sink_buffer_frames > 0 {
1794 if pre_sink_buffer.len() == policy.pre_sink_buffer_frames {
1795 pre_sink_buffer.pop_front();
1796 stats.dropped_frames = stats.dropped_frames.saturating_add(1);
1797 metrics::counter!(
1798 "rvoip_media_graph_drops_total",
1799 "reason" => "pre-sink-buffer-full"
1800 )
1801 .increment(1);
1802 }
1803 pre_sink_buffer.push_back(frame);
1804 } else {
1805 stats.dropped_frames = stats.dropped_frames.saturating_add(1);
1806 metrics::counter!(
1807 "rvoip_media_graph_drops_total",
1808 "reason" => "pre-sink-buffer-disabled"
1809 )
1810 .increment(1);
1811 }
1812 } else {
1813 let terminal = route_source_frame(
1814 frame,
1815 source_pt,
1816 &policy,
1817 &mut sinks,
1818 &mut groups,
1819 &mut stats,
1820 );
1821 aggregate_metrics.set_sink_count(sinks.len());
1822 aggregate_metrics.set_codec_group_count(groups.len());
1823 if !terminal.is_empty() {
1824 publish_actor_snapshot(
1825 &snapshot_for_task,
1826 &graph_id_for_task,
1827 source_state,
1828 &source_codec,
1829 source_pt,
1830 &sinks,
1831 &groups,
1832 &stats,
1833 );
1834 snapshot_dirty = false;
1835 for (route_id, reason) in terminal {
1836 terminate_route(
1837 &route_statuses_for_actor,
1838 &route_id,
1839 reason,
1840 );
1841 }
1842 }
1843 }
1844 }
1845 _ = snapshot_tick.tick() => {
1846 let owner_removed =
1849 prune_cancelled_owner_sinks(&mut sinks, &mut groups);
1850 prune_sink_tasks(&sink_tasks_for_actor);
1851 if !owner_removed.is_empty() {
1852 aggregate_metrics.set_sink_count(sinks.len());
1853 aggregate_metrics.set_codec_group_count(groups.len());
1854 }
1855 if snapshot_dirty || !owner_removed.is_empty() {
1856 publish_snapshot(
1857 &snapshot_for_task,
1858 Arc::new(build_snapshot(
1859 &graph_id_for_task,
1860 source_state,
1861 &source_codec,
1862 source_pt,
1863 &sinks,
1864 &groups,
1865 &stats,
1866 )),
1867 );
1868 snapshot_dirty = false;
1869 }
1870 for route_id in owner_removed {
1871 terminate_route(
1872 &route_statuses_for_actor,
1873 &route_id,
1874 MediaGraphRouteTerminalReason::OwnerRemoved,
1875 );
1876 }
1877 }
1878 _ = activity_tick.tick() => {
1879 publish_pending_activity(
1880 &activity_tx,
1881 &mut activity_pending,
1882 last_activity_at,
1883 stats.source_frames,
1884 );
1885 }
1886 }
1887 }
1888
1889 publish_pending_activity(
1894 &activity_tx,
1895 &mut activity_pending,
1896 last_activity_at,
1897 stats.source_frames,
1898 );
1899 sinks.clear();
1900 groups.clear();
1901 aggregate_metrics.set_sink_count(0);
1902 aggregate_metrics.set_codec_group_count(0);
1903 publish_snapshot(
1904 &snapshot_for_task,
1905 Arc::new(build_snapshot(
1906 &graph_id_for_task,
1907 source_state,
1908 &source_codec,
1909 source_pt,
1910 &sinks,
1911 &groups,
1912 &stats,
1913 )),
1914 );
1915 debug!(graph_id = %graph_id_for_task, ?source_state, "rvoip media graph stopped");
1916 });
1917
1918 Ok(MediaGraphHandle {
1919 graph_id,
1920 commands: command_tx,
1921 abort: task.abort_handle(),
1922 latest_snapshot,
1923 snapshot_in_flight,
1924 completion: completion_rx,
1925 activity: activity_rx,
1926 route_statuses,
1927 sink_admission,
1928 })
1929}
1930
1931fn publish_pending_activity(
1932 sender: &watch::Sender<Option<MediaGraphActivityObservation>>,
1933 pending: &mut bool,
1934 observed_at: Option<DateTime<Utc>>,
1935 source_frames: u64,
1936) {
1937 if !std::mem::replace(pending, false) {
1938 return;
1939 }
1940 let observed_at = observed_at.expect("pending media activity has an observation time");
1941 sender.send_replace(Some(MediaGraphActivityObservation {
1942 source_frames,
1943 observed_at,
1944 }));
1945}
1946
1947fn record_source_activity(
1948 last_observed_at: &mut Option<DateTime<Utc>>,
1949 pending: &mut bool,
1950 observed_at: DateTime<Utc>,
1951) {
1952 *last_observed_at =
1953 Some(last_observed_at.map_or(observed_at, |previous| previous.max(observed_at)));
1954 *pending = true;
1955}
1956
1957fn publish_actor_snapshot(
1958 retained: &RetainedSnapshot,
1959 graph_id: &MediaGraphId,
1960 source_state: MediaGraphSourceState,
1961 source_codec: &CodecInfo,
1962 source_pt: u8,
1963 sinks: &HashMap<MediaRouteId, SinkRuntime>,
1964 groups: &HashMap<CodecGroupKey, CodecGroup>,
1965 stats: &GraphStats,
1966) {
1967 publish_snapshot(
1968 retained,
1969 Arc::new(build_snapshot(
1970 graph_id,
1971 source_state,
1972 source_codec,
1973 source_pt,
1974 sinks,
1975 groups,
1976 stats,
1977 )),
1978 );
1979}
1980
1981fn route_source_frame(
1985 frame: MediaFrame,
1986 source_pt: u8,
1987 policy: &MediaGraphPolicy,
1988 sinks: &mut HashMap<MediaRouteId, SinkRuntime>,
1989 groups: &mut HashMap<CodecGroupKey, CodecGroup>,
1990 stats: &mut GraphStats,
1991) -> Vec<(MediaRouteId, MediaGraphRouteTerminalReason)> {
1992 let now = Instant::now();
1993 let mut evict = Vec::new();
1994 let mut closed = Vec::new();
1995 let is_telephone_event = frame.payload_type == Some(DEFAULT_TELEPHONE_EVENT_PT);
1996
1997 for group in groups.values_mut() {
1998 group.source_frames_routed = group.source_frames_routed.saturating_add(1);
1999 let mut grouped = frame.clone();
2000 if !is_telephone_event {
2001 if let Some(transcoder) = group.transcoder.as_mut() {
2002 group.transcode_operations = group.transcode_operations.saturating_add(1);
2003 stats.transcode_operations = stats.transcode_operations.saturating_add(1);
2004 metrics::counter!(
2005 "rvoip_media_graph_transcodes_total",
2006 "target_payload_type" => group.target_pt.to_string()
2007 )
2008 .increment(1);
2009 match transcoder.transcode(&frame.payload) {
2010 Ok(payload) => {
2011 grouped.payload = payload.into();
2012 grouped.payload_type = Some(group.target_pt);
2013 }
2014 Err(error) => {
2015 stats.transcode_errors = stats.transcode_errors.saturating_add(1);
2016 warn!(
2017 %error,
2018 source_pt,
2019 target_pt = group.target_pt,
2020 "media graph transcode failed"
2021 );
2022 metrics::counter!("rvoip_media_graph_transcode_errors_total").increment(1);
2023 continue;
2024 }
2025 }
2026 }
2027 }
2028
2029 for route_id in &group.sinks {
2030 let Some(sink) = sinks.get_mut(route_id) else {
2031 continue;
2032 };
2033 let mut routed = grouped.clone();
2034 if !is_telephone_event {
2035 routed.timestamp_rtp = sink.clock.translate(frame.timestamp_rtp);
2036 }
2037 let offer = sink.queue.offer(routed);
2038 if offer == OfferResult::Closed {
2039 closed.push(route_id.clone());
2040 continue;
2041 }
2042 let dropped = offer == OfferResult::DroppedOldest;
2043 stats.sink_offers = stats.sink_offers.saturating_add(1);
2044 metrics::counter!("rvoip_media_graph_frames_total").increment(1);
2045 if dropped {
2046 stats.dropped_frames = stats.dropped_frames.saturating_add(1);
2047 metrics::counter!(
2048 "rvoip_media_graph_drops_total",
2049 "reason" => "queue-full"
2050 )
2051 .increment(1);
2052 }
2053 if sink.record_offer(now, dropped, policy) {
2054 evict.push(route_id.clone());
2055 }
2056 }
2057 }
2058
2059 closed.sort();
2060 closed.dedup();
2061 evict.sort();
2062 evict.dedup();
2063 let mut terminal = Vec::with_capacity(closed.len() + evict.len());
2064 for route_id in closed {
2065 if remove_sink(&route_id, sinks, groups) {
2066 terminal.push((route_id, MediaGraphRouteTerminalReason::TargetClosed));
2067 }
2068 }
2069 for route_id in evict {
2070 if let Some(sink) = sinks.get(&route_id) {
2071 stats.record_eviction(route_id.clone(), sink);
2072 metrics::counter!(
2073 "rvoip_media_graph_evictions_total",
2074 "reason" => "slow-consumer"
2075 )
2076 .increment(1);
2077 }
2078 if remove_sink(&route_id, sinks, groups) {
2079 terminal.push((route_id, MediaGraphRouteTerminalReason::SlowConsumerEvicted));
2080 }
2081 }
2082 terminal
2083}
2084
2085fn update_sink_group(
2086 route_id: &MediaRouteId,
2087 codec: CodecInfo,
2088 target_pt: u8,
2089 source_codec: &CodecInfo,
2090 source_pt: u8,
2091 sinks: &mut HashMap<MediaRouteId, SinkRuntime>,
2092 groups: &mut HashMap<CodecGroupKey, CodecGroup>,
2093) -> Result<()> {
2094 let Some(sink) = sinks.get_mut(route_id) else {
2095 return Err(RvoipError::InvalidState("media graph sink not found"));
2096 };
2097 if let Some(group) = groups.get_mut(&sink.group_key) {
2098 group.sinks.remove(route_id);
2099 }
2100 let group_key = CodecGroupKey::new(&codec, target_pt);
2101 sink.clock
2102 .reconfigure(source_codec.clock_rate_hz, codec.clock_rate_hz);
2103 sink.target_codec = codec.clone();
2104 sink.target_pt = target_pt;
2105 sink.group_key = group_key.clone();
2106 groups.retain(|_, group| !group.sinks.is_empty());
2107 groups
2108 .entry(group_key)
2109 .or_insert_with(|| CodecGroup::new(source_codec, source_pt, codec, target_pt))
2110 .sinks
2111 .insert(route_id.clone());
2112 Ok(())
2113}
2114
2115fn rebuild_transcoders(
2116 source_codec: &CodecInfo,
2117 source_pt: u8,
2118 sinks: &mut HashMap<MediaRouteId, SinkRuntime>,
2119 groups: &mut HashMap<CodecGroupKey, CodecGroup>,
2120) {
2121 for sink in sinks.values_mut() {
2122 sink.clock
2123 .reconfigure(source_codec.clock_rate_hz, sink.target_codec.clock_rate_hz);
2124 }
2125 for group in groups.values_mut() {
2126 group.transcoder = make_transcoder(
2127 source_codec,
2128 source_pt,
2129 &group.target_codec,
2130 group.target_pt,
2131 );
2132 }
2133}
2134
2135fn remove_sink(
2136 route_id: &MediaRouteId,
2137 sinks: &mut HashMap<MediaRouteId, SinkRuntime>,
2138 groups: &mut HashMap<CodecGroupKey, CodecGroup>,
2139) -> bool {
2140 let Some(sink) = sinks.remove(route_id) else {
2141 return false;
2142 };
2143 if let Some(group) = groups.get_mut(&sink.group_key) {
2144 group.sinks.remove(route_id);
2145 }
2146 groups.retain(|_, group| !group.sinks.is_empty());
2147 true
2148}
2149
2150fn prune_cancelled_owner_sinks(
2154 sinks: &mut HashMap<MediaRouteId, SinkRuntime>,
2155 groups: &mut HashMap<CodecGroupKey, CodecGroup>,
2156) -> Vec<MediaRouteId> {
2157 let mut cancelled = sinks
2158 .iter()
2159 .filter_map(|(route_id, sink)| sink.owner_liveness.is_cancelled().then(|| route_id.clone()))
2160 .collect::<Vec<_>>();
2161 cancelled.sort();
2162 cancelled.retain(|route_id| remove_sink(route_id, sinks, groups));
2163 if !cancelled.is_empty() {
2164 metrics::counter!(
2165 "rvoip_media_graph_owner_prunes_total",
2166 "phase" => "installed"
2167 )
2168 .increment(cancelled.len() as u64);
2169 }
2170 cancelled
2171}
2172
2173fn codec_for_payload_type(payload_type: u8) -> Option<CodecInfo> {
2174 let (name, clock_rate_hz) = match payload_type {
2175 0 => ("pcmu", 8_000),
2176 8 => ("pcma", 8_000),
2177 18 => ("g729", 8_000),
2178 111 => ("opus", 48_000),
2179 PCM_S16LE => ("pcm_s16le", 16_000),
2180 _ => return None,
2181 };
2182 Some(CodecInfo {
2183 name: name.into(),
2184 clock_rate_hz,
2185 channels: 1,
2186 fmtp: None,
2187 })
2188}
2189
2190fn build_snapshot(
2191 graph_id: &MediaGraphId,
2192 source_state: MediaGraphSourceState,
2193 source_codec: &CodecInfo,
2194 source_pt: u8,
2195 sinks: &HashMap<MediaRouteId, SinkRuntime>,
2196 groups: &HashMap<CodecGroupKey, CodecGroup>,
2197 stats: &GraphStats,
2198) -> MediaGraphSnapshot {
2199 let mut sink_snapshots: Vec<_> = sinks
2200 .iter()
2201 .map(|(route_id, sink)| {
2202 let (rolling_samples, rolling_drops) = sink.rolling_drop_counts();
2203 MediaGraphSinkSnapshot {
2204 route_id: route_id.clone(),
2205 target_codec: sink.target_codec.clone(),
2206 target_payload_type: sink.target_pt,
2207 queue_depth: sink.queue.depth(),
2208 queue_capacity: sink.queue.capacity,
2209 offered_frames: sink.offered_frames,
2210 dropped_frames: sink.dropped_frames,
2211 rolling_samples,
2212 rolling_drops,
2213 rolling_drop_ratio: if rolling_samples == 0 {
2214 0.0
2215 } else {
2216 rolling_drops as f32 / rolling_samples as f32
2217 },
2218 }
2219 })
2220 .collect();
2221 sink_snapshots.sort_by(|a, b| a.route_id.cmp(&b.route_id));
2222
2223 let mut codec_groups: Vec<_> = groups
2224 .values()
2225 .map(|group| {
2226 let mut sink_routes: Vec<_> = group.sinks.iter().cloned().collect();
2227 sink_routes.sort();
2228 MediaGraphCodecGroupSnapshot {
2229 target_codec: group.target_codec.clone(),
2230 target_payload_type: group.target_pt,
2231 sink_routes,
2232 transcoding: group.transcoder.is_some(),
2233 source_frames_routed: group.source_frames_routed,
2234 transcode_operations: group.transcode_operations,
2235 }
2236 })
2237 .collect();
2238 codec_groups.sort_by(|a, b| {
2239 (
2240 a.target_payload_type,
2241 a.target_codec.name.as_str(),
2242 a.target_codec.clock_rate_hz,
2243 a.target_codec.channels,
2244 a.target_codec.fmtp.as_deref(),
2245 )
2246 .cmp(&(
2247 b.target_payload_type,
2248 b.target_codec.name.as_str(),
2249 b.target_codec.clock_rate_hz,
2250 b.target_codec.channels,
2251 b.target_codec.fmtp.as_deref(),
2252 ))
2253 });
2254
2255 MediaGraphSnapshot {
2256 graph_id: graph_id.clone(),
2257 source_state,
2258 source_codec: source_codec.clone(),
2259 source_payload_type: source_pt,
2260 source_frames: stats.source_frames,
2261 sink_offers: stats.sink_offers,
2262 dropped_frames: stats.dropped_frames,
2263 evictions: stats.evictions,
2264 transcode_operations: stats.transcode_operations,
2265 transcode_errors: stats.transcode_errors,
2266 sinks: sink_snapshots,
2267 codec_groups,
2268 recent_evictions: stats.recent_evictions.iter().cloned().collect(),
2269 }
2270}
2271
2272fn publish_snapshot(shared: &RetainedSnapshot, snapshot: Arc<MediaGraphSnapshot>) {
2273 *shared
2274 .write()
2275 .unwrap_or_else(|poisoned| poisoned.into_inner()) = snapshot;
2276}
2277
2278fn read_snapshot(snapshot: &RetainedSnapshot) -> Arc<MediaGraphSnapshot> {
2279 Arc::clone(
2280 &snapshot
2281 .read()
2282 .unwrap_or_else(|poisoned| poisoned.into_inner()),
2283 )
2284}
2285
2286#[cfg(test)]
2287mod tests {
2288 use std::sync::Arc;
2289
2290 use bytes::Bytes;
2291 use chrono::{Duration as ChronoDuration, Utc};
2292
2293 use super::*;
2294 use crate::ids::StreamId;
2295 use crate::stream::StreamKind;
2296
2297 fn codec(name: &str, clock_rate: u32) -> CodecInfo {
2298 CodecInfo {
2299 name: name.into(),
2300 clock_rate_hz: clock_rate,
2301 channels: 1,
2302 fmtp: None,
2303 }
2304 }
2305
2306 #[test]
2307 fn graph_diagnostics_never_render_ids_or_codec_parameters() {
2308 const CANARY: &str = "media-graph-canary\r\nAuthorization: exposed";
2309 let codec = CodecInfo {
2310 name: CANARY.into(),
2311 clock_rate_hz: 48_000,
2312 channels: 1,
2313 fmtp: Some(CANARY.into()),
2314 };
2315 let graph_id = MediaGraphId::from_string(CANARY);
2316 let key = CodecGroupKey::new(&codec, 111);
2317 let snapshot = MediaGraphSnapshot {
2318 graph_id,
2319 source_state: MediaGraphSourceState::Open,
2320 source_codec: codec,
2321 source_payload_type: 111,
2322 source_frames: 1,
2323 sink_offers: 0,
2324 dropped_frames: 0,
2325 evictions: 0,
2326 transcode_operations: 0,
2327 transcode_errors: 0,
2328 sinks: Vec::new(),
2329 codec_groups: Vec::new(),
2330 recent_evictions: Vec::new(),
2331 };
2332 for debug in [format!("{key:?}"), format!("{snapshot:?}")] {
2333 assert!(!debug.contains(CANARY), "graph value leaked: {debug}");
2334 }
2335 }
2336
2337 fn frame(value: u8) -> MediaFrame {
2338 frame_at(value, value as u32 * 160)
2339 }
2340
2341 fn frame_at(value: u8, timestamp_rtp: u32) -> MediaFrame {
2342 frame_at_pt(value, timestamp_rtp, 0)
2343 }
2344
2345 fn frame_at_pt(value: u8, timestamp_rtp: u32, payload_type: u8) -> MediaFrame {
2346 MediaFrame {
2347 stream_id: StreamId::from_string("strm_media_graph_test"),
2348 kind: StreamKind::Audio,
2349 payload: Bytes::from(vec![value; 160]),
2350 timestamp_rtp,
2351 captured_at: Utc::now(),
2352 payload_type: Some(payload_type),
2353 }
2354 }
2355
2356 async fn wait_until(mut predicate: impl FnMut() -> bool) {
2357 tokio::time::timeout(Duration::from_secs(2), async {
2358 while !predicate() {
2359 tokio::task::yield_now().await;
2360 }
2361 })
2362 .await
2363 .expect("condition did not become true");
2364 }
2365
2366 async fn wait_for_route_state(status: &MediaGraphRouteStatus, expected: MediaGraphRouteState) {
2367 tokio::time::timeout(Duration::from_secs(2), async {
2368 while status.state() != expected {
2369 tokio::task::yield_now().await;
2370 }
2371 })
2372 .await
2373 .expect("route state did not converge");
2374 }
2375
2376 #[tokio::test]
2377 async fn one_source_reaches_multiple_sinks() {
2378 let (source_tx, source_rx) = mpsc::channel(4);
2379 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2380 let graph_id = graph.id().clone();
2381 let (a_tx, mut a_rx) = mpsc::channel(4);
2382 let (b_tx, mut b_rx) = mpsc::channel(4);
2383 graph.add_sink(codec("pcmu", 8_000), a_tx).unwrap();
2384 graph.add_sink(codec("pcmu", 8_000), b_tx).unwrap();
2385 let snapshot = graph.snapshot().await;
2386 assert_eq!(snapshot.graph_id, graph_id);
2387 assert_eq!(snapshot.sinks.len(), 2);
2388 assert_eq!(snapshot.codec_groups.len(), 1);
2389
2390 source_tx.send(frame(7)).await.unwrap();
2391 assert_eq!(a_rx.recv().await.unwrap().payload[0], 7);
2392 assert_eq!(b_rx.recv().await.unwrap().payload[0], 7);
2393 graph.shutdown();
2394 }
2395
2396 #[tokio::test]
2397 async fn source_activity_is_retained_and_coalesced_per_interval() {
2398 let (source_tx, source_rx) = mpsc::channel(64);
2399 for value in 0..32 {
2400 source_tx.try_send(frame(value)).unwrap();
2401 }
2402 let graph = start_media_graph_with_activity_interval(
2403 source_rx,
2404 codec("pcmu", 8_000),
2405 MediaGraphPolicy::default(),
2406 Duration::from_millis(40),
2407 )
2408 .unwrap();
2409 let mut activity = graph.subscribe_activity();
2410 tokio::time::timeout(Duration::from_secs(1), activity.changed())
2411 .await
2412 .expect("coalesced observation deadline")
2413 .expect("activity publisher remains live");
2414 let first = activity
2415 .borrow_and_update()
2416 .clone()
2417 .expect("activity observation");
2418 assert_eq!(first.source_frames, 32);
2419 assert!(
2420 tokio::time::timeout(Duration::from_millis(60), activity.changed())
2421 .await
2422 .is_err(),
2423 "idle ticks do not manufacture activity"
2424 );
2425
2426 source_tx.send(frame(33)).await.unwrap();
2427 tokio::time::timeout(Duration::from_secs(1), activity.changed())
2428 .await
2429 .expect("next observation deadline")
2430 .expect("activity publisher remains live");
2431 let second = activity
2432 .borrow_and_update()
2433 .clone()
2434 .expect("next activity observation");
2435 assert_eq!(second.source_frames, 33);
2436 assert!(second.observed_at >= first.observed_at);
2437 graph.shutdown_and_wait().await.unwrap();
2438 }
2439
2440 #[test]
2441 fn zero_activity_observation_interval_is_rejected() {
2442 let (_source_tx, source_rx) = mpsc::channel(1);
2443 assert!(matches!(
2444 start_media_graph_with_activity_interval(
2445 source_rx,
2446 codec("pcmu", 8_000),
2447 MediaGraphPolicy::default(),
2448 Duration::ZERO,
2449 ),
2450 Err(RvoipError::InvalidState(
2451 "media graph activity observation interval is invalid"
2452 ))
2453 ));
2454 }
2455
2456 #[test]
2457 fn activity_observation_time_does_not_move_backward() {
2458 let mut last_observed_at = None;
2459 let mut pending = false;
2460 let later = Utc::now();
2461 record_source_activity(&mut last_observed_at, &mut pending, later);
2462 assert!(pending);
2463 pending = false;
2464
2465 record_source_activity(
2466 &mut last_observed_at,
2467 &mut pending,
2468 later - ChronoDuration::seconds(10),
2469 );
2470 assert!(pending);
2471 assert_eq!(last_observed_at, Some(later));
2472 }
2473
2474 #[tokio::test]
2475 async fn buffered_first_frame_waits_for_initial_sink_registration() {
2476 let (source_tx, source_rx) = mpsc::channel(1);
2477 source_tx.send(frame(42)).await.unwrap();
2478 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2479 for _ in 0..10 {
2482 tokio::task::yield_now().await;
2483 }
2484 let (target_tx, mut target_rx) = mpsc::channel(1);
2485 graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
2486
2487 let received = tokio::time::timeout(Duration::from_secs(1), target_rx.recv())
2488 .await
2489 .expect("initial media was not routed")
2490 .expect("initial sink closed");
2491 assert_eq!(received.payload[0], 42);
2492 graph.shutdown_and_wait().await.unwrap();
2493 }
2494
2495 #[tokio::test]
2496 async fn pre_sink_buffer_is_bounded_drop_oldest_and_flushes_in_order() {
2497 let policy = MediaGraphPolicy {
2498 sink_queue_frames: 4,
2499 pre_sink_buffer_frames: 3,
2500 ..Default::default()
2501 };
2502 let (source_tx, source_rx) = mpsc::channel(8);
2503 for value in 0..6 {
2504 source_tx.send(frame(value)).await.unwrap();
2505 }
2506 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), policy).unwrap();
2507
2508 tokio::time::timeout(Duration::from_secs(2), async {
2511 loop {
2512 if graph.snapshot().await.source_frames == 6 {
2513 break;
2514 }
2515 tokio::task::yield_now().await;
2516 }
2517 })
2518 .await
2519 .expect("source did not drain into the pre-sink buffer");
2520
2521 let (target_tx, mut target_rx) = mpsc::channel(8);
2522 let route = graph
2523 .add_managed_sink(codec("pcmu", 8_000), target_tx)
2524 .unwrap();
2525 wait_for_route_state(&route.status(), MediaGraphRouteState::Active).await;
2526 for expected in 3..6 {
2527 assert_eq!(target_rx.recv().await.unwrap().payload[0], expected);
2528 }
2529 assert!(
2530 tokio::time::timeout(Duration::from_millis(20), target_rx.recv())
2531 .await
2532 .is_err()
2533 );
2534 assert_eq!(graph.snapshot().await.dropped_frames, 3);
2535 graph.shutdown_and_wait().await.unwrap();
2536 }
2537
2538 #[tokio::test]
2539 async fn closed_source_with_buffered_media_converges_and_accounts_for_drops() {
2540 let (source_tx, source_rx) = mpsc::channel(4);
2541 for value in 0..3 {
2542 source_tx.send(frame(value)).await.unwrap();
2543 }
2544 drop(source_tx);
2545 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2546
2547 assert_eq!(
2548 graph.wait_closed().await.unwrap(),
2549 MediaGraphSourceState::Closed
2550 );
2551 let snapshot = graph.latest_snapshot();
2552 assert_eq!(snapshot.source_frames, 3);
2553 assert_eq!(snapshot.dropped_frames, 3);
2554 assert_eq!(snapshot.source_state, MediaGraphSourceState::Closed);
2555 }
2556
2557 #[tokio::test]
2558 async fn removing_one_sink_does_not_stop_others() {
2559 let (source_tx, source_rx) = mpsc::channel(4);
2560 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2561 let (a_tx, mut a_rx) = mpsc::channel(4);
2562 let (b_tx, mut b_rx) = mpsc::channel(4);
2563 let a = graph.add_sink(codec("pcmu", 8_000), a_tx).unwrap();
2564 graph.add_sink(codec("pcmu", 8_000), b_tx).unwrap();
2565 graph.snapshot().await;
2566 graph.remove_sink(a);
2567 assert_eq!(graph.snapshot().await.sinks.len(), 1);
2568
2569 source_tx.send(frame(9)).await.unwrap();
2570 let removed = tokio::time::timeout(Duration::from_millis(50), a_rx.recv()).await;
2571 assert!(matches!(removed, Ok(None) | Err(_)));
2572 assert_eq!(b_rx.recv().await.unwrap().payload[0], 9);
2573 graph.shutdown();
2574 }
2575
2576 #[tokio::test]
2577 async fn acknowledged_removal_reports_route_existence() {
2578 let (_source_tx, source_rx) = mpsc::channel(1);
2579 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2580 let (target_tx, mut target_rx) = mpsc::channel(1);
2581 let route = graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
2582 graph.snapshot().await;
2583
2584 assert!(graph.remove_sink_and_wait(route.clone()).await.unwrap());
2585 assert!(graph.latest_snapshot().sinks.is_empty());
2586 assert!(!graph.remove_sink_and_wait(route).await.unwrap());
2587 assert!(target_rx.recv().await.is_none());
2588 graph.shutdown_and_wait().await.unwrap();
2589 }
2590
2591 #[tokio::test]
2592 async fn managed_owner_drop_removes_route_without_status_clone_ownership() {
2593 let (_source_tx, source_rx) = mpsc::channel(1);
2594 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2595 let (target_tx, _target_rx) = mpsc::channel(1);
2596 let route = graph
2597 .add_managed_sink(codec("pcmu", 8_000), target_tx)
2598 .unwrap();
2599 let status = route.status();
2600 let observer = status.clone();
2601 wait_for_route_state(&status, MediaGraphRouteState::Active).await;
2602 assert_eq!(graph.latest_snapshot().sinks.len(), 1);
2603
2604 drop(route);
2605 assert_eq!(
2606 status.wait_terminal().await,
2607 MediaGraphRouteTerminalReason::OwnerRemoved
2608 );
2609 assert_eq!(
2610 observer.state(),
2611 MediaGraphRouteState::Terminal(MediaGraphRouteTerminalReason::OwnerRemoved)
2612 );
2613 assert!(graph.latest_snapshot().sinks.is_empty());
2614 assert!(graph
2615 .route_statuses
2616 .lock()
2617 .unwrap_or_else(|poisoned| poisoned.into_inner())
2618 .is_empty());
2619 graph.shutdown_and_wait().await.unwrap();
2620 }
2621
2622 #[tokio::test]
2623 async fn managed_owner_drop_converges_when_control_queue_is_saturated() {
2624 let (source_tx, source_rx) = mpsc::channel(1);
2625 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2626 let (target_tx, mut target_rx) = mpsc::channel(1);
2627 let route = graph
2628 .add_managed_sink(codec("pcmu", 8_000), target_tx)
2629 .unwrap();
2630 let status = route.status();
2631 wait_for_route_state(&status, MediaGraphRouteState::Active).await;
2632
2633 let mut snapshot_receivers = Vec::with_capacity(CONTROL_QUEUE_CAPACITY);
2637 for _ in 0..CONTROL_QUEUE_CAPACITY {
2638 let (reply, receive) = oneshot::channel();
2639 assert!(graph.commands.try_send(Command::Snapshot(reply)).is_ok());
2640 snapshot_receivers.push(receive);
2641 }
2642 let (overflow, _receive) = oneshot::channel();
2643 assert!(matches!(
2644 graph.commands.try_send(Command::Snapshot(overflow)),
2645 Err(mpsc::error::TrySendError::Full(_))
2646 ));
2647 assert_eq!(graph.commands.capacity(), 0);
2648
2649 drop(route);
2653 source_tx.try_send(frame(77)).unwrap();
2654 assert_eq!(
2655 tokio::time::timeout(Duration::from_secs(2), status.wait_terminal())
2656 .await
2657 .expect("cancelled route did not converge"),
2658 MediaGraphRouteTerminalReason::OwnerRemoved
2659 );
2660
2661 let snapshot = graph.snapshot().await;
2662 assert!(snapshot.sinks.is_empty());
2663 assert!(snapshot.codec_groups.is_empty());
2664 assert_eq!(snapshot.sink_offers, 0);
2665 assert_eq!(graph.sink_admission.in_use.load(Ordering::Acquire), 0);
2666 assert!(graph
2667 .route_statuses
2668 .lock()
2669 .unwrap_or_else(|poisoned| poisoned.into_inner())
2670 .is_empty());
2671 assert!(matches!(
2672 tokio::time::timeout(Duration::from_secs(1), target_rx.recv()).await,
2673 Ok(None)
2674 ));
2675 drop(snapshot_receivers);
2676 graph.shutdown_and_wait().await.unwrap();
2677 }
2678
2679 #[test]
2680 fn default_sink_limit_reserves_direct_fanout_headroom() {
2681 assert_eq!(
2682 MediaGraphPolicy::default().max_sinks,
2683 DEFAULT_MEDIA_GRAPH_MAX_SINKS
2684 );
2685 assert_eq!(DEFAULT_MEDIA_GRAPH_MAX_SINKS, 1_024);
2686 }
2687
2688 #[tokio::test]
2689 async fn sink_admission_accepts_the_boundary_rejects_one_more_and_recovers() {
2690 let policy = MediaGraphPolicy {
2691 max_sinks: 2,
2692 ..Default::default()
2693 };
2694 let (_source_tx, source_rx) = mpsc::channel(1);
2695 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), policy).unwrap();
2696 let (first_tx, _first_rx) = mpsc::channel(1);
2697 let (second_tx, _second_rx) = mpsc::channel(1);
2698 let first = graph.add_sink(codec("pcmu", 8_000), first_tx).unwrap();
2699 graph.add_sink(codec("pcmu", 8_000), second_tx).unwrap();
2700
2701 let (rejected_tx, _rejected_rx) = mpsc::channel(1);
2702 assert!(matches!(
2703 graph.add_sink(codec("pcmu", 8_000), rejected_tx),
2704 Err(RvoipError::AdmissionRejected(
2705 "media graph maximum sink count reached"
2706 ))
2707 ));
2708 let at_limit = graph.snapshot().await;
2709 assert_eq!(at_limit.sinks.len(), 2);
2710 assert_eq!(at_limit.codec_groups.len(), 1);
2711 assert_eq!(graph.sink_admission.in_use.load(Ordering::Acquire), 2);
2712 assert_eq!(
2713 graph
2714 .route_statuses
2715 .lock()
2716 .unwrap_or_else(|poisoned| poisoned.into_inner())
2717 .len(),
2718 2
2719 );
2720
2721 assert!(graph.remove_sink_and_wait(first).await.unwrap());
2722 assert_eq!(graph.sink_admission.in_use.load(Ordering::Acquire), 1);
2723 let (replacement_tx, _replacement_rx) = mpsc::channel(1);
2724 graph
2725 .add_sink(codec("pcmu", 8_000), replacement_tx)
2726 .unwrap();
2727 let recovered = graph.snapshot().await;
2728 assert_eq!(recovered.sinks.len(), 2);
2729 assert_eq!(recovered.codec_groups.len(), 1);
2730 assert_eq!(graph.sink_admission.in_use.load(Ordering::Acquire), 2);
2731 graph.shutdown_and_wait().await.unwrap();
2732 }
2733
2734 #[tokio::test]
2735 async fn repeated_managed_add_remove_does_not_retain_status_registry_entries() {
2736 let (_source_tx, source_rx) = mpsc::channel(1);
2737 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2738
2739 for iteration in 0..32 {
2740 let (target_tx, _target_rx) = mpsc::channel(1);
2741 let route = graph
2742 .add_managed_sink(codec("pcmu", 8_000), target_tx)
2743 .unwrap();
2744 let status = route.status();
2745 wait_for_route_state(&status, MediaGraphRouteState::Active).await;
2746 if iteration % 2 == 0 {
2747 drop(route);
2748 } else {
2749 assert!(route.remove().await.unwrap());
2750 }
2751 assert_eq!(
2752 status.wait_terminal().await,
2753 MediaGraphRouteTerminalReason::OwnerRemoved
2754 );
2755 assert!(graph
2756 .route_statuses
2757 .lock()
2758 .unwrap_or_else(|poisoned| poisoned.into_inner())
2759 .is_empty());
2760 }
2761 assert!(graph.latest_snapshot().sinks.is_empty());
2762 graph.shutdown_and_wait().await.unwrap();
2763 }
2764
2765 #[tokio::test]
2766 async fn target_close_reports_terminal_reason_and_late_status_clone() {
2767 let (source_tx, source_rx) = mpsc::channel(2);
2768 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2769 let (target_tx, target_rx) = mpsc::channel(1);
2770 let route = graph
2771 .add_managed_sink(codec("pcmu", 8_000), target_tx)
2772 .unwrap();
2773 let status = route.status();
2774 wait_for_route_state(&status, MediaGraphRouteState::Active).await;
2775 drop(target_rx);
2776 source_tx.send(frame(1)).await.unwrap();
2777
2778 assert_eq!(
2779 status.wait_terminal().await,
2780 MediaGraphRouteTerminalReason::TargetClosed
2781 );
2782 let late = status.clone();
2783 assert_eq!(
2784 late.state(),
2785 MediaGraphRouteState::Terminal(MediaGraphRouteTerminalReason::TargetClosed)
2786 );
2787 assert!(graph.latest_snapshot().sinks.is_empty());
2788 graph.shutdown_and_wait().await.unwrap();
2789 }
2790
2791 #[tokio::test]
2792 async fn full_queue_drops_oldest_frame() {
2793 let queue = SinkQueue::new(2);
2794 assert_eq!(queue.offer(frame(1)), OfferResult::Enqueued);
2795 assert_eq!(queue.offer(frame(2)), OfferResult::Enqueued);
2796 assert_eq!(queue.offer(frame(3)), OfferResult::DroppedOldest);
2797 assert_eq!(queue.depth(), 2);
2798 assert_eq!(queue.receive().await.unwrap().payload[0], 2);
2799 assert_eq!(queue.depth(), 1);
2800 assert_eq!(queue.receive().await.unwrap().payload[0], 3);
2801 queue.close();
2802 assert!(queue.receive().await.is_none());
2803 }
2804
2805 #[tokio::test]
2806 async fn eviction_is_strictly_greater_than_twenty_five_percent_over_ten_seconds() {
2807 let policy = MediaGraphPolicy {
2808 max_sinks: DEFAULT_MEDIA_GRAPH_MAX_SINKS,
2809 sink_queue_frames: 1,
2810 pre_sink_buffer_frames: 10,
2811 eviction_window: Duration::from_secs(10),
2812 eviction_drop_ratio: 0.25,
2813 minimum_eviction_samples: 4,
2814 };
2815 let (target, _receiver) = mpsc::channel::<MediaFrame>(1);
2816 let task = tokio::spawn(async move { drop(target) });
2817 let target_codec = codec("pcmu", 8_000);
2818 let mut sink = SinkRuntime {
2819 group_key: CodecGroupKey::new(&target_codec, 0),
2820 target_codec,
2821 target_pt: 0,
2822 owner_liveness: Arc::new(RouteOwnerLiveness::default()),
2823 _admission: Arc::new(SinkAdmissionState::new(1)).try_acquire().unwrap(),
2824 clock: RtpClockTranslator::new(8_000, 8_000),
2825 queue: Arc::new(SinkQueue::new(1)),
2826 task: task.abort_handle(),
2827 history: VecDeque::new(),
2828 rolling_drops: 0,
2829 offered_frames: 0,
2830 dropped_frames: 0,
2831 };
2832 let start = Instant::now();
2833 assert!(!sink.record_offer(start, true, &policy));
2834 assert!(!sink.record_offer(start + Duration::from_secs(1), false, &policy));
2835 assert!(!sink.record_offer(start + Duration::from_secs(2), false, &policy));
2836 assert!(!sink.record_offer(start + Duration::from_secs(3), false, &policy));
2838 assert_eq!(sink.rolling_drop_counts(), (4, 1));
2839 assert!(!sink.record_offer(start + Duration::from_secs(10), false, &policy));
2841 assert!(!sink.record_offer(
2844 start + Duration::from_secs(10) + Duration::from_nanos(1),
2845 true,
2846 &policy,
2847 ));
2848 assert_eq!(sink.rolling_drop_counts(), (5, 1));
2849 assert!(sink.record_offer(
2851 start + Duration::from_secs(10) + Duration::from_nanos(2),
2852 true,
2853 &policy,
2854 ));
2855 assert_eq!(sink.rolling_drop_counts(), (6, 2));
2856 }
2857
2858 #[tokio::test]
2859 async fn slow_sink_is_evicted_and_reported() {
2860 let policy = MediaGraphPolicy {
2861 max_sinks: DEFAULT_MEDIA_GRAPH_MAX_SINKS,
2862 sink_queue_frames: 1,
2863 pre_sink_buffer_frames: 10,
2864 eviction_window: Duration::from_secs(10),
2865 eviction_drop_ratio: 0.25,
2866 minimum_eviction_samples: 4,
2867 };
2868 let (source_tx, source_rx) = mpsc::channel(64);
2869 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), policy).unwrap();
2870 let (target_tx, _target_rx) = mpsc::channel(1);
2871 let route = graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
2872 graph.snapshot().await;
2873 for value in 0..32 {
2874 source_tx.send(frame(value)).await.unwrap();
2875 }
2876 wait_until(|| graph.latest_snapshot().evictions == 1).await;
2877 let snapshot = graph.snapshot().await;
2878 assert!(snapshot.sinks.is_empty());
2879 assert!(snapshot.dropped_frames > 0);
2880 assert_eq!(snapshot.recent_evictions[0].route_id, route);
2881 assert_eq!(
2882 snapshot.recent_evictions[0].reason,
2883 MediaGraphEvictionReason::SlowConsumer
2884 );
2885 graph.shutdown();
2886 }
2887
2888 #[tokio::test]
2889 async fn managed_slow_sink_reports_eviction_terminal_reason() {
2890 let policy = MediaGraphPolicy {
2891 max_sinks: DEFAULT_MEDIA_GRAPH_MAX_SINKS,
2892 sink_queue_frames: 1,
2893 pre_sink_buffer_frames: 10,
2894 eviction_window: Duration::from_secs(10),
2895 eviction_drop_ratio: 0.25,
2896 minimum_eviction_samples: 4,
2897 };
2898 let (source_tx, source_rx) = mpsc::channel(64);
2899 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), policy).unwrap();
2900 let (target_tx, _target_rx) = mpsc::channel(1);
2901 let route = graph
2902 .add_managed_sink(codec("pcmu", 8_000), target_tx)
2903 .unwrap();
2904 let status = route.status();
2905 wait_for_route_state(&status, MediaGraphRouteState::Active).await;
2906 for value in 0..32 {
2907 source_tx.send(frame(value)).await.unwrap();
2908 }
2909
2910 assert_eq!(
2911 tokio::time::timeout(Duration::from_secs(2), status.wait_terminal())
2912 .await
2913 .expect("slow sink was not evicted"),
2914 MediaGraphRouteTerminalReason::SlowConsumerEvicted
2915 );
2916 assert!(graph.latest_snapshot().sinks.is_empty());
2917 graph.shutdown_and_wait().await.unwrap();
2918 }
2919
2920 #[tokio::test]
2921 async fn transcodes_once_per_codec_group_and_preserves_timestamps() {
2922 let (source_tx, source_rx) = mpsc::channel(4);
2923 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2924 let (a_tx, mut a_rx) = mpsc::channel(4);
2925 let (b_tx, mut b_rx) = mpsc::channel(4);
2926 graph.add_sink(codec("pcma", 8_000), a_tx).unwrap();
2927 graph.add_sink(codec("pcma", 8_000), b_tx).unwrap();
2928 graph.snapshot().await;
2929
2930 source_tx.send(frame_at(0x7f, 42_424)).await.unwrap();
2931 let a = a_rx.recv().await.unwrap();
2932 let b = b_rx.recv().await.unwrap();
2933 assert_eq!(a.timestamp_rtp, 42_424);
2934 assert_eq!(b.timestamp_rtp, 42_424);
2935 assert_eq!(a.stream_id, b.stream_id);
2936 assert_eq!(a.payload, b.payload);
2937 let snapshot = graph.snapshot().await;
2938 assert_eq!(snapshot.transcode_operations, 1);
2939 assert_eq!(snapshot.codec_groups[0].transcode_operations, 1);
2940 assert_eq!(snapshot.codec_groups[0].sink_routes.len(), 2);
2941 graph.shutdown();
2942 }
2943
2944 #[tokio::test]
2945 async fn frame_order_and_rtp_timestamp_continuity_are_preserved() {
2946 let (source_tx, source_rx) = mpsc::channel(8);
2947 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2948 let (target_tx, mut target_rx) = mpsc::channel(8);
2949 graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
2950 graph.snapshot().await;
2951
2952 let timestamps = [u32::MAX - 159, 0, 160, 320];
2953 for (value, timestamp) in timestamps.into_iter().enumerate() {
2954 source_tx
2955 .send(frame_at(value as u8, timestamp))
2956 .await
2957 .unwrap();
2958 }
2959 for (value, timestamp) in timestamps.into_iter().enumerate() {
2960 let received = target_rx.recv().await.unwrap();
2961 assert_eq!(received.payload[0], value as u8);
2962 assert_eq!(received.timestamp_rtp, timestamp);
2963 }
2964 graph.shutdown();
2965 }
2966
2967 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2968 async fn concurrent_add_remove_leaves_no_routes_or_groups() {
2969 let (_source_tx, source_rx) = mpsc::channel(4);
2970 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2971 let mut tasks = Vec::new();
2972 for _ in 0..64 {
2973 let graph = graph.clone();
2974 tasks.push(tokio::spawn(async move {
2975 let (target, _receiver) = mpsc::channel(1);
2976 let route = graph.add_sink(codec("pcmu", 8_000), target).unwrap();
2977 assert!(graph.remove_sink(route));
2978 }));
2979 }
2980 for task in tasks {
2981 task.await.unwrap();
2982 }
2983 let snapshot = graph.snapshot().await;
2984 assert!(snapshot.sinks.is_empty());
2985 assert!(snapshot.codec_groups.is_empty());
2986 graph.shutdown();
2987 }
2988
2989 #[tokio::test]
2990 async fn source_close_closes_sinks_and_retains_final_snapshot() {
2991 let (source_tx, source_rx) = mpsc::channel(1);
2992 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
2993 let (target_tx, mut target_rx) = mpsc::channel(1);
2994 graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
2995 graph.snapshot().await;
2996 drop(source_tx);
2997
2998 wait_until(|| graph.abort_handle().is_finished()).await;
2999 assert!(target_rx.recv().await.is_none());
3000 let snapshot = graph.snapshot().await;
3001 assert_eq!(snapshot.source_state, MediaGraphSourceState::Closed);
3002 assert!(snapshot.sinks.is_empty());
3003 assert!(snapshot.codec_groups.is_empty());
3004 }
3005
3006 #[tokio::test]
3007 async fn empty_source_close_is_observed_before_first_sink() {
3008 let (source_tx, source_rx) = mpsc::channel(1);
3009 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3010 drop(source_tx);
3011 assert_eq!(
3012 graph.wait_closed().await.unwrap(),
3013 MediaGraphSourceState::Closed
3014 );
3015 assert!(graph.abort_handle().is_finished());
3016 }
3017
3018 #[tokio::test]
3019 async fn shutdown_and_abort_cleanup_all_sink_tasks() {
3020 let (_source_tx, source_rx) = mpsc::channel(1);
3021 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3022 let (target_tx, mut target_rx) = mpsc::channel(1);
3023 graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
3024 graph.snapshot().await;
3025 assert_eq!(
3026 graph.shutdown_and_wait().await.unwrap(),
3027 MediaGraphSourceState::Shutdown
3028 );
3029 assert!(graph.abort_handle().is_finished());
3030 assert!(target_rx.recv().await.is_none());
3031 assert_eq!(
3032 graph.snapshot().await.source_state,
3033 MediaGraphSourceState::Shutdown
3034 );
3035
3036 let (_source_tx, source_rx) = mpsc::channel(1);
3037 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3038 let (target_tx, mut target_rx) = mpsc::channel(1);
3039 graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
3040 graph.snapshot().await;
3041 graph.abort_handle().abort();
3042 assert_eq!(
3043 graph.wait_closed().await.unwrap(),
3044 MediaGraphSourceState::Aborted
3045 );
3046 assert!(graph.abort_handle().is_finished());
3047 assert!(target_rx.recv().await.is_none());
3048 let snapshot = graph.snapshot().await;
3049 assert_eq!(snapshot.source_state, MediaGraphSourceState::Aborted);
3050 assert!(snapshot.sinks.is_empty());
3051 assert!(snapshot.codec_groups.is_empty());
3052 }
3053
3054 #[tokio::test]
3055 async fn managed_routes_report_shutdown_source_close_and_abort() {
3056 let (_source_tx, source_rx) = mpsc::channel(1);
3058 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3059 let (target_tx, _target_rx) = mpsc::channel(1);
3060 let route = graph
3061 .add_managed_sink(codec("pcmu", 8_000), target_tx)
3062 .unwrap();
3063 let status = route.status();
3064 wait_for_route_state(&status, MediaGraphRouteState::Active).await;
3065 assert_eq!(
3066 graph.shutdown_and_wait().await.unwrap(),
3067 MediaGraphSourceState::Shutdown
3068 );
3069 assert_eq!(
3070 status.wait_terminal().await,
3071 MediaGraphRouteTerminalReason::GraphShutdown
3072 );
3073
3074 let (source_tx, source_rx) = mpsc::channel(1);
3076 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3077 let (target_tx, _target_rx) = mpsc::channel(1);
3078 let route = graph
3079 .add_managed_sink(codec("pcmu", 8_000), target_tx)
3080 .unwrap();
3081 let status = route.status();
3082 wait_for_route_state(&status, MediaGraphRouteState::Active).await;
3083 drop(source_tx);
3084 assert_eq!(
3085 status.wait_terminal().await,
3086 MediaGraphRouteTerminalReason::SourceClosed
3087 );
3088 assert_eq!(
3089 graph.wait_closed().await.unwrap(),
3090 MediaGraphSourceState::Closed
3091 );
3092
3093 let (_source_tx, source_rx) = mpsc::channel(1);
3096 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3097 let (target_tx, _target_rx) = mpsc::channel(1);
3098 let route = graph
3099 .add_managed_sink(codec("pcmu", 8_000), target_tx)
3100 .unwrap();
3101 let status = route.status();
3102 wait_for_route_state(&status, MediaGraphRouteState::Active).await;
3103 graph.abort_handle().abort();
3104 assert_eq!(
3105 status.wait_terminal().await,
3106 MediaGraphRouteTerminalReason::GraphAborted
3107 );
3108 let late = status.clone();
3109 assert_eq!(
3110 late.state(),
3111 MediaGraphRouteState::Terminal(MediaGraphRouteTerminalReason::GraphAborted)
3112 );
3113 assert_eq!(
3114 graph.wait_closed().await.unwrap(),
3115 MediaGraphSourceState::Aborted
3116 );
3117 }
3118
3119 #[tokio::test]
3120 async fn cloned_handles_observe_the_same_terminal_convergence() {
3121 let (_source_tx, source_rx) = mpsc::channel(1);
3122 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3123 let waiter_a = graph.clone();
3124 let waiter_b = graph.clone();
3125 let a = tokio::spawn(async move { waiter_a.wait_closed().await.unwrap() });
3126 let b = tokio::spawn(async move { waiter_b.wait_closed().await.unwrap() });
3127
3128 assert_eq!(
3129 graph.shutdown_and_wait().await.unwrap(),
3130 MediaGraphSourceState::Shutdown
3131 );
3132 assert_eq!(a.await.unwrap(), MediaGraphSourceState::Shutdown);
3133 assert_eq!(b.await.unwrap(), MediaGraphSourceState::Shutdown);
3134 }
3135
3136 #[tokio::test]
3137 async fn source_and_sink_codec_updates_are_independent() {
3138 let (_source_tx, source_rx) = mpsc::channel(1);
3139 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3140 let (target_tx, _target_rx) = mpsc::channel(1);
3141 let route = graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
3142 graph.snapshot().await;
3143
3144 graph
3145 .update_source_codec(codec("pcma", 8_000))
3146 .await
3147 .unwrap();
3148 let snapshot = graph.latest_snapshot();
3150 assert_eq!(snapshot.source_payload_type, 8);
3151 assert_eq!(snapshot.sinks[0].target_payload_type, 0);
3152 assert!(snapshot.codec_groups[0].transcoding);
3153
3154 graph
3155 .update_sink_codec(route.clone(), codec("opus", 48_000))
3156 .await
3157 .unwrap();
3158 let snapshot = graph.latest_snapshot();
3159 assert_eq!(snapshot.source_payload_type, 8);
3160 assert_eq!(snapshot.sinks[0].target_payload_type, 111);
3161 assert_eq!(snapshot.codec_groups[0].target_payload_type, 111);
3162
3163 graph.update_route(route, 0, 8).await.unwrap();
3164 let snapshot = graph.latest_snapshot();
3165 assert_eq!(snapshot.source_payload_type, 0);
3166 assert_eq!(snapshot.sinks[0].target_payload_type, 8);
3167 graph.shutdown();
3168 }
3169
3170 #[test]
3171 fn rtp_clock_translation_is_wrap_safe_in_both_directions() {
3172 let mut upsample = RtpClockTranslator::new(8_000, 48_000);
3173 let first = u32::MAX - 159;
3174 let upsampled = [
3175 upsample.translate(first),
3176 upsample.translate(0),
3177 upsample.translate(160),
3178 ];
3179 assert_eq!(upsampled[0], first);
3180 assert_eq!(upsampled[1].wrapping_sub(upsampled[0]), 960);
3181 assert_eq!(upsampled[2].wrapping_sub(upsampled[1]), 960);
3182
3183 let mut downsample = RtpClockTranslator::new(48_000, 8_000);
3184 let first = u32::MAX - 959;
3185 let downsampled = [
3186 downsample.translate(first),
3187 downsample.translate(0),
3188 downsample.translate(960),
3189 ];
3190 assert_eq!(downsampled[0], first);
3191 assert_eq!(downsampled[1].wrapping_sub(downsampled[0]), 160);
3192 assert_eq!(downsampled[2].wrapping_sub(downsampled[1]), 160);
3193 }
3194
3195 #[tokio::test]
3196 async fn each_codec_group_uses_its_own_rtp_clock() {
3197 let (source_tx, source_rx) = mpsc::channel(4);
3198 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3199 let (pcmu_tx, mut pcmu_rx) = mpsc::channel(4);
3200 let (opus_tx, mut opus_rx) = mpsc::channel(4);
3201 graph.add_sink(codec("pcmu", 8_000), pcmu_tx).unwrap();
3202 graph.add_sink(codec("opus", 48_000), opus_tx).unwrap();
3203 graph.snapshot().await;
3204
3205 source_tx.send(frame_at(0xff, 10_000)).await.unwrap();
3206 source_tx.send(frame_at(0xff, 10_160)).await.unwrap();
3207 let pcmu_first = pcmu_rx.recv().await.unwrap();
3208 let pcmu_second = pcmu_rx.recv().await.unwrap();
3209 let opus_first = opus_rx.recv().await.unwrap();
3210 let opus_second = opus_rx.recv().await.unwrap();
3211 assert_eq!(
3212 pcmu_second
3213 .timestamp_rtp
3214 .wrapping_sub(pcmu_first.timestamp_rtp),
3215 160
3216 );
3217 assert_eq!(
3218 opus_second
3219 .timestamp_rtp
3220 .wrapping_sub(opus_first.timestamp_rtp),
3221 960
3222 );
3223 graph.shutdown();
3224 }
3225
3226 #[tokio::test]
3227 async fn sink_rekey_preserves_timestamp_epoch_across_8k_and_48k_targets() {
3228 let (source_tx, source_rx) = mpsc::channel(8);
3229 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3230 let (target_tx, mut target_rx) = mpsc::channel(8);
3231 let route = graph.add_sink(codec("opus", 48_000), target_tx).unwrap();
3232 graph.snapshot().await;
3233
3234 source_tx.send(frame_at(0x7f, 10_000)).await.unwrap();
3235 source_tx.send(frame_at(0x7f, 10_160)).await.unwrap();
3236 let first = target_rx.recv().await.unwrap();
3237 let second = target_rx.recv().await.unwrap();
3238 assert_eq!(second.timestamp_rtp.wrapping_sub(first.timestamp_rtp), 960);
3239
3240 let mut rekeyed_opus = codec("opus", 48_000);
3241 rekeyed_opus.fmtp = Some("minptime=20;useinbandfec=1".into());
3242 graph
3243 .update_sink_codec(route.clone(), rekeyed_opus)
3244 .await
3245 .unwrap();
3246 source_tx.send(frame_at(0x7f, 10_320)).await.unwrap();
3247 let after_fmtp_rekey = target_rx.recv().await.unwrap();
3248 assert_eq!(
3249 after_fmtp_rekey
3250 .timestamp_rtp
3251 .wrapping_sub(second.timestamp_rtp),
3252 960
3253 );
3254
3255 graph
3256 .update_sink_codec(route, codec("pcmu", 8_000))
3257 .await
3258 .unwrap();
3259 source_tx.send(frame_at(0x7f, 10_480)).await.unwrap();
3260 let after_48k_to_8k = target_rx.recv().await.unwrap();
3261 assert_eq!(
3262 after_48k_to_8k
3263 .timestamp_rtp
3264 .wrapping_sub(after_fmtp_rekey.timestamp_rtp),
3265 160
3266 );
3267 graph.shutdown_and_wait().await.unwrap();
3268 }
3269
3270 #[tokio::test]
3271 async fn source_reconfiguration_preserves_target_clock_across_8k_48k_8k() {
3272 let (source_tx, source_rx) = mpsc::channel(8);
3273 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3274 let (target_tx, mut target_rx) = mpsc::channel(8);
3275 graph.add_sink(codec("opus", 48_000), target_tx).unwrap();
3276 graph.snapshot().await;
3277
3278 source_tx.send(frame_at(0x7f, 10_000)).await.unwrap();
3279 source_tx.send(frame_at(0x7f, 10_160)).await.unwrap();
3280 let first = target_rx.recv().await.unwrap();
3281 let second = target_rx.recv().await.unwrap();
3282 assert_eq!(second.timestamp_rtp.wrapping_sub(first.timestamp_rtp), 960);
3283
3284 graph
3285 .update_source_codec(codec("opus", 48_000))
3286 .await
3287 .unwrap();
3288 source_tx
3289 .send(frame_at_pt(0x11, 11_120, 111))
3290 .await
3291 .unwrap();
3292 let after_8k_to_48k = target_rx.recv().await.unwrap();
3293 assert_eq!(
3294 after_8k_to_48k
3295 .timestamp_rtp
3296 .wrapping_sub(second.timestamp_rtp),
3297 960
3298 );
3299
3300 graph
3301 .update_source_codec(codec("pcmu", 8_000))
3302 .await
3303 .unwrap();
3304 source_tx.send(frame_at(0x7f, 11_280)).await.unwrap();
3305 let after_48k_to_8k = target_rx.recv().await.unwrap();
3306 assert_eq!(
3307 after_48k_to_8k
3308 .timestamp_rtp
3309 .wrapping_sub(after_8k_to_48k.timestamp_rtp),
3310 960
3311 );
3312 graph.shutdown_and_wait().await.unwrap();
3313 }
3314
3315 #[tokio::test]
3316 async fn normalized_codec_identity_controls_grouping() {
3317 let (_source_tx, source_rx) = mpsc::channel(1);
3318 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3319 let mut opus_a = codec("OPUS", 48_000);
3320 opus_a.fmtp = Some("minptime=10; useinbandfec=1".into());
3321 let mut opus_a_reordered = codec("opus", 48_000);
3322 opus_a_reordered.fmtp = Some("useinbandfec=1;minptime=10".into());
3323 let mut opus_b = codec("opus", 48_000);
3324 opus_b.fmtp = Some("minptime=20;useinbandfec=1".into());
3325 let (a_tx, _a_rx) = mpsc::channel(1);
3326 let (b_tx, _b_rx) = mpsc::channel(1);
3327 let (c_tx, _c_rx) = mpsc::channel(1);
3328 graph.add_sink(opus_a, a_tx).unwrap();
3329 graph.add_sink(opus_a_reordered, b_tx).unwrap();
3330 graph.add_sink(opus_b, c_tx).unwrap();
3331
3332 let mut group_sizes: Vec<_> = graph
3333 .snapshot()
3334 .await
3335 .codec_groups
3336 .into_iter()
3337 .map(|group| group.sink_routes.len())
3338 .collect();
3339 group_sizes.sort_unstable();
3340 assert_eq!(group_sizes, vec![1, 2]);
3341 graph.shutdown();
3342 }
3343
3344 #[test]
3345 fn configured_transcoder_honors_canonical_opus_mono() {
3346 let source = codec("pcmu", 8_000);
3347 let mut target = codec("opus", 48_000);
3348 target.fmtp = Some("maxaveragebitrate=32000;cbr=1".into());
3349 let session = ConfiguredTranscodingSession::new(&source, 0, &target, 111).unwrap();
3350 let source_info = session.source_codec.get_info();
3351 let target_info = session.target_codec.get_info();
3352 assert_eq!(source_info.sample_rate, 8_000);
3353 assert_eq!(source_info.channels, 1);
3354 assert_eq!(target_info.sample_rate, 48_000);
3355 assert_eq!(target_info.channels, 1);
3356 }
3357
3358 #[test]
3359 fn configured_transcoder_converts_pcm_s16le_and_g711() {
3360 let pcm = codec("pcm_s16le", 16_000);
3361 let linear = (0..320)
3362 .map(|sample| ((sample as f32 * 0.2).sin() * 10_000.0) as i16)
3363 .flat_map(i16::to_le_bytes)
3364 .collect::<Vec<_>>();
3365
3366 for (g711_name, g711_pt) in [("pcmu", 0), ("pcma", 8)] {
3367 let g711 = codec(g711_name, 8_000);
3368 let mut to_g711 =
3369 ConfiguredTranscodingSession::new(&pcm, PCM_S16LE, &g711, g711_pt).unwrap();
3370 let encoded = to_g711.transcode(&linear).unwrap();
3371 assert_eq!(encoded.len(), 160);
3372
3373 let mut to_pcm =
3374 ConfiguredTranscodingSession::new(&g711, g711_pt, &pcm, PCM_S16LE).unwrap();
3375 assert_eq!(to_pcm.transcode(&encoded).unwrap().len(), 640);
3376 }
3377 }
3378
3379 #[test]
3380 fn configured_transcoder_preserves_pcm_wideband_path_through_opus() {
3381 let pcm = codec("pcm_s16le", 16_000);
3382 let opus = codec("opus", 48_000);
3383 let linear = (0..320)
3384 .map(|sample| ((sample as f32 * 0.2).sin() * 10_000.0) as i16)
3385 .flat_map(i16::to_le_bytes)
3386 .collect::<Vec<_>>();
3387
3388 let mut to_opus = ConfiguredTranscodingSession::new(&pcm, PCM_S16LE, &opus, 111).unwrap();
3389 let encoded = to_opus.transcode(&linear).unwrap();
3390 assert!(!encoded.is_empty());
3391
3392 let mut to_pcm = ConfiguredTranscodingSession::new(&opus, 111, &pcm, PCM_S16LE).unwrap();
3393 assert_eq!(to_pcm.transcode(&encoded).unwrap().len(), 640);
3394 }
3395
3396 #[tokio::test]
3397 async fn compatibility_update_rejects_unknown_payload_types_atomically() {
3398 let (_source_tx, source_rx) = mpsc::channel(1);
3399 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3400 let (target_tx, _target_rx) = mpsc::channel(1);
3401 let route = graph.add_sink(codec("pcma", 8_000), target_tx).unwrap();
3402 let before = graph.snapshot().await;
3403
3404 assert!(matches!(
3405 graph.update_route(route.clone(), 127, 8).await,
3406 Err(RvoipError::UnsupportedCodec(_))
3407 ));
3408 assert!(matches!(
3409 graph.update_route(route, 0, 127).await,
3410 Err(RvoipError::UnsupportedCodec(_))
3411 ));
3412 let after = graph.snapshot().await;
3413 assert_eq!(after.source_codec, before.source_codec);
3414 assert_eq!(after.source_payload_type, before.source_payload_type);
3415 assert_eq!(after.sinks[0].target_codec, before.sinks[0].target_codec);
3416 assert_eq!(
3417 after.sinks[0].target_payload_type,
3418 before.sinks[0].target_payload_type
3419 );
3420 graph.shutdown();
3421 }
3422
3423 #[tokio::test]
3424 async fn frame_hot_path_does_not_rebuild_retained_snapshot() {
3425 let (source_tx, source_rx) = mpsc::channel(2);
3426 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), Default::default()).unwrap();
3427 let (target_tx, mut target_rx) = mpsc::channel(2);
3428 graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
3429 let baseline = graph.snapshot_arc().await;
3430
3431 source_tx.send(frame(1)).await.unwrap();
3432 target_rx.recv().await.unwrap();
3433 let retained = graph.latest_snapshot_arc();
3434 assert!(Arc::ptr_eq(&baseline, &retained));
3435
3436 let explicit = graph.snapshot_arc().await;
3437 assert!(!Arc::ptr_eq(&baseline, &explicit));
3438 assert_eq!(explicit.source_frames, 1);
3439 graph.shutdown();
3440 }
3441
3442 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3443 async fn snapshot_flood_is_coalesced_without_starving_media() {
3444 let policy = MediaGraphPolicy {
3445 sink_queue_frames: 64,
3446 ..Default::default()
3447 };
3448 let (source_tx, source_rx) = mpsc::channel(64);
3449 let graph = start_media_graph(source_rx, codec("pcmu", 8_000), policy).unwrap();
3450 let (target_tx, mut target_rx) = mpsc::channel(64);
3451 graph.add_sink(codec("pcmu", 8_000), target_tx).unwrap();
3452 graph.snapshot().await;
3453
3454 let mut readers = Vec::new();
3455 for _ in 0..32 {
3456 let graph = graph.clone();
3457 readers.push(tokio::spawn(async move {
3458 for _ in 0..100 {
3459 let _ = graph.snapshot_arc().await;
3460 }
3461 }));
3462 }
3463 for value in 0..20 {
3464 source_tx.send(frame(value)).await.unwrap();
3465 }
3466 tokio::time::timeout(Duration::from_secs(2), async {
3467 for _ in 0..20 {
3468 target_rx.recv().await.expect("media sink closed");
3469 }
3470 })
3471 .await
3472 .expect("snapshot traffic starved media");
3473 for reader in readers {
3474 reader.await.unwrap();
3475 }
3476 assert_eq!(graph.snapshot().await.source_frames, 20);
3477 graph.shutdown();
3478 }
3479}