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        self.attach_recorder_tap(&mut track).await;
361        match track
362            .start(self.event_sender.clone(), self.packet_sender.clone())
363            .await
364        {
365            Ok(_) => {
366                info!(session_id = self.id, track_id = track.id(), "track started");
367                let track_id = track.id().clone();
368                if track_id.as_str() == self.id.as_str() {
369                    let pending = std::mem::take(&mut *self.pending_ice_candidates.lock().await);
370                    for (candidate, sdp_mid, sdp_mline_index) in pending {
371                        if let Err(e) =
372                            track.add_ice_candidate(&candidate, sdp_mid.as_deref(), sdp_mline_index)
373                        {
374                            warn!(
375                                session_id = self.id,
376                                track_id = track.id(),
377                                "failed to apply buffered ICE candidate: {}",
378                                e
379                            );
380                        }
381                    }
382                }
383                self.tracks
384                    .lock()
385                    .await
386                    .insert(track_id.clone(), (track, DtmfDetector::new()));
387                self.event_sender
388                    .send(SessionEvent::TrackStart {
389                        track_id,
390                        timestamp: crate::media::get_timestamp(),
391                        play_id,
392                    })
393                    .ok();
394            }
395            Err(e) => {
396                warn!(
397                    session_id = self.id,
398                    track_id = track.id(),
399                    play_id = play_id.as_deref(),
400                    "Failed to start track: {}",
401                    e
402                );
403            }
404        }
405    }
406
407    pub async fn mute_track(&self, id: Option<TrackId>) {
408        if let Some(id) = id {
409            if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
410                MuteProcessor::mute_track(track.as_mut());
411            }
412        } else {
413            for (track, _) in self.tracks.lock().await.values_mut() {
414                MuteProcessor::mute_track(track.as_mut());
415            }
416        }
417    }
418
419    pub async fn unmute_track(&self, id: Option<TrackId>) {
420        if let Some(id) = id {
421            if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
422                MuteProcessor::unmute_track(track.as_mut());
423            }
424        } else {
425            for (track, _) in self.tracks.lock().await.values_mut() {
426                MuteProcessor::unmute_track(track.as_mut());
427            }
428        }
429    }
430
431    /// Trickle ICE: feed a remote candidate into the ICE-backed (WebRTC)
432    /// track, if it's up yet, or buffer it for `update_track` to replay
433    /// otherwise.
434    pub async fn add_ice_candidate(
435        &self,
436        candidate: &str,
437        sdp_mid: Option<&str>,
438        sdp_mline_index: Option<u32>,
439    ) -> Result<()> {
440        let tracks = self.tracks.lock().await;
441        if let Some((track, _)) = tracks.get(self.id.as_str()) {
442            track.add_ice_candidate(candidate, sdp_mid, sdp_mline_index)?;
443            return Ok(());
444        }
445        drop(tracks);
446        self.pending_ice_candidates.lock().await.push((
447            candidate.to_string(),
448            sdp_mid.map(|s| s.to_string()),
449            sdp_mline_index,
450        ));
451        Ok(())
452    }
453
454    pub async fn pause_playback(&self, id: TrackId) -> Result<()> {
455        self.set_playback_paused(id, true).await
456    }
457
458    pub async fn resume_playback(&self, id: TrackId) -> Result<()> {
459        self.set_playback_paused(id, false).await
460    }
461
462    async fn set_playback_paused(&self, id: TrackId, paused: bool) -> Result<()> {
463        if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
464            if track.set_paused(paused) {
465                Ok(())
466            } else {
467                warn!(
468                    session_id = self.id,
469                    track_id = %id,
470                    paused,
471                    "pause state requested for track that does not support pausing"
472                );
473                Err(anyhow::anyhow!("track does not support pausing: {}", id))
474            }
475        } else {
476            warn!(
477                session_id = self.id,
478                track_id = %id,
479                paused,
480                "pause state requested for unknown track"
481            );
482            Err(anyhow::anyhow!("track not found: {}", id))
483        }
484    }
485
486    pub async fn hold_track(&self, id: Option<TrackId>) {
487        if let Some(id) = id {
488            if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
489                HoldTrack::hold_track(track.as_mut());
490            }
491        } else {
492            for (track, _) in self.tracks.lock().await.values_mut() {
493                HoldTrack::hold_track(track.as_mut());
494            }
495        }
496    }
497
498    pub async fn resume_track(&self, id: Option<TrackId>) {
499        if let Some(id) = id {
500            if let Some((track, _)) = self.tracks.lock().await.get_mut(&id) {
501                HoldTrack::resume_track(track.as_mut());
502            }
503        } else {
504            for (track, _) in self.tracks.lock().await.values_mut() {
505                HoldTrack::resume_track(track.as_mut());
506            }
507        }
508    }
509
510    pub async fn suppress_forwarding(&self, track_id: &TrackId) {
511        self.suppressed_sources
512            .lock()
513            .await
514            .insert(track_id.clone());
515    }
516
517    pub async fn resume_forwarding(&self, track_id: &TrackId) {
518        self.suppressed_sources.lock().await.remove(track_id);
519    }
520
521    pub async fn remove_processor<T: 'static>(&self, track_id: &TrackId) -> Result<()> {
522        if let Some((track, _)) = self.tracks.lock().await.get_mut(track_id) {
523            track.as_mut().processor_chain().remove_processor::<T>();
524            Ok(())
525        } else {
526            Err(anyhow::anyhow!("Track {} not found", track_id))
527        }
528    }
529
530    pub async fn append_processor(
531        &self,
532        track_id: &TrackId,
533        processor: Box<dyn crate::media::processor::Processor>,
534    ) -> Result<()> {
535        if let Some((track, _)) = self.tracks.lock().await.get_mut(track_id) {
536            track.as_mut().processor_chain().append_processor(processor);
537            Ok(())
538        } else {
539            Err(anyhow::anyhow!("Track {} not found", track_id))
540        }
541    }
542}
543
544#[derive(Clone)]
545pub struct RecorderProcessor {
546    sender: mpsc::UnboundedSender<AudioFrame>,
547}
548
549impl RecorderProcessor {
550    pub fn new(sender: mpsc::UnboundedSender<AudioFrame>) -> Self {
551        Self { sender }
552    }
553}
554
555impl Processor for RecorderProcessor {
556    fn process_frame(&mut self, frame: &mut AudioFrame) -> Result<()> {
557        let frame_clone = frame.clone();
558        let _ = self.sender.send(frame_clone);
559        Ok(())
560    }
561}
562
563impl MediaStream {
564    pub async fn start_recorder(&self) -> Result<()> {
565        let recorder_option = self.recorder_option.lock().await.clone();
566        if let Some(recorder_option) = recorder_option {
567            if recorder_option.recorder_file.is_empty() {
568                warn!(
569                    session_id = self.id,
570                    "recorder file is empty, skipping recorder start"
571                );
572                return Ok(());
573            }
574            let recorder_receiver = match self.recorder_receiver.lock().await.take() {
575                Some(receiver) => receiver,
576                None => {
577                    return Ok(());
578                }
579            };
580            let cancel_token = self.cancel_token.child_token();
581            let session_id_clone = self.id.clone();
582
583            info!(
584                session_id = session_id_clone,
585                sample_rate = recorder_option.samplerate,
586                native_samplerate = recorder_option.native_samplerate.unwrap_or(false),
587                ptime = recorder_option.ptime,
588                "start recorder",
589            );
590
591            let native = recorder_option.native_samplerate.unwrap_or(false);
592            let recorder_handle = crate::spawn(async move {
593                let recorder_file = recorder_option.recorder_file.clone();
594                let recorder =
595                    Recorder::new(cancel_token, session_id_clone.clone(), recorder_option);
596                match recorder
597                    .process_recording(Path::new(&recorder_file), recorder_receiver)
598                    .await
599                {
600                    Ok(_) => {}
601                    Err(e) => {
602                        warn!(
603                            session_id = session_id_clone,
604                            "Failed to process recorder: {}", e
605                        );
606                    }
607                }
608            });
609            *self.recorder_handle.lock().await = Some(recorder_handle);
610            self.recording_active.store(true, Ordering::SeqCst);
611
612            // Inject the recorder tap into tracks that were added before the
613            // recorder started. `set_raw_tap` replaces any previous tap, so
614            // this stays idempotent.
615            for (track, _) in self.tracks.lock().await.values_mut() {
616                if native {
617                    track.set_raw_tap(Some(self.recorder_sender.clone()));
618                } else if !track.processor_chain().has_processor::<RecorderProcessor>() {
619                    track.insert_processor(Box::new(RecorderProcessor::new(
620                        self.recorder_sender.clone(),
621                    )));
622                }
623            }
624        }
625        Ok(())
626    }
627
628    /// Attach the recorder tap to a track according to the recorder mode:
629    /// native-samplerate mode mirrors pre-resample frames via the chain's
630    /// raw tap, otherwise a post-pipeline RecorderProcessor is appended.
631    /// No-op when no recorder is configured.
632    async fn attach_recorder_tap(&self, track: &mut Box<dyn Track>) {
633        let Some(option) = self.recorder_option.lock().await.clone() else {
634            return;
635        };
636        if option.native_samplerate.unwrap_or(false) {
637            track.set_raw_tap(Some(self.recorder_sender.clone()));
638        } else {
639            track.append_processor(Box::new(RecorderProcessor::new(
640                self.recorder_sender.clone(),
641            )));
642        }
643    }
644
645    pub async fn set_track_refer(&self, track_id: &TrackId, refer: Option<bool>) {
646        if let Some((_, dtmf)) = self.tracks.lock().await.get_mut(track_id) {
647            dtmf.refer = refer;
648        }
649    }
650
651    pub async fn set_track_dtmf_forward(&self, track_id: &TrackId, forward: bool) {
652        if let Some((_, dtmf)) = self.tracks.lock().await.get_mut(track_id) {
653            dtmf.suppress_dtmf_forward = !forward;
654        }
655    }
656
657    async fn handle_forward_track(&self, mut packet_receiver: TrackPacketReceiver) {
658        let event_sender = self.event_sender.clone();
659        while let Some(packet) = packet_receiver.recv().await {
660            if self
661                .ambiance_source_id
662                .lock()
663                .ok()
664                .and_then(|id| id.clone())
665                .as_ref()
666                == Some(&packet.track_id)
667            {
668                self.last_server_packet_ts
669                    .store(crate::media::get_timestamp(), Ordering::Relaxed);
670            }
671
672            let suppressed = {
673                self.suppressed_sources
674                    .lock()
675                    .await
676                    .contains(&packet.track_id)
677            };
678
679            let is_dtmf = matches!(&packet.samples,
680                Samples::RTP { payload_type, .. } if *payload_type >= 96 && *payload_type <= 127);
681
682            let mut tracks = self.tracks.lock().await;
683
684            // Check once whether the source track suppresses DTMF forwarding.
685            let source_suppresses_dtmf = is_dtmf
686                && tracks
687                    .get(&packet.track_id)
688                    .map(|(_, d)| d.suppress_dtmf_forward)
689                    .unwrap_or(false);
690
691            for (track, dtmf_detector) in tracks.values_mut() {
692                if track.id() == &packet.track_id {
693                    if let Samples::RTP {
694                        payload_type,
695                        payload,
696                        ..
697                    } = &packet.samples
698                    {
699                        if let Some(digit) = dtmf_detector.detect_rtp(*payload_type, payload) {
700                            debug!(track_id = track.id(), digit, "DTMF detected");
701                            event_sender
702                                .send(SessionEvent::Dtmf {
703                                    track_id: packet.track_id.to_string(),
704                                    timestamp: packet.timestamp,
705                                    digit,
706                                    refer: dtmf_detector.refer,
707                                })
708                                .ok();
709                        }
710                    }
711                    continue;
712                }
713                if suppressed {
714                    continue;
715                }
716                // Skip DTMF forwarding if source or destination has it suppressed.
717                if source_suppresses_dtmf || (is_dtmf && dtmf_detector.suppress_dtmf_forward) {
718                    continue;
719                }
720                if packet.track_id == QUEUE_HOLD_TRACK_ID && track.id() == CALLEE_TRACK_ID {
721                    continue;
722                }
723                if let Err(e) = track.send_packet(&packet).await {
724                    warn!(
725                        id = track.id(),
726                        "media_stream: Failed to send packet to track: {}", e
727                    );
728                }
729            }
730        }
731    }
732}
733
734pub struct MuteProcessor;
735
736impl MuteProcessor {
737    pub fn mute_track(track: &mut dyn Track) {
738        let chain = track.processor_chain();
739        if !chain.has_processor::<MuteProcessor>() {
740            chain.insert_processor(Box::new(MuteProcessor));
741        }
742    }
743
744    pub fn unmute_track(track: &mut dyn Track) {
745        let chain = track.processor_chain();
746        chain.remove_processor::<MuteProcessor>();
747    }
748}
749
750impl Processor for MuteProcessor {
751    fn process_frame(&mut self, frame: &mut AudioFrame) -> Result<()> {
752        match &mut frame.samples {
753            Samples::PCM { samples } => {
754                samples.fill(0);
755            }
756            // discard DTMF frames
757            Samples::RTP { payload_type, .. } if *payload_type >= 96 && *payload_type <= 127 => {
758                frame.samples = Samples::Empty;
759            }
760            _ => {}
761        }
762        Ok(())
763    }
764}
765
766pub struct HoldTrack;
767
768impl HoldTrack {
769    pub fn hold_track(track: &mut dyn Track) {
770        let chain = track.processor_chain();
771        // Remove existing processor if present
772        chain.remove_processor::<HoldProcessor>();
773        // Add a new processor with hold state set to true
774        let processor = HoldProcessor::new();
775        processor.set_hold(true);
776        chain.insert_processor(Box::new(processor));
777    }
778
779    pub fn resume_track(track: &mut dyn Track) {
780        let chain = track.processor_chain();
781        // Simply remove the hold processor to resume normal operation
782        chain.remove_processor::<HoldProcessor>();
783    }
784}