Skip to main content

active_call/media/track/
forwarding.rs

1use crate::event::EventSender;
2use crate::media::processor::ProcessorChain;
3use crate::media::track::{Track, TrackConfig, TrackPacketSender};
4use crate::media::{AudioFrame, Samples, TrackId};
5use anyhow::Result;
6use async_trait::async_trait;
7use std::sync::{
8    Arc,
9    atomic::{AtomicBool, Ordering},
10};
11use tokio::sync::mpsc;
12use tokio::time::Duration;
13use tokio_util::sync::CancellationToken;
14use tracing::{info, warn};
15
16pub struct ForwardingTrack {
17    track_id: TrackId,
18    source_peer_track_id: TrackId,
19    peer_sender: mpsc::Sender<AudioFrame>,
20    inbound_receiver: Option<mpsc::Receiver<AudioFrame>>,
21    processor_chain: ProcessorChain,
22    config: TrackConfig,
23    cancel_token: CancellationToken,
24    ssrc: u32,
25    paused: Arc<AtomicBool>,
26}
27
28impl ForwardingTrack {
29    pub fn new(
30        track_id: TrackId,
31        source_peer_track_id: TrackId,
32        peer_sender: mpsc::Sender<AudioFrame>,
33        inbound_receiver: mpsc::Receiver<AudioFrame>,
34        config: TrackConfig,
35        cancel_token: CancellationToken,
36        ssrc: u32,
37        paused: Arc<AtomicBool>,
38    ) -> Self {
39        Self {
40            processor_chain: ProcessorChain::new(config.samplerate),
41            track_id,
42            source_peer_track_id,
43            peer_sender,
44            inbound_receiver: Some(inbound_receiver),
45            config,
46            cancel_token,
47            ssrc,
48            paused,
49        }
50    }
51}
52
53#[async_trait]
54impl Track for ForwardingTrack {
55    fn ssrc(&self) -> u32 {
56        self.ssrc
57    }
58
59    fn id(&self) -> &TrackId {
60        &self.track_id
61    }
62
63    fn config(&self) -> &TrackConfig {
64        &self.config
65    }
66
67    fn processor_chain(&mut self) -> &mut ProcessorChain {
68        &mut self.processor_chain
69    }
70
71    async fn handshake(&mut self, _offer: String, _timeout: Option<Duration>) -> Result<String> {
72        Ok(String::new())
73    }
74
75    async fn update_remote_description(&mut self, _answer: &String) -> Result<()> {
76        Ok(())
77    }
78
79    async fn start(
80        &mut self,
81        _event_sender: EventSender,
82        packet_sender: TrackPacketSender,
83    ) -> Result<()> {
84        let mut inbound_receiver = self
85            .inbound_receiver
86            .take()
87            .ok_or_else(|| anyhow::anyhow!("forwarding track already started"))?;
88        let track_id = self.track_id.clone();
89        let cancel_token = self.cancel_token.clone();
90        let mut processor_chain = self.processor_chain.clone();
91
92        crate::spawn(async move {
93            let stop_reason = loop {
94                tokio::select! {
95                    _ = cancel_token.cancelled() => {
96                        break "track stopped";
97                    }
98                    packet = inbound_receiver.recv() => {
99                        match packet {
100                            Some(mut packet) => {
101                                packet.track_id = track_id.clone();
102                                if let Err(e) = processor_chain.process_frame(&mut packet) {
103                                    warn!(track_id, "processor_chain process_frame error: {:?}", e);
104                                }
105                                if packet_sender.send(packet).is_err() {
106                                    break "media stream closed";
107                                }
108                            }
109                            None => {
110                                break "peer bridge channel closed";
111                            }
112                        }
113                    }
114                }
115            };
116            cancel_token.cancel();
117            info!(
118                track_id,
119                reason = stop_reason,
120                "audio bridge forwarding task stopped"
121            );
122        });
123        Ok(())
124    }
125
126    async fn stop(&self) -> Result<()> {
127        self.cancel_token.cancel();
128        Ok(())
129    }
130
131    async fn send_packet(&mut self, packet: &AudioFrame) -> Result<()> {
132        if self.cancel_token.is_cancelled()
133            || self.paused.load(Ordering::Relaxed)
134            || packet.track_id != self.source_peer_track_id
135        {
136            return Ok(());
137        }
138
139        if let Samples::RTP { payload_type, .. } = &packet.samples {
140            if *payload_type >= 96 && *payload_type <= 127 {
141                return Ok(());
142            }
143        }
144
145        match self.peer_sender.try_send(packet.clone()) {
146            Ok(_) => {}
147            Err(mpsc::error::TrySendError::Full(_)) => {}
148            Err(mpsc::error::TrySendError::Closed(_)) => {
149                self.cancel_token.cancel();
150            }
151        }
152
153        Ok(())
154    }
155}