diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index 471aea927..c6a018972 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -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 diff --git a/pkg/sfu/buffer/dependencydescriptorparser.go b/pkg/sfu/buffer/dependencydescriptorparser.go index 05b675ad0..6c91af260 100644 --- a/pkg/sfu/buffer/dependencydescriptorparser.go +++ b/pkg/sfu/buffer/dependencydescriptorparser.go @@ -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 } diff --git a/pkg/sfu/streamtracker/streamtracker_dd.go b/pkg/sfu/streamtracker/streamtracker_dd.go index b0f74ef44..be8009eb1 100644 --- a/pkg/sfu/streamtracker/streamtracker_dd.go +++ b/pkg/sfu/streamtracker/streamtracker_dd.go @@ -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 diff --git a/pkg/sfu/videolayerselector/dependencydescriptor.go b/pkg/sfu/videolayerselector/dependencydescriptor.go index f03f78c6c..f458928d8 100644 --- a/pkg/sfu/videolayerselector/dependencydescriptor.go +++ b/pkg/sfu/videolayerselector/dependencydescriptor.go @@ -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() { diff --git a/pkg/sfu/videolayerselector/dependencydescriptor_test.go b/pkg/sfu/videolayerselector/dependencydescriptor_test.go index c013e46af..863c6755f 100644 --- a/pkg/sfu/videolayerselector/dependencydescriptor_test.go +++ b/pkg/sfu/videolayerselector/dependencydescriptor_test.go @@ -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{