Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 25 additions & 8 deletions pkg/twcc/twcc.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ func (r *Recorder) Record(mediaSSRC uint32, sequenceNumber uint16, arrivalTime i
unwrappedSN := r.sequenceUnwrapper.Unwrap(sequenceNumber)
r.maybeCullOldPackets(unwrappedSN, arrivalTime)
if r.startSequenceNumber == nil || unwrappedSN < *r.startSequenceNumber {
r.startSequenceNumber = &unwrappedSN
r.setStartSequenceNumber(unwrappedSN)
}

// We are only interested in the first time a packet is received.
Expand All @@ -65,11 +65,20 @@ func (r *Recorder) Record(mediaSSRC uint32, sequenceNumber uint16, arrivalTime i

// Limit the range of sequence numbers to send feedback for.
if *r.startSequenceNumber < r.arrivalTimeMap.BeginSequenceNumber() {
sn := r.arrivalTimeMap.BeginSequenceNumber()
r.startSequenceNumber = &sn
r.setStartSequenceNumber(r.arrivalTimeMap.BeginSequenceNumber())
}
}

// setStartSequenceNumber stores sequenceNumber in a single reused allocation;
// storing the address of a stack value instead would heap-allocate it on every
// call.
func (r *Recorder) setStartSequenceNumber(sequenceNumber int64) {
if r.startSequenceNumber == nil {
r.startSequenceNumber = new(int64)
}
*r.startSequenceNumber = sequenceNumber
}

func (r *Recorder) maybeCullOldPackets(sequenceNumber int64, arrivalTime int64) {
if r.startSequenceNumber != nil && *r.startSequenceNumber >= r.arrivalTimeMap.EndSequenceNumber() &&
arrivalTime >= packetWindowMicroseconds {
Expand Down Expand Up @@ -147,7 +156,7 @@ func (r *Recorder) maybeBuildFeedbackPacket(beginSeqNumInclusive, endSeqNumExclu
// try again after skipping any missing packets.
// NOTE: It's fine that we already incremented fbPktCnt, as in essence
// we did actually "skip" a feedback (and this matches Chrome's behavior).
r.startSequenceNumber = &seq
r.setStartSequenceNumber(seq)

return nil
}
Expand All @@ -160,7 +169,7 @@ func (r *Recorder) maybeBuildFeedbackPacket(beginSeqNumInclusive, endSeqNumExclu
nextSequenceNumber = seq + 1
}

r.startSequenceNumber = &nextSequenceNumber
r.setStartSequenceNumber(nextSequenceNumber)

return fb
}
Expand All @@ -175,7 +184,7 @@ type feedback struct {
len int
lastChunk chunk
chunks []rtcp.PacketStatusChunk
deltas []*rtcp.RecvDelta
deltas []rtcp.RecvDelta
}

func newFeedback(senderSSRC, mediaSSRC uint32, count uint8) *feedback {
Expand All @@ -195,6 +204,8 @@ func (f *feedback) setBase(sequenceNumber uint16, timeUS int64) {
f.lastTimestampUS = f.refTimestamp64MS * 64e3
}

// getRTCP finalizes the feedback and returns it as an RTCP packet. It consumes
// the accumulated state and must be called at most once.
func (f *feedback) getRTCP() *rtcp.TransportLayerCC {
f.rtcp.PacketStatusCount = f.sequenceNumberCount
f.rtcp.ReferenceTime = uint32(f.refTimestamp64MS) //nolint:gosec // G115
Expand All @@ -203,7 +214,13 @@ func (f *feedback) getRTCP() *rtcp.TransportLayerCC {
f.chunks = append(f.chunks, f.lastChunk.encode())
}
f.rtcp.PacketChunks = append(f.rtcp.PacketChunks, f.chunks...)
f.rtcp.RecvDeltas = f.deltas
// The pointers alias f.deltas entries, which is cleared so that the packet
// becomes their sole owner.
f.rtcp.RecvDeltas = make([]*rtcp.RecvDelta, len(f.deltas))
for i := range f.deltas {
f.rtcp.RecvDeltas[i] = &f.deltas[i]
}
f.deltas = nil

// 4 bytes header + 16 bytes twcc header + 2 bytes for each chunk + length of deltas
padLen := 20 + len(f.rtcp.PacketChunks)*2 + f.len
Expand Down Expand Up @@ -257,7 +274,7 @@ func (f *feedback) addReceived(sequenceNumber uint16, timestampUS int64) bool {
f.chunks = append(f.chunks, f.lastChunk.encode())
}
f.lastChunk.add(recvDelta)
f.deltas = append(f.deltas, &rtcp.RecvDelta{
f.deltas = append(f.deltas, rtcp.RecvDelta{
Type: recvDelta,
Delta: deltaUSRounded,
})
Expand Down
53 changes: 53 additions & 0 deletions pkg/twcc/twcc_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -921,3 +921,56 @@ func TestPacketsHheld(t *testing.T) {
recorder.BuildFeedbackPacket()
assert.Zero(t, recorder.PacketsHeld())
}

func TestRecordDoesNotAllocateInSteadyState(t *testing.T) {
recorder := NewRecorder(5000)

arrivalTime := int64(0)
sequenceNumber := uint16(0)
record := func() {
arrivalTime += 1000
recorder.Record(5000, sequenceNumber, arrivalTime)
sequenceNumber++
}

// Warm up until the arrival-time map reaches its maximum size, after which
// in-order packets no longer resize it. In normal operation the map stays
// small because old entries are culled, but culling only runs once
// BuildFeedbackPacket has consumed the recorded range, and building
// feedback allocates. Growing the map to its cap instead gives Record a
// steady state without any builds in the measured loop.
for range maxNumberOfPackets + minCapacity {
record()
}

assert.Zero(t, testing.AllocsPerRun(1000, record))
}

func BenchmarkRecord(b *testing.B) {
recorder := NewRecorder(5000)

arrivalTime := int64(0)
sequenceNumber := uint16(0)
b.ReportAllocs()
for b.Loop() {
arrivalTime += 1000
recorder.Record(5000, sequenceNumber, arrivalTime)
sequenceNumber++
}
}

func BenchmarkRecordAndBuildFeedbackPacket(b *testing.B) {
recorder := NewRecorder(5000)

arrivalTime := int64(0)
sequenceNumber := uint16(0)
b.ReportAllocs()
for b.Loop() {
for range 100 {
arrivalTime += 1000
recorder.Record(5000, sequenceNumber, arrivalTime)
sequenceNumber++
}
recorder.BuildFeedbackPacket()
}
}
Loading