Skip to main content

active_call/media/
stream.rs

1use crate::event::{EventSender, SessionEvent};
2use crate::media::ambiance::{AmbianceOption, AmbianceProcessor};
3use crate::media::dtmf::DtmfDetector;
4use crate::media::volume_control::HoldProcessor;
5use crate::media::{AudioFrame, INTERNAL_SAMPLERATE, Samples, TrackId};
6use crate::media::{
7    processor::Processor,
8    recorder::{Recorder, RecorderOption},
9    track::{Track, TrackPacketReceiver, TrackPacketSender},
10};
11use anyhow::Result;
12use std::collections::{HashMap, HashSet};
13use std::path::Path;
14use std::sync::{
15    Arc, Mutex as StdMutex,
16    atomic::{AtomicBool, AtomicU64, Ordering},
17};
18use std::time::Duration;
19use tokio::task::JoinHandle;
20use tokio::{
21    select,
22    sync::{Mutex, mpsc},
23};
24use tokio_util::sync::CancellationToken;
25use tracing::{debug, info, warn};
26use uuid;
27
28pub struct MediaStream {
29    id: String,
30    pub cancel_token: CancellationToken,
31    recorder_option: Mutex<Option<RecorderOption>>,
32    tracks: Mutex<HashMap<TrackId, (Box<dyn Track>, DtmfDetector)>>,
33    /// Trickle ICE candidates that arrived before any track existed to feed
34    /// them to. Drained into the next track started.
35    pending_ice_candidates: Mutex<Vec<(String, Option<String>, Option<u32>)>>,
36    suppressed_sources: Mutex<HashSet<TrackId>>,
37    event_sender: EventSender,
38    pub packet_sender: TrackPacketSender,
39    packet_receiver: Mutex<Option<TrackPacketReceiver>>,
40    recorder_sender: mpsc::UnboundedSender<AudioFrame>,
41    recorder_receiver: Mutex<Option<mpsc::UnboundedReceiver<AudioFrame>>>,
42    recorder_handle: Mutex<Option<JoinHandle<()>>>,
43    /// True once a recorder task is actually consuming `recorder_sender`.
44    /// The ambiance idle loop mirrors its frames to the recorder only when
45    /// this is set, avoiding unbounded buffering when recording starts late.
46    recording_active: Arc<AtomicBool>,
47    ambiance: Mutex<Option<Arc<StdMutex<AmbianceProcessor>>>>,
48    ambiance_source_id: StdMutex<Option<TrackId>>,
49    last_server_packet_ts: Arc<AtomicU64>,
50    ambiance_idle_started: AtomicBool,
51}
52
53const CALLEE_TRACK_ID: &str = "callee-track";
54const QUEUE_HOLD_TRACK_ID: &str = "queue-hold-track";
55pub const SERVER_SIDE_TRACK_ID: &str = "server-side-track";
56const AMBIANCE_IDLE_TRACK_ID: &str = "ambiance-track";
57const AMBIANCE_IDLE_PTIME: Duration = Duration::from_millis(20);
58// Skip idle fill if server-side audio arrived within this window (TTS ptime is 20ms).
59const AMBIANCE_IDLE_GAP_MS: u64 = 25;
60
61pub struct MediaStreamBuilder {
62    cancel_token: Option<CancellationToken>,
63    id: Option<String>,
64    event_sender: EventSender,
65    recorder_config: Option<RecorderOption>,
66}
67
68impl MediaStreamBuilder {
69    pub fn new(event_sender: EventSender) -> Self {
70        Self {
71            id: Some(format!("ms:{}", uuid::Uuid::new_v4())),
72            cancel_token: None,
73            event_sender,
74            recorder_config: None,
75        }
76    }
77    pub fn with_id(mut self, id: String) -> Self {
78        self.id = Some(id);
79        self
80    }
81
82    pub fn with_cancel_token(mut self, cancel_token: CancellationToken) -> Self {
83        self.cancel_token = Some(cancel_token);
84        self
85    }
86
87    pub fn with_recorder_config(mut self, recorder_config: RecorderOption) -> Self {
88        self.recorder_config = Some(recorder_config);
89        self
90    }
91
92    pub fn build(self) -> MediaStream {
93        let cancel_token = self
94            .cancel_token
95            .unwrap_or_else(|| CancellationToken::new());
96        let tracks = Mutex::new(HashMap::new());
97        let (track_packet_sender, track_packet_receiver) = mpsc::unbounded_channel();
98        let (recorder_sender, recorder_receiver) = mpsc::unbounded_channel();
99        MediaStream {
100            id: self.id.unwrap_or_default(),
101            cancel_token,
102            recorder_option: Mutex::new(self.recorder_config),
103            tracks,
104            pending_ice_candidates: Mutex::new(Vec::new()),
105            suppressed_sources: Mutex::new(HashSet::new()),
106            event_sender: self.event_sender,
107            packet_sender: track_packet_sender,
108            packet_receiver: Mutex::new(Some(track_packet_receiver)),
109            recorder_sender,
110            recorder_receiver: Mutex::new(Some(recorder_receiver)),
111            recorder_handle: Mutex::new(None),
112            recording_active: Arc::new(AtomicBool::new(false)),
113            ambiance: Mutex::new(None),
114            ambiance_source_id: StdMutex::new(None),
115            last_server_packet_ts: Arc::new(AtomicU64::new(0)),
116            ambiance_idle_started: AtomicBool::new(false),
117        }
118    }
119}
120
121impl MediaStream {
122    pub async fn serve(&self) -> Result<()> {
123        let packet_receiver = match self.packet_receiver.lock().await.take() {
124            Some(receiver) => receiver,
125            None => {
126                warn!(
127                    session_id = self.id,
128                    "MediaStream::serve() called multiple times, stream already serving"
129                );
130                return Ok(());
131            }
132        };
133        self.start_recorder().await.ok();
134        info!(session_id = self.id, "mediastream serving");
135        select! {
136            _ = self.cancel_token.cancelled() => {}
137            r = self.handle_forward_track(packet_receiver) => {
138                info!(session_id = self.id, "track packet receiver stopped {:?}", r);
139            }
140        }
141        Ok(())
142    }
143
144    /// Load ambiance once per call and keep mixing it while TTS/file playback is silent.
145    pub async fn ensure_ambiance(
146        &self,
147        option: AmbianceOption,
148        source_track_id: TrackId,
149    ) -> Result<Option<Arc<StdMutex<AmbianceProcessor>>>> {
150        let mut slot = self.ambiance.lock().await;
151        if let Some(existing) = slot.as_ref() {
152            return Ok(Some(existing.clone()));
153        }
154
155        if option.path.is_none() || option.enabled == Some(false) {
156            return Ok(None);
157        }
158
159        let processor = AmbianceProcessor::new(option).await?;
160        let shared = Arc::new(StdMutex::new(processor));
161        *slot = Some(shared.clone());
162        drop(slot);
163
164        *self.ambiance_source_id.lock().unwrap() = Some(source_track_id);
165        self.start_ambiance_idle_loop(shared.clone());
166        info!(session_id = self.id, "ambiance idle mixer started");
167        Ok(Some(shared))
168    }
169
170    fn start_ambiance_idle_loop(&self, processor: Arc<StdMutex<AmbianceProcessor>>) {
171        if self.ambiance_idle_started.swap(true, Ordering::SeqCst) {
172            return;
173        }
174
175        let cancel_token = self.cancel_token.clone();
176        let packet_sender = self.packet_sender.clone();
177        let last_server_packet_ts = self.last_server_packet_ts.clone();
178        let session_id = self.id.clone();
179        let recorder_sender = self.recorder_sender.clone();
180        let recording_active = self.recording_active.clone();
181
182        crate::spawn(async move {
183            let mut ticker = tokio::time::interval(AMBIANCE_IDLE_PTIME);
184            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
185            loop {
186                tokio::select! {
187                    _ = cancel_token.cancelled() => break,
188                    _ = ticker.tick() => {
189                        let now = crate::media::get_timestamp();
190                        let last = last_server_packet_ts.load(Ordering::Relaxed);
191                        if last != 0 && now.saturating_sub(last) < AMBIANCE_IDLE_GAP_MS {
192                            continue;
193                        }
194
195                        let mut frame = AudioFrame {
196                            track_id: AMBIANCE_IDLE_TRACK_ID.to_string(),
197                            samples: Samples::Empty,
198                            timestamp: now,
199                            sample_rate: INTERNAL_SAMPLERATE,
200                            channels: 1,
201                            ..Default::default()
202                        };
203                        {
204                            let mut ambiance = match processor.lock() {
205                                Ok(guard) => guard,
206                                Err(_) => break,
207                            };
208                            if let Err(e) = ambiance.process_frame(&mut frame) {
209                                warn!(session_id, "ambiance idle mix failed: {}", e);
210                                continue;
211                            }
212                        }
213                        if matches!(frame.samples, Samples::Empty) {
214                            continue;
215                        }
216                        // Mirror the mixed idle ambiance into the recording so the
217                        // stereo WAV's server-side channel reflects what the user
218                        // actually heard, including silent-period background audio.
219                        if recording_active.load(Ordering::SeqCst) {
220                            let mut recorded = frame.clone();
221                            recorded.track_id = SERVER_SIDE_TRACK_ID.to_string();
222                            let _ = recorder_sender.send(recorded);
223                        }
224                        if packet_sender.send(frame).is_err() {
225                            debug!(session_id, "ambiance idle sender closed");
226                            break;
227                        }
228                    }
229                }
230            }
231        });
232    }
233
234    pub fn stop(&self, _reason: Option<String>, _initiator: Option<String>) {
235        self.cancel_token.cancel()
236    }
237
238    pub async fn cleanup(&self) -> Result<()> {
239        self.cancel_token.cancel();
240        {
241            let mut tracks = self.tracks.lock().await;
242            for (id, (track, _)) in tracks.drain() {
243                if let Err(e) = track.stop().await {
244                    warn!(session_id = self.id, track_id = %id, "failed to stop track during cleanup: {}", e);
245                }
246            }
247        }
248        self.suppressed_sources.lock().await.clear();
249
250        if let Some(recorder_handle) = self.recorder_handle.lock().await.take() {
251            if let Ok(Ok(_)) = tokio::time::timeout(Duration::from_secs(30), recorder_handle).await
252            {
253                info!(session_id = self.id, "recorder stopped");
254            } else {
255                warn!(session_id = self.id, "recorder timeout");
256            }
257        }
258        Ok(())
259    }
260    pub async fn track_count(&self) -> usize {
261        self.tracks.lock().await.len()
262    }
263
264    pub async fn update_recorder_option(&self, recorder_config: RecorderOption) {
265        *self.recorder_option.lock().await = Some(recorder_config);
266        self.start_recorder().await.ok();
267    }
268
269    pub async fn remove_track(&self, id: &TrackId, graceful: bool) {
270        let track_entry = { self.tracks.lock().await.remove(id) };
271        if let Some((track, _)) = track_entry {
272            self.suppressed_sources.lock().await.remove(id);
273            let res = if !graceful {
274                track.stop().await
275            } else {
276                track.stop_graceful().await
277            };
278            match res {
279                Ok(_) => {}
280                Err(e) => {
281                    warn!(session_id = self.id, "failed to stop track: {}", e);
282                }
283            }
284        }
285    }
286    pub async fn update_remote_description(
287        &self,
288        track_id: &TrackId,
289        answer: &String,
290    ) -> Result<()> {
291        let track_entry = { self.tracks.lock().await.remove(track_id) };
292        if let Some((mut track, dtmf)) = track_entry {
293            let res = track.update_remote_description(answer).await;
294            self.tracks
295                .lock()
296                .await
297                .insert(track_id.clone(), (track, dtmf));
298            res?;
299        }
300        Ok(())
301    }
302
303    pub async fn update_remote_description_force(
304        &self,
305        track_id: &TrackId,
306        answer: &String,
307    ) -> Result<()> {
308        let track_entry = { self.tracks.lock().await.remove(track_id) };
309        if let Some((mut track, dtmf)) = track_entry {
310            let res = track.update_remote_description_force(answer).await;
311            self.tracks
312                .lock()
313                .await
314                .insert(track_id.clone(), (track, dtmf));
315            res?;
316        }
317        Ok(())
318    }
319
320    /// Apply a provisional remote description (SIP 183 early media). See
321    /// `Track::update_remote_description_provisional`.
322    pub async fn update_remote_description_provisional(
323        &self,
324        track_id: &TrackId,
325        answer: &String,
326    ) -> Result<()> {
327        let track_entry = { self.tracks.lock().await.remove(track_id) };
328        if let Some((mut track, dtmf)) = track_entry {
329            let res = track.update_remote_description_provisional(answer).await;
330            self.tracks
331                .lock()
332                .await
333                .insert(track_id.clone(), (track, dtmf));
334            res?;
335        }
336        Ok(())
337    }
338
339    pub async fn handshake(
340        &self,
341        track_id: &TrackId,
342        offer: String,
343        timeout: Option<Duration>,
344    ) -> Result<String> {
345        let track_entry = { self.tracks.lock().await.remove(track_id) };
346        if let Some((mut track, dtmf)) = track_entry {
347            let res = track.handshake(offer, timeout).await;
348            self.tracks
349                .lock()
350                .await
351                .insert(track_id.clone(), (track, dtmf));
352            res
353        } else {
354            anyhow::bail!("track not found: {}", track_id)
355        }
356    }
357
358    pub async fn update_track(&self, mut track: Box<dyn Track>, play_id: Option<String>) {
359        self.remove_track(track.id(), false).await;
360        if self.recorder_option.lock().await.is_some() {
361            track.append_processor(Box::new(RecorderProcessor::new(
362                self.recorder_sender.clone(),
363            )));
364        }
365        match track
366            .start(self.event_sender.clone(), self.packet_sender.clone())
367            .await
368        {
369            Ok(_) => {
370                info!(session_id = self.id, track_id = track.id(), "track started");
371                let track_id = track.id().clone();
372                if track_id.as_str() == self.id.as_str() {
373                    let pending = std::mem::take(&mut *self.pending_ice_candidates.lock().await);
374                    for (candidate, sdp_mid, sdp_mline_index) in pending {
375                        if let Err(e) =
376                            track.add_ice_candidate(&candidate, sdp_mid.as_deref(), sdp_mline_index)
377                        {
378                            warn!(
379                                session_id = self.id,
380                                track_id = track.id(),
381                                "failed to apply buffered ICE candidate: {}",
382                                e
383                            );
384                        }
385                    }
386                }
387                self.tracks
388                    .lock()
389                    .await
390                    .insert(track_id.clone(), (track, DtmfDetector::new()));
391                self.event_sender
392                    .send(SessionEvent::TrackStart {
393                        track_id,
394                        timestamp: crate::media::get_timestamp(),
395                        play_id,
396                    })
397                    .ok();
398            }
399            Err(e) => {
400                warn!(
401                    session_id = self.id,
402                    track_id = track.id(),
403                    play_id = play_id.as_deref(),
404                    "Failed to start track: {}",
405                    e
406                );
407            }
408        }
409    }
410
411    pub async fn mute_track(&self, id: Option<TrackId>) {
412        if let Some(id) = id {
413            if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
414                MuteProcessor::mute_track(track.as_mut());
415            }
416        } else {
417            for (track, _) in self.tracks.lock().await.values_mut() {
418                MuteProcessor::mute_track(track.as_mut());
419            }
420        }
421    }
422
423    pub async fn unmute_track(&self, id: Option<TrackId>) {
424        if let Some(id) = id {
425            if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
426                MuteProcessor::unmute_track(track.as_mut());
427            }
428        } else {
429            for (track, _) in self.tracks.lock().await.values_mut() {
430                MuteProcessor::unmute_track(track.as_mut());
431            }
432        }
433    }
434
435    /// Trickle ICE: feed a remote candidate into the ICE-backed (WebRTC)
436    /// track, if it's up yet, or buffer it for `update_track` to replay
437    /// otherwise.
438    pub async fn add_ice_candidate(
439        &self,
440        candidate: &str,
441        sdp_mid: Option<&str>,
442        sdp_mline_index: Option<u32>,
443    ) -> Result<()> {
444        let tracks = self.tracks.lock().await;
445        if let Some((track, _)) = tracks.get(self.id.as_str()) {
446            track.add_ice_candidate(candidate, sdp_mid, sdp_mline_index)?;
447            return Ok(());
448        }
449        drop(tracks);
450        self.pending_ice_candidates.lock().await.push((
451            candidate.to_string(),
452            sdp_mid.map(|s| s.to_string()),
453            sdp_mline_index,
454        ));
455        Ok(())
456    }
457
458    pub async fn pause_playback(&self, id: TrackId) -> Result<()> {
459        self.set_playback_paused(id, true).await
460    }
461
462    pub async fn resume_playback(&self, id: TrackId) -> Result<()> {
463        self.set_playback_paused(id, false).await
464    }
465
466    async fn set_playback_paused(&self, id: TrackId, paused: bool) -> Result<()> {
467        if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
468            if track.set_paused(paused) {
469                Ok(())
470            } else {
471                warn!(
472                    session_id = self.id,
473                    track_id = %id,
474                    paused,
475                    "pause state requested for track that does not support pausing"
476                );
477                Err(anyhow::anyhow!("track does not support pausing: {}", id))
478            }
479        } else {
480            warn!(
481                session_id = self.id,
482                track_id = %id,
483                paused,
484                "pause state requested for unknown track"
485            );
486            Err(anyhow::anyhow!("track not found: {}", id))
487        }
488    }
489
490    pub async fn hold_track(&self, id: Option<TrackId>) {
491        if let Some(id) = id {
492            if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
493                HoldTrack::hold_track(track.as_mut());
494            }
495        } else {
496            for (track, _) in self.tracks.lock().await.values_mut() {
497                HoldTrack::hold_track(track.as_mut());
498            }
499        }
500    }
501
502    pub async fn resume_track(&self, id: Option<TrackId>) {
503        if let Some(id) = id {
504            if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
505                HoldTrack::resume_track(track.as_mut());
506            }
507        } else {
508            for (track, _) in self.tracks.lock().await.values_mut() {
509                HoldTrack::resume_track(track.as_mut());
510            }
511        }
512    }
513
514    pub async fn suppress_forwarding(&self, track_id: &TrackId) {
515        self.suppressed_sources
516            .lock()
517            .await
518            .insert(track_id.clone());
519    }
520
521    pub async fn resume_forwarding(&self, track_id: &TrackId) {
522        self.suppressed_sources.lock().await.remove(track_id);
523    }
524
525    pub async fn remove_processor<T: 'static>(&self, track_id: &TrackId) -> Result<()> {
526        if let Some((track, _)) = self.tracks.lock().await.get_mut(track_id) {
527            track.as_mut().processor_chain().remove_processor::<T>();
528            Ok(())
529        } else {
530            Err(anyhow::anyhow!("Track {} not found", track_id))
531        }
532    }
533
534    pub async fn append_processor(
535        &self,
536        track_id: &TrackId,
537        processor: Box<dyn crate::media::processor::Processor>,
538    ) -> Result<()> {
539        if let Some((track, _)) = self.tracks.lock().await.get_mut(track_id) {
540            track.as_mut().processor_chain().append_processor(processor);
541            Ok(())
542        } else {
543            Err(anyhow::anyhow!("Track {} not found", track_id))
544        }
545    }
546}
547
548#[derive(Clone)]
549pub struct RecorderProcessor {
550    sender: mpsc::UnboundedSender<AudioFrame>,
551}
552
553impl RecorderProcessor {
554    pub fn new(sender: mpsc::UnboundedSender<AudioFrame>) -> Self {
555        Self { sender }
556    }
557}
558
559impl Processor for RecorderProcessor {
560    fn process_frame(&mut self, frame: &mut AudioFrame) -> Result<()> {
561        let frame_clone = frame.clone();
562        let _ = self.sender.send(frame_clone);
563        Ok(())
564    }
565}
566
567impl MediaStream {
568    pub async fn start_recorder(&self) -> Result<()> {
569        let recorder_option = self.recorder_option.lock().await.clone();
570        if let Some(recorder_option) = recorder_option {
571            if recorder_option.recorder_file.is_empty() {
572                warn!(
573                    session_id = self.id,
574                    "recorder file is empty, skipping recorder start"
575                );
576                return Ok(());
577            }
578            let recorder_receiver = match self.recorder_receiver.lock().await.take() {
579                Some(receiver) => receiver,
580                None => {
581                    return Ok(());
582                }
583            };
584            let cancel_token = self.cancel_token.child_token();
585            let session_id_clone = self.id.clone();
586
587            info!(
588                session_id = session_id_clone,
589                sample_rate = recorder_option.samplerate,
590                ptime = recorder_option.ptime,
591                "start recorder",
592            );
593
594            let recorder_handle = crate::spawn(async move {
595                let recorder_file = recorder_option.recorder_file.clone();
596                let recorder =
597                    Recorder::new(cancel_token, session_id_clone.clone(), recorder_option);
598                match recorder
599                    .process_recording(Path::new(&recorder_file), recorder_receiver)
600                    .await
601                {
602                    Ok(_) => {}
603                    Err(e) => {
604                        warn!(
605                            session_id = session_id_clone,
606                            "Failed to process recorder: {}", e
607                        );
608                    }
609                }
610            });
611            *self.recorder_handle.lock().await = Some(recorder_handle);
612            self.recording_active.store(true, Ordering::SeqCst);
613
614            // Inject RecorderProcessor into tracks that were added before the recorder started
615            for (track, _) in self.tracks.lock().await.values_mut() {
616                track.insert_processor(Box::new(RecorderProcessor::new(
617                    self.recorder_sender.clone(),
618                )));
619            }
620        }
621        Ok(())
622    }
623
624    pub async fn set_track_refer(&self, track_id: &TrackId, refer: Option<bool>) {
625        if let Some((_, dtmf)) = self.tracks.lock().await.get_mut(track_id) {
626            dtmf.refer = refer;
627        }
628    }
629
630    pub async fn set_track_dtmf_forward(&self, track_id: &TrackId, forward: bool) {
631        if let Some((_, dtmf)) = self.tracks.lock().await.get_mut(track_id) {
632            dtmf.suppress_dtmf_forward = !forward;
633        }
634    }
635
636    async fn handle_forward_track(&self, mut packet_receiver: TrackPacketReceiver) {
637        let event_sender = self.event_sender.clone();
638        while let Some(packet) = packet_receiver.recv().await {
639            if self
640                .ambiance_source_id
641                .lock()
642                .ok()
643                .and_then(|id| id.clone())
644                .as_ref()
645                == Some(&packet.track_id)
646            {
647                self.last_server_packet_ts
648                    .store(crate::media::get_timestamp(), Ordering::Relaxed);
649            }
650
651            let suppressed = {
652                self.suppressed_sources
653                    .lock()
654                    .await
655                    .contains(&packet.track_id)
656            };
657
658            let is_dtmf = matches!(&packet.samples,
659                Samples::RTP { payload_type, .. } if *payload_type >= 96 && *payload_type <= 127);
660
661            let mut tracks = self.tracks.lock().await;
662
663            // Check once whether the source track suppresses DTMF forwarding.
664            let source_suppresses_dtmf = is_dtmf
665                && tracks
666                    .get(&packet.track_id)
667                    .map(|(_, d)| d.suppress_dtmf_forward)
668                    .unwrap_or(false);
669
670            for (track, dtmf_detector) in tracks.values_mut() {
671                if track.id() == &packet.track_id {
672                    if let Samples::RTP {
673                        payload_type,
674                        payload,
675                        ..
676                    } = &packet.samples
677                    {
678                        if let Some(digit) = dtmf_detector.detect_rtp(*payload_type, payload) {
679                            debug!(track_id = track.id(), digit, "DTMF detected");
680                            event_sender
681                                .send(SessionEvent::Dtmf {
682                                    track_id: packet.track_id.to_string(),
683                                    timestamp: packet.timestamp,
684                                    digit,
685                                    refer: dtmf_detector.refer,
686                                })
687                                .ok();
688                        }
689                    }
690                    continue;
691                }
692                if suppressed {
693                    continue;
694                }
695                // Skip DTMF forwarding if source or destination has it suppressed.
696                if source_suppresses_dtmf || (is_dtmf && dtmf_detector.suppress_dtmf_forward) {
697                    continue;
698                }
699                if packet.track_id == QUEUE_HOLD_TRACK_ID && track.id() == CALLEE_TRACK_ID {
700                    continue;
701                }
702                if let Err(e) = track.send_packet(&packet).await {
703                    warn!(
704                        id = track.id(),
705                        "media_stream: Failed to send packet to track: {}", e
706                    );
707                }
708            }
709        }
710    }
711}
712
713pub struct MuteProcessor;
714
715impl MuteProcessor {
716    pub fn mute_track(track: &mut dyn Track) {
717        let chain = track.processor_chain();
718        if !chain.has_processor::<MuteProcessor>() {
719            chain.insert_processor(Box::new(MuteProcessor));
720        }
721    }
722
723    pub fn unmute_track(track: &mut dyn Track) {
724        let chain = track.processor_chain();
725        chain.remove_processor::<MuteProcessor>();
726    }
727}
728
729impl Processor for MuteProcessor {
730    fn process_frame(&mut self, frame: &mut AudioFrame) -> Result<()> {
731        match &mut frame.samples {
732            Samples::PCM { samples } => {
733                samples.fill(0);
734            }
735            // discard DTMF frames
736            Samples::RTP { payload_type, .. } if *payload_type >= 96 && *payload_type <= 127 => {
737                frame.samples = Samples::Empty;
738            }
739            _ => {}
740        }
741        Ok(())
742    }
743}
744
745pub struct HoldTrack;
746
747impl HoldTrack {
748    pub fn hold_track(track: &mut dyn Track) {
749        let chain = track.processor_chain();
750        // Remove existing processor if present
751        chain.remove_processor::<HoldProcessor>();
752        // Add a new processor with hold state set to true
753        let processor = HoldProcessor::new();
754        processor.set_hold(true);
755        chain.insert_processor(Box::new(processor));
756    }
757
758    pub fn resume_track(track: &mut dyn Track) {
759        let chain = track.processor_chain();
760        // Simply remove the hold processor to resume normal operation
761        chain.remove_processor::<HoldProcessor>();
762    }
763}