Skip to content

Commit f189294

Browse files
committed
improve: use receive timestamp for rtp capture
1 parent a0e1b9f commit f189294

7 files changed

Lines changed: 200 additions & 84 deletions

File tree

src/bin/sipflow.rs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -224,20 +224,24 @@ mod tests {
224224
let mut packets = Vec::new(); // should use Vec<(i32, u64, Vec<u8>)>
225225
let payload = vec![0x7F; 160]; // Silence
226226

227-
// 12 bytes dummy header
227+
// 12 bytes RTP header
228228
let mut header = vec![0u8; 12];
229+
header[0] = 0x80; // RTP v2
229230
header[1] = 0; // PCMU
230231

231232
let mut p1 = header.clone();
233+
p1[4..8].copy_from_slice(&1000u32.to_be_bytes());
232234
p1.extend_from_slice(&payload);
233235
packets.push((0, 1000u64, p1));
234236

235237
let mut p2 = header.clone();
238+
p2[4..8].copy_from_slice(&1000u32.to_be_bytes());
236239
p2.extend_from_slice(&payload);
237240
packets.push((1, 1000u64, p2));
238241

239242
// Next 20ms
240243
let mut p3 = header.clone();
244+
p3[4..8].copy_from_slice(&1160u32.to_be_bytes());
241245
p3.extend_from_slice(&payload);
242246
packets.push((0, 1160u64, p3));
243247

src/media/bridge.rs

Lines changed: 19 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@
4141
//! 3. The bridge's WebRTC side connects to the WebRTC client
4242
//! 4. The bridge's RTP side connects to the SIP/RTP endpoint
4343
44+
use crate::media::ReceiveTimestampClock;
4445
use crate::media::recorder::{Leg as RecLeg, Recorder};
4546
use crate::media::transcoder::{RtpTiming, Transcoder, rewrite_dtmf_duration};
4647
use anyhow::Result;
@@ -340,7 +341,8 @@ pub struct BridgePeer {
340341
/// When true, the bridge skips writing to the recorder (recording paused).
341342
recording_paused: Arc<AtomicBool>,
342343
/// Non-blocking RTP capture tee for SipFlow.
343-
sipflow_tx: Option<mpsc::Sender<(RecLeg, MediaSample)>>,
344+
sipflow_tx: Option<mpsc::Sender<(RecLeg, MediaSample, u64)>>,
345+
receive_clock: ReceiveTimestampClock,
344346
dtmf_sink: Arc<parking_lot::RwLock<Option<BridgeDtmfSink>>>,
345347
/// Audio sender channels for forwarding — fast-path aliases
346348
caller_send: Arc<AsyncMutex<Option<MediaSender>>>,
@@ -411,6 +413,7 @@ impl BridgePeer {
411413
recorder: None,
412414
recording_paused: Arc::new(AtomicBool::new(false)),
413415
sipflow_tx: None,
416+
receive_clock: ReceiveTimestampClock::new(),
414417
dtmf_sink: Arc::new(parking_lot::RwLock::new(None)),
415418
caller_send: Arc::new(AsyncMutex::new(None)),
416419
callee_send: Arc::new(AsyncMutex::new(None)),
@@ -1314,6 +1317,7 @@ impl BridgePeer {
13141317
let recorder = self.recorder.clone();
13151318
let recording_paused = self.recording_paused.clone();
13161319
let sipflow_tx = self.sipflow_tx.clone();
1320+
let receive_clock = self.receive_clock.clone();
13171321
let dtmf_sink = Arc::clone(&self.dtmf_sink);
13181322
let caller_to_callee_transcoder = Arc::clone(&self.caller_to_callee_transcoder);
13191323
let caller_to_callee_timing = Arc::clone(&self.caller_to_callee_timing);
@@ -1432,6 +1436,7 @@ impl BridgePeer {
14321436
if !is_video { recorder.clone() } else { None },
14331437
if !is_video { Some(RecLeg::A) } else { None },
14341438
if !is_video { sipflow_tx.clone() } else { None },
1439+
receive_clock.clone(),
14351440
recording_paused.clone(),
14361441
Arc::clone(&dtmf_sink),
14371442
Some(Arc::clone(&caller_to_callee_transcoder)),
@@ -1556,6 +1561,7 @@ impl BridgePeer {
15561561
if !is_video { recorder.clone() } else { None },
15571562
if !is_video { Some(RecLeg::B) } else { None },
15581563
if !is_video { sipflow_tx.clone() } else { None },
1564+
receive_clock.clone(),
15591565
recording_paused.clone(),
15601566
Arc::clone(&dtmf_sink),
15611567
Some(Arc::clone(&callee_to_caller_transcoder)),
@@ -1603,7 +1609,8 @@ impl BridgePeer {
16031609
leg_stats: Arc<LegStats>,
16041610
recorder: Option<Arc<parking_lot::RwLock<Option<Recorder>>>>,
16051611
recorder_leg: Option<RecLeg>,
1606-
sipflow_tx: Option<mpsc::Sender<(RecLeg, MediaSample)>>,
1612+
sipflow_tx: Option<mpsc::Sender<(RecLeg, MediaSample, u64)>>,
1613+
receive_clock: ReceiveTimestampClock,
16071614
recording_paused: Arc<AtomicBool>,
16081615
dtmf_sink: Arc<parking_lot::RwLock<Option<BridgeDtmfSink>>>,
16091616
transcoder: Option<Arc<parking_lot::Mutex<Option<Transcoder>>>>,
@@ -1661,6 +1668,7 @@ impl BridgePeer {
16611668
match sample_result {
16621669
Ok(sample) => {
16631670
packet_count += 1;
1671+
let received_at_micros = receive_clock.now_micros();
16641672
if !is_video {
16651673
if let (Some(rec), Some(leg)) = (&recorder, recorder_leg)
16661674
&& !recording_paused.load(std::sync::atomic::Ordering::Relaxed)
@@ -1677,7 +1685,7 @@ impl BridgePeer {
16771685
}
16781686

16791687
if let (Some(tx), Some(leg)) = (&sipflow_tx, recorder_leg) {
1680-
let _ = tx.try_send((leg, sample.clone()));
1688+
let _ = tx.try_send((leg, sample.clone(), received_at_micros));
16811689
}
16821690
}
16831691
if !is_video {
@@ -2005,7 +2013,8 @@ impl BridgePeer {
20052013
leg_stats: Arc<LegStats>,
20062014
recorder: Option<Arc<parking_lot::RwLock<Option<Recorder>>>>,
20072015
recorder_leg: Option<RecLeg>,
2008-
sipflow_tx: Option<mpsc::Sender<(RecLeg, MediaSample)>>,
2016+
sipflow_tx: Option<mpsc::Sender<(RecLeg, MediaSample, u64)>>,
2017+
receive_clock: ReceiveTimestampClock,
20092018
recording_paused: Arc<AtomicBool>,
20102019
dtmf_sink: Arc<parking_lot::RwLock<Option<BridgeDtmfSink>>>,
20112020
transcoder: Option<Arc<parking_lot::Mutex<Option<Transcoder>>>>,
@@ -2025,6 +2034,7 @@ impl BridgePeer {
20252034
recorder,
20262035
recorder_leg,
20272036
sipflow_tx,
2037+
receive_clock,
20282038
recording_paused,
20292039
dtmf_sink,
20302040
transcoder,
@@ -2102,7 +2112,7 @@ pub struct BridgePeerBuilder {
21022112
ice_servers: Vec<IceServer>,
21032113
recorder: Option<Arc<parking_lot::RwLock<Option<Recorder>>>>,
21042114
recording_paused: Arc<AtomicBool>,
2105-
sipflow_tx: Option<mpsc::Sender<(RecLeg, MediaSample)>>,
2115+
sipflow_tx: Option<mpsc::Sender<(RecLeg, MediaSample, u64)>>,
21062116
cname: Option<String>,
21072117
}
21082118

@@ -2251,7 +2261,7 @@ impl BridgePeerBuilder {
22512261

22522262
pub fn with_sipflow_capture(
22532263
mut self,
2254-
sipflow_tx: mpsc::Sender<(RecLeg, MediaSample)>,
2264+
sipflow_tx: mpsc::Sender<(RecLeg, MediaSample, u64)>,
22552265
) -> Self {
22562266
self.sipflow_tx = Some(sipflow_tx);
22572267
self
@@ -3632,6 +3642,7 @@ mod tests {
36323642
None, // no recorder
36333643
None, // no recorder leg
36343644
None, // no sipflow capture
3645+
ReceiveTimestampClock::new(),
36353646
Arc::new(AtomicBool::new(false)), // not paused
36363647
ds,
36373648
tr,
@@ -3736,6 +3747,7 @@ mod tests {
37363747
None,
37373748
None,
37383749
None,
3750+
ReceiveTimestampClock::new(),
37393751
Arc::new(AtomicBool::new(false)),
37403752
ds,
37413753
tr,
@@ -4051,6 +4063,7 @@ mod tests {
40514063
None,
40524064
None,
40534065
None,
4066+
ReceiveTimestampClock::new(),
40544067
Arc::new(AtomicBool::new(false)),
40554068
ds,
40564069
None, // no transcoder

src/media/forwarding_track.rs

Lines changed: 17 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
use crate::media::negotiate::NegotiatedLegProfile;
22
use crate::media::transcoder::{RtpTiming, Transcoder, rewrite_dtmf_duration};
3-
use crate::media::{Track, recorder::Leg};
3+
use crate::media::{ReceiveTimestampClock, Track, recorder::Leg};
44
use anyhow::{Result, anyhow};
55
use async_trait::async_trait;
66
use parking_lot::Mutex;
@@ -43,7 +43,8 @@ pub struct ForwardingTrack {
4343
audio_timing: Mutex<Option<RtpTiming>>,
4444
dtmf_timing: Mutex<Option<RtpTiming>>,
4545
recorder_tx: Option<mpsc::Sender<(Leg, MediaSample)>>,
46-
sipflow_tx: Option<mpsc::Sender<(Leg, MediaSample)>>,
46+
sipflow_tx: Option<mpsc::Sender<(Leg, MediaSample, u64)>>,
47+
receive_clock: ReceiveTimestampClock,
4748
recorder_leg: Leg,
4849
dtmf_mapping: Mutex<Option<DtmfMapping>>,
4950
}
@@ -60,7 +61,7 @@ impl ForwardingTrack {
6061
track_id: String,
6162
inner: Arc<dyn MediaStreamTrack>,
6263
recorder_tx: Option<mpsc::Sender<(Leg, MediaSample)>>,
63-
sipflow_tx: Option<mpsc::Sender<(Leg, MediaSample)>>,
64+
sipflow_tx: Option<mpsc::Sender<(Leg, MediaSample, u64)>>,
6465
recorder_leg: Leg,
6566
ingress_profile: NegotiatedLegProfile,
6667
egress_profile: NegotiatedLegProfile,
@@ -78,6 +79,7 @@ impl ForwardingTrack {
7879
dtmf_timing: Mutex::new(None),
7980
recorder_tx,
8081
sipflow_tx,
82+
receive_clock: ReceiveTimestampClock::new(),
8183
recorder_leg,
8284
dtmf_mapping: Mutex::new(None),
8385
}
@@ -250,14 +252,15 @@ impl MediaStreamTrack for ForwardingTrack {
250252
let audio_mapping = self.audio_mapping.lock().clone();
251253
let dtmf_mapping = self.dtmf_mapping.lock().clone();
252254
let sample = self.inner.recv().await?;
255+
let received_at_micros = self.receive_clock.now_micros();
253256

254257
if let Some(tx) = &self.recorder_tx {
255258
let _ = tx.try_send((self.recorder_leg, sample.clone()));
256259
}
257260

258261
// SipFlow RTP recording: non-blocking tee, drops if consumer falls behind.
259262
if let Some(tx) = &self.sipflow_tx {
260-
let _ = tx.try_send((self.recorder_leg, sample.clone()));
263+
let _ = tx.try_send((self.recorder_leg, sample.clone(), received_at_micros));
261264
}
262265

263266
if let MediaSample::Audio(ref frame) = sample {
@@ -471,7 +474,7 @@ mod tests {
471474
/// blocking the hot path and without interfering with the recorder_tx.
472475
#[tokio::test]
473476
async fn sipflow_tx_receives_sample() {
474-
let (sf_tx, mut sf_rx) = mpsc::channel::<(Leg, MediaSample)>(256);
477+
let (sf_tx, mut sf_rx) = mpsc::channel::<(Leg, MediaSample, u64)>(256);
475478
let sample = audio_sample(0 /* PCMU */);
476479
let track = OneShotTrack::new(sample.clone());
477480

@@ -494,16 +497,18 @@ mod tests {
494497
assert!(matches!(result, MediaSample::Audio(_)));
495498

496499
// sipflow channel must also have received the sample.
497-
let (leg, _sf_sample) = sf_rx.try_recv().expect("sample must be in sipflow channel");
500+
let (leg, _sf_sample, received_at_micros) =
501+
sf_rx.try_recv().expect("sample must be in sipflow channel");
498502
assert_eq!(leg, Leg::A);
503+
assert!(received_at_micros > 0);
499504
}
500505

501506
/// Both recorder_tx AND sipflow_tx can be active simultaneously; each
502507
/// must receive its own copy of the sample.
503508
#[tokio::test]
504509
async fn both_recorder_and_sipflow_receive_sample() {
505510
let (rec_tx, mut rec_rx) = mpsc::channel::<(Leg, MediaSample)>(256);
506-
let (sf_tx, mut sf_rx) = mpsc::channel::<(Leg, MediaSample)>(256);
511+
let (sf_tx, mut sf_rx) = mpsc::channel::<(Leg, MediaSample, u64)>(256);
507512
let sample = audio_sample(0 /* PCMU */);
508513
let track = OneShotTrack::new(sample.clone());
509514

@@ -527,16 +532,18 @@ mod tests {
527532
let (rec_leg, _) = rec_rx
528533
.try_recv()
529534
.expect("recorder channel must have sample");
530-
let (sf_leg, _) = sf_rx.try_recv().expect("sipflow channel must have sample");
535+
let (sf_leg, _, received_at_micros) =
536+
sf_rx.try_recv().expect("sipflow channel must have sample");
531537
assert_eq!(rec_leg, Leg::B);
532538
assert_eq!(sf_leg, Leg::B);
539+
assert!(received_at_micros > 0);
533540
}
534541

535542
#[tokio::test]
536543
async fn sipflow_full_channel_does_not_block() {
537-
let (sf_tx, _sf_rx) = mpsc::channel::<(Leg, MediaSample)>(1);
544+
let (sf_tx, _sf_rx) = mpsc::channel::<(Leg, MediaSample, u64)>(1);
538545

539-
let _ = sf_tx.try_send((Leg::A, audio_sample(0)));
546+
let _ = sf_tx.try_send((Leg::A, audio_sample(0), 1));
540547

541548
let track = OneShotTrack::new(audio_sample(0));
542549
let ft = ForwardingTrack::new(

src/media/mod.rs

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,37 @@ pub fn get_timestamp() -> u64 {
5050
.as_millis() as u64
5151
}
5252

53+
#[derive(Debug, Clone)]
54+
pub struct ReceiveTimestampClock {
55+
base_instant: std::time::Instant,
56+
base_epoch_micros: u64,
57+
}
58+
59+
impl ReceiveTimestampClock {
60+
pub fn new() -> Self {
61+
let base_epoch_micros = std::time::SystemTime::now()
62+
.duration_since(std::time::UNIX_EPOCH)
63+
.map(|duration| duration.as_micros() as u64)
64+
.unwrap_or_default();
65+
66+
Self {
67+
base_instant: std::time::Instant::now(),
68+
base_epoch_micros,
69+
}
70+
}
71+
72+
pub fn now_micros(&self) -> u64 {
73+
self.base_epoch_micros
74+
.saturating_add(self.base_instant.elapsed().as_micros() as u64)
75+
}
76+
}
77+
78+
impl Default for ReceiveTimestampClock {
79+
fn default() -> Self {
80+
Self::new()
81+
}
82+
}
83+
5384
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5485
struct AudioFrameTiming {
5586
pcm_sample_rate: u32,

src/proxy/proxy_call/sip_session.rs

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,7 @@ use tracing::{debug, error, info, warn};
7777
type SipFlowRtpCaptureTx = mpsc::Sender<(
7878
crate::media::recorder::Leg,
7979
rustrtc::media::frame::MediaSample,
80+
u64,
8081
)>;
8182

8283
mod conference;
@@ -4581,19 +4582,20 @@ impl SipSession {
45814582
let (tx, mut rx) = mpsc::channel::<(
45824583
crate::media::recorder::Leg,
45834584
rustrtc::media::frame::MediaSample,
4585+
u64,
45844586
)>(crate::media::forwarding_track::ForwardingTrack::DEFAULT_SIPFLOW_CHANNEL_CAPACITY);
45854587

45864588
let call_id = self.context.session_id.clone();
45874589
crate::utils::spawn(async move {
45884590
use crate::sipflow::{SipFlowItem, SipFlowMsgType};
45894591

4590-
while let Some((leg, sample)) = rx.recv().await {
4592+
while let Some((leg, sample, received_at_micros)) = rx.recv().await {
45914593
if let rustrtc::media::frame::MediaSample::Audio(ref frame) = sample
45924594
&& let Some(ref rtp_packet) = frame.raw_packet
45934595
&& let Ok(rtp_bytes) = rtp_packet.marshal()
45944596
{
45954597
let item = SipFlowItem {
4596-
timestamp: frame.rtp_timestamp as u64,
4598+
timestamp: received_at_micros,
45974599
seq: frame.sequence_number.unwrap_or(0) as u64,
45984600
msg_type: SipFlowMsgType::Rtp,
45994601
src_addr: format!("{leg:?}"),
@@ -4947,6 +4949,7 @@ impl SipSession {
49474949
tokio::sync::mpsc::Sender<(
49484950
crate::media::recorder::Leg,
49494951
rustrtc::media::frame::MediaSample,
4952+
u64,
49504953
)>,
49514954
>,
49524955
session_id: &str,

0 commit comments

Comments
 (0)