1use super::track_codec::TrackCodec;
2use crate::{
3 event::{EventSender, SessionEvent},
4 media::AudioFrame,
5 media::{
6 processor::ProcessorChain,
7 track::{Track, TrackConfig, TrackId, TrackPacketSender},
8 },
9};
10use anyhow::Result;
11use async_trait::async_trait;
12use audio_codec::CodecType;
13use bytes::Bytes;
14use futures::{FutureExt, StreamExt, stream::FuturesUnordered};
15use rustrtc::{
16 AudioCapability, IceServer, MediaKind, PeerConnection, PeerConnectionEvent,
17 PeerConnectionState, RtcConfiguration, RtpCodecParameters, SdpType, TransportMode,
18 config::MediaCapabilities,
19 media::{
20 MediaStreamTrack, SampleStreamSource, frame::AudioFrame as RtcAudioFrame, sample_track,
21 track::SampleStreamTrack,
22 },
23};
24use std::{
25 sync::{
26 Arc,
27 atomic::{AtomicBool, Ordering},
28 },
29 time::{Duration, Instant},
30};
31use tokio::sync::Mutex;
32use tokio_util::sync::CancellationToken;
33use tracing::{debug, info};
34
35#[derive(Clone)]
36pub struct RtcTrackConfig {
37 pub mode: TransportMode,
38 pub ice_servers: Option<Vec<IceServer>>,
39 pub external_ip: Option<String>,
40 pub rtp_port_range: Option<(u16, u16)>,
41 pub bind_ip: Option<String>,
42 pub preferred_codec: Option<CodecType>,
43 pub codecs: Vec<CodecType>,
44 pub payload_type: Option<u8>,
45 pub enable_latching: Option<bool>,
46 pub enable_ice_lite: Option<bool>,
47}
48
49impl Default for RtcTrackConfig {
50 fn default() -> Self {
51 Self {
52 mode: TransportMode::WebRtc, ice_servers: None,
54 external_ip: None,
55 rtp_port_range: None,
56 bind_ip: None,
57 preferred_codec: None,
58 codecs: Vec::new(),
59 payload_type: None,
60 enable_latching: None,
61 enable_ice_lite: None,
62 }
63 }
64}
65
66pub struct RtcTrack {
67 track_id: TrackId,
68 track_config: TrackConfig,
69 rtc_config: RtcTrackConfig,
70 processor_chain: ProcessorChain,
71 packet_sender: Arc<Mutex<Option<TrackPacketSender>>>,
72 event_sender: Arc<Mutex<Option<EventSender>>>,
73 media_ready_sent: Arc<AtomicBool>,
74 cancel_token: CancellationToken,
75 local_source: Option<Arc<SampleStreamSource>>,
76 encoder: TrackCodec,
77 ssrc: u32,
78 payload_type: Option<u8>,
79 pub peer_connection: Option<Arc<PeerConnection>>,
80 next_rtp_timestamp: u32,
81 next_rtp_sequence_number: u16,
82 last_packet_time: Option<Instant>,
83 last_remote_sdp: Option<String>,
84 need_marker: bool,
85}
86
87impl RtcTrack {
88 pub fn new(
89 cancel_token: CancellationToken,
90 id: TrackId,
91 track_config: TrackConfig,
92 rtc_config: RtcTrackConfig,
93 ) -> Self {
94 let processor_chain = ProcessorChain::new(track_config.samplerate);
95 Self {
96 track_id: id,
97 track_config,
98 rtc_config,
99 processor_chain,
100 packet_sender: Arc::new(Mutex::new(None)),
101 event_sender: Arc::new(Mutex::new(None)),
102 media_ready_sent: Arc::new(AtomicBool::new(false)),
103 cancel_token,
104 local_source: None,
105 encoder: TrackCodec::new(),
106 ssrc: 0,
107 payload_type: None,
108 peer_connection: None,
109 next_rtp_timestamp: 0,
110 next_rtp_sequence_number: 0,
111 last_packet_time: None,
112 last_remote_sdp: None,
113 need_marker: false,
114 }
115 }
116
117 pub fn with_ssrc(mut self, ssrc: u32) -> Self {
118 self.ssrc = ssrc;
119 self
120 }
121
122 pub fn create_audio_track(
123 _codec: CodecType,
124 _stream_id: Option<String>,
125 ) -> (Arc<SampleStreamSource>, Arc<SampleStreamTrack>) {
126 let (source, track, _) = sample_track(rustrtc::media::MediaKind::Audio, 100);
127 (Arc::new(source), track)
128 }
129
130 pub async fn local_description(&self) -> Result<String> {
131 let pc = self
132 .peer_connection
133 .as_ref()
134 .ok_or_else(|| anyhow::anyhow!("No PeerConnection"))?;
135 let offer = pc.create_offer().await?;
136 pc.set_local_description(offer.clone())?;
137 Ok(offer.to_sdp_string())
138 }
139
140 pub async fn create(&mut self) -> Result<()> {
141 if self.peer_connection.is_some() {
142 return Ok(());
143 }
144
145 let mut config = RtcConfiguration::default();
146 if self.ssrc != 0 {
147 config.ssrc_start = self.ssrc;
148 }
149 config.transport_mode = self.rtc_config.mode.clone();
150
151 if let Some(ice_servers) = &self.rtc_config.ice_servers {
152 config.ice_servers = ice_servers.clone();
153 }
154
155 if let Some(external_ip) = &self.rtc_config.external_ip {
156 config.external_ip = Some(external_ip.clone());
157 }
158 if let Some(bind_ip) = &self.rtc_config.bind_ip {
159 config.bind_ip = Some(bind_ip.clone());
160 }
161 if let Some((rtp_start_port, rtp_end_port)) = self.rtc_config.rtp_port_range {
162 config.rtp_start_port = Some(rtp_start_port);
163 config.rtp_end_port = Some(rtp_end_port);
164 }
165 config.enable_ice_lite = self.rtc_config.enable_ice_lite.unwrap_or(false);
166 config.enable_latching = self
167 .rtc_config
168 .enable_latching
169 .unwrap_or_else(|| self.rtc_config.mode == TransportMode::Rtp);
170
171 if !self.rtc_config.codecs.is_empty() {
172 let mut caps = MediaCapabilities::default();
173 caps.audio.clear();
174
175 for codec in &self.rtc_config.codecs {
176 let cap = match codec {
177 CodecType::PCMU => AudioCapability::pcmu(),
178 CodecType::PCMA => AudioCapability::pcma(),
179 CodecType::G722 => AudioCapability::g722(),
180 CodecType::G729 => AudioCapability::g729(),
181 CodecType::TelephoneEvent => AudioCapability::telephone_event(),
182 #[cfg(feature = "opus")]
183 CodecType::Opus => AudioCapability::opus(),
184 };
185 caps.audio.push(cap);
186 }
187 config.media_capabilities = Some(caps);
188 }
189
190 let peer_connection = Arc::new(PeerConnection::new(config));
191 self.peer_connection = Some(peer_connection.clone());
192
193 let default_codec = CodecType::G722;
194 let codec = self.rtc_config.preferred_codec.unwrap_or(default_codec);
195
196 let (source, track) = Self::create_audio_track(codec, Some(self.track_id.clone()));
197 self.local_source = Some(source);
198
199 let payload_type = self
200 .rtc_config
201 .payload_type
202 .unwrap_or_else(|| codec.payload_type());
203
204 self.payload_type = Some(payload_type);
205
206 let params = RtpCodecParameters {
207 clock_rate: codec.clock_rate(),
208 channels: codec.channels() as u8,
209 payload_type,
210 ..Default::default()
211 };
212
213 peer_connection.add_track_with_stream_id(track, self.track_id.clone(), params)?;
214
215 self.spawn_handlers(
217 peer_connection.clone(),
218 self.track_id.clone(),
219 self.processor_chain.clone(),
220 payload_type,
221 self.event_sender.clone(),
222 self.media_ready_sent.clone(),
223 );
224
225 Ok(())
226 }
227
228 fn spawn_handlers(
229 &self,
230 pc: Arc<PeerConnection>,
231 track_id: TrackId,
232 processor_chain: ProcessorChain,
233 default_payload_type: u8,
234 event_sender: Arc<Mutex<Option<EventSender>>>,
235 media_ready_sent: Arc<AtomicBool>,
236 ) {
237 let cancel_token = self.cancel_token.clone();
238 let packet_sender = self.packet_sender.clone();
239 let pc_event = pc.clone();
240 let pc_stats = pc.clone();
241 let pc_state = pc.clone();
242 let track_id_log = track_id.clone();
243 let is_rtp_media = matches!(
244 self.rtc_config.mode,
245 TransportMode::Rtp | TransportMode::Srtp
246 );
247 let is_webrtc = self.rtc_config.mode != TransportMode::Rtp;
248
249 crate::spawn(async move {
250 info!(track_id=%track_id_log, "RtcTrack event/stats loop started");
251
252 let mut events = futures::stream::unfold(pc_event, |pc| async move {
253 pc.recv().await.map(|ev| (ev, pc))
254 })
255 .boxed();
256
257 let mut state_rx = if is_webrtc {
258 Some(pc_state.subscribe_peer_state())
259 } else {
260 None
261 };
262
263 let mut stats_interval = tokio::time::interval(Duration::from_secs(5));
264 let mut event_count = 0;
265 let mut workers = FuturesUnordered::new();
266
267 loop {
268 tokio::select! {
269 _ = cancel_token.cancelled() => {
270 debug!(track_id=%track_id_log, "RtcTrack loop cancelled");
271 break;
272 }
273
274 Some(event) = events.next() => {
275 event_count += 1;
276 let event_type = match &event {
277 PeerConnectionEvent::Track(_) => "Track",
278 PeerConnectionEvent::DataChannel(_) => "DataChannel",
279 };
280 debug!(track_id=%track_id_log, "Received PeerConnectionEvent #{}: {}", event_count, event_type);
281
282 if let PeerConnectionEvent::Track(transceiver) = event {
283 if let Some(receiver) = transceiver.receiver() {
284 let track = receiver.track();
285 if is_rtp_media {
286 let maybe_sender = event_sender.lock().await.clone();
287 if let Some(sender) = maybe_sender {
288 if media_ready_sent
289 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
290 .is_ok()
291 {
292 let result = sender.send(SessionEvent::MediaReady {
293 track_id: track_id_log.clone(),
294 timestamp: crate::media::get_timestamp(),
295 });
296 if result.is_err() {
297 media_ready_sent.store(false, Ordering::SeqCst);
298 }
299 }
300 }
301 }
302 info!(track_id=%track_id_log, "New track received");
303
304 let (f1, f2) = Self::create_track_workers(
305 track,
306 packet_sender.clone(),
307 track_id_log.clone(),
308 processor_chain.clone(),
309 default_payload_type,
310 );
311 workers.push(f1);
312 workers.push(f2);
313 }
314 }
315 }
316
317 _ = workers.next(), if !workers.is_empty() => {}
318
319 _ = stats_interval.tick() => {
320 match pc_stats.get_stats().await {
321 Ok(stats) => {
322 info!(track_id=%track_id_log, %stats, "RTCP Stats");
323 }
324 Err(e) => {
325 debug!(track_id=%track_id_log, "Failed to get stats: {:?}", e);
326 }
327 }
328 }
329
330 res = async {
332 if let Some(rx) = state_rx.as_mut() {
333 rx.changed().await
334 } else {
335 std::future::pending().await
336 }
337 } => {
338 if res.is_ok() {
339 if let Some(rx) = state_rx.as_ref() {
340 let s = *rx.borrow();
341 debug!(track_id=%track_id_log, "peer connection state changed: {:?}", s);
342 match s {
343 PeerConnectionState::Disconnected
344 | PeerConnectionState::Closed
345 | PeerConnectionState::Failed => {
346 info!(
347 track_id = %track_id_log,
348 "peer connection is {:?}, try to close", s
349 );
350 cancel_token.cancel();
351 pc_state.close();
352 break;
353 }
354 _ => {}
355 }
356 }
357 }
358 }
359 }
360 }
361 debug!(track_id=%track_id_log, "RtcTrack event/stats loop ended, total events: {}", event_count);
362 });
363 }
364
365 fn create_track_workers(
366 track: Arc<SampleStreamTrack>,
367 packet_sender_arc: Arc<Mutex<Option<TrackPacketSender>>>,
368 track_id: TrackId,
369 processor_chain: ProcessorChain,
370 default_payload_type: u8,
371 ) -> (
372 std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>,
373 std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>,
374 ) {
375 let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<rustrtc::media::frame::AudioFrame>();
376
377 let track_id_proc = track_id.clone();
379 let packet_sender_proc = packet_sender_arc.clone();
380 let processor_chain_proc = processor_chain.clone();
381 let proc_fut = Self::run_processing_worker(
382 rx,
383 track_id_proc,
384 packet_sender_proc,
385 processor_chain_proc,
386 default_payload_type,
387 );
388
389 let track_id_recv = track_id.clone();
391 let recv_fut = Self::run_receiving_worker(track, tx, track_id_recv);
392
393 (proc_fut.boxed(), recv_fut.boxed())
394 }
395
396 async fn run_processing_worker(
397 mut rx: tokio::sync::mpsc::UnboundedReceiver<rustrtc::media::frame::AudioFrame>,
398 track_id: TrackId,
399 packet_sender: Arc<Mutex<Option<TrackPacketSender>>>,
400 mut processor_chain: ProcessorChain,
401 default_payload_type: u8,
402 ) {
403 info!(track_id=%track_id, "RtcTrack processing worker started");
404 while let Some(frame) = rx.recv().await {
405 let res = std::panic::AssertUnwindSafe(Self::process_audio_frame(
406 frame,
407 &track_id,
408 &packet_sender,
409 &mut processor_chain,
410 default_payload_type,
411 ))
412 .catch_unwind()
413 .await;
414
415 if let Err(cause) = res {
416 let msg = if let Some(s) = cause.downcast_ref::<&str>() {
417 *s
418 } else if let Some(s) = cause.downcast_ref::<String>() {
419 &s[..]
420 } else {
421 "Unknown panic"
422 };
423 tracing::error!(track_id=%track_id, "RtcTrack processing worker PANIC: {}", msg);
424 break;
425 }
426 }
427 info!(track_id=%track_id, "RtcTrack processing worker stopped");
428 }
429
430 async fn run_receiving_worker(
431 track: Arc<SampleStreamTrack>,
432 tx: tokio::sync::mpsc::UnboundedSender<rustrtc::media::frame::AudioFrame>,
433 track_id: TrackId,
434 ) {
435 let mut samples =
436 futures::stream::unfold(
437 track,
438 |t| async move { t.recv().await.ok().map(|s| (s, t)) },
439 )
440 .boxed();
441
442 while let Some(sample) = samples.next().await {
443 if let rustrtc::media::frame::MediaSample::Audio(frame) = sample {
444 if let Err(_) = tx.send(frame) {
445 break;
446 }
447 } else {
448 debug!(track_id=%track_id, "Received non-audio sample");
449 }
450 }
451 info!(track_id=%track_id, "RtcTrack receiving worker stopped");
452 }
453
454 async fn process_audio_frame(
455 frame: rustrtc::media::frame::AudioFrame,
456 track_id: &TrackId,
457 packet_sender: &Arc<Mutex<Option<TrackPacketSender>>>,
458 processor_chain: &mut ProcessorChain,
459 default_payload_type: u8,
460 ) {
461 let packet_sender = packet_sender.lock().await;
462 if let Some(sender) = packet_sender.as_ref() {
463 let payload_type = frame.payload_type.unwrap_or(default_payload_type);
464 let src_codec = match processor_chain.codec.get_codec_for_pt(payload_type) {
465 Some(c) => c,
466 None => {
467 debug!(track_id=%track_id, "Unknown payload type {}, skipping frame", payload_type);
468 return;
469 }
470 };
471
472 let mut af = AudioFrame {
473 track_id: track_id.clone(),
474 samples: crate::media::Samples::RTP {
475 payload_type,
476 payload: frame.data.to_vec(),
477 sequence_number: frame.sequence_number.unwrap_or(0),
478 },
479 timestamp: crate::media::get_timestamp(),
480 sample_rate: src_codec.samplerate(),
481 channels: src_codec.channels(),
482 ..Default::default()
483 };
484 if let Err(e) = processor_chain.process_frame(&mut af) {
485 debug!(track_id=%track_id, "processor_chain process_frame error: {:?}", e);
486 }
487
488 sender.send(af).ok();
489 }
490 }
491
492 pub fn parse_sdp_payload_types(&mut self, sdp_type: SdpType, sdp_str: &str) -> Result<()> {
493 use crate::media::negotiate::parse_rtpmap;
494 let sdp = rustrtc::SessionDescription::parse(sdp_type, sdp_str)?;
495
496 if let Some(media) = sdp
497 .media_sections
498 .iter()
499 .find(|m| m.kind == MediaKind::Audio)
500 {
501 for attr in &media.attributes {
502 if attr.key == "rtpmap" {
503 if let Some(value) = &attr.value {
504 if let Ok((pt, codec, _, _)) = parse_rtpmap(value) {
505 self.encoder.set_payload_type(pt, codec.clone());
506 self.processor_chain.codec.set_payload_type(pt, codec);
507 }
508 }
509 }
510 }
511
512 let mut negotiated = None;
514
515 if sdp_type == rustrtc::sdp::SdpType::Answer && !self.rtc_config.codecs.is_empty() {
518 for preferred_codec in &self.rtc_config.codecs {
519 if *preferred_codec == CodecType::TelephoneEvent {
520 continue;
521 }
522 for fmt in &media.formats {
523 if let Ok(pt) = fmt.parse::<u8>() {
524 let codec = self.encoder.get_codec_for_pt(pt);
525 if let Some(c) = codec {
526 if c == *preferred_codec {
527 negotiated = Some((pt, c));
528 break;
529 }
530 }
531 }
532 }
533 if negotiated.is_some() {
534 break;
535 }
536 }
537 }
538
539 if negotiated.is_none() {
541 for fmt in &media.formats {
542 if let Ok(pt) = fmt.parse::<u8>() {
543 let codec = self.encoder.get_codec_for_pt(pt);
544 if let Some(codec) = codec {
545 if codec != CodecType::TelephoneEvent {
546 negotiated = Some((pt, codec));
547 break;
548 }
549 }
550 }
551 }
552 }
553
554 if let Some((pt, codec)) = negotiated {
555 info!(track_id=%self.track_id, "Negotiated primary audio PT {} ({:?})", pt, codec);
556 self.payload_type = Some(pt);
557 }
558 }
559 Ok(())
560 }
561
562 fn normalize_sdp(sdp: &str) -> String {
563 sdp.lines()
564 .map(|line| {
565 if line.starts_with("o=") {
566 let parts: Vec<&str> = line.split_whitespace().collect();
567 if parts.len() >= 3 {
568 return format!("o= {} {}", parts[1], parts[2]);
569 }
570 }
571 line.to_string()
572 })
573 .filter(|line| {
574 !line.starts_with("t=") && !line.starts_with("a=ssrc:") && !line.starts_with("a=msid:") && !line.trim().is_empty()
578 })
579 .collect::<Vec<_>>()
580 .join("\n")
581 }
582
583 async fn update_remote_description_internal(
584 &mut self,
585 answer: &String,
586 force_update: bool,
587 ) -> Result<()> {
588 info!(
589 track_id=%self.track_id,
590 "update_remote_description_internal called. force={}, last_sdp_is_some={}, mode={:?}",
591 force_update,
592 self.last_remote_sdp.is_some(),
593 self.rtc_config.mode
594 );
595
596 if let Some(pc) = &self.peer_connection {
597 if !force_update {
598 if let Some(ref last_sdp) = self.last_remote_sdp {
599 if Self::normalize_sdp(last_sdp) == Self::normalize_sdp(answer) {
600 debug!(track_id=%self.track_id, "SDP unchanged, skipping update_remote_description");
601 return Ok(());
602 }
603 }
604 } else {
605 debug!(track_id=%self.track_id, "Force update requested, skipping SDP comparison");
606 }
607
608 let _is_first_remote_sdp = self.last_remote_sdp.is_none();
609
610 let sdp_obj = rustrtc::SessionDescription::parse(rustrtc::SdpType::Answer, answer)?;
611 match pc.set_remote_description(sdp_obj.clone()).await {
612 Ok(_) => {
613 debug!(track_id=%self.track_id, "set_remote_description succeeded");
614 self.last_remote_sdp = Some(answer.clone());
615 }
616 Err(e) => {
617 if self.rtc_config.mode == TransportMode::Rtp {
618 info!(track_id=%self.track_id, "set_remote_description failed ({}), attempting to re-sync state for SIP update", e);
619
620 if let Some(current_local) = pc.local_description() {
621 let sdp = current_local.to_sdp_string();
622 for line in sdp.lines() {
623 if line.starts_with("a=ssrc:") {
624 info!(track_id=%self.track_id, "SSRC before re-sync: {}", line);
625 }
626 }
627 }
628
629 let offer = pc.create_offer().await?;
630
631 let sdp = offer.to_sdp_string();
632 for line in sdp.lines() {
633 if line.starts_with("a=ssrc:") {
634 info!(track_id=%self.track_id, "SSRC in new offer (re-sync): {}", line);
635 }
636 }
637
638 pc.set_local_description(offer)?;
639 pc.set_remote_description(sdp_obj).await?;
640 self.last_remote_sdp = Some(answer.clone());
641 info!(track_id=%self.track_id, "successfully re-synced WebRTC state for SIP update");
642 } else {
643 return Err(e.into());
644 }
645 }
646 }
647
648 self.parse_sdp_payload_types(rustrtc::SdpType::Answer, answer)?;
652 }
653 Ok(())
654 }
655}
656
657#[async_trait]
658impl Track for RtcTrack {
659 fn ssrc(&self) -> u32 {
660 self.ssrc
661 }
662 fn id(&self) -> &TrackId {
663 &self.track_id
664 }
665 fn config(&self) -> &TrackConfig {
666 &self.track_config
667 }
668 fn processor_chain(&mut self) -> &mut ProcessorChain {
669 &mut self.processor_chain
670 }
671
672 async fn handshake(&mut self, offer: String, _: Option<Duration>) -> Result<String> {
673 info!(track_id=%self.track_id, "rtc handshake start");
674 self.create().await?;
675
676 let pc = self.peer_connection.clone().ok_or_else(|| {
677 anyhow::anyhow!("No PeerConnection available for track {}", self.track_id)
678 })?;
679
680 debug!(track_id=%self.track_id, "Before set_remote_description: transceivers count = {}", pc.get_transceivers().len());
681 for (i, t) in pc.get_transceivers().iter().enumerate() {
682 debug!(track_id=%self.track_id, " Transceiver #{}: kind={:?}, mid={:?}, direction={:?}",
683 i, t.kind(), t.mid(), t.direction());
684 }
685
686 let sdp = rustrtc::SessionDescription::parse(rustrtc::SdpType::Offer, &offer)?;
687 pc.set_remote_description(sdp.clone()).await?;
688
689 debug!(track_id=%self.track_id, "After set_remote_description: transceivers count = {}", pc.get_transceivers().len());
690 for (i, t) in pc.get_transceivers().iter().enumerate() {
691 debug!(track_id=%self.track_id, " Transceiver #{}: kind={:?}, mid={:?}, direction={:?}, has_receiver={}",
692 i, t.kind(), t.mid(), t.direction(), t.receiver().is_some());
693 }
694
695 info!(track_id=%self.track_id, "Waiting for Track events (SSRC latching for RTP mode)");
698
699 self.parse_sdp_payload_types(rustrtc::SdpType::Offer, &offer)?;
700
701 let mut answer = pc.create_answer().await?;
702 crate::media::negotiate::intersect_answer(&sdp, &mut answer);
703 self.parse_sdp_payload_types(rustrtc::SdpType::Answer, &answer.to_sdp_string())?;
704
705 pc.set_local_description(answer.clone())?;
706
707 if self.rtc_config.mode != TransportMode::Rtp {
708 pc.wait_for_gathering_complete().await;
709 }
710
711 let final_answer = pc
712 .local_description()
713 .ok_or(anyhow::anyhow!("No local description"))?;
714
715 Ok(final_answer.to_sdp_string())
716 }
717
718 async fn update_remote_description(&mut self, answer: &String) -> Result<()> {
719 self.update_remote_description_internal(answer, false).await
720 }
721
722 async fn update_remote_description_force(&mut self, answer: &String) -> Result<()> {
723 self.update_remote_description_internal(answer, true).await
724 }
725
726 async fn start(
727 &mut self,
728 event_sender: EventSender,
729 packet_sender: TrackPacketSender,
730 ) -> Result<()> {
731 *self.packet_sender.lock().await = Some(packet_sender.clone());
732 *self.event_sender.lock().await = Some(event_sender.clone());
733 let token_clone = self.cancel_token.clone();
734 let event_sender_clone = event_sender.clone();
735 let track_id = self.track_id.clone();
736 let ssrc = self.ssrc;
737
738 if self.rtc_config.mode != TransportMode::Rtp {
739 let start_time = crate::media::get_timestamp();
740 crate::spawn(async move {
741 token_clone.cancelled().await;
742 let _ = event_sender_clone.send(SessionEvent::TrackEnd {
743 track_id,
744 timestamp: crate::media::get_timestamp(),
745 duration: crate::media::get_timestamp() - start_time,
746 ssrc,
747 play_id: None,
748 });
749 });
750 }
751
752 Ok(())
753 }
754
755 async fn stop(&self) -> Result<()> {
756 self.cancel_token.cancel();
757 if let Some(pc) = &self.peer_connection {
758 pc.close();
759 }
760 Ok(())
761 }
762
763 async fn send_packet(&mut self, packet: &AudioFrame) -> Result<()> {
764 let packet = packet.clone();
765
766 if let Some(source) = &self.local_source {
767 match &packet.samples {
768 crate::media::Samples::PCM { samples } => {
769 let payload_type = self.get_payload_type();
770 let (_, encoded) = self.encoder.encode(payload_type, packet.clone());
771 let target_codec = self
772 .encoder
773 .get_codec_for_pt(payload_type)
774 .ok_or_else(|| anyhow::anyhow!("Invalid codec type: {}", payload_type))?;
775 if !encoded.is_empty() {
776 let clock_rate = target_codec.clock_rate();
777
778 let now = Instant::now();
779 if let Some(last_time) = self.last_packet_time {
780 let elapsed = now.duration_since(last_time);
781 if elapsed.as_millis() > 50 {
782 let gap_increment =
783 (elapsed.as_millis() as u32 * clock_rate) / 1000;
784 self.next_rtp_timestamp += gap_increment;
785 self.need_marker = true;
786 }
787 }
788
789 self.last_packet_time = Some(now);
790
791 let timestamp_increment = (samples.len() as u64 * clock_rate as u64
792 / packet.sample_rate as u64
793 / self.track_config.channels as u64)
794 as u32;
795 let rtp_timestamp = self.next_rtp_timestamp;
796 self.next_rtp_timestamp += timestamp_increment;
797 let sequence_number = self.next_rtp_sequence_number;
798 self.next_rtp_sequence_number += 1;
799
800 let mut marker = false;
801 if self.need_marker {
802 marker = true;
803 self.need_marker = false;
804 }
805
806 let frame = RtcAudioFrame {
807 data: Bytes::from(encoded),
808 clock_rate,
809 payload_type: Some(payload_type),
810 sequence_number: Some(sequence_number),
811 rtp_timestamp,
812 marker,
813 ..Default::default()
814 };
815 source.try_send_audio(frame).ok();
816 }
817 }
818 crate::media::Samples::RTP {
819 payload,
820 payload_type,
821 sequence_number,
822 } => {
823 let target_codec = self
824 .encoder
825 .get_codec_for_pt(*payload_type)
826 .ok_or_else(|| anyhow::anyhow!("Invalid codec type: {}", payload_type))?;
827 let clock_rate = target_codec.clock_rate();
828
829 let now = Instant::now();
830 if let Some(last_time) = self.last_packet_time {
831 let elapsed = now.duration_since(last_time);
832 if elapsed.as_millis() > 50 {
833 let gap_increment = (elapsed.as_millis() as u32 * clock_rate) / 1000;
834 self.next_rtp_timestamp += gap_increment;
835 self.need_marker = true;
836 }
837 }
838 self.last_packet_time = Some(now);
839
840 let increment = match *payload_type {
841 0 | 8 | 18 => payload.len() as u32,
842 9 => payload.len() as u32,
843 111 => (clock_rate / 50) as u32,
844 _ => (clock_rate / 50) as u32,
845 };
846
847 let rtp_timestamp = self.next_rtp_timestamp;
848 self.next_rtp_timestamp += increment;
849 let sequence_number = *sequence_number;
850
851 let mut marker = false;
852 if self.need_marker {
853 marker = true;
854 self.need_marker = false;
855 }
856
857 let frame = RtcAudioFrame {
858 data: Bytes::from(payload.clone()),
859 clock_rate,
860 payload_type: Some(*payload_type),
861 sequence_number: Some(sequence_number),
862 rtp_timestamp,
863 marker,
864 ..Default::default()
865 };
866 source.try_send_audio(frame).ok();
867 }
868 _ => {}
869 }
870 }
871 Ok(())
872 }
873}
874
875impl RtcTrack {
876 fn get_payload_type(&self) -> u8 {
877 if let Some(pt) = self.payload_type {
878 return pt;
879 }
880
881 self.rtc_config.payload_type.unwrap_or_else(|| {
882 match self.rtc_config.preferred_codec.unwrap_or(CodecType::G722) {
883 CodecType::PCMU => 0,
884 CodecType::PCMA => 8,
885 #[cfg(feature = "opus")]
886 CodecType::Opus => 111,
887 CodecType::G722 => 9,
888 CodecType::G729 => 18,
889 _ => 111,
890 }
891 })
892 }
893}
894
895#[cfg(test)]
896mod tests {
897 use super::*;
898 use crate::media::track::TrackConfig;
899
900 #[test]
901 fn test_parse_sdp_payload_types() {
902 let track_id = "test-track".to_string();
903 let cancel_token = CancellationToken::new();
904 let mut track = RtcTrack::new(
905 cancel_token,
906 track_id,
907 TrackConfig::default(),
908 RtcTrackConfig::default(),
909 );
910
911 let sdp1 = "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=-\r\nc=IN IP4 127.0.0.1\r\nt=0 0\r\nm=audio 1234 RTP/AVP 8 0 101\r\na=rtpmap:8 PCMA/8000\r\na=rtpmap:0 PCMU/8000\r\na=rtpmap:101 telephone-event/8000\r\n";
913 track
914 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp1)
915 .expect("parse offer");
916 assert_eq!(track.get_payload_type(), 8);
917
918 let mut rtc_config = RtcTrackConfig::default();
920 rtc_config.preferred_codec = Some(CodecType::PCMU);
921 let mut track2 = RtcTrack::new(
922 CancellationToken::new(),
923 "test-track-2".to_string(),
924 TrackConfig::default(),
925 rtc_config,
926 );
927
928 let sdp2 = "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=-\r\nc=IN IP4 127.0.0.1\r\nt=0 0\r\nm=audio 1234 RTP/AVP 101 0 8\r\na=rtpmap:101 telephone-event/8000\r\na=rtpmap:0 PCMU/8000\r\na=rtpmap:8 PCMA/8000\r\n";
929 track2
930 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp2)
931 .expect("parse offer");
932 assert_eq!(track2.get_payload_type(), 0);
933
934 let sdp3 = "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=-\r\nc=IN IP4 127.0.0.1\r\nt=0 0\r\nm=audio 1234 RTP/AVP 111 101\r\na=rtpmap:111 opus/48000/2\r\na=rtpmap:101 telephone-event/8000\r\n";
936 track
937 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp3)
938 .expect("parse offer");
939 assert_eq!(track.get_payload_type(), 111);
940
941 let mut rtc_config = RtcTrackConfig::default();
944 rtc_config.preferred_codec = Some(CodecType::PCMU);
945 rtc_config.codecs = vec![CodecType::PCMU, CodecType::PCMA];
946 let mut track4 = RtcTrack::new(
947 CancellationToken::new(),
948 "test-track-4".to_string(),
949 TrackConfig::default(),
950 rtc_config,
951 );
952
953 let sdp4 = "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=-\r\nc=IN IP4 127.0.0.1\r\nt=0 0\r\nm=audio 1234 RTP/AVP 18 0 101\r\na=fmtp:18 annexb=yes\r\na=rtpmap:101 telephone-event/8000\r\n";
954 track4
955 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp4)
956 .expect("parse offer");
957 assert_eq!(track4.get_payload_type(), 18);
958
959 let answer4 = "v=0\r\no=- 0 0 IN IP4 127.0.0.1\r\ns=-\r\nc=IN IP4 127.0.0.1\r\nt=0 0\r\nm=audio 1234 RTP/AVP 0\r\na=rtpmap:0 PCMU/8000\r\n";
960 track4
961 .parse_sdp_payload_types(rustrtc::SdpType::Answer, answer4)
962 .expect("parse answer");
963 assert_eq!(track4.get_payload_type(), 0);
964 }
965
966 #[tokio::test]
967 async fn test_rtp_mode_handshake_spawns_handler() {
968 use rustrtc::TransportMode;
969
970 let track_id = "test-track-sip".to_string();
971 let cancel = CancellationToken::new();
972 let track_config = TrackConfig::default();
973 let mut rtc_config = RtcTrackConfig::default();
974 rtc_config.mode = TransportMode::Rtp;
975 rtc_config.preferred_codec = Some(CodecType::PCMU);
976 rtc_config.codecs = vec![CodecType::PCMU, CodecType::TelephoneEvent];
977
978 let mut track = RtcTrack::new(cancel, track_id, track_config, rtc_config);
979
980 let offer = "v=0\r\n\
982o=- 123456 123456 IN IP4 172.0.0.1\r\n\
983s=-\r\n\
984c=IN IP4 172.0.0.1\r\n\
985t=0 0\r\n\
986m=audio 10000 RTP/AVP 0 101\r\n\
987a=rtpmap:0 PCMU/8000\r\n\
988a=rtpmap:101 telephone-event/8000\r\n\
989a=sendrecv\r\n";
990
991 let res = track.handshake(offer.to_string(), None).await;
993 assert!(res.is_ok(), "handshake failed: {res:?}");
994
995 if let Some(pc) = &track.peer_connection {
997 let transceivers = pc.get_transceivers();
998 assert_eq!(transceivers.len(), 1);
1001 assert!(transceivers[0].receiver().is_some());
1002 } else {
1003 panic!("PeerConnection not initialized");
1004 }
1005 }
1006}