active_call/media/track/
forwarding.rs1use 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}