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
704 pc.set_local_description(answer.clone())?;
705
706 if self.rtc_config.mode != TransportMode::Rtp {
707 pc.wait_for_gathering_complete().await;
708 }
709
710 let final_answer = pc
711 .local_description()
712 .ok_or(anyhow::anyhow!("No local description"))?;
713
714 Ok(final_answer.to_sdp_string())
715 }
716
717 async fn update_remote_description(&mut self, answer: &String) -> Result<()> {
718 self.update_remote_description_internal(answer, false).await
719 }
720
721 async fn update_remote_description_force(&mut self, answer: &String) -> Result<()> {
722 self.update_remote_description_internal(answer, true).await
723 }
724
725 async fn start(
726 &mut self,
727 event_sender: EventSender,
728 packet_sender: TrackPacketSender,
729 ) -> Result<()> {
730 *self.packet_sender.lock().await = Some(packet_sender.clone());
731 *self.event_sender.lock().await = Some(event_sender.clone());
732 let token_clone = self.cancel_token.clone();
733 let event_sender_clone = event_sender.clone();
734 let track_id = self.track_id.clone();
735 let ssrc = self.ssrc;
736
737 if self.rtc_config.mode != TransportMode::Rtp {
738 let start_time = crate::media::get_timestamp();
739 crate::spawn(async move {
740 token_clone.cancelled().await;
741 let _ = event_sender_clone.send(SessionEvent::TrackEnd {
742 track_id,
743 timestamp: crate::media::get_timestamp(),
744 duration: crate::media::get_timestamp() - start_time,
745 ssrc,
746 play_id: None,
747 });
748 });
749 }
750
751 Ok(())
752 }
753
754 async fn stop(&self) -> Result<()> {
755 self.cancel_token.cancel();
756 if let Some(pc) = &self.peer_connection {
757 pc.close();
758 }
759 Ok(())
760 }
761
762 async fn send_packet(&mut self, packet: &AudioFrame) -> Result<()> {
763 let packet = packet.clone();
764
765 if let Some(source) = &self.local_source {
766 match &packet.samples {
767 crate::media::Samples::PCM { samples } => {
768 let payload_type = self.get_payload_type();
769 let (_, encoded) = self.encoder.encode(payload_type, packet.clone());
770 let target_codec = self
771 .encoder
772 .get_codec_for_pt(payload_type)
773 .ok_or_else(|| anyhow::anyhow!("Invalid codec type: {}", payload_type))?;
774 if !encoded.is_empty() {
775 let clock_rate = target_codec.clock_rate();
776
777 let now = Instant::now();
778 if let Some(last_time) = self.last_packet_time {
779 let elapsed = now.duration_since(last_time);
780 if elapsed.as_millis() > 50 {
781 let gap_increment =
782 (elapsed.as_millis() as u32 * clock_rate) / 1000;
783 self.next_rtp_timestamp += gap_increment;
784 self.need_marker = true;
785 }
786 }
787
788 self.last_packet_time = Some(now);
789
790 let timestamp_increment = (samples.len() as u64 * clock_rate as u64
791 / packet.sample_rate as u64
792 / self.track_config.channels as u64)
793 as u32;
794 let rtp_timestamp = self.next_rtp_timestamp;
795 self.next_rtp_timestamp += timestamp_increment;
796 let sequence_number = self.next_rtp_sequence_number;
797 self.next_rtp_sequence_number += 1;
798
799 let mut marker = false;
800 if self.need_marker {
801 marker = true;
802 self.need_marker = false;
803 }
804
805 let frame = RtcAudioFrame {
806 data: Bytes::from(encoded),
807 clock_rate,
808 payload_type: Some(payload_type),
809 sequence_number: Some(sequence_number),
810 rtp_timestamp,
811 marker,
812 ..Default::default()
813 };
814 source.try_send_audio(frame).ok();
815 }
816 }
817 crate::media::Samples::RTP {
818 payload,
819 payload_type,
820 sequence_number,
821 } => {
822 let target_codec = self
823 .encoder
824 .get_codec_for_pt(*payload_type)
825 .ok_or_else(|| anyhow::anyhow!("Invalid codec type: {}", payload_type))?;
826 let clock_rate = target_codec.clock_rate();
827
828 let now = Instant::now();
829 if let Some(last_time) = self.last_packet_time {
830 let elapsed = now.duration_since(last_time);
831 if elapsed.as_millis() > 50 {
832 let gap_increment = (elapsed.as_millis() as u32 * clock_rate) / 1000;
833 self.next_rtp_timestamp += gap_increment;
834 self.need_marker = true;
835 }
836 }
837 self.last_packet_time = Some(now);
838
839 let increment = match *payload_type {
840 0 | 8 | 18 => payload.len() as u32,
841 9 => payload.len() as u32,
842 111 => (clock_rate / 50) as u32,
843 _ => (clock_rate / 50) as u32,
844 };
845
846 let rtp_timestamp = self.next_rtp_timestamp;
847 self.next_rtp_timestamp += increment;
848 let sequence_number = *sequence_number;
849
850 let mut marker = false;
851 if self.need_marker {
852 marker = true;
853 self.need_marker = false;
854 }
855
856 let frame = RtcAudioFrame {
857 data: Bytes::from(payload.clone()),
858 clock_rate,
859 payload_type: Some(*payload_type),
860 sequence_number: Some(sequence_number),
861 rtp_timestamp,
862 marker,
863 ..Default::default()
864 };
865 source.try_send_audio(frame).ok();
866 }
867 _ => {}
868 }
869 }
870 Ok(())
871 }
872}
873
874impl RtcTrack {
875 fn get_payload_type(&self) -> u8 {
876 if let Some(pt) = self.payload_type {
877 return pt;
878 }
879
880 self.rtc_config.payload_type.unwrap_or_else(|| {
881 match self.rtc_config.preferred_codec.unwrap_or(CodecType::G722) {
882 CodecType::PCMU => 0,
883 CodecType::PCMA => 8,
884 #[cfg(feature = "opus")]
885 CodecType::Opus => 111,
886 CodecType::G722 => 9,
887 CodecType::G729 => 18,
888 _ => 111,
889 }
890 })
891 }
892}
893
894#[cfg(test)]
895mod tests {
896 use super::*;
897 use crate::media::track::TrackConfig;
898
899 #[test]
900 fn test_parse_sdp_payload_types() {
901 let track_id = "test-track".to_string();
902 let cancel_token = CancellationToken::new();
903 let mut track = RtcTrack::new(
904 cancel_token,
905 track_id,
906 TrackConfig::default(),
907 RtcTrackConfig::default(),
908 );
909
910 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";
912 track
913 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp1)
914 .expect("parse offer");
915 assert_eq!(track.get_payload_type(), 8);
916
917 let mut rtc_config = RtcTrackConfig::default();
919 rtc_config.preferred_codec = Some(CodecType::PCMU);
920 let mut track2 = RtcTrack::new(
921 CancellationToken::new(),
922 "test-track-2".to_string(),
923 TrackConfig::default(),
924 rtc_config,
925 );
926
927 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";
928 track2
929 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp2)
930 .expect("parse offer");
931 assert_eq!(track2.get_payload_type(), 0);
932
933 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";
935 track
936 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp3)
937 .expect("parse offer");
938 assert_eq!(track.get_payload_type(), 111);
939 }
940
941 #[tokio::test]
942 async fn test_rtp_mode_handshake_spawns_handler() {
943 use rustrtc::TransportMode;
944
945 let track_id = "test-track-sip".to_string();
946 let cancel = CancellationToken::new();
947 let track_config = TrackConfig::default();
948 let mut rtc_config = RtcTrackConfig::default();
949 rtc_config.mode = TransportMode::Rtp;
950
951 let mut track = RtcTrack::new(cancel, track_id, track_config, rtc_config);
952
953 let offer = "v=0\r\n\
955o=- 123456 123456 IN IP4 172.0.0.1\r\n\
956s=-\r\n\
957c=IN IP4 172.0.0.1\r\n\
958t=0 0\r\n\
959m=audio 10000 RTP/AVP 0 101\r\n\
960a=rtpmap:0 PCMU/8000\r\n\
961a=rtpmap:101 telephone-event/8000\r\n\
962a=sendrecv\r\n";
963
964 let res = track.handshake(offer.to_string(), None).await;
966 assert!(res.is_ok());
967
968 if let Some(pc) = &track.peer_connection {
970 let transceivers = pc.get_transceivers();
971 assert_eq!(transceivers.len(), 1);
974 assert!(transceivers[0].receiver().is_some());
975 } else {
976 panic!("PeerConnection not initialized");
977 }
978 }
979}