mirror of
https://github.com/livekit/livekit.git
synced 2026-07-31 09:29:55 +00:00
Make sure dd selector uses correct keyframe to select packets (#2218)
* Make sure dd selector uses correct keyframe to select packets * Fix test case * remove unsed field
This commit is contained in:
@@ -614,9 +614,6 @@ func (b *Buffer) getExtPacket(rtpPacket *rtp.Packet, arrivalTime time.Time, flow
|
||||
if b.ddParser != nil {
|
||||
ddVal, videoLayer, err := b.ddParser.Parse(ep.Packet)
|
||||
if err != nil {
|
||||
if !errors.Is(err, ErrFrameEarlierThanKeyFrame) && !errors.Is(err, dd.ErrDDReaderNoStructure) {
|
||||
b.logger.Warnw("could not parse dependency descriptor", err)
|
||||
}
|
||||
return nil
|
||||
} else if ddVal != nil {
|
||||
ep.DependencyDescriptor = ddVal
|
||||
|
||||
@@ -27,7 +27,8 @@ import (
|
||||
)
|
||||
|
||||
var (
|
||||
ErrFrameEarlierThanKeyFrame = fmt.Errorf("frame is earlier than current keyframe")
|
||||
ErrFrameEarlierThanKeyFrame = fmt.Errorf("frame is earlier than current keyframe")
|
||||
ErrDDStructureAttachedToNonFirstPacket = fmt.Errorf("dependency descriptor structure is attached to non-first packet of a frame")
|
||||
)
|
||||
|
||||
type DependencyDescriptorParser struct {
|
||||
@@ -39,7 +40,6 @@ type DependencyDescriptorParser struct {
|
||||
|
||||
seqWrapAround *utils.WrapAround[uint16, uint64]
|
||||
frameWrapAround *utils.WrapAround[uint16, uint64]
|
||||
structureExtSeq uint64
|
||||
structureExtFrameNum uint64
|
||||
activeDecodeTargetsExtSeq uint64
|
||||
activeDecodeTargetsMask uint32
|
||||
@@ -66,12 +66,15 @@ type ExtDependencyDescriptor struct {
|
||||
ActiveDecodeTargetsUpdated bool
|
||||
Integrity bool
|
||||
ExtFrameNum uint64
|
||||
// the frame number of the keyframe which the current frame depends on
|
||||
ExtKeyFrameNum uint64
|
||||
}
|
||||
|
||||
func (r *DependencyDescriptorParser) Parse(pkt *rtp.Packet) (*ExtDependencyDescriptor, VideoLayer, error) {
|
||||
var videoLayer VideoLayer
|
||||
ddBuf := pkt.GetExtension(r.ddExtID)
|
||||
if ddBuf == nil {
|
||||
r.logger.Warnw("dependency descriptor extension is not present", nil, "seq", pkt.SequenceNumber)
|
||||
return nil, videoLayer, nil
|
||||
}
|
||||
|
||||
@@ -82,7 +85,9 @@ func (r *DependencyDescriptorParser) Parse(pkt *rtp.Packet) (*ExtDependencyDescr
|
||||
}
|
||||
_, err := ext.Unmarshal(ddBuf)
|
||||
if err != nil {
|
||||
// r.logger.Debugw("failed to parse generic dependency descriptor", "err", err, "payload", pkt.PayloadType, "ddbufLen", len(ddBuf))
|
||||
if err != dd.ErrDDReaderNoStructure {
|
||||
r.logger.Warnw("failed to parse generic dependency descriptor", err, "payload", pkt.PayloadType, "ddbufLen", len(ddBuf))
|
||||
}
|
||||
return nil, videoLayer, err
|
||||
}
|
||||
|
||||
@@ -108,17 +113,21 @@ func (r *DependencyDescriptorParser) Parse(pkt *rtp.Packet) (*ExtDependencyDescr
|
||||
}
|
||||
|
||||
if ddVal.AttachedStructure != nil {
|
||||
r.logger.Debugw("parsed dependency descriptor", "extSeq", extSeq, "extFN", extFN, "structureID", ddVal.AttachedStructure.StructureId, "descriptor", ddVal.String())
|
||||
if extSeq > r.structureExtSeq {
|
||||
r.structure = ddVal.AttachedStructure
|
||||
r.decodeTargets = ProcessFrameDependencyStructure(ddVal.AttachedStructure)
|
||||
r.structureExtSeq = extSeq
|
||||
r.structureExtFrameNum = extFN
|
||||
extDD.StructureUpdated = true
|
||||
extDD.ActiveDecodeTargetsUpdated = true
|
||||
// The dependency descriptor reader will always set ActiveDecodeTargetsBitmask for TemplateDependencyStructure is present,
|
||||
// so don't need to notify max layer change here.
|
||||
if !ddVal.FirstPacketInFrame {
|
||||
r.logger.Warnw("attached structure is not the first packet in frame", nil, "extSeq", extSeq, "extFN", extFN)
|
||||
return nil, videoLayer, ErrDDStructureAttachedToNonFirstPacket
|
||||
}
|
||||
|
||||
if r.structure == nil || ddVal.AttachedStructure.StructureId != r.structure.StructureId {
|
||||
r.logger.Infow("structure updated", "structureID", ddVal.AttachedStructure.StructureId, "extSeq", extSeq, "extFN", extFN, "descriptor", ddVal.String())
|
||||
}
|
||||
r.structure = ddVal.AttachedStructure
|
||||
r.decodeTargets = ProcessFrameDependencyStructure(ddVal.AttachedStructure)
|
||||
r.structureExtFrameNum = extFN
|
||||
extDD.StructureUpdated = true
|
||||
extDD.ActiveDecodeTargetsUpdated = true
|
||||
// The dependency descriptor reader will always set ActiveDecodeTargetsBitmask for TemplateDependencyStructure is present,
|
||||
// so don't need to notify max layer change here.
|
||||
}
|
||||
|
||||
if mask := ddVal.ActiveDecodeTargetsBitmask; mask != nil && extSeq > r.activeDecodeTargetsExtSeq {
|
||||
@@ -143,6 +152,7 @@ func (r *DependencyDescriptorParser) Parse(pkt *rtp.Packet) (*ExtDependencyDescr
|
||||
}
|
||||
|
||||
extDD.DecodeTargets = r.decodeTargets
|
||||
extDD.ExtKeyFrameNum = r.structureExtFrameNum
|
||||
|
||||
return extDD, videoLayer, nil
|
||||
}
|
||||
|
||||
@@ -193,7 +193,7 @@ func (s *StreamTrackerDependencyDescriptor) Observe(temporalLayer int32, pktSize
|
||||
|
||||
for _, dt := range ddVal.DecodeTargets {
|
||||
if len(dtis) <= dt.Target {
|
||||
s.params.Logger.Errorw("len(dtis) less than target", nil, "target", dt.Target, "dtls", dtis)
|
||||
s.params.Logger.Errorw("len(dtis) less than target", nil, "target", dt.Target, "dtis", dtis)
|
||||
continue
|
||||
}
|
||||
// we are not dropping discardable frames now, so only ingore not present frames
|
||||
|
||||
@@ -16,6 +16,7 @@ package videolayerselector
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"runtime/debug"
|
||||
"sync"
|
||||
|
||||
"github.com/livekit/livekit-server/pkg/sfu/buffer"
|
||||
@@ -31,6 +32,8 @@ type DependencyDescriptor struct {
|
||||
previousActiveDecodeTargetsBitmask *uint32
|
||||
activeDecodeTargetsBitmask *uint32
|
||||
structure *dede.FrameDependencyStructure
|
||||
extKeyFrameNum uint64
|
||||
keyFrameValid bool
|
||||
|
||||
chains []*FrameChain
|
||||
|
||||
@@ -79,38 +82,76 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
|
||||
Temporal: int32(fd.TemporalId),
|
||||
}
|
||||
|
||||
if !d.keyFrameValid && dd.AttachedStructure == nil {
|
||||
return
|
||||
}
|
||||
|
||||
// early return if this frame is already forwarded or dropped
|
||||
sd, err := d.decisions.GetDecision(extFrameNum)
|
||||
if err != nil {
|
||||
// do not mark as dropped as only error is an old frame
|
||||
d.logger.Debugw(fmt.Sprintf("drop packet on decision error, incoming %v, fn: %d/%d, sn: %d",
|
||||
incomingLayer,
|
||||
dd.FrameNumber,
|
||||
extFrameNum,
|
||||
extPkt.Packet.SequenceNumber,
|
||||
), "err", err)
|
||||
// d.logger.Debugw(fmt.Sprintf("drop packet on decision error, incoming %v, fn: %d/%d, sn: %d",
|
||||
// incomingLayer,
|
||||
// dd.FrameNumber,
|
||||
// extFrameNum,
|
||||
// extPkt.Packet.SequenceNumber,
|
||||
// ), "err", err)
|
||||
return
|
||||
}
|
||||
switch sd {
|
||||
case selectorDecisionDropped:
|
||||
// a packet of an alreadty dropped frame, maintain decision
|
||||
d.logger.Debugw(fmt.Sprintf("drop packet already dropped, incoming %v, fn: %d/%d, sn: %d",
|
||||
incomingLayer,
|
||||
dd.FrameNumber,
|
||||
extFrameNum,
|
||||
extPkt.Packet.SequenceNumber,
|
||||
))
|
||||
// d.logger.Debugw(fmt.Sprintf("drop packet already dropped, incoming %v, fn: %d/%d, sn: %d",
|
||||
// incomingLayer,
|
||||
// dd.FrameNumber,
|
||||
// extFrameNum,
|
||||
// extPkt.Packet.SequenceNumber,
|
||||
// ))
|
||||
return
|
||||
}
|
||||
|
||||
if ddwdt.StructureUpdated {
|
||||
d.updateDependencyStructure(dd.AttachedStructure, ddwdt.DecodeTargets)
|
||||
// TODO-REMOVE: remove this log after stable
|
||||
d.logger.Infow("update dependency structure",
|
||||
"structureID", dd.AttachedStructure.StructureId,
|
||||
"structure", dd.AttachedStructure,
|
||||
"decodeTargets", ddwdt.DecodeTargets,
|
||||
"efn", extFrameNum,
|
||||
"sn", extPkt.Packet.SequenceNumber,
|
||||
"isKeyFrame", extPkt.KeyFrame,
|
||||
"currentKeyframe", d.extKeyFrameNum,
|
||||
)
|
||||
|
||||
d.updateDependencyStructure(dd.AttachedStructure, ddwdt.DecodeTargets, extFrameNum)
|
||||
}
|
||||
|
||||
if ddwdt.ExtKeyFrameNum != d.extKeyFrameNum {
|
||||
// keyframe mismatch, drop and reset chains
|
||||
// TODO-REMOVE: remove this log after stable
|
||||
d.logger.Infow("drop packet for keyframe mismatch", "incoming", incomingLayer, "efn", extFrameNum, "sn", extPkt.Packet.SequenceNumber, "requiredKeyFrame", ddwdt.ExtKeyFrameNum, "structureKeyFrame", d.extKeyFrameNum)
|
||||
d.decisions.AddDropped(extFrameNum)
|
||||
d.invalidateKeyFrame()
|
||||
return
|
||||
}
|
||||
|
||||
if ddwdt.ActiveDecodeTargetsUpdated {
|
||||
d.updateActiveDecodeTargets(*dd.ActiveDecodeTargetsBitmask)
|
||||
}
|
||||
|
||||
// TODO-REMOVE: remove this log after stable
|
||||
if len(fd.ChainDiffs) != len(d.chains) {
|
||||
d.logger.Warnw("frame chain diff length mismatch", nil,
|
||||
"incoming", incomingLayer,
|
||||
"efn", extFrameNum,
|
||||
"sn", extPkt.Packet.SequenceNumber,
|
||||
"chainDiffs", fd.ChainDiffs,
|
||||
"chains", len(d.chains),
|
||||
"requiredKeyFrame", ddwdt.ExtKeyFrameNum,
|
||||
"structureKeyFrame", d.extKeyFrameNum)
|
||||
d.decisions.AddDropped(extFrameNum)
|
||||
return
|
||||
}
|
||||
|
||||
for _, chain := range d.chains {
|
||||
chain.OnFrame(extFrameNum, fd)
|
||||
}
|
||||
@@ -133,7 +174,7 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
|
||||
if err != nil {
|
||||
d.decodeTargetsLock.RUnlock()
|
||||
// dtis error, dependency descriptor might lost
|
||||
d.logger.Debugw(fmt.Sprintf("drop packet for frame detection error, incoming: %v", incomingLayer), "err", err)
|
||||
d.logger.Warnw(fmt.Sprintf("drop packet for frame detection error, incoming: %v", incomingLayer), err)
|
||||
d.decisions.AddDropped(extFrameNum)
|
||||
return
|
||||
}
|
||||
@@ -148,34 +189,34 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
|
||||
|
||||
if highestDecodeTarget.Target < 0 {
|
||||
// no active decode target, do not select
|
||||
d.logger.Debugw(
|
||||
"drop packet for no target found",
|
||||
"highestDecodeTarget", highestDecodeTarget,
|
||||
"decodeTargets", d.decodeTargets,
|
||||
"tagetLayer", d.targetLayer,
|
||||
"incoming", incomingLayer,
|
||||
"fn", dd.FrameNumber,
|
||||
"efn", extFrameNum,
|
||||
"sn", extPkt.Packet.SequenceNumber,
|
||||
"isKeyFrame", extPkt.KeyFrame,
|
||||
)
|
||||
// d.logger.Debugw(
|
||||
// "drop packet for no target found",
|
||||
// "highestDecodeTarget", highestDecodeTarget,
|
||||
// "decodeTargets", d.decodeTargets,
|
||||
// "tagetLayer", d.targetLayer,
|
||||
// "incoming", incomingLayer,
|
||||
// "fn", dd.FrameNumber,
|
||||
// "efn", extFrameNum,
|
||||
// "sn", extPkt.Packet.SequenceNumber,
|
||||
// "isKeyFrame", extPkt.KeyFrame,
|
||||
// )
|
||||
d.decisions.AddDropped(extFrameNum)
|
||||
return
|
||||
}
|
||||
|
||||
// DD-TODO : if bandwidth in congest, could drop the 'Discardable' frame
|
||||
if dti == dede.DecodeTargetNotPresent {
|
||||
d.logger.Debugw(
|
||||
"drop packet for decode target not present",
|
||||
"highestDecodeTarget", highestDecodeTarget,
|
||||
"decodeTargets", d.decodeTargets,
|
||||
"tagetLayer", d.targetLayer,
|
||||
"incoming", incomingLayer,
|
||||
"fn", dd.FrameNumber,
|
||||
"efn", extFrameNum,
|
||||
"sn", extPkt.Packet.SequenceNumber,
|
||||
"isKeyFrame", extPkt.KeyFrame,
|
||||
)
|
||||
// d.logger.Debugw(
|
||||
// "drop packet for decode target not present",
|
||||
// "highestDecodeTarget", highestDecodeTarget,
|
||||
// "decodeTargets", d.decodeTargets,
|
||||
// "tagetLayer", d.targetLayer,
|
||||
// "incoming", incomingLayer,
|
||||
// "fn", dd.FrameNumber,
|
||||
// "efn", extFrameNum,
|
||||
// "sn", extPkt.Packet.SequenceNumber,
|
||||
// "isKeyFrame", extPkt.KeyFrame,
|
||||
// )
|
||||
d.decisions.AddDropped(extFrameNum)
|
||||
return
|
||||
}
|
||||
@@ -195,17 +236,17 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
|
||||
}
|
||||
}
|
||||
if !isDecodable {
|
||||
d.logger.Debugw(
|
||||
"drop packet for not decodable",
|
||||
"highestDecodeTarget", highestDecodeTarget,
|
||||
"decodeTargets", d.decodeTargets,
|
||||
"tagetLayer", d.targetLayer,
|
||||
"incoming", incomingLayer,
|
||||
"fn", dd.FrameNumber,
|
||||
"efn", extFrameNum,
|
||||
"sn", extPkt.Packet.SequenceNumber,
|
||||
"isKeyFrame", extPkt.KeyFrame,
|
||||
)
|
||||
// d.logger.Debugw(
|
||||
// "drop packet for not decodable",
|
||||
// "highestDecodeTarget", highestDecodeTarget,
|
||||
// "decodeTargets", d.decodeTargets,
|
||||
// "tagetLayer", d.targetLayer,
|
||||
// "incoming", incomingLayer,
|
||||
// "fn", dd.FrameNumber,
|
||||
// "efn", extFrameNum,
|
||||
// "sn", extPkt.Packet.SequenceNumber,
|
||||
// "isKeyFrame", extPkt.KeyFrame,
|
||||
// )
|
||||
d.decisions.AddDropped(extFrameNum)
|
||||
return
|
||||
}
|
||||
@@ -263,11 +304,33 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
|
||||
// d.logger.Debugw("set active decode targets bitmask", "activeDecodeTargetsBitmask", d.activeDecodeTargetsBitmask)
|
||||
}
|
||||
}
|
||||
bytes, err := ddExtension.Marshal()
|
||||
if err != nil {
|
||||
d.logger.Warnw("error marshalling dependency descriptor extension", err)
|
||||
} else {
|
||||
result.DependencyDescriptorExtension = bytes
|
||||
|
||||
var ddMarshaled bool
|
||||
func() {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
d.logger.Errorw("panic marshalling dependency descriptor extension", nil,
|
||||
"efn", extFrameNum,
|
||||
"sn", extPkt.Packet.SequenceNumber,
|
||||
"keyframeRequired", ddwdt.ExtKeyFrameNum,
|
||||
"currentKeyframe", d.extKeyFrameNum,
|
||||
"panic", r,
|
||||
"stack", string(debug.Stack()))
|
||||
}
|
||||
}()
|
||||
bytes, err := ddExtension.Marshal()
|
||||
if err != nil {
|
||||
d.logger.Warnw("error marshalling dependency descriptor extension", err)
|
||||
} else {
|
||||
result.DependencyDescriptorExtension = bytes
|
||||
ddMarshaled = true
|
||||
}
|
||||
}()
|
||||
|
||||
if !ddMarshaled {
|
||||
// drop packet if we can't marshal dependency descriptor
|
||||
d.decisions.AddDropped(extFrameNum)
|
||||
return
|
||||
}
|
||||
|
||||
if ddwdt.Integrity {
|
||||
@@ -284,8 +347,10 @@ func (d *DependencyDescriptor) Rollback() {
|
||||
d.Base.Rollback()
|
||||
}
|
||||
|
||||
func (d *DependencyDescriptor) updateDependencyStructure(structure *dede.FrameDependencyStructure, decodeTargets []buffer.DependencyDescriptorDecodeTarget) {
|
||||
func (d *DependencyDescriptor) updateDependencyStructure(structure *dede.FrameDependencyStructure, decodeTargets []buffer.DependencyDescriptorDecodeTarget, extFrameNum uint64) {
|
||||
d.structure = structure
|
||||
d.extKeyFrameNum = extFrameNum
|
||||
d.keyFrameValid = true
|
||||
|
||||
d.chains = d.chains[:0]
|
||||
|
||||
@@ -329,6 +394,14 @@ func (d *DependencyDescriptor) updateActiveDecodeTargets(activeDecodeTargetsBitm
|
||||
}
|
||||
}
|
||||
|
||||
func (d *DependencyDescriptor) invalidateKeyFrame() {
|
||||
d.keyFrameValid = false
|
||||
d.chains = d.chains[:0]
|
||||
d.decodeTargetsLock.Lock()
|
||||
d.decodeTargets = d.decodeTargets[:0]
|
||||
d.decodeTargetsLock.Unlock()
|
||||
}
|
||||
|
||||
func (d *DependencyDescriptor) CheckSync() (locked bool, layer int32) {
|
||||
layer = d.GetRequestSpatial()
|
||||
if !d.currentLayer.IsValid() {
|
||||
|
||||
@@ -258,6 +258,14 @@ func TestDependencyDescriptor(t *testing.T) {
|
||||
locked, layer := ddSelector.CheckSync()
|
||||
require.True(t, locked)
|
||||
require.Equal(t, targetLayer.Spatial, layer)
|
||||
|
||||
// should drop frame that relies on a keyframe is not present in current selection
|
||||
framesPrevious := createDDFrames(buffer.VideoLayer{Spatial: 2, Temporal: 2}, 1000)
|
||||
ret = ddSelector.Select(framesPrevious[1], 0)
|
||||
require.False(t, ret.IsSelected)
|
||||
// keyframe lost, out of sync
|
||||
locked, _ = ddSelector.CheckSync()
|
||||
require.False(t, locked)
|
||||
}
|
||||
|
||||
func createDDFrames(maxLayer buffer.VideoLayer, startFrameNumber uint16) []*buffer.ExtPacket {
|
||||
@@ -279,7 +287,7 @@ func createDDFrames(maxLayer buffer.VideoLayer, startFrameNumber uint16) []*buff
|
||||
return decodeTargets[i].Layer.GreaterThan(decodeTargets[j].Layer)
|
||||
})
|
||||
|
||||
chainDiffs := make([]int, len(decodeTargets))
|
||||
chainDiffs := make([]int, int(maxLayer.Spatial)+1)
|
||||
dtis := make([]dd.DecodeTargetIndication, len(decodeTargets))
|
||||
for _, dt := range decodeTargets {
|
||||
dtis[dt.Target] = dd.DecodeTargetSwitch
|
||||
@@ -319,6 +327,7 @@ func createDDFrames(maxLayer buffer.VideoLayer, startFrameNumber uint16) []*buff
|
||||
ActiveDecodeTargetsUpdated: true,
|
||||
Integrity: true,
|
||||
ExtFrameNum: uint64(startFrameNumber),
|
||||
ExtKeyFrameNum: uint64(startFrameNumber),
|
||||
},
|
||||
Packet: &rtp.Packet{
|
||||
Header: rtp.Header{
|
||||
@@ -356,7 +365,6 @@ func createDDFrames(maxLayer buffer.VideoLayer, startFrameNumber uint16) []*buff
|
||||
}
|
||||
|
||||
frame := &buffer.ExtPacket{
|
||||
KeyFrame: true,
|
||||
DependencyDescriptor: &buffer.ExtDependencyDescriptor{
|
||||
Descriptor: &dd.DependencyDescriptor{
|
||||
FrameNumber: startFrameNumber,
|
||||
@@ -367,9 +375,10 @@ func createDDFrames(maxLayer buffer.VideoLayer, startFrameNumber uint16) []*buff
|
||||
DecodeTargetIndications: frameDtis,
|
||||
},
|
||||
},
|
||||
DecodeTargets: decodeTargets,
|
||||
Integrity: true,
|
||||
ExtFrameNum: uint64(startFrameNumber),
|
||||
DecodeTargets: decodeTargets,
|
||||
Integrity: true,
|
||||
ExtFrameNum: uint64(startFrameNumber),
|
||||
ExtKeyFrameNum: keyFrame.DependencyDescriptor.ExtFrameNum,
|
||||
},
|
||||
Packet: &rtp.Packet{
|
||||
Header: rtp.Header{
|
||||
|
||||
Reference in New Issue
Block a user