Skip to content
72 changes: 72 additions & 0 deletions receiver/capture_latency.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
// SPDX-FileCopyrightText: 2026 The Pion community <https://pion.ly>
// SPDX-License-Identifier: MIT

//go:build !js

package receiver

import "time"

// CaptureTimeUsFromRTP recovers the capture instant (unix microseconds) that the
// sender encoded into an RTP timestamp via captureUs*ClockRate/1e6 (mod 2^32)
// (see sender's captureTimestampInterceptor). nowUnixUs is a recent wall-clock
// reading near packet receipt, used to disambiguate the 32-bit wrap; the capture
// instant is taken to be the most recent time whose encoded low 32 bits match
// rtpTS and which is not in the future relative to nowUnixUs. clockRate 0 falls
// back to 90 kHz, matching the sender. Mirrors the browser recovery
// getSynchronizationSources()[].rtpTimestamp * 1e6/ClockRate.
//
// The recovered value carries a quantization error of at most 1e6/ClockRate
// microseconds (~11 µs at 90 kHz, ~21 µs at 48 kHz) from the sender's and this
// function's integer divisions.
func CaptureTimeUsFromRTP(rtpTS uint32, clockRate uint32, nowUnixUs int64) int64 {
rate := int64(clockRate)
if rate == 0 {
rate = 90000
}

// Reduce ClockRate/1e6 to lowest terms so nowUnixUs*num and ticks*den stay
// within int64: nowUnixUs is unix micros (~1.7e15) and nowUnixUs*ClockRate
// would overflow. 90 kHz -> 9/100, 48 kHz -> 6/125. Mirrors the sender.
g := gcd(rate, 1_000_000)
num := rate / g
den := int64(1_000_000) / g

// Full (non-wrapped) tick count at nowUnixUs; may exceed 2^32.
nowTicks := nowUnixUs * num / den

// Splice rtpTS into the high bits of nowTicks, then snap to the wrap nearest
// nowTicks. Capture is within half a wrap period of receipt for any realistic
// latency or cross-machine clock skew (the wrap period is 2^32 ticks, ~13.25 h
// at 90 kHz / ~24.9 h at 48 kHz), so a candidate more than 2^31 ticks away is a
// wrap artifact. Using a half-period guard band (rather than "pull back
// whenever candidate > nowTicks") keeps a sender clock that runs slightly ahead
// mapping to a small negative latency instead of collapsing a full 2^32 wrap
// (~24.9 h) onto it.
candidate := (nowTicks &^ 0xFFFFFFFF) | int64(rtpTS)
switch {
case candidate-nowTicks > 1<<31:
candidate -= 1 << 32
case nowTicks-candidate > 1<<31:
candidate += 1 << 32
}

return candidate * den / num
}

// GlassToGlassLatency returns nowUnixUs minus the capture instant recovered from
// the RTP timestamp: the elapsed time from frame capture to this observation.
func GlassToGlassLatency(rtpTS uint32, clockRate uint32, nowUnixUs int64) time.Duration {
captureUs := CaptureTimeUsFromRTP(rtpTS, clockRate, nowUnixUs)

return time.Duration(nowUnixUs-captureUs) * time.Microsecond
}

// gcd returns the greatest common divisor of a and b (both assumed positive).
func gcd(a, b int64) int64 {
for b != 0 {
a, b = b, a%b
}

return a
}
127 changes: 127 additions & 0 deletions receiver/capture_latency_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,127 @@
// SPDX-FileCopyrightText: 2026 The Pion community <https://pion.ly>
// SPDX-License-Identifier: MIT

//go:build !js

package receiver

import (
"testing"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// encodeCaptureRTP reproduces the sender's captureTimestampInterceptor encoding:
// the outgoing RTP timestamp is captureUs*ClockRate/1e6 truncated to 32 bits.
// It uses the same reduced fraction as the interceptor so captureUs*num does not
// overflow int64 (captureUs*ClockRate would).
func encodeCaptureRTP(captureUs int64, clockRate uint32) uint32 {
rate := int64(clockRate)
g := gcd(rate, 1_000_000)
num := rate / g
den := int64(1_000_000) / g

return uint32(captureUs * num / den) //nolint:gosec // intentional 32-bit wrap
}

// quantUs is the maximum round-trip error introduced by the sender's and the
// recovery's integer divisions: one tick, i.e. 1e6/ClockRate microseconds.
func quantUs(clockRate uint32) int64 {
return 1_000_000/int64(clockRate) + 1
}

// TestCaptureTimeFromRTP_RoundTrip asserts that a capture time encoded with the
// sender's formula is recovered within one RTP tick, and that the resulting
// glass-to-glass latency matches the injected delay, across clock rates and
// delays.
func TestCaptureTimeFromRTP_RoundTrip(t *testing.T) {
// A recent, realistic capture instant (unix micros).
const baseCaptureUs = int64(1_751_000_000_000_000)

clockRates := []uint32{90000, 48000}
latencies := []time.Duration{
0,
5 * time.Millisecond,
50 * time.Millisecond,
250 * time.Millisecond,
time.Second,
}

for _, rate := range clockRates {
for _, lat := range latencies {
captureUs := baseCaptureUs
rtpTS := encodeCaptureRTP(captureUs, rate)

latUs := lat.Microseconds()
nowUs := captureUs + latUs

recovered := CaptureTimeUsFromRTP(rtpTS, rate, nowUs)
assert.InDeltaf(t, captureUs, recovered, float64(quantUs(rate)),
"rate=%d lat=%s: recovered capture time out of tolerance", rate, lat)

gotLatency := GlassToGlassLatency(rtpTS, rate, nowUs)
assert.InDeltaf(t, latUs, gotLatency.Microseconds(), float64(quantUs(rate)),
"rate=%d lat=%s: recovered latency out of tolerance", rate, lat)
}
}
}

// TestCaptureTimeFromRTP_WrapBoundary asserts recovery is correct when the
// capture instant's encoded ticks and nowUnixUs's ticks straddle a 2^32
// boundary, exercising the wrap-correction branch.
func TestCaptureTimeFromRTP_WrapBoundary(t *testing.T) {
const rate = uint32(90000)

// Find a capture time whose encoded 90 kHz ticks sit just below a 2^32
// multiple, so that a small added latency pushes nowUnixUs's ticks across
// the boundary. 2^32 ticks / 90000 ticks-per-second = wrap period; work in
// microseconds: wrapUs = 2^32 * 1e6 / 90000.
wrapUs := int64(1) << 32 * 1_000_000 / int64(rate)
// Capture 1 ms of ticks before the 5th wrap boundary.
captureUs := 5*wrapUs - time.Millisecond.Microseconds()

rtpTS := encodeCaptureRTP(captureUs, rate)
// 50 ms later — nowUnixUs is past the boundary while rtpTS is from before it.
nowUs := captureUs + 50*time.Millisecond.Microseconds()

recovered := CaptureTimeUsFromRTP(rtpTS, rate, nowUs)
assert.InDelta(t, captureUs, recovered, float64(quantUs(rate)),
"recovery must cross the 2^32 wrap boundary correctly")
}

// TestCaptureTimeFromRTP_ClockRateZeroFallsBackTo90k asserts a zero clock rate
// is treated as 90 kHz, matching the sender's fallback.
func TestCaptureTimeFromRTP_ClockRateZeroFallsBackTo90k(t *testing.T) {
const baseCaptureUs = int64(1_751_000_000_000_000)

rtpTS := encodeCaptureRTP(baseCaptureUs, 90000)
nowUs := baseCaptureUs + 30*time.Millisecond.Microseconds()

withZero := CaptureTimeUsFromRTP(rtpTS, 0, nowUs)
with90k := CaptureTimeUsFromRTP(rtpTS, 90000, nowUs)
assert.Equal(t, with90k, withZero, "clockRate 0 must behave like 90 kHz")
}

// TestReportGlassToGlassLatency asserts the receiver's read-loop helper recovers
// a recent capture instant from a live RTP timestamp and reports a latency close
// to the true elapsed time, at both the audio (48 kHz) and video (90 kHz) rates.
func TestReportGlassToGlassLatency(t *testing.T) {
r, err := NewReceiver()
require.NoError(t, err)

const elapsed = 40 * time.Millisecond
// Generous upper bound: recovery quantization plus scheduler/test slack.
const tol = 25 * time.Millisecond

for _, rate := range []uint32{48000, 90000} {
// A frame captured `elapsed` ago, stamped as the sender would.
captureUs := time.Now().UnixMicro() - elapsed.Microseconds()
rtpTS := encodeCaptureRTP(captureUs, rate)

got := r.reportGlassToGlassLatency("track-test", rtpTS, rate)
assert.InDeltaf(t, elapsed.Microseconds(), got.Microseconds(), float64(tol.Microseconds()),
"rate=%d: reported latency %s should track elapsed %s", rate, got, elapsed)
}
}
29 changes: 26 additions & 3 deletions receiver/receiver.go
Original file line number Diff line number Diff line change
Expand Up @@ -282,11 +282,29 @@ func (r *Receiver) createOutputFile(trackIdentifier string) {
r.log.Infof("Created output file: %s", cleanFilename)
}

// handleNonVP8Track handles non-VP8 tracks by simply reading packets.
// reportGlassToGlassLatency recovers the capture instant the sender encoded into
// rtpTS (see the sender's captureTimestampInterceptor) using the track's
// negotiated clock rate (90 kHz VP8, 48 kHz Opus), logs the resulting
// glass-to-glass latency, and returns it. Shared by the video and audio read
// loops so both surface the capture timestamp identically.
func (r *Receiver) reportGlassToGlassLatency(trackID string, rtpTS, clockRate uint32) time.Duration {
lat := GlassToGlassLatency(rtpTS, clockRate, time.Now().UnixMicro())
r.log.Debugf("glass-to-glass latency: track=%s rtpTS=%d clockRate=%d latency=%v",
trackID, rtpTS, clockRate, lat)

return lat
}

// handleNonVP8Track reads non-VP8 tracks (in practice, encoded Opus audio) and
// recovers each packet's capture timestamp from its RTP timestamp, mirroring the
// VP8 video path in processPackets. The sender stamps the capture instant into
// the RTP timestamp for audio (48 kHz) exactly as it does for video (90 kHz), so
// the same recovery applies here using the track's negotiated clock rate.
func (r *Receiver) handleNonVP8Track(
ctx context.Context, trackRemote *webrtc.TrackRemote,
rtpReceiver *webrtc.RTPReceiver, _ *trackInfo,
rtpReceiver *webrtc.RTPReceiver, trackInfo *trackInfo,
) {
clockRate := trackRemote.Codec().ClockRate
for {
select {
case <-ctx.Done():
Expand All @@ -296,7 +314,7 @@ func (r *Receiver) handleNonVP8Track(
continue
}

_, _, err := trackRemote.ReadRTP()
packet, _, err := trackRemote.ReadRTP()
if errors.Is(err, io.EOF) {
r.log.Infof("trackRemote.ReadRTP received EOF")

Expand All @@ -307,6 +325,8 @@ func (r *Receiver) handleNonVP8Track(

continue
}

r.reportGlassToGlassLatency(trackInfo.identifier, packet.Timestamp, clockRate)
}
}
}
Expand Down Expand Up @@ -378,6 +398,7 @@ func (r *Receiver) processPackets(ctx context.Context, trackRemote *webrtc.Track
trackInfo *trackInfo, frameAssembler *VP8FrameAssembler,
bytesReceivedChan chan int, stats *trackStats,
) {
clockRate := trackRemote.Codec().ClockRate
for {
select {
case <-ctx.Done():
Expand All @@ -402,6 +423,8 @@ func (r *Receiver) processPackets(ctx context.Context, trackRemote *webrtc.Track
bytesReceivedChan <- packet.MarshalSize()
stats.rtpPacketsReceived++

r.reportGlassToGlassLatency(trackInfo.identifier, packet.Timestamp, clockRate)

r.processVP8Packet(packet, trackInfo, frameAssembler, stats)
}
}
Expand Down
Loading
Loading