mirror of
https://github.com/livekit/livekit.git
synced 2026-10-06 03:27:46 +00:00
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 <agent@perfloop.ai> * Use require in the descriptor marshal benchmark test Signed-off-by: Perfloop Agent <agent@perfloop.ai> * 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 <agent@perfloop.ai>
This commit is contained in:
+24
-16
@@ -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 {
|
||||
|
||||
+11
-2
@@ -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)
|
||||
|
||||
|
||||
@@ -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)))
|
||||
}
|
||||
|
||||
@@ -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)))
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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"
|
||||
|
||||
+247
@@ -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)
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}()
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user