Skip to main content

rvoip_core/
media_graph.rs

1//! Bounded, one-source-to-many real-time media routing.
2//!
3//! A `MediaStream::frames_in()` receiver is intentionally single-take. The
4//! media graph owns that receiver once and exposes dynamic sink routes so a
5//! call peer, recorder, UCTP publisher, and MOQT publisher can observe the
6//! same source without racing for frames.
7
8use 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);
37/// Maximum publication cadence for retained source-activity observations.
38pub 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
43/// Default graph fanout ceiling. This admits the 1,000-listener direct UCTP
44/// target while retaining 24 routes for the call peer, recorders, and other
45/// operational observers.
46pub const DEFAULT_MEDIA_GRAPH_MAX_SINKS: usize = 1_024;
47
48/// Stable identifier for the lifetime of a media graph.
49///
50/// It is intentionally defined here rather than in the shared ID vocabulary:
51/// a graph is an rvoip-core runtime concern, not a cross-adapter wire type.
52#[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    /// Maximum installed or queued sink routes. Admission is reserved before
90    /// an add command enters the bounded control queue, so concurrent callers
91    /// cannot oversubscribe the graph.
92    pub max_sinks: usize,
93    pub sink_queue_frames: usize,
94    /// Frames retained before the first sink is registered. The buffer is
95    /// always bounded and drops its oldest frame on overflow.
96    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/// Latest coalesced proof that the graph consumed media from its single
116/// authoritative source receiver.
117///
118/// `source_frames` is diagnostic only. Consumers that persist activity should
119/// assign their own consecutive delivery generation because this retained
120/// value can skip counts while they are backpressured.
121#[derive(Clone, Debug, Eq, PartialEq)]
122pub struct MediaGraphActivityObservation {
123    pub source_frames: u64,
124    pub observed_at: DateTime<Utc>,
125}
126
127/// Current state of the graph's single source receiver.
128#[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/// Why a managed media route reached its terminal state.
144#[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/// Latest lifecycle state of a managed media route.
156///
157/// This is delivered through a Tokio watch channel, so observers retain one
158/// bounded latest value rather than an unbounded event backlog.
159#[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/// Cloneable lifecycle observer for one graph sink route.
168///
169/// Status observers do not own graph membership. Cloning this value never
170/// extends the route's lifetime after its [`ManagedMediaRoute`] owner drops.
171#[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    /// Wait until the actor has installed the route. If the graph terminates
196    /// first, return the retained terminal reason instead of hanging.
197    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/// Owning lease for a managed graph route.
225///
226/// Dropping this lease signals owner cancellation and best-effort queues a
227/// fast removal command. The actor also observes cancellation during bounded
228/// frame/periodic maintenance, so queue saturation cannot retain the route.
229/// Callers may clone [`Self::status`] freely without keeping the route alive.
230#[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    /// Remove this route with actor acknowledgement. Clone [`Self::status`]
260    /// first when the caller also wants to observe the retained terminal
261    /// reason after consuming the owning lease.
262    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            // The bounded control queue is only a fast path. The actor also
291            // observes this signal during frame processing and its periodic
292            // maintenance tick, so a full queue cannot retain an orphan.
293            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/// Last-known state for an active sink route.
303#[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/// Codec-group state. A source frame is transcoded at most once per group and
318/// the resulting immutable payload is cloned cheaply into every member sink.
319#[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/// Point-in-time operational view of a graph.
338///
339/// The leading concurrent `MediaGraphHandle::snapshot` request places a
340/// barrier on the graph command queue; concurrent followers coalesce onto the
341/// retained view. `latest_snapshot_arc` is a cheap non-blocking last-known view
342/// and remains available after the graph has stopped.
343#[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    /// Compatibility command for the original payload-type based API.
384    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    /// Subscribe to a bounded retained source-activity observation.
481    ///
482    /// The graph publishes at most one value per configured observation
483    /// interval and overwrites an unread value with the newest one. A slow
484    /// observer therefore cannot stall or grow memory in the media path.
485    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    /// Add a sink and retain a bounded, cloneable lifecycle observer for it.
499    /// The legacy [`Self::add_sink`] API remains a route-ID-only wrapper.
500    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    /// Queue a removal command. The return value reports whether the command
549    /// was accepted, not whether the route existed. Use
550    /// `remove_sink_and_wait` when route-existence acknowledgement matters.
551    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    /// Update the source codec and rebuild every codec group's transcoder.
574    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    /// Move one sink to the codec group represented by `codec`.
588    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    /// Compatibility wrapper for callers that still renegotiate with RTP
603    /// payload types. New code should call `update_source_codec` and
604    /// `update_sink_codec` independently.
605    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    /// Return a command-barrier-consistent snapshot when this call leads a
629    /// snapshot batch, a coalesced retained snapshot when another request is
630    /// already pending, or the final retained snapshot after shutdown.
631    pub async fn snapshot(&self) -> MediaGraphSnapshot {
632        (*self.snapshot_arc().await).clone()
633    }
634
635    /// Arc-returning snapshot API for high-frequency diagnostics. Concurrent
636    /// callers coalesce behind at most one actor request; followers receive the
637    /// retained snapshot instead of growing the control queue.
638    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    /// Return the most recently published snapshot without waiting on the
657    /// graph actor.
658    pub fn latest_snapshot(&self) -> MediaGraphSnapshot {
659        (*self.latest_snapshot_arc()).clone()
660    }
661
662    /// Cheap retained-snapshot access. The read lock is held only long enough
663    /// to clone an Arc; callers never contend while inspecting its contents.
664    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            // A saturated control plane must not make shutdown impossible.
674            self.abort.abort();
675        }
676    }
677
678    /// Request graceful shutdown and wait for both the graph actor and every
679    /// sink-forwarding task to converge on a terminal state.
680    pub async fn shutdown_and_wait(&self) -> Result<MediaGraphSourceState> {
681        self.shutdown();
682        self.wait_closed().await
683    }
684
685    /// Wait for graph and sink-task convergence without initiating shutdown.
686    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
741/// Validate codec identity before transferring ownership of a stream's
742/// single-consumer receiver into a graph.
743pub 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    /// Enqueue without awaiting a slow sink. The oldest queued frame is
788    /// discarded when full so the sink always sees the freshest media.
789    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    /// Releasing the runtime route returns one slot to the graph-wide
848    /// admission budget. Pending add commands hold the same kind of permit.
849    _admission: SinkAdmissionPermit,
850    /// RTP clock ownership is per route, not per codec group. This lets a
851    /// route retain its timestamp epoch when fmtp/payload changes move it
852    /// between groups while payload transcoding remains shared by the group.
853    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
1205/// Balances process-wide gauges by delta. A graph must never `set` a gauge to
1206/// its local count because doing so corrupts the aggregate when graphs overlap.
1207struct 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
1248/// Ensures an aborted actor retains an accurate terminal snapshot.
1249struct 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            // Runtime teardown itself guarantees the sink tasks can no longer
1298            // execute; publish terminal state for any remaining synchronous
1299            // observers.
1300            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
1348/// Start a media graph task that owns `source` for its lifetime.
1349pub 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    // Construct the guard before spawning. If the task is aborted before its
1408    // first poll, dropping the unpolled future still terminalizes every route.
1409    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        // Once a first sink has existed, later zero-sink periods deliberately
1426        // discard media rather than replaying stale RTP to a future attachment.
1427        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        // Consume interval's immediate first tick; the initial retained
1436        // snapshot was already published before spawning the actor.
1437        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                            // Preserve the original API's success-on-missing-route behavior.
1667                            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                            // The waiting caller may have timed out while media
1694                            // was being processed. Skip stale diagnostic work.
1695                            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                    // A managed owner may have been dropped while the bounded
1744                    // control queue was full. Prune before routing the next
1745                    // frame so an orphan never receives further media.
1746                    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                    // Periodic maintenance guarantees convergence even when
1847                    // the source is idle and Drop could not enqueue Remove.
1848                    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        // Preserve the last accepted frame even when the source closes or a
1890        // graceful graph shutdown happens before the next coalescing tick.
1891        // Lifecycle ownership is revalidated by the Orchestrator before this
1892        // retained observation can become an authoritative event.
1893        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
1981/// Route one already-accounted source frame. Payload transcoding remains once
1982/// per codec group, while each sink advances its own RTP clock so a group
1983/// re-key cannot reset that route's timestamp epoch.
1984fn 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
2150/// Remove managed routes whose sole owner has cancelled them. The scan is
2151/// bounded by `MediaGraphPolicy::max_sinks`; it never creates a fallback task
2152/// or channel when the control queue is saturated.
2153fn 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        // Give the actor opportunities to race the already-ready source. The
2480        // initial-route gate must keep the frame buffered.
2481        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        // Establish that all six ready source frames reached the bounded
2509        // pre-sink buffer before registering the first sink.
2510        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        // This test runs on Tokio's current-thread runtime. With no await in
2634        // this loop, the actor cannot drain commands while every bounded slot
2635        // is filled with a live snapshot request.
2636        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's Remove fast path is necessarily rejected, but cancellation
2650        // remains visible to the actor. Frame activity must prune the route
2651        // before routing this next packet.
2652        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        // Exactly 25% is retained, not evicted.
2837        assert!(!sink.record_offer(start + Duration::from_secs(3), false, &policy));
2838        assert_eq!(sink.rolling_drop_counts(), (4, 1));
2839        // The sample at exactly the ten-second boundary remains in the window.
2840        assert!(!sink.record_offer(start + Duration::from_secs(10), false, &policy));
2841        // One nanosecond later the original drop is pruned, so the new drop is
2842        // only one of the five current samples.
2843        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        // Two of six current samples is 33%, which is strictly over policy.
2850        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        // Graceful graph shutdown.
3057        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        // Natural source closure.
3075        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        // Forced actor abort, including a status clone subscribing after the
3094        // terminal update has already been retained.
3095        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        // Update acknowledgement is also a retained-snapshot barrier.
3149        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}