From 34e99f02a9b649ce8bf8705170994b16906170fc Mon Sep 17 00:00:00 2001 From: Kevin Caffrey Date: Tue, 14 Jul 2026 15:45:53 -0400 Subject: [PATCH] Reduce TWCC sender allocations Record stored the unwrapped sequence number by taking its address, which heap-allocates the local on every call, whether or not the branch is taken. Updates now go through a setter that writes into a single reused allocation. addReceived allocated one rtcp.RecvDelta per received packet. Store values and build the pointer slice once per feedback in getRTCP. BenchmarkRecord: 21.1ns/op 1 alloc -> 6.9ns/op 0 allocs. BenchmarkRecordAndBuildFeedbackPacket (100 packets per feedback): 4546ns/op 322 allocs -> 2270ns/op 21 allocs. --- pkg/twcc/twcc.go | 33 ++++++++++++++++++++------- pkg/twcc/twcc_test.go | 53 +++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 78 insertions(+), 8 deletions(-) diff --git a/pkg/twcc/twcc.go b/pkg/twcc/twcc.go index 2648e852..1427864b 100644 --- a/pkg/twcc/twcc.go +++ b/pkg/twcc/twcc.go @@ -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. @@ -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 { @@ -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 } @@ -160,7 +169,7 @@ func (r *Recorder) maybeBuildFeedbackPacket(beginSeqNumInclusive, endSeqNumExclu nextSequenceNumber = seq + 1 } - r.startSequenceNumber = &nextSequenceNumber + r.setStartSequenceNumber(nextSequenceNumber) return fb } @@ -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 { @@ -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 @@ -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 @@ -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, }) diff --git a/pkg/twcc/twcc_test.go b/pkg/twcc/twcc_test.go index ece6a148..4a77e537 100644 --- a/pkg/twcc/twcc_test.go +++ b/pkg/twcc/twcc_test.go @@ -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() + } +}