From 7ef0cd014e950373acd789557099305716eea196 Mon Sep 17 00:00:00 2001 From: Perfloop Agent Date: Wed, 30 Sep 2026 08:11:49 -0700 Subject: [PATCH] Reduce descriptor marshal allocations on Go 1.26.4/linux-amd64 (#4681) * Reduce descriptor marshal allocations on Go 1.26.4/linux-amd64 Signed-off-by: Perfloop Agent * Use require in the descriptor marshal benchmark test Signed-off-by: Perfloop Agent * sfu: marshal dependency descriptors into caller-owned scratch * sfu: reuse the selector's descriptor clone * sfu: share one inline header-extension size constant --------- Signed-off-by: Perfloop Agent --- pkg/sfu/downtrack.go | 40 +-- pkg/sfu/forwarder.go | 13 +- pkg/sfu/forwarder_test.go | 11 + pkg/sfu/pacer/base_test.go | 22 ++ pkg/sfu/pacer/pacer.go | 10 + .../dependencydescriptorextension.go | 46 +++- ...dencydescriptorextension_benchmark_test.go | 247 ++++++++++++++++++ .../dependencydescriptorwriter.go | 21 +- .../dependencydescriptor.go | 16 +- .../dependencydescriptor_test.go | 70 +++++ .../videolayerselector/videolayerselector.go | 37 ++- 11 files changed, 490 insertions(+), 43 deletions(-) create mode 100644 pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorextension_benchmark_test.go diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 1ff4ac7aa..e01c96748 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -1079,9 +1079,28 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) int32 { SSRC: d.ssrc, }) + pacerPacket := pacer.PacketFactory.Get().(*pacer.Packet) + *pacerPacket = pacer.Packet{ + Header: hdr, + HeaderPool: RTPHeaderFactory, + Payload: payload, + ProbeClusterId: ccutils.ProbeClusterId(d.probeClusterId.Load()), + AbsSendTimeExtID: uint8(d.absSendTimeExtID), + TransportWideExtID: uint8(d.transportWideExtID), + WriteStream: d.writeStream, + Pool: PacketFactory, + PoolEntity: poolEntity, + } + // add extensions - if d.dependencyDescriptorExtID != 0 && tp.ddBytes != nil { - hdr.SetExtension(uint8(d.dependencyDescriptorExtID), tp.ddBytes) + // + // the header refers to the descriptor until the pacer marshals it, so hold it in the pacer packet, which one send owns + ddBytes := tp.ddBytesSpill + if inline := tp.ddInline(); ddBytes == nil && len(inline) != 0 { + ddBytes = pacerPacket.HoldExtension(inline) + } + if d.dependencyDescriptorExtID != 0 && len(ddBytes) != 0 { + hdr.SetExtension(uint8(d.dependencyDescriptorExtID), ddBytes) } if d.playoutDelayExtID != 0 && d.playoutDelay != nil { if val := d.playoutDelay.GetDelayExtension(hdr.SequenceNumber); val != nil { @@ -1128,6 +1147,7 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) int32 { } d.addDummyExtensions(hdr) + // the sequencer copies the descriptor, which has to happen before Enqueue frees the scratch if d.sequencer != nil { d.sequencer.push( extPkt.Arrival, @@ -1138,13 +1158,14 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) int32 { int8(layer), payload[:len(codecBytes)], tp.incomingHeaderSize, - tp.ddBytes, + ddBytes, actBytes, trailerStripped, ) } headerSize := hdr.MarshalSize() + pacerPacket.HeaderSize = headerSize d.rtpStats.Update( extPkt.Arrival, tp.rtp.extSequenceNumber, @@ -1155,19 +1176,6 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) int32 { 0, extPkt.IsOutOfOrder, ) - pacerPacket := pacer.PacketFactory.Get().(*pacer.Packet) - *pacerPacket = pacer.Packet{ - Header: hdr, - HeaderPool: RTPHeaderFactory, - HeaderSize: headerSize, - Payload: payload, - ProbeClusterId: ccutils.ProbeClusterId(d.probeClusterId.Load()), - AbsSendTimeExtID: uint8(d.absSendTimeExtID), - TransportWideExtID: uint8(d.transportWideExtID), - WriteStream: d.writeStream, - Pool: PacketFactory, - PoolEntity: poolEntity, - } d.pacer.Enqueue(pacerPacket) if extPkt.IsKeyFrame { diff --git a/pkg/sfu/forwarder.go b/pkg/sfu/forwarder.go index 9c87b0653..fbdc28092 100644 --- a/pkg/sfu/forwarder.go +++ b/pkg/sfu/forwarder.go @@ -199,7 +199,9 @@ type TranslationParams struct { isResuming bool isSwitching bool rtp TranslationParamsRTP - ddBytes []byte + ddBytes [dd.MaxInlineExtensionSize]byte + ddBytesLen int + ddBytesSpill []byte incomingHeaderSize int codecBytes [codecmunger.MaxHeaderSize]byte codecBytesLen int @@ -213,6 +215,13 @@ func (tp *TranslationParams) codecHeader() []byte { return tp.codecBytes[:tp.codecBytesLen] } +// ddInline returns the dependency descriptor extension held inline. It points +// into tp, so it has to be copied before it is attached to anything that +// outlives the write. +func (tp *TranslationParams) ddInline() []byte { + return tp.ddBytes[:tp.ddBytesLen] +} + // ------------------------------------------------------------------- type refInfo struct { @@ -2190,7 +2199,7 @@ func (f *Forwarder) getTranslationParamsVideo(extPkt *buffer.ExtPacket, layer in } tp.isResuming = result.IsResuming tp.isSwitching = result.IsSwitching - tp.ddBytes = result.DependencyDescriptorExtension + tp.ddBytes, tp.ddBytesLen, tp.ddBytesSpill = result.DDBytes, result.DDBytesLen, result.DDBytesSpill tp.marker = result.RTPMarker tp.isEndOfLayerFrame = f.isEndOfLayerFrame(extPkt) diff --git a/pkg/sfu/forwarder_test.go b/pkg/sfu/forwarder_test.go index defec62db..931c79e7e 100644 --- a/pkg/sfu/forwarder_test.go +++ b/pkg/sfu/forwarder_test.go @@ -28,6 +28,7 @@ import ( "github.com/livekit/livekit-server/pkg/sfu/buffer" "github.com/livekit/livekit-server/pkg/sfu/codecmunger" + "github.com/livekit/livekit-server/pkg/sfu/pacer" dd "github.com/livekit/livekit-server/pkg/sfu/rtpextension/dependencydescriptor" "github.com/livekit/livekit-server/pkg/sfu/testutils" ) @@ -2261,3 +2262,13 @@ func codecHeaderOf(b []byte) [codecmunger.MaxHeaderSize]byte { copy(hdr[:], b) return hdr } + +// WriteRTP holds the dependency descriptor in pacer scratch sized to pion's one +// byte extension profile payload cap, which has to be the inline size +// TranslationParams carries, or descriptors would take the heap copy fallback +func TestPacerScratchFitsInlineDependencyDescriptor(t *testing.T) { + p := &pacer.Packet{} + + require.Len(t, p.HoldExtension(make([]byte, dd.MaxInlineExtensionSize)), dd.MaxInlineExtensionSize) + require.Nil(t, p.HoldExtension(make([]byte, dd.MaxInlineExtensionSize+1))) +} diff --git a/pkg/sfu/pacer/base_test.go b/pkg/sfu/pacer/base_test.go index 3842203a6..82599f049 100644 --- a/pkg/sfu/pacer/base_test.go +++ b/pkg/sfu/pacer/base_test.go @@ -15,6 +15,7 @@ package pacer import ( + "bytes" "sync" "testing" @@ -24,6 +25,7 @@ import ( "github.com/livekit/protocol/logger" "github.com/livekit/livekit-server/pkg/sfu/bwe" + dd "github.com/livekit/livekit-server/pkg/sfu/rtpextension/dependencydescriptor" "github.com/livekit/livekit-server/pkg/sfu/utils" ) @@ -85,3 +87,23 @@ func TestSendPacketHeaderExtensionsNoAlloc(t *testing.T) { require.Equal(t, 0.0, allocs, "allocations per SendPacket") } } + +func TestHoldExtension(t *testing.T) { + p := &Packet{} + + ext := []byte{1, 2, 3} + held := p.HoldExtension(ext) + require.Equal(t, ext, held) + require.Same(t, &p.extBuf[0], &held[0], "the extension has to be held in the packet's scratch") + + ext[0] = 9 + require.EqualValues(t, 1, held[0], "the held copy has to be independent of the caller's slice") + + // exactly the one byte extension profile payload cap + exact := bytes.Repeat([]byte{7}, dd.MaxInlineExtensionSize) + held = p.HoldExtension(exact) + require.Equal(t, exact, held) + + // one byte more does not fit + require.Nil(t, p.HoldExtension(bytes.Repeat([]byte{7}, dd.MaxInlineExtensionSize+1))) +} diff --git a/pkg/sfu/pacer/pacer.go b/pkg/sfu/pacer/pacer.go index 199a8bf17..6b89b6d1f 100644 --- a/pkg/sfu/pacer/pacer.go +++ b/pkg/sfu/pacer/pacer.go @@ -19,6 +19,7 @@ import ( "time" "github.com/livekit/livekit-server/pkg/sfu/ccutils" + dd "github.com/livekit/livekit-server/pkg/sfu/rtpextension/dependencydescriptor" "github.com/pion/rtp" "github.com/pion/webrtc/v4" ) @@ -59,6 +60,15 @@ type Packet struct { // the Packet is owned by one send until SendPacket returns absSendTimeBuf [3]byte twccBuf [2]byte + extBuf [dd.MaxInlineExtensionSize]byte +} + +// HoldExtension copies ext into scratch the header can point at until SendPacket returns, nil if it does not fit +func (p *Packet) HoldExtension(ext []byte) []byte { + if len(ext) > len(p.extBuf) { + return nil + } + return p.extBuf[:copy(p.extBuf[:], ext)] } type Pacer interface { diff --git a/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorextension.go b/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorextension.go index 293482396..28812d09a 100644 --- a/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorextension.go +++ b/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorextension.go @@ -15,6 +15,7 @@ package dependencydescriptor import ( + "errors" "fmt" "math" "strconv" @@ -23,6 +24,10 @@ import ( // DependencyDescriptorExtension is a extension payload format in // https://aomediacodec.github.io/av1-rtp-spec/#dependency-descriptor-rtp-header-extension +// ErrBufferTooSmall is returned by MarshalTo when the caller's storage cannot +// hold the marshaled extension. +var ErrBufferTooSmall = errors.New("dependency descriptor buffer too small") + func formatBitmask(b *uint32) string { if b == nil { return "-" @@ -42,18 +47,48 @@ func (d *DependencyDescriptorExtension) Marshal() ([]byte, error) { } func (d *DependencyDescriptorExtension) MarshalWithActiveChains(activeChains uint32) ([]byte, error) { - writer, err := NewDependencyDescriptorWriter(nil, d.Structure, activeChains, d.Descriptor) + writer, size, err := d.newWriter(activeChains) if err != nil { return nil, err } - buf := make([]byte, int(math.Ceil(float64(writer.ValueSizeBits())/8))) - writer.ResetBuf(buf) - if err = writer.Write(); err != nil { + + buf := make([]byte, size) + writer = writer.withBuf(buf) + if err := writer.Write(); err != nil { return nil, err } return buf, nil } +// MarshalTo marshals into caller owned storage and returns the number of bytes +// written, or ErrBufferTooSmall, leaving buf untouched, if the extension does not fit. +func (d *DependencyDescriptorExtension) MarshalTo(buf []byte) (int, error) { + writer, size, err := d.newWriter(^uint32(0)) + if err != nil { + return 0, err + } + if len(buf) < size { + // bare sentinel, key frames take this path and formatting an error there would allocate + return 0, ErrBufferTooSmall + } + + writer = writer.withBuf(buf[:size]) + if err := writer.Write(); err != nil { + return 0, err + } + return size, nil +} + +// newWriter returns a writer for the extension, by value so that it stays on the +// caller's stack, along with the number of bytes the marshaled extension needs. +func (d *DependencyDescriptorExtension) newWriter(activeChains uint32) (DependencyDescriptorWriter, int, error) { + writer := newDependencyDescriptorWriter(nil, d.Structure, activeChains, d.Descriptor) + if err := writer.findBestTemplate(); err != nil { + return writer, 0, err + } + return writer, int(math.Ceil(float64(writer.ValueSizeBits()) / 8)), nil +} + func (d *DependencyDescriptorExtension) Unmarshal(buf []byte) (int, error) { reader := NewDependencyDescriptorReader(buf, d.Structure, d.Descriptor) return reader.Parse() @@ -67,6 +102,9 @@ const ( MaxDecodeTargets = 32 MaxTemplates = 64 + // MaxInlineExtensionSize is the inline storage the forwarding path and the pacer reserve for a marshaled extension, pion's one byte profile cap + MaxInlineExtensionSize = 16 + AllChainsAreActive = uint32(0) ExtensionURI = "https://aomediacodec.github.io/av1-rtp-spec/#dependency-descriptor-rtp-header-extension" diff --git a/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorextension_benchmark_test.go b/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorextension_benchmark_test.go new file mode 100644 index 000000000..6e9509c74 --- /dev/null +++ b/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorextension_benchmark_test.go @@ -0,0 +1,247 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package dependencydescriptor + +import ( + "encoding/hex" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/livekit/livekit-server/pkg/sfu/utils" +) + +// dependencyDescriptorMarshalFixture is a dependency descriptor captured from +// traffic. The fixture helper decodes its attached structure once, then returns +// a regular per-packet descriptor that references that structure. +const dependencyDescriptorMarshalFixture = "c1017280081485214eafffaaaa863cf0430c10c302afc0aaa0063c00430010c002a000a80006000040001d954926e082b04a0941b820ac1282503157f974000ca864330e222222eca8655304224230eca877530077004200ef008601df010d" + +func newDependencyDescriptorMarshalFixture(tb testing.TB) *DependencyDescriptorExtension { + tb.Helper() + + buf, err := hex.DecodeString(dependencyDescriptorMarshalFixture) + require.NoError(tb, err) + + descriptor := DependencyDescriptor{} + parser := DependencyDescriptorExtension{Descriptor: &descriptor} + _, err = parser.Unmarshal(buf) + require.NoError(tb, err) + require.NotNil(tb, descriptor.AttachedStructure, "fixture did not contain a dependency structure") + + structure := descriptor.AttachedStructure + descriptor.AttachedStructure = nil + descriptor.ActiveDecodeTargetsBitmask = nil + + return &DependencyDescriptorExtension{ + Descriptor: &descriptor, + Structure: structure, + } +} + +func checkDependencyDescriptorMarshal(tb testing.TB, buf []byte, structure *FrameDependencyStructure, frameNumber uint16) { + tb.Helper() + + require.NotEmpty(tb, buf, "marshal returned an empty dependency descriptor") + + decoded := DependencyDescriptor{} + parser := DependencyDescriptorExtension{ + Descriptor: &decoded, + Structure: structure, + } + _, err := parser.Unmarshal(buf) + require.NoError(tb, err) + require.Equal(tb, frameNumber, decoded.FrameNumber) +} + +func TestDependencyDescriptorMarshalRoundTrip(t *testing.T) { + extension := newDependencyDescriptorMarshalFixture(t) + + first, err := extension.Marshal() + require.NoError(t, err) + firstCopy := append([]byte(nil), first...) + checkDependencyDescriptorMarshal(t, first, extension.Structure, extension.Descriptor.FrameNumber) + + extension.Descriptor.FrameNumber++ + second, err := extension.Marshal() + require.NoError(t, err) + checkDependencyDescriptorMarshal(t, second, extension.Structure, extension.Descriptor.FrameNumber) + + require.Equal(t, firstCopy, first, "a later marshal modified a previously returned buffer") + require.NotEqual(t, first, second, "different frame numbers produced identical descriptors") +} + +func TestDependencyDescriptorMarshalTo(t *testing.T) { + extension := newDependencyDescriptorMarshalFixture(t) + + // every template of the fixture structure, i. e. the per-packet descriptors + // the forwarding path marshals, has to match Marshal byte for byte and fit + // the inline size the forwarding path reserves + for i, template := range extension.Structure.Templates { + extension.Descriptor.FrameDependencies = template + extension.Descriptor.FrameNumber = uint16(i) + + want, err := extension.Marshal() + require.NoError(t, err) + require.LessOrEqual(t, len(want), MaxInlineExtensionSize) + + var scratch [MaxInlineExtensionSize]byte + n, err := extension.MarshalTo(scratch[:]) + require.NoError(t, err) + require.Equal(t, want, scratch[:n]) + checkDependencyDescriptorMarshal(t, scratch[:n], extension.Structure, uint16(i)) + } +} + +func TestDependencyDescriptorMarshalToBufferTooSmall(t *testing.T) { + extension := newDependencyDescriptorMarshalFixture(t) + + want, err := extension.Marshal() + require.NoError(t, err) + + _, err = extension.MarshalTo(make([]byte, len(want)-1)) + require.ErrorIs(t, err, ErrBufferTooSmall) + + n, err := extension.MarshalTo(make([]byte, len(want))) + require.NoError(t, err) + require.Equal(t, len(want), n) + + // a descriptor carrying the full dependency structure does not fit inline + extension.Descriptor.AttachedStructure = extension.Structure + _, err = extension.MarshalTo(make([]byte, MaxInlineExtensionSize)) + require.ErrorIs(t, err, ErrBufferTooSmall) + + keyFrame, err := extension.Marshal() + require.NoError(t, err) + require.Greater(t, len(keyFrame), MaxInlineExtensionSize) +} + +// TestDependencyDescriptorMarshalToLargerDescriptors covers the 8 to 16 byte +// band, where the active decode targets bitmask and custom frame dependencies +// land and where the sequencer's own inline copy already spills. +func TestDependencyDescriptorMarshalToLargerDescriptors(t *testing.T) { + extension := newDependencyDescriptorMarshalFixture(t) + bitmask := uint32(0x1ff) + + // the widest per-packet descriptor of the fixture, with the bitmask attached + widest := 0 + for _, template := range extension.Structure.Templates { + extension.Descriptor.FrameDependencies = template + extension.Descriptor.ActiveDecodeTargetsBitmask = &bitmask + + want, err := extension.Marshal() + require.NoError(t, err) + widest = max(widest, len(want)) + + var scratch [MaxInlineExtensionSize]byte + n, err := extension.MarshalTo(scratch[:]) + require.NoError(t, err) + require.Equal(t, want, scratch[:n]) + } + require.GreaterOrEqual(t, widest, 8) + + // custom frame diffs and chain diffs, i. e. a frame that matches no template + custom := extension.Structure.Templates[0].Clone() + custom.FrameDiffs = []int{1, 300, 4000} + for i := range custom.ChainDiffs { + custom.ChainDiffs[i] = 200 + i + } + extension.Descriptor.FrameDependencies = custom + extension.Descriptor.ActiveDecodeTargetsBitmask = nil + + want, err := extension.Marshal() + require.NoError(t, err) + require.GreaterOrEqual(t, len(want), 8) + require.LessOrEqual(t, len(want), MaxInlineExtensionSize) + + var scratch [MaxInlineExtensionSize]byte + n, err := extension.MarshalTo(scratch[:]) + require.NoError(t, err) + require.Equal(t, want, scratch[:n]) + checkDependencyDescriptorMarshal(t, scratch[:n], extension.Structure, extension.Descriptor.FrameNumber) +} + +func TestDependencyDescriptorMarshalAllocs(t *testing.T) { + if utils.RaceEnabled { + // the race detector perturbs allocation counts, see #4874 + t.Skip("allocation count is not meaningful under the race detector") + } + + extension := newDependencyDescriptorMarshalFixture(t) + + var ( + scratch [MaxInlineExtensionSize]byte + n int + err error + ) + allocs := testing.AllocsPerRun(100, func() { + n, err = extension.MarshalTo(scratch[:]) + }) + require.NoError(t, err) + require.NotZero(t, n) + require.Zero(t, allocs, "a per-packet descriptor has to marshal without allocating") + + // a key frame descriptor carries the full dependency structure, does not fit + // the inline scratch and spills to exactly one allocation, the Marshal slice + extension.Descriptor.AttachedStructure = extension.Structure + + var buf []byte + allocs = testing.AllocsPerRun(100, func() { + if n, err = extension.MarshalTo(scratch[:]); err != nil { + buf, err = extension.Marshal() + } + }) + require.NoError(t, err) + require.Greater(t, len(buf), MaxInlineExtensionSize) + require.Equal(t, float64(1), allocs, "the spill path has to allocate only the Marshal slice") +} + +func BenchmarkDependencyDescriptorMarshal(b *testing.B) { + extension := newDependencyDescriptorMarshalFixture(b) + frameNumber := extension.Descriptor.FrameNumber + + b.ReportAllocs() + var buf []byte + for b.Loop() { + extension.Descriptor.FrameNumber = frameNumber + frameNumber++ + + var err error + buf, err = extension.Marshal() + require.NoError(b, err) + } + + checkDependencyDescriptorMarshal(b, buf, extension.Structure, frameNumber-1) +} + +func BenchmarkDependencyDescriptorMarshalTo(b *testing.B) { + extension := newDependencyDescriptorMarshalFixture(b) + frameNumber := extension.Descriptor.FrameNumber + + b.ReportAllocs() + var ( + scratch [MaxInlineExtensionSize]byte + n int + ) + for b.Loop() { + extension.Descriptor.FrameNumber = frameNumber + frameNumber++ + + var err error + n, err = extension.MarshalTo(scratch[:]) + require.NoError(b, err) + } + + checkDependencyDescriptorMarshal(b, scratch[:n], extension.Structure, frameNumber-1) +} diff --git a/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorwriter.go b/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorwriter.go index 3c6ba720d..fac78f036 100644 --- a/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorwriter.go +++ b/pkg/sfu/rtpextension/dependencydescriptor/dependencydescriptorwriter.go @@ -33,23 +33,28 @@ type DependencyDescriptorWriter struct { descriptor *DependencyDescriptor structure *FrameDependencyStructure activeChains uint32 - writer *BitStreamWriter + writer BitStreamWriter bestTemplate TemplateMatch } -func NewDependencyDescriptorWriter(buf []byte, structure *FrameDependencyStructure, activeChains uint32, descriptor *DependencyDescriptor) (*DependencyDescriptorWriter, error) { - writer := NewBitStreamWriter(buf) - w := &DependencyDescriptorWriter{ +func newDependencyDescriptorWriter(buf []byte, structure *FrameDependencyStructure, activeChains uint32, descriptor *DependencyDescriptor) DependencyDescriptorWriter { + return DependencyDescriptorWriter{ descriptor: descriptor, structure: structure, activeChains: activeChains, - writer: writer, + writer: BitStreamWriter{buf: buf}, } - return w, w.findBestTemplate() } -func (w *DependencyDescriptorWriter) ResetBuf(buf []byte) { - w.writer = NewBitStreamWriter(buf) +func NewDependencyDescriptorWriter(buf []byte, structure *FrameDependencyStructure, activeChains uint32, descriptor *DependencyDescriptor) (*DependencyDescriptorWriter, error) { + w := newDependencyDescriptorWriter(buf, structure, activeChains, descriptor) + return &w, w.findBestTemplate() +} + +// withBuf returns a copy of the writer that writes into buf, value receiver so a caller owned buffer does not escape +func (w DependencyDescriptorWriter) withBuf(buf []byte) DependencyDescriptorWriter { + w.writer = BitStreamWriter{buf: buf} + return w } func (w *DependencyDescriptorWriter) Write() error { diff --git a/pkg/sfu/videolayerselector/dependencydescriptor.go b/pkg/sfu/videolayerselector/dependencydescriptor.go index 3e2e090a2..2909e1da5 100644 --- a/pkg/sfu/videolayerselector/dependencydescriptor.go +++ b/pkg/sfu/videolayerselector/dependencydescriptor.go @@ -46,6 +46,10 @@ type DependencyDescriptor struct { decodeTargets []*DecodeTarget fnWrapper FrameNumberWrapper + // per packet copy of the descriptor when the frame number or the active + // decode targets are rewritten, owned here so it does not allocate + ddClone dede.DependencyDescriptor + restartGeneration int } @@ -322,8 +326,8 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r unWrapFn := uint16(d.fnWrapper.UpdateAndGet(extFrameNum, ddwdt.StructureUpdated)) var ddClone *dede.DependencyDescriptor if unWrapFn != dd.FrameNumber { - clone := *dd - ddClone = &clone + d.ddClone = *dd + ddClone = &d.ddClone ddClone.FrameNumber = unWrapFn ddExtension.Descriptor = ddClone } @@ -333,8 +337,8 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r if ddClone == nil { // clone and override activebitmask // DD-TODO: if the packet that contains the bitmask is acknowledged by RR, then we don't need it until it changed. - clone := *dd - ddClone = &clone + d.ddClone = *dd + ddClone = &d.ddClone ddExtension.Descriptor = ddClone } ddClone.ActiveDecodeTargetsBitmask = d.activeDecodeTargetsBitmask @@ -355,11 +359,9 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r "stack", string(debug.Stack())) } }() - bytes, err := ddExtension.Marshal() - if err != nil { + if err := result.marshalDependencyDescriptorExtension(ddExtension); err != nil { d.logger.Warnw("error marshalling dependency descriptor extension", err) } else { - result.DependencyDescriptorExtension = bytes ddMarshaled = true } }() diff --git a/pkg/sfu/videolayerselector/dependencydescriptor_test.go b/pkg/sfu/videolayerselector/dependencydescriptor_test.go index 925c017dd..b5bfc4d36 100644 --- a/pkg/sfu/videolayerselector/dependencydescriptor_test.go +++ b/pkg/sfu/videolayerselector/dependencydescriptor_test.go @@ -15,6 +15,7 @@ package videolayerselector import ( + "encoding/hex" "slices" "testing" @@ -23,6 +24,7 @@ import ( "github.com/livekit/livekit-server/pkg/sfu/buffer" dd "github.com/livekit/livekit-server/pkg/sfu/rtpextension/dependencydescriptor" + "github.com/livekit/livekit-server/pkg/sfu/utils" "github.com/livekit/protocol/logger" ) @@ -409,3 +411,71 @@ func createDDFrames(maxLayer buffer.VideoLayer, startFrameNumber uint16) []*buff return frames } + +// same capture as the dependency descriptor package's marshal fixture, an L3T3 +// key frame descriptor with the full dependency structure attached +const dependencyDescriptorFixture = "c1017280081485214eafffaaaa863cf0430c10c302afc0aaa0063c00430010c002a000a80006000040001d954926e082b04a0941b820ac1282503157f974000ca864330e222222eca8655304224230eca877530077004200ef008601df010d" + +func TestVideoLayerSelectorResultDependencyDescriptor(t *testing.T) { + raw, err := hex.DecodeString(dependencyDescriptorFixture) + require.NoError(t, err) + + descriptor := dd.DependencyDescriptor{} + _, err = (&dd.DependencyDescriptorExtension{Descriptor: &descriptor}).Unmarshal(raw) + require.NoError(t, err) + require.NotNil(t, descriptor.AttachedStructure) + structure := descriptor.AttachedStructure + + // a key frame descriptor carries the full structure and spills to the heap + keyFrame := VideoLayerSelectorResult{} + require.NoError(t, keyFrame.marshalDependencyDescriptorExtension(&dd.DependencyDescriptorExtension{ + Descriptor: &descriptor, + Structure: structure, + })) + require.Zero(t, keyFrame.DDBytesLen) + require.Greater(t, len(keyFrame.DDBytesSpill), dd.MaxInlineExtensionSize) + + // a per-packet descriptor is held inline, byte for byte what Marshal returns + perPacket := descriptor + perPacket.AttachedStructure = nil + perPacket.ActiveDecodeTargetsBitmask = nil + ddExtension := &dd.DependencyDescriptorExtension{Descriptor: &perPacket, Structure: structure} + + result := VideoLayerSelectorResult{} + require.NoError(t, result.marshalDependencyDescriptorExtension(ddExtension)) + require.Nil(t, result.DDBytesSpill) + require.NotZero(t, result.DDBytesLen) + require.LessOrEqual(t, result.DDBytesLen, dd.MaxInlineExtensionSize) + + want, err := ddExtension.Marshal() + require.NoError(t, err) + require.Equal(t, want, result.DDBytes[:result.DDBytesLen]) +} + +func TestDependencyDescriptorSelectAllocations(t *testing.T) { + for _, rewriteFrameNumber := range []bool{false, true} { + selector := NewDependencyDescriptor(logger.GetLogger()) + selector.SetTarget(buffer.VideoLayer{Spatial: 0, Temporal: 0}) + selector.SetRequestSpatial(0) + frames := createDDFrames(buffer.VideoLayer{Spatial: 0, Temporal: 0}, 100) + require.True(t, selector.Select(frames[0], 0).IsSelected) + require.NotNil(t, selector.activeDecodeTargetsBitmask) + if rewriteFrameNumber { + selector.fnWrapper.offset = 6000 + } + packet := frames[1] + packet.DependencyDescriptor.ExtFrameNum-- + packet.DependencyDescriptor.Descriptor.FrameNumber-- + + selected := true + allocs := testing.AllocsPerRun(1000, func() { + packet.DependencyDescriptor.ExtFrameNum++ + packet.DependencyDescriptor.Descriptor.FrameNumber++ + selected = selected && selector.Select(packet, 0).IsSelected + }) + require.True(t, selected) + if !utils.RaceEnabled { + require.Equal(t, 0.0, allocs, "allocations per selected packet, rewriteFrameNumber=%v", rewriteFrameNumber) + } + } +} diff --git a/pkg/sfu/videolayerselector/videolayerselector.go b/pkg/sfu/videolayerselector/videolayerselector.go index 72a6e2c81..795a058e4 100644 --- a/pkg/sfu/videolayerselector/videolayerselector.go +++ b/pkg/sfu/videolayerselector/videolayerselector.go @@ -15,18 +15,43 @@ package videolayerselector import ( + "errors" + "github.com/livekit/livekit-server/pkg/sfu/buffer" + dd "github.com/livekit/livekit-server/pkg/sfu/rtpextension/dependencydescriptor" "github.com/livekit/livekit-server/pkg/sfu/videolayerselector/temporallayerselector" "github.com/livekit/protocol/logger" ) type VideoLayerSelectorResult struct { - IsSelected bool - IsRelevant bool - IsSwitching bool - IsResuming bool - RTPMarker bool - DependencyDescriptorExtension []byte + IsSelected bool + IsRelevant bool + IsSwitching bool + IsResuming bool + RTPMarker bool + + // marshaled dependency descriptor extension, held inline when it fits and + // spilled to the heap when it does not, the way the sequencer keeps its copy + DDBytes [dd.MaxInlineExtensionSize]byte + DDBytesLen int + DDBytesSpill []byte +} + +// marshalDependencyDescriptorExtension marshals ddExtension into the result, +// inline if it fits and on the heap otherwise. +func (v *VideoLayerSelectorResult) marshalDependencyDescriptorExtension(ddExtension *dd.DependencyDescriptorExtension) error { + ddBytesLen, err := ddExtension.MarshalTo(v.DDBytes[:]) + if err == nil { + v.DDBytesLen = ddBytesLen + return nil + } + if !errors.Is(err, dd.ErrBufferTooSmall) { + return err + } + + // descriptors carrying the full dependency structure do not fit inline + v.DDBytesSpill, err = ddExtension.Marshal() + return err } type VideoLayerSelector interface {