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 CodecType::Opus => AudioCapability::opus(),
183 };
184 caps.audio.push(cap);
185 }
186 config.media_capabilities = Some(caps);
187 }
188
189 let peer_connection = Arc::new(PeerConnection::new(config));
190 self.peer_connection = Some(peer_connection.clone());
191
192 let default_codec = CodecType::G722;
193 let codec = self.rtc_config.preferred_codec.unwrap_or(default_codec);
194
195 let (source, track) = Self::create_audio_track(codec, Some(self.track_id.clone()));
196 self.local_source = Some(source);
197
198 let payload_type = self
199 .rtc_config
200 .payload_type
201 .unwrap_or_else(|| codec.payload_type());
202
203 self.payload_type = Some(payload_type);
204
205 let params = RtpCodecParameters {
206 clock_rate: codec.clock_rate(),
207 channels: codec.channels() as u8,
208 payload_type,
209 ..Default::default()
210 };
211
212 peer_connection.add_track_with_stream_id(track, self.track_id.clone(), params)?;
213
214 self.spawn_handlers(
216 peer_connection.clone(),
217 self.track_id.clone(),
218 self.processor_chain.clone(),
219 payload_type,
220 self.event_sender.clone(),
221 self.media_ready_sent.clone(),
222 );
223
224 Ok(())
225 }
226
227 fn spawn_handlers(
228 &self,
229 pc: Arc<PeerConnection>,
230 track_id: TrackId,
231 processor_chain: ProcessorChain,
232 default_payload_type: u8,
233 event_sender: Arc<Mutex<Option<EventSender>>>,
234 media_ready_sent: Arc<AtomicBool>,
235 ) {
236 let cancel_token = self.cancel_token.clone();
237 let packet_sender = self.packet_sender.clone();
238 let pc_event = pc.clone();
239 let pc_stats = pc.clone();
240 let pc_state = pc.clone();
241 let track_id_log = track_id.clone();
242 let is_rtp_media = matches!(
243 self.rtc_config.mode,
244 TransportMode::Rtp | TransportMode::Srtp
245 );
246 let is_webrtc = self.rtc_config.mode != TransportMode::Rtp;
247
248 crate::spawn(async move {
249 info!(track_id=%track_id_log, "RtcTrack event/stats loop started");
250
251 let mut events = futures::stream::unfold(pc_event, |pc| async move {
252 pc.recv().await.map(|ev| (ev, pc))
253 })
254 .boxed();
255
256 let mut state_rx = if is_webrtc {
257 Some(pc_state.subscribe_peer_state())
258 } else {
259 None
260 };
261
262 let mut stats_interval = tokio::time::interval(Duration::from_secs(5));
263 let mut event_count = 0;
264 let mut workers = FuturesUnordered::new();
265
266 loop {
267 tokio::select! {
268 _ = cancel_token.cancelled() => {
269 debug!(track_id=%track_id_log, "RtcTrack loop cancelled");
270 break;
271 }
272
273 Some(event) = events.next() => {
274 event_count += 1;
275 let event_type = match &event {
276 PeerConnectionEvent::Track(_) => "Track",
277 PeerConnectionEvent::DataChannel(_) => "DataChannel",
278 };
279 debug!(track_id=%track_id_log, "Received PeerConnectionEvent #{}: {}", event_count, event_type);
280
281 if let PeerConnectionEvent::Track(transceiver) = event {
282 if let Some(receiver) = transceiver.receiver() {
283 let track = receiver.track();
284 if is_rtp_media {
285 let maybe_sender = event_sender.lock().await.clone();
286 if let Some(sender) = maybe_sender {
287 if media_ready_sent
288 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
289 .is_ok()
290 {
291 let result = sender.send(SessionEvent::MediaReady {
292 track_id: track_id_log.clone(),
293 timestamp: crate::media::get_timestamp(),
294 });
295 if result.is_err() {
296 media_ready_sent.store(false, Ordering::SeqCst);
297 }
298 }
299 }
300 }
301 info!(track_id=%track_id_log, "New track received");
302
303 let (f1, f2) = Self::create_track_workers(
304 track,
305 packet_sender.clone(),
306 track_id_log.clone(),
307 processor_chain.clone(),
308 default_payload_type,
309 );
310 workers.push(f1);
311 workers.push(f2);
312 }
313 }
314 }
315
316 _ = workers.next(), if !workers.is_empty() => {}
317
318 _ = stats_interval.tick() => {
319 match pc_stats.get_stats().await {
320 Ok(stats) => {
321 info!(track_id=%track_id_log, %stats, "RTCP Stats");
322 }
323 Err(e) => {
324 debug!(track_id=%track_id_log, "Failed to get stats: {:?}", e);
325 }
326 }
327 }
328
329 res = async {
331 if let Some(rx) = state_rx.as_mut() {
332 rx.changed().await
333 } else {
334 std::future::pending().await
335 }
336 } => {
337 if res.is_ok() {
338 if let Some(rx) = state_rx.as_ref() {
339 let s = *rx.borrow();
340 debug!(track_id=%track_id_log, "peer connection state changed: {:?}", s);
341 match s {
342 PeerConnectionState::Disconnected
343 | PeerConnectionState::Closed
344 | PeerConnectionState::Failed => {
345 info!(
346 track_id = %track_id_log,
347 "peer connection is {:?}, try to close", s
348 );
349 cancel_token.cancel();
350 pc_state.close();
351 break;
352 }
353 _ => {}
354 }
355 }
356 }
357 }
358 }
359 }
360 debug!(track_id=%track_id_log, "RtcTrack event/stats loop ended, total events: {}", event_count);
361 });
362 }
363
364 fn create_track_workers(
365 track: Arc<SampleStreamTrack>,
366 packet_sender_arc: Arc<Mutex<Option<TrackPacketSender>>>,
367 track_id: TrackId,
368 processor_chain: ProcessorChain,
369 default_payload_type: u8,
370 ) -> (
371 std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>,
372 std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>,
373 ) {
374 let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<rustrtc::media::frame::AudioFrame>();
375
376 let track_id_proc = track_id.clone();
378 let packet_sender_proc = packet_sender_arc.clone();
379 let processor_chain_proc = processor_chain.clone();
380 let proc_fut = Self::run_processing_worker(
381 rx,
382 track_id_proc,
383 packet_sender_proc,
384 processor_chain_proc,
385 default_payload_type,
386 );
387
388 let track_id_recv = track_id.clone();
390 let recv_fut = Self::run_receiving_worker(track, tx, track_id_recv);
391
392 (proc_fut.boxed(), recv_fut.boxed())
393 }
394
395 async fn run_processing_worker(
396 mut rx: tokio::sync::mpsc::UnboundedReceiver<rustrtc::media::frame::AudioFrame>,
397 track_id: TrackId,
398 packet_sender: Arc<Mutex<Option<TrackPacketSender>>>,
399 mut processor_chain: ProcessorChain,
400 default_payload_type: u8,
401 ) {
402 info!(track_id=%track_id, "RtcTrack processing worker started");
403 while let Some(frame) = rx.recv().await {
404 let res = std::panic::AssertUnwindSafe(Self::process_audio_frame(
405 frame,
406 &track_id,
407 &packet_sender,
408 &mut processor_chain,
409 default_payload_type,
410 ))
411 .catch_unwind()
412 .await;
413
414 if let Err(cause) = res {
415 let msg = if let Some(s) = cause.downcast_ref::<&str>() {
416 *s
417 } else if let Some(s) = cause.downcast_ref::<String>() {
418 &s[..]
419 } else {
420 "Unknown panic"
421 };
422 tracing::error!(track_id=%track_id, "RtcTrack processing worker PANIC: {}", msg);
423 break;
424 }
425 }
426 info!(track_id=%track_id, "RtcTrack processing worker stopped");
427 }
428
429 async fn run_receiving_worker(
430 track: Arc<SampleStreamTrack>,
431 tx: tokio::sync::mpsc::UnboundedSender<rustrtc::media::frame::AudioFrame>,
432 track_id: TrackId,
433 ) {
434 let mut samples =
435 futures::stream::unfold(
436 track,
437 |t| async move { t.recv().await.ok().map(|s| (s, t)) },
438 )
439 .boxed();
440
441 while let Some(sample) = samples.next().await {
442 if let rustrtc::media::frame::MediaSample::Audio(frame) = sample {
443 if let Err(_) = tx.send(frame) {
444 break;
445 }
446 } else {
447 debug!(track_id=%track_id, "Received non-audio sample");
448 }
449 }
450 info!(track_id=%track_id, "RtcTrack receiving worker stopped");
451 }
452
453 async fn process_audio_frame(
454 frame: rustrtc::media::frame::AudioFrame,
455 track_id: &TrackId,
456 packet_sender: &Arc<Mutex<Option<TrackPacketSender>>>,
457 processor_chain: &mut ProcessorChain,
458 default_payload_type: u8,
459 ) {
460 let packet_sender = packet_sender.lock().await;
461 if let Some(sender) = packet_sender.as_ref() {
462 let payload_type = frame.payload_type.unwrap_or(default_payload_type);
463 let src_codec = match processor_chain.codec.get_codec_for_pt(payload_type) {
464 Some(c) => c,
465 None => {
466 debug!(track_id=%track_id, "Unknown payload type {}, skipping frame", payload_type);
467 return;
468 }
469 };
470
471 let mut af = AudioFrame {
472 track_id: track_id.clone(),
473 samples: crate::media::Samples::RTP {
474 payload_type,
475 payload: frame.data.to_vec(),
476 sequence_number: frame.sequence_number.unwrap_or(0),
477 },
478 timestamp: crate::media::get_timestamp(),
479 sample_rate: src_codec.samplerate(),
480 channels: src_codec.channels(),
481 ..Default::default()
482 };
483 if let Err(e) = processor_chain.process_frame(&mut af) {
484 debug!(track_id=%track_id, "processor_chain process_frame error: {:?}", e);
485 }
486
487 sender.send(af).ok();
488 }
489 }
490
491 pub fn parse_sdp_payload_types(&mut self, sdp_type: SdpType, sdp_str: &str) -> Result<()> {
492 use crate::media::negotiate::parse_rtpmap;
493 let sdp = rustrtc::SessionDescription::parse(sdp_type, sdp_str)?;
494
495 if let Some(media) = sdp
496 .media_sections
497 .iter()
498 .find(|m| m.kind == MediaKind::Audio)
499 {
500 for attr in &media.attributes {
501 if attr.key == "rtpmap" {
502 if let Some(value) = &attr.value {
503 if let Ok((pt, codec, _, _)) = parse_rtpmap(value) {
504 self.encoder.set_payload_type(pt, codec.clone());
505 self.processor_chain.codec.set_payload_type(pt, codec);
506 }
507 }
508 }
509 }
510
511 let mut negotiated = None;
513
514 if sdp_type == rustrtc::sdp::SdpType::Answer && !self.rtc_config.codecs.is_empty() {
517 for preferred_codec in &self.rtc_config.codecs {
518 if *preferred_codec == CodecType::TelephoneEvent {
519 continue;
520 }
521 for fmt in &media.formats {
522 if let Ok(pt) = fmt.parse::<u8>() {
523 let codec = self.encoder.get_codec_for_pt(pt);
524 if let Some(c) = codec {
525 if c == *preferred_codec {
526 negotiated = Some((pt, c));
527 break;
528 }
529 }
530 }
531 }
532 if negotiated.is_some() {
533 break;
534 }
535 }
536 }
537
538 if negotiated.is_none() {
540 for fmt in &media.formats {
541 if let Ok(pt) = fmt.parse::<u8>() {
542 let codec = self.encoder.get_codec_for_pt(pt);
543 if let Some(codec) = codec {
544 if codec != CodecType::TelephoneEvent {
545 negotiated = Some((pt, codec));
546 break;
547 }
548 }
549 }
550 }
551 }
552
553 if let Some((pt, codec)) = negotiated {
554 info!(track_id=%self.track_id, "Negotiated primary audio PT {} ({:?})", pt, codec);
555 self.payload_type = Some(pt);
556 }
557 }
558 Ok(())
559 }
560
561 fn normalize_sdp(sdp: &str) -> String {
562 sdp.lines()
563 .map(|line| {
564 if line.starts_with("o=") {
565 let parts: Vec<&str> = line.split_whitespace().collect();
566 if parts.len() >= 3 {
567 return format!("o= {} {}", parts[1], parts[2]);
568 }
569 }
570 line.to_string()
571 })
572 .filter(|line| {
573 !line.starts_with("t=") && !line.starts_with("a=ssrc:") && !line.starts_with("a=msid:") && !line.trim().is_empty()
577 })
578 .collect::<Vec<_>>()
579 .join("\n")
580 }
581
582 async fn update_remote_description_internal(
583 &mut self,
584 answer: &String,
585 force_update: bool,
586 ) -> Result<()> {
587 info!(
588 track_id=%self.track_id,
589 "update_remote_description_internal called. force={}, last_sdp_is_some={}, mode={:?}",
590 force_update,
591 self.last_remote_sdp.is_some(),
592 self.rtc_config.mode
593 );
594
595 if let Some(pc) = &self.peer_connection {
596 if !force_update {
597 if let Some(ref last_sdp) = self.last_remote_sdp {
598 if Self::normalize_sdp(last_sdp) == Self::normalize_sdp(answer) {
599 debug!(track_id=%self.track_id, "SDP unchanged, skipping update_remote_description");
600 return Ok(());
601 }
602 }
603 } else {
604 debug!(track_id=%self.track_id, "Force update requested, skipping SDP comparison");
605 }
606
607 let _is_first_remote_sdp = self.last_remote_sdp.is_none();
608
609 let sdp_obj = rustrtc::SessionDescription::parse(rustrtc::SdpType::Answer, answer)?;
610 match pc.set_remote_description(sdp_obj.clone()).await {
611 Ok(_) => {
612 debug!(track_id=%self.track_id, "set_remote_description succeeded");
613 self.last_remote_sdp = Some(answer.clone());
614 }
615 Err(e) => {
616 if self.rtc_config.mode == TransportMode::Rtp {
617 info!(track_id=%self.track_id, "set_remote_description failed ({}), attempting to re-sync state for SIP update", e);
618
619 if let Some(current_local) = pc.local_description() {
620 let sdp = current_local.to_sdp_string();
621 for line in sdp.lines() {
622 if line.starts_with("a=ssrc:") {
623 info!(track_id=%self.track_id, "SSRC before re-sync: {}", line);
624 }
625 }
626 }
627
628 let offer = pc.create_offer().await?;
629
630 let sdp = offer.to_sdp_string();
631 for line in sdp.lines() {
632 if line.starts_with("a=ssrc:") {
633 info!(track_id=%self.track_id, "SSRC in new offer (re-sync): {}", line);
634 }
635 }
636
637 pc.set_local_description(offer)?;
638 pc.set_remote_description(sdp_obj).await?;
639 self.last_remote_sdp = Some(answer.clone());
640 info!(track_id=%self.track_id, "successfully re-synced WebRTC state for SIP update");
641 } else {
642 return Err(e.into());
643 }
644 }
645 }
646
647 self.parse_sdp_payload_types(rustrtc::SdpType::Answer, answer)?;
651 }
652 Ok(())
653 }
654}
655
656#[async_trait]
657impl Track for RtcTrack {
658 fn ssrc(&self) -> u32 {
659 self.ssrc
660 }
661 fn id(&self) -> &TrackId {
662 &self.track_id
663 }
664 fn config(&self) -> &TrackConfig {
665 &self.track_config
666 }
667 fn processor_chain(&mut self) -> &mut ProcessorChain {
668 &mut self.processor_chain
669 }
670
671 async fn handshake(&mut self, offer: String, _: Option<Duration>) -> Result<String> {
672 info!(track_id=%self.track_id, "rtc handshake start");
673 self.create().await?;
674
675 let pc = self.peer_connection.clone().ok_or_else(|| {
676 anyhow::anyhow!("No PeerConnection available for track {}", self.track_id)
677 })?;
678
679 debug!(track_id=%self.track_id, "Before set_remote_description: transceivers count = {}", pc.get_transceivers().len());
680 for (i, t) in pc.get_transceivers().iter().enumerate() {
681 debug!(track_id=%self.track_id, " Transceiver #{}: kind={:?}, mid={:?}, direction={:?}",
682 i, t.kind(), t.mid(), t.direction());
683 }
684
685 let sdp = rustrtc::SessionDescription::parse(rustrtc::SdpType::Offer, &offer)?;
686 pc.set_remote_description(sdp.clone()).await?;
687
688 debug!(track_id=%self.track_id, "After set_remote_description: transceivers count = {}", pc.get_transceivers().len());
689 for (i, t) in pc.get_transceivers().iter().enumerate() {
690 debug!(track_id=%self.track_id, " Transceiver #{}: kind={:?}, mid={:?}, direction={:?}, has_receiver={}",
691 i, t.kind(), t.mid(), t.direction(), t.receiver().is_some());
692 }
693
694 info!(track_id=%self.track_id, "Waiting for Track events (SSRC latching for RTP mode)");
697
698 self.parse_sdp_payload_types(rustrtc::SdpType::Offer, &offer)?;
699
700 let mut answer = pc.create_answer().await?;
701 crate::media::negotiate::intersect_answer(&sdp, &mut answer);
702 self.parse_sdp_payload_types(rustrtc::SdpType::Answer, &answer.to_sdp_string())?;
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 CodecType::Opus => 111,
885 CodecType::G722 => 9,
886 CodecType::G729 => 18,
887 _ => 111,
888 }
889 })
890 }
891}
892
893#[cfg(test)]
894mod tests {
895 use super::*;
896 use crate::media::track::TrackConfig;
897
898 #[test]
899 fn test_parse_sdp_payload_types() {
900 let track_id = "test-track".to_string();
901 let cancel_token = CancellationToken::new();
902 let mut track = RtcTrack::new(
903 cancel_token,
904 track_id,
905 TrackConfig::default(),
906 RtcTrackConfig::default(),
907 );
908
909 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";
911 track
912 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp1)
913 .expect("parse offer");
914 assert_eq!(track.get_payload_type(), 8);
915
916 let mut rtc_config = RtcTrackConfig::default();
918 rtc_config.preferred_codec = Some(CodecType::PCMU);
919 let mut track2 = RtcTrack::new(
920 CancellationToken::new(),
921 "test-track-2".to_string(),
922 TrackConfig::default(),
923 rtc_config,
924 );
925
926 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";
927 track2
928 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp2)
929 .expect("parse offer");
930 assert_eq!(track2.get_payload_type(), 0);
931
932 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";
934 track
935 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp3)
936 .expect("parse offer");
937 assert_eq!(track.get_payload_type(), 111);
938
939 let mut rtc_config = RtcTrackConfig::default();
942 rtc_config.preferred_codec = Some(CodecType::PCMU);
943 rtc_config.codecs = vec![CodecType::PCMU, CodecType::PCMA];
944 let mut track4 = RtcTrack::new(
945 CancellationToken::new(),
946 "test-track-4".to_string(),
947 TrackConfig::default(),
948 rtc_config,
949 );
950
951 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";
952 track4
953 .parse_sdp_payload_types(rustrtc::SdpType::Offer, sdp4)
954 .expect("parse offer");
955 assert_eq!(track4.get_payload_type(), 18);
956
957 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";
958 track4
959 .parse_sdp_payload_types(rustrtc::SdpType::Answer, answer4)
960 .expect("parse answer");
961 assert_eq!(track4.get_payload_type(), 0);
962 }
963
964 #[tokio::test]
965 async fn test_rtp_mode_handshake_spawns_handler() {
966 use rustrtc::TransportMode;
967
968 let track_id = "test-track-sip".to_string();
969 let cancel = CancellationToken::new();
970 let track_config = TrackConfig::default();
971 let mut rtc_config = RtcTrackConfig::default();
972 rtc_config.mode = TransportMode::Rtp;
973 rtc_config.preferred_codec = Some(CodecType::PCMU);
974 rtc_config.codecs = vec![CodecType::PCMU, CodecType::TelephoneEvent];
975
976 let mut track = RtcTrack::new(cancel, track_id, track_config, rtc_config);
977
978 let offer = "v=0\r\n\
980o=- 123456 123456 IN IP4 172.0.0.1\r\n\
981s=-\r\n\
982c=IN IP4 172.0.0.1\r\n\
983t=0 0\r\n\
984m=audio 10000 RTP/AVP 0 101\r\n\
985a=rtpmap:0 PCMU/8000\r\n\
986a=rtpmap:101 telephone-event/8000\r\n\
987a=sendrecv\r\n";
988
989 let res = track.handshake(offer.to_string(), None).await;
991 assert!(res.is_ok(), "handshake failed: {res:?}");
992
993 if let Some(pc) = &track.peer_connection {
995 let transceivers = pc.get_transceivers();
996 assert_eq!(transceivers.len(), 1);
999 assert!(transceivers[0].receiver().is_some());
1000 } else {
1001 panic!("PeerConnection not initialized");
1002 }
1003 }
1004}