Introducing frame based stream tracker. (#1267)

* Split stream tracker impl from base

* slight re-arrangement of code

* fps based stream tracker

* MinFPS config

* switch back to packet based tracker

* use video config by default to handle sources without type
This commit is contained in:
Raja Subramanian
2022-12-28 13:00:21 +05:30
committed by GitHub
parent c9ccff0a80
commit 2b031a5112
10 changed files with 630 additions and 239 deletions
+86 -42
View File
@@ -23,6 +23,7 @@ var DefaultStunServers = []string{
}
type CongestionControlProbeMode string
type StreamTrackerType string
const (
generatedCLIFlagUsage = "generated"
@@ -30,6 +31,9 @@ const (
CongestionControlProbeModePadding CongestionControlProbeMode = "padding"
CongestionControlProbeModeMedia CongestionControlProbeMode = "media"
StreamTrackerTypePacket StreamTrackerType = "packet"
StreamTrackerTypeFrame StreamTrackerType = "frame"
StatsUpdateInterval = time.Second * 10
)
@@ -143,22 +147,31 @@ type AudioConfig struct {
SmoothIntervals uint32 `yaml:"smooth_intervals"`
}
type StreamTrackerPacketConfig struct {
SamplesRequired uint32 `yaml:"samples_required"` // number of samples needed per cycle
CyclesRequired uint32 `yaml:"cycles_required"` // number of cycles needed to be active
CycleDuration time.Duration `yaml:"cycle_duration"`
}
type StreamTrackerFrameConfig struct {
MinFPS float64 `yaml:"min_fps"`
}
type StreamTrackerConfig struct {
SamplesRequired uint32 `yaml:"samples_required"`
CyclesRequired uint32 `yaml:"cycles_required"`
CycleDuration time.Duration `yaml:"cycle_duration"`
BitrateReportInterval time.Duration `yaml:"bitrate_report_interval"`
BitrateReportInterval map[int32]time.Duration `yaml:"bitrate_report_interval,omitempty"`
ExemptedLayers []int32 `yaml:"exempted_layers,omitempty"`
PacketTracker map[int32]StreamTrackerPacketConfig `yaml:"packet_tracker,omitempty"`
FrameTracker map[int32]StreamTrackerFrameConfig `yaml:"frame_tracker,omitempty"`
}
type StreamTrackersConfig struct {
Video []StreamTrackerConfig `yaml:"video"`
Screenshare []StreamTrackerConfig `yaml:"screenshare"`
ExemptedLayersVideo []int32 `yaml:"exempted_layers_video"`
ExemptedLayersScreenshare []int32 `yaml:"exempted_layers_screenshare"`
Video StreamTrackerConfig `yaml:"video"`
Screenshare StreamTrackerConfig `yaml:"screenshare"`
}
type VideoConfig struct {
DynacastPauseDelay time.Duration `yaml:"dynacast_pause_delay,omitempty"`
StreamTrackerType StreamTrackerType `yaml:"stream_tracker_type,omitempty"`
StreamTracker StreamTrackersConfig `yaml:"stream_tracker,omitempty"`
}
@@ -272,47 +285,78 @@ func NewConfig(confString string, strictMode bool, c *cli.Context, baseFlags []c
},
Video: VideoConfig{
DynacastPauseDelay: 5 * time.Second,
StreamTrackerType: StreamTrackerTypePacket,
StreamTracker: StreamTrackersConfig{
ExemptedLayersScreenshare: []int32{0},
ExemptedLayersVideo: []int32{},
Screenshare: []StreamTrackerConfig{
{
SamplesRequired: 1,
CyclesRequired: 1,
CycleDuration: 2 * time.Second,
BitrateReportInterval: 4 * time.Second,
Video: StreamTrackerConfig{
BitrateReportInterval: map[int32]time.Duration{
0: 1 * time.Second,
1: 1 * time.Second,
2: 1 * time.Second,
},
{
SamplesRequired: 1,
CyclesRequired: 1,
CycleDuration: 2 * time.Second,
BitrateReportInterval: 4 * time.Second,
ExemptedLayers: []int32{},
PacketTracker: map[int32]StreamTrackerPacketConfig{
0: StreamTrackerPacketConfig{
SamplesRequired: 1,
CyclesRequired: 4,
CycleDuration: 500 * time.Millisecond,
},
1: StreamTrackerPacketConfig{
SamplesRequired: 5,
CyclesRequired: 20,
CycleDuration: 500 * time.Millisecond,
},
2: StreamTrackerPacketConfig{
SamplesRequired: 5,
CyclesRequired: 20,
CycleDuration: 500 * time.Millisecond,
},
},
{
SamplesRequired: 1,
CyclesRequired: 1,
CycleDuration: 2 * time.Second,
BitrateReportInterval: 4 * time.Second,
FrameTracker: map[int32]StreamTrackerFrameConfig{
0: StreamTrackerFrameConfig{
MinFPS: 5.0,
},
1: StreamTrackerFrameConfig{
MinFPS: 5.0,
},
2: StreamTrackerFrameConfig{
MinFPS: 5.0,
},
},
},
Video: []StreamTrackerConfig{
{
SamplesRequired: 1,
CyclesRequired: 4,
CycleDuration: 500 * time.Millisecond,
BitrateReportInterval: 1 * time.Second,
Screenshare: StreamTrackerConfig{
BitrateReportInterval: map[int32]time.Duration{
0: 4 * time.Second,
1: 4 * time.Second,
2: 4 * time.Second,
},
{
SamplesRequired: 5,
CyclesRequired: 20,
CycleDuration: 500 * time.Millisecond,
BitrateReportInterval: 1 * time.Second,
ExemptedLayers: []int32{0},
PacketTracker: map[int32]StreamTrackerPacketConfig{
0: StreamTrackerPacketConfig{
SamplesRequired: 1,
CyclesRequired: 1,
CycleDuration: 2 * time.Second,
},
1: StreamTrackerPacketConfig{
SamplesRequired: 1,
CyclesRequired: 1,
CycleDuration: 2 * time.Second,
},
2: StreamTrackerPacketConfig{
SamplesRequired: 1,
CyclesRequired: 1,
CycleDuration: 2 * time.Second,
},
},
{
SamplesRequired: 5,
CyclesRequired: 20,
CycleDuration: 500 * time.Millisecond,
BitrateReportInterval: 1 * time.Second,
FrameTracker: map[int32]StreamTrackerFrameConfig{
0: StreamTrackerFrameConfig{
MinFPS: 0.5,
},
1: StreamTrackerFrameConfig{
MinFPS: 0.5,
},
2: StreamTrackerFrameConfig{
MinFPS: 0.5,
},
},
},
},
+1
View File
@@ -234,6 +234,7 @@ func (t *MediaTrack) AddReceiver(receiver *webrtc.RTPReceiver, track *webrtc.Tra
t.params.TrackInfo,
LoggerWithCodecMime(t.params.Logger, mime),
twcc,
t.params.VideoConfig.StreamTrackerType,
t.params.VideoConfig.StreamTracker,
sfu.WithPliThrottleConfig(t.params.PLIThrottleConfig),
sfu.WithAudioConfig(t.params.AudioConfig),
+9 -1
View File
@@ -95,6 +95,7 @@ type Buffer struct {
ddParser *DependencyDescriptorParser
maxLayerChangedCB func(int32, int32)
paused bool
frameRateCalculator [DefaultMaxLayerSpatial + 1]FrameRateCalculator
frameRateCalculated bool
}
@@ -123,6 +124,13 @@ func (b *Buffer) SetLogger(logger logger.Logger) {
}
}
func (b *Buffer) SetPaused(paused bool) {
b.Lock()
defer b.Unlock()
b.paused = paused
}
func (b *Buffer) SetTWCC(twcc *twcc.Responder) {
b.Lock()
defer b.Unlock()
@@ -431,7 +439,7 @@ func (b *Buffer) patchExtPacket(ep *ExtPacket, buf []byte) *ExtPacket {
}
func (b *Buffer) doFpsCalc(ep *ExtPacket) {
if b.frameRateCalculated || len(ep.Packet.Payload) == 0 {
if b.paused || b.frameRateCalculated || len(ep.Packet.Payload) == 0 {
return
}
spatial := ep.Spatial
+20 -2
View File
@@ -171,6 +171,7 @@ func NewWebRTCReceiver(
trackInfo *livekit.TrackInfo,
logger logger.Logger,
twcc *twcc.Responder,
trackerType config.StreamTrackerType,
trackerConfig config.StreamTrackersConfig,
opts ...ReceiverOpts,
) *WebRTCReceiver {
@@ -189,7 +190,7 @@ func NewWebRTCReceiver(
isRED: IsRedCodec(track.Codec().MimeType),
}
w.streamTrackerManager = NewStreamTrackerManager(logger, trackInfo, w.isSVC, trackerConfig)
w.streamTrackerManager = NewStreamTrackerManager(logger, trackInfo, w.isSVC, w.codec.ClockRate, trackerType, trackerConfig)
w.streamTrackerManager.OnAvailableLayersChanged(w.downTrackLayerChange)
w.streamTrackerManager.OnBitrateAvailabilityChanged(w.downTrackBitrateAvailabilityChange)
@@ -335,6 +336,7 @@ func (w *WebRTCReceiver) AddUpTrack(track *webrtc.TrackRemote, buff *buffer.Buff
rtt := w.rtt
w.bufferMu.Unlock()
buff.SetRTT(rtt)
buff.SetPaused(w.streamTrackerManager.IsPaused())
if w.Kind() == webrtc.RTPCodecTypeVideo && w.useTrackers {
w.streamTrackerManager.AddTracker(layer)
@@ -348,6 +350,16 @@ func (w *WebRTCReceiver) AddUpTrack(track *webrtc.TrackRemote, buff *buffer.Buff
// the layer
func (w *WebRTCReceiver) SetUpTrackPaused(paused bool) {
w.streamTrackerManager.SetPaused(paused)
w.bufferMu.RLock()
for _, buff := range w.buffers {
if buff == nil {
continue
}
buff.SetPaused(paused)
}
w.bufferMu.RUnlock()
}
func (w *WebRTCReceiver) AddDownTrack(track TrackSender) error {
@@ -560,7 +572,13 @@ func (w *WebRTCReceiver) forwardRTP(layer int32) {
}
if spatialTracker != nil {
spatialTracker.Observe(pkt.Temporal, len(pkt.RawPacket), len(pkt.Packet.Payload))
spatialTracker.Observe(
pkt.Temporal,
len(pkt.RawPacket),
len(pkt.Packet.Payload),
pkt.Packet.Marker,
pkt.Packet.Timestamp,
)
}
w.downTrackSpreader.Broadcast(func(dt TrackSender) {
+42
View File
@@ -0,0 +1,42 @@
package streamtracker
import (
"fmt"
"time"
)
// ------------------------------------------------------------
type StreamStatusChange int32
func (s StreamStatusChange) String() string {
switch s {
case StreamStatusChangeNone:
return "none"
case StreamStatusChangeStopped:
return "stopped"
case StreamStatusChangeActive:
return "active"
default:
return fmt.Sprintf("unknown: %d", int(s))
}
}
const (
StreamStatusChangeNone StreamStatusChange = iota
StreamStatusChangeStopped
StreamStatusChangeActive
)
// ------------------------------------------------------------
type StreamTrackerImpl interface {
Start()
Stop()
Reset()
GetCheckInterval() time.Duration
Observe(hasMarker bool, ts uint32) StreamStatusChange
CheckStatus() StreamStatusChange
}
@@ -1,14 +1,16 @@
package sfu
package streamtracker
import (
"fmt"
"sync"
"time"
"go.uber.org/atomic"
"github.com/livekit/protocol/logger"
"go.uber.org/atomic"
)
// ------------------------------------------------------------
type StreamStatus int32
func (s StreamStatus) String() string {
@@ -18,31 +20,24 @@ func (s StreamStatus) String() string {
case StreamStatusActive:
return "active"
default:
return "unknown"
return fmt.Sprintf("unknown: %d", int(s))
}
}
const (
StreamStatusStopped StreamStatus = 0
StreamStatusActive StreamStatus = 1
StreamStatusStopped StreamStatus = iota
StreamStatusActive
)
// ------------------------------------------------------------
type StreamTrackerParams struct {
// number of samples needed per cycle
SamplesRequired uint32
// number of cycles needed to be active
CyclesRequired uint32
CycleDuration time.Duration
StreamTrackerImpl StreamTrackerImpl
BitrateReportInterval time.Duration
Logger logger.Logger
}
// StreamTracker keeps track of packet flow and ensures a particular up track is consistently producing
// It runs its own goroutine for detection, and fires OnStatusChanged callback
type StreamTracker struct {
params StreamTrackerParams
@@ -51,16 +46,11 @@ type StreamTracker struct {
lock sync.RWMutex
paused bool
countSinceLast uint32 // number of packets received since last check
generation atomic.Uint32
paused bool
generation atomic.Uint32
initialized bool
status StreamStatus
// only access within detectWorker
cycleCount uint32
status StreamStatus
lastNotifiedStatus StreamStatus
lastBitrateReport time.Time
bytesForBitrate [4]int64
@@ -91,33 +81,31 @@ func (s *StreamTracker) Status() StreamStatus {
return s.status
}
func (s *StreamTracker) maybeSetStatus(status StreamStatus) (StreamStatus, bool) {
changed := false
if s.status != status {
s.status = status
changed = true
}
return status, changed
func (s *StreamTracker) setStatusLocked(status StreamStatus) {
s.status = status
}
func (s *StreamTracker) maybeNotifyStatus(status StreamStatus, changed bool) {
if changed && s.onStatusChanged != nil {
func (s *StreamTracker) maybeNotifyStatus() {
var status StreamStatus
notify := false
s.lock.Lock()
if s.status != s.lastNotifiedStatus {
notify = true
status = s.status
s.lastNotifiedStatus = s.status
}
s.lock.Unlock()
if notify && s.onStatusChanged != nil {
s.onStatusChanged(status)
}
}
func (s *StreamTracker) init() {
s.lock.Lock()
status, changed := s.maybeSetStatus(StreamStatusActive)
s.lock.Unlock()
s.maybeNotifyStatus(status, changed)
go s.detectWorker(s.generation.Load())
}
func (s *StreamTracker) Start() {
s.lock.Lock()
defer s.lock.Unlock()
s.params.StreamTrackerImpl.Start()
}
func (s *StreamTracker) Stop() {
@@ -131,29 +119,28 @@ func (s *StreamTracker) Stop() {
// bump generation to trigger exit of worker
s.generation.Inc()
s.params.StreamTrackerImpl.Stop()
}
func (s *StreamTracker) Reset() {
s.lock.Lock()
defer s.lock.Unlock()
if s.isStopped {
s.lock.Unlock()
return
}
s.resetLocked()
s.lock.Unlock()
s.maybeNotifyStatus()
}
func (s *StreamTracker) resetLocked() {
// bump generation to trigger exit of current worker
s.generation.Inc()
s.countSinceLast = 0
s.cycleCount = 0
s.initialized = false
s.status = StreamStatusStopped
s.setStatusLocked(StreamStatusStopped)
for i := 0; i < len(s.bytesForBitrate); i++ {
s.bytesForBitrate[i] = 0
@@ -161,59 +148,56 @@ func (s *StreamTracker) resetLocked() {
for i := 0; i < len(s.bitrate); i++ {
s.bitrate[i] = 0
}
s.params.StreamTrackerImpl.Reset()
}
func (s *StreamTracker) SetPaused(paused bool) {
s.lock.Lock()
s.paused = paused
status := s.status
changed := false
if !paused {
s.resetLocked()
} else {
// bump generation to trigger exit of current worker
s.generation.Inc()
status, changed = s.maybeSetStatus(StreamStatusStopped)
s.setStatusLocked(StreamStatusStopped)
}
s.lock.Unlock()
s.maybeNotifyStatus(status, changed)
s.maybeNotifyStatus()
}
// Observe a packet that's received
func (s *StreamTracker) Observe(temporalLayer int32, pktSize int, payloadSize int) {
func (s *StreamTracker) Observe(
temporalLayer int32,
pktSize int,
payloadSize int,
hasMarker bool,
ts uint32,
) {
s.lock.Lock()
defer s.lock.Unlock()
if s.isStopped || s.paused || payloadSize == 0 {
s.lock.Unlock()
return
}
if !s.initialized {
// first packet
s.initialized = true
s.countSinceLast = 1
statusChange := s.params.StreamTrackerImpl.Observe(hasMarker, ts)
if statusChange == StreamStatusChangeActive {
s.setStatusLocked(StreamStatusActive)
s.lastBitrateReport = time.Now()
if temporalLayer >= 0 {
s.bytesForBitrate[temporalLayer] += int64(pktSize)
}
// declare stream active and start the detection worker
go s.init()
return
go s.worker(s.generation.Load())
}
s.countSinceLast++
if temporalLayer >= 0 {
s.bytesForBitrate[temporalLayer] += int64(pktSize)
}
s.lock.Unlock()
if statusChange != StreamStatusChangeNone {
s.maybeNotifyStatus()
}
}
// BitrateTemporalCumulative returns the current stream bitrate temporal layer accumulated with lower temporal layers.
@@ -236,8 +220,8 @@ func (s *StreamTracker) BitrateTemporalCumulative() []int64 {
return brs
}
func (s *StreamTracker) detectWorker(generation uint32) {
ticker := time.NewTicker(s.params.CycleDuration)
func (s *StreamTracker) worker(generation uint32) {
ticker := time.NewTicker(s.params.StreamTrackerImpl.GetCheckInterval())
defer ticker.Stop()
tickerBitrate := time.NewTicker(s.params.BitrateReportInterval)
@@ -249,7 +233,7 @@ func (s *StreamTracker) detectWorker(generation uint32) {
if generation != s.generation.Load() {
return
}
s.detectChanges()
s.updateStatus()
case <-tickerBitrate.C:
if generation != s.generation.Load() {
@@ -260,28 +244,17 @@ func (s *StreamTracker) detectWorker(generation uint32) {
}
}
func (s *StreamTracker) detectChanges() {
func (s *StreamTracker) updateStatus() {
s.lock.Lock()
if s.countSinceLast >= s.params.SamplesRequired {
s.cycleCount++
} else {
s.cycleCount = 0
switch s.params.StreamTrackerImpl.CheckStatus() {
case StreamStatusChangeStopped:
s.setStatusLocked(StreamStatusStopped)
case StreamStatusChangeActive:
s.setStatusLocked(StreamStatusActive)
}
status := s.status
changed := false
if s.cycleCount == 0 {
// flip to stopped
status, changed = s.maybeSetStatus(StreamStatusStopped)
} else if s.cycleCount >= s.params.CyclesRequired {
// flip to active
status, changed = s.maybeSetStatus(StreamStatusActive)
}
s.countSinceLast = 0
s.lock.Unlock()
s.maybeNotifyStatus(status, changed)
s.maybeNotifyStatus()
}
func (s *StreamTracker) bitrateReport() {
@@ -0,0 +1,177 @@
package streamtracker
import (
"math"
"time"
"github.com/livekit/livekit-server/pkg/config"
"github.com/livekit/protocol/logger"
)
const (
checkInterval = 500 * time.Millisecond
staleWindowFactor = 5
frameRateResolution = float64(0.01) // 1 frame every 100 seconds
)
type StreamTrackerFrameParams struct {
Config config.StreamTrackerFrameConfig
ClockRate uint32
Logger logger.Logger
}
type StreamTrackerFrame struct {
params StreamTrackerFrameParams
initialized bool
tsInitialized bool
oldestTS uint32
newestTS uint32
numFrames int
lowestFrameRate float64
evalInterval time.Duration
lastStatusCheckAt time.Time
}
func NewStreamTrackerFrame(params StreamTrackerFrameParams) StreamTrackerImpl {
s := &StreamTrackerFrame{
params: params,
}
s.Reset()
return s
}
func (s *StreamTrackerFrame) Start() {
}
func (s *StreamTrackerFrame) Stop() {
}
func (s *StreamTrackerFrame) Reset() {
s.initialized = false
s.tsInitialized = false
s.oldestTS = 0
s.newestTS = 0
s.numFrames = 0
s.lowestFrameRate = 0.0
s.updateEvalInterval()
s.lastStatusCheckAt = time.Time{}
}
func (s *StreamTrackerFrame) GetCheckInterval() time.Duration {
return checkInterval
}
func (s *StreamTrackerFrame) Observe(hasMarker bool, ts uint32) StreamStatusChange {
if !s.initialized {
s.initialized = true
if hasMarker {
s.tsInitialized = true
s.oldestTS = ts
s.newestTS = ts
s.numFrames = 1
}
return StreamStatusChangeActive
}
if hasMarker {
if !s.tsInitialized {
s.tsInitialized = true
s.oldestTS = ts
s.newestTS = ts
s.numFrames = 1
} else {
diff := ts - s.oldestTS
if diff > (1 << 31) {
s.oldestTS = ts
}
diff = ts - s.newestTS
if diff < (1 << 31) {
s.newestTS = ts
}
s.numFrames++
}
}
return StreamStatusChangeNone
}
func (s *StreamTrackerFrame) CheckStatus() StreamStatusChange {
if !s.initialized {
// should not be getting called when not initialized, but be safe
return StreamStatusChangeNone
}
// calculate frame rate since last check
frameRate := float64(0.0)
diff := s.newestTS - s.oldestTS
if diff > 0 || s.numFrames > 1 {
if diff > s.params.ClockRate*staleWindowFactor {
s.params.Logger.Infow("eval window might be stale", "numFrames", s.numFrames, "timeElapsed", float64(diff)/float64(s.params.ClockRate))
// STREAM-TRACKER-FRAME-TODO: might need to protect against one frame, long pause and then one or more frames, i. e. window getting stale.
// One possible option is to reset the fps measurement variables (tsInitialized, oldestTS, newestTS, numFrames, lowestFrameRate, evelInterval)
// and restart the lowest frame rate calulation process.
}
frameRate = float64(s.params.ClockRate) / float64(diff) * float64(s.numFrames-1)
frameRate = math.Round(frameRate/frameRateResolution) * frameRateResolution
}
if s.lowestFrameRate == 0.0 {
if frameRate == 0.0 {
// need at least two frames to kick things off
return StreamStatusChangeNone
}
s.lowestFrameRate = frameRate
s.updateEvalInterval()
s.params.Logger.Infow("initializing lowest frame rate", "lowestFPS", s.lowestFrameRate, "evalInterval", s.evalInterval)
} else {
// check only at intervals based on lowest seen frame rate
if s.lastStatusCheckAt.IsZero() {
s.lastStatusCheckAt = time.Now()
}
if time.Since(s.lastStatusCheckAt) < s.evalInterval {
return StreamStatusChangeNone
}
s.lastStatusCheckAt = time.Now()
}
// reset for next evaluation interval
s.oldestTS = s.newestTS
s.numFrames = 1
// STREAM-TRACKER-FRAME-TODO: this will run into challenges for frame rate falling steeply, how to address that
// look at some referential rules (between layers) for possibilities to solve it. Currently, this is addressed
// by setting a source aware min FPS to ensure evaluation window in long enough
// update lowest seen frame rate
if frameRate > 0.0 && s.lowestFrameRate > frameRate {
s.lowestFrameRate = frameRate
s.updateEvalInterval()
s.params.Logger.Infow("updating lowest frame rate", "lowestFPS", s.lowestFrameRate, "evalInterval", s.evalInterval)
}
if frameRate == 0.0 {
return StreamStatusChangeStopped
}
return StreamStatusChangeActive
}
func (s *StreamTrackerFrame) updateEvalInterval() {
s.evalInterval = checkInterval
if s.lowestFrameRate > 0 {
lowestFrameRateInterval := time.Duration(float64(time.Second) / s.lowestFrameRate)
if lowestFrameRateInterval > s.evalInterval {
s.evalInterval = lowestFrameRateInterval
}
}
if s.params.Config.MinFPS > 0 {
minFPSInterval := time.Duration(float64(time.Second) / s.params.Config.MinFPS)
if minFPSInterval > s.evalInterval {
s.evalInterval = minFPSInterval
}
}
}
@@ -0,0 +1,83 @@
package streamtracker
import (
"time"
"github.com/livekit/livekit-server/pkg/config"
"github.com/livekit/protocol/logger"
)
type StreamTrackerPacketParams struct {
Config config.StreamTrackerPacketConfig
Logger logger.Logger
}
type StreamTrackerPacket struct {
params StreamTrackerPacketParams
countSinceLast uint32 // number of packets received since last check
initialized bool
cycleCount uint32
}
func NewStreamTrackerPacket(params StreamTrackerPacketParams) StreamTrackerImpl {
return &StreamTrackerPacket{
params: params,
}
}
func (s *StreamTrackerPacket) Start() {
}
func (s *StreamTrackerPacket) Stop() {
}
func (s *StreamTrackerPacket) Reset() {
s.countSinceLast = 0
s.cycleCount = 0
s.initialized = false
}
func (s *StreamTrackerPacket) GetCheckInterval() time.Duration {
return s.params.Config.CycleDuration
}
func (s *StreamTrackerPacket) Observe(_hasMarker bool, _ts uint32) StreamStatusChange {
if !s.initialized {
// first packet
s.initialized = true
s.countSinceLast = 1
return StreamStatusChangeActive
}
s.countSinceLast++
return StreamStatusChangeNone
}
func (s *StreamTrackerPacket) CheckStatus() StreamStatusChange {
if !s.initialized {
// should not be getting called when not initialized, but be safe
return StreamStatusChangeNone
}
if s.countSinceLast >= s.params.Config.SamplesRequired {
s.cycleCount++
} else {
s.cycleCount = 0
}
statusChange := StreamStatusChangeNone
if s.cycleCount == 0 {
// no packets seen for a period, flip to stopped
statusChange = StreamStatusChangeStopped
} else if s.cycleCount >= s.params.Config.CyclesRequired {
// packets seen for some time after resume, flip to active
statusChange = StreamStatusChangeActive
}
s.countSinceLast = 0
return statusChange
}
@@ -1,4 +1,4 @@
package sfu
package streamtracker
import (
"fmt"
@@ -9,15 +9,23 @@ import (
"github.com/stretchr/testify/require"
"go.uber.org/atomic"
"github.com/livekit/livekit-server/pkg/config"
"github.com/livekit/livekit-server/pkg/testutils"
"github.com/livekit/protocol/logger"
)
func newStreamTracker(samplesRequired uint32, cyclesRequired uint32, cycleDuration time.Duration) *StreamTracker {
func newStreamTrackerPacket(samplesRequired uint32, cyclesRequired uint32, cycleDuration time.Duration) *StreamTracker {
stp := NewStreamTrackerPacket(StreamTrackerPacketParams{
Config: config.StreamTrackerPacketConfig{
SamplesRequired: samplesRequired,
CyclesRequired: cyclesRequired,
CycleDuration: cycleDuration,
},
Logger: logger.GetLogger(),
})
return NewStreamTracker(StreamTrackerParams{
SamplesRequired: samplesRequired,
CyclesRequired: cyclesRequired,
CycleDuration: cycleDuration,
StreamTrackerImpl: stp,
BitrateReportInterval: 1 * time.Second,
Logger: logger.GetLogger(),
})
@@ -26,7 +34,7 @@ func newStreamTracker(samplesRequired uint32, cyclesRequired uint32, cycleDurati
func TestStreamTracker(t *testing.T) {
t.Run("flips to active on first observe", func(t *testing.T) {
callbackCalled := atomic.NewBool(false)
tracker := newStreamTracker(5, 60, 500*time.Millisecond)
tracker := newStreamTrackerPacket(5, 60, 500*time.Millisecond)
tracker.Start()
tracker.OnStatusChanged(func(status StreamStatus) {
callbackCalled.Store(true)
@@ -34,14 +42,14 @@ func TestStreamTracker(t *testing.T) {
require.Equal(t, StreamStatusStopped, tracker.Status())
// observe first packet
tracker.Observe(0, 20, 10)
tracker.Observe(0, 20, 10, false, 0)
testutils.WithTimeout(t, func() string {
if callbackCalled.Load() {
return ""
} else {
return "first packet didn't activate stream"
}
return "first packet didn't activate stream"
})
require.Equal(t, StreamStatusActive, tracker.Status())
@@ -51,7 +59,7 @@ func TestStreamTracker(t *testing.T) {
})
t.Run("flips to inactive immediately", func(t *testing.T) {
tracker := newStreamTracker(5, 60, 500*time.Millisecond)
tracker := newStreamTrackerPacket(5, 60, 500*time.Millisecond)
tracker.Start()
require.Equal(t, StreamStatusStopped, tracker.Status())
@@ -65,7 +73,7 @@ func TestStreamTracker(t *testing.T) {
callbackStatusMu.Unlock()
})
tracker.Observe(0, 20, 10)
tracker.Observe(0, 20, 10, false, 0)
testutils.WithTimeout(t, func() string {
callbackStatusMu.RLock()
defer callbackStatusMu.RUnlock()
@@ -79,7 +87,7 @@ func TestStreamTracker(t *testing.T) {
require.Equal(t, StreamStatusActive, tracker.Status())
// run a single iteration
tracker.detectChanges()
tracker.updateStatus()
testutils.WithTimeout(t, func() string {
callbackStatusMu.RLock()
@@ -98,46 +106,46 @@ func TestStreamTracker(t *testing.T) {
})
t.Run("flips back to active after iterations", func(t *testing.T) {
tracker := newStreamTracker(1, 2, 500*time.Millisecond)
tracker := newStreamTrackerPacket(1, 2, 500*time.Millisecond)
tracker.Start()
require.Equal(t, StreamStatusStopped, tracker.Status())
tracker.Observe(0, 20, 10)
tracker.Observe(0, 20, 10, false, 0)
testutils.WithTimeout(t, func() string {
if tracker.Status() == StreamStatusActive {
return ""
} else {
return "first packet did not activate stream"
}
return "first packet did not activate stream"
})
tracker.maybeSetStatus(StreamStatusStopped)
tracker.setStatusLocked(StreamStatusStopped)
tracker.Observe(0, 20, 10)
tracker.detectChanges()
tracker.Observe(0, 20, 10, false, 0)
tracker.updateStatus()
require.Equal(t, StreamStatusStopped, tracker.Status())
tracker.Observe(0, 20, 10)
tracker.detectChanges()
tracker.Observe(0, 20, 10, false, 0)
tracker.updateStatus()
require.Equal(t, StreamStatusActive, tracker.Status())
tracker.Stop()
})
t.Run("changes to inactive when paused", func(t *testing.T) {
tracker := newStreamTracker(5, 60, 500*time.Millisecond)
tracker := newStreamTrackerPacket(5, 60, 500*time.Millisecond)
tracker.Start()
tracker.Observe(0, 20, 10)
tracker.Observe(0, 20, 10, false, 0)
testutils.WithTimeout(t, func() string {
if tracker.Status() == StreamStatusActive {
return ""
} else {
return "first packet did not activate stream"
}
return "first packet did not activate stream"
})
tracker.SetPaused(true)
tracker.detectChanges()
tracker.updateStatus()
require.Equal(t, StreamStatusStopped, tracker.Status())
tracker.Stop()
@@ -145,7 +153,7 @@ func TestStreamTracker(t *testing.T) {
t.Run("flips back to active on first observe after reset", func(t *testing.T) {
callbackCalled := atomic.NewUint32(0)
tracker := newStreamTracker(5, 60, 500*time.Millisecond)
tracker := newStreamTrackerPacket(5, 60, 500*time.Millisecond)
tracker.Start()
tracker.OnStatusChanged(func(status StreamStatus) {
callbackCalled.Inc()
@@ -153,46 +161,48 @@ func TestStreamTracker(t *testing.T) {
require.Equal(t, StreamStatusStopped, tracker.Status())
// observe first packet
tracker.Observe(0, 20, 10)
tracker.Observe(0, 20, 10, false, 0)
testutils.WithTimeout(t, func() string {
if callbackCalled.Load() == 1 {
return ""
} else {
return fmt.Sprintf("expected onStatusChanged to be called once, actual: %d", callbackCalled.Load())
}
return fmt.Sprintf("expected onStatusChanged to be called once, actual: %d", callbackCalled.Load())
})
require.Equal(t, StreamStatusActive, tracker.Status())
require.Equal(t, uint32(1), callbackCalled.Load())
// observe a few more
tracker.Observe(0, 20, 10)
tracker.Observe(0, 20, 10)
tracker.Observe(0, 20, 10)
tracker.Observe(0, 20, 10)
tracker.detectChanges()
tracker.Observe(0, 20, 10, false, 0)
tracker.Observe(0, 20, 10, false, 0)
tracker.Observe(0, 20, 10, false, 0)
tracker.Observe(0, 20, 10, false, 0)
tracker.updateStatus()
// should still be active
require.Equal(t, StreamStatusActive, tracker.Status())
require.Equal(t, uint32(1), callbackCalled.Load())
// Reset. The first packet after reset should flip state again
tracker.Reset()
require.Equal(t, StreamStatusStopped, tracker.Status())
require.Equal(t, uint32(2), callbackCalled.Load())
// first packet after reset
tracker.Observe(0, 20, 10)
tracker.Observe(0, 20, 10, false, 0)
testutils.WithTimeout(t, func() string {
if callbackCalled.Load() == 2 {
if callbackCalled.Load() == 3 {
return ""
} else {
return fmt.Sprintf("expected onStatusChanged to be called twice, actual %d", callbackCalled.Load())
}
return fmt.Sprintf("expected onStatusChanged to be called thrice, actual %d", callbackCalled.Load())
})
require.Equal(t, StreamStatusActive, tracker.Status())
require.Equal(t, uint32(2), callbackCalled.Load())
require.Equal(t, uint32(3), callbackCalled.Load())
tracker.Stop()
})
+91 -56
View File
@@ -6,6 +6,7 @@ import (
"github.com/livekit/livekit-server/pkg/config"
"github.com/livekit/livekit-server/pkg/sfu/buffer"
"github.com/livekit/livekit-server/pkg/sfu/streamtracker"
"github.com/livekit/protocol/livekit"
"github.com/livekit/protocol/logger"
)
@@ -15,15 +16,14 @@ type StreamTrackerManager struct {
trackInfo *livekit.TrackInfo
isSVC bool
maxPublishedLayer int32
clockRate uint32
configVideo []StreamTrackerParams
configScreenshare []StreamTrackerParams
exemptedLayersVideo []int32
exemptedLayersScreenshare []int32
trackerType config.StreamTrackerType
trackerConfig config.StreamTrackerConfig
lock sync.RWMutex
trackers [DefaultMaxLayerSpatial + 1]*StreamTracker
trackers [DefaultMaxLayerSpatial + 1]*streamtracker.StreamTracker
availableLayers []int32
exemptedLayers []int32
@@ -35,32 +35,30 @@ type StreamTrackerManager struct {
onMaxLayerChanged func(maxLayer int32)
}
func NewStreamTrackerManager(logger logger.Logger, trackInfo *livekit.TrackInfo, isSVC bool, trackerConfig config.StreamTrackersConfig) *StreamTrackerManager {
func NewStreamTrackerManager(
logger logger.Logger,
trackInfo *livekit.TrackInfo,
isSVC bool,
clockRate uint32,
trackerType config.StreamTrackerType,
trackersConfig config.StreamTrackersConfig,
) *StreamTrackerManager {
s := &StreamTrackerManager{
logger: logger,
trackInfo: trackInfo,
isSVC: isSVC,
maxPublishedLayer: 0,
exemptedLayersVideo: trackerConfig.ExemptedLayersVideo,
exemptedLayersScreenshare: trackerConfig.ExemptedLayersScreenshare,
logger: logger,
trackInfo: trackInfo,
isSVC: isSVC,
maxPublishedLayer: 0,
clockRate: clockRate,
trackerType: trackerType,
}
for _, layer := range trackerConfig.Video {
s.configVideo = append(s.configVideo, StreamTrackerParams{
SamplesRequired: layer.SamplesRequired,
CyclesRequired: layer.CyclesRequired,
CycleDuration: layer.CycleDuration,
BitrateReportInterval: layer.BitrateReportInterval,
})
}
for _, layer := range trackerConfig.Screenshare {
s.configScreenshare = append(s.configScreenshare, StreamTrackerParams{
SamplesRequired: layer.SamplesRequired,
CyclesRequired: layer.CyclesRequired,
CycleDuration: layer.CycleDuration,
BitrateReportInterval: layer.BitrateReportInterval,
})
switch s.trackInfo.Source {
case livekit.TrackSource_SCREEN_SHARE:
s.trackerConfig = trackersConfig.Screenshare
case livekit.TrackSource_CAMERA:
s.trackerConfig = trackersConfig.Video
default:
s.trackerConfig = trackersConfig.Video
}
for _, layer := range s.trackInfo.Layers {
@@ -86,27 +84,60 @@ func (s *StreamTrackerManager) OnMaxLayerChanged(f func(maxLayer int32)) {
s.onMaxLayerChanged = f
}
func (s *StreamTrackerManager) AddTracker(layer int32) *StreamTracker {
var params StreamTrackerParams
if s.trackInfo.Source == livekit.TrackSource_SCREEN_SHARE {
if int(layer) >= len(s.configScreenshare) {
return nil
}
params = s.configScreenshare[layer]
} else {
if int(layer) >= len(s.configVideo) {
return nil
}
params = s.configVideo[layer]
func (s *StreamTrackerManager) createStreamTrackerPacket(layer int32) streamtracker.StreamTrackerImpl {
packetTrackerConfig, ok := s.trackerConfig.PacketTracker[layer]
if !ok {
return nil
}
params.Logger = s.logger.WithValues("layer", layer)
tracker := NewStreamTracker(params)
params := streamtracker.StreamTrackerPacketParams{
Config: packetTrackerConfig,
Logger: s.logger.WithValues("layer", layer),
}
return streamtracker.NewStreamTrackerPacket(params)
}
func (s *StreamTrackerManager) createStreamTrackerFrame(layer int32) streamtracker.StreamTrackerImpl {
frameTrackerConfig, ok := s.trackerConfig.FrameTracker[layer]
if !ok {
return nil
}
params := streamtracker.StreamTrackerFrameParams{
Config: frameTrackerConfig,
ClockRate: s.clockRate,
Logger: s.logger.WithValues("layer", layer),
}
return streamtracker.NewStreamTrackerFrame(params)
}
func (s *StreamTrackerManager) AddTracker(layer int32) *streamtracker.StreamTracker {
bitrateInterval, ok := s.trackerConfig.BitrateReportInterval[layer]
if !ok {
return nil
}
var trackerImpl streamtracker.StreamTrackerImpl
switch s.trackerType {
case config.StreamTrackerTypePacket:
trackerImpl = s.createStreamTrackerPacket(layer)
case config.StreamTrackerTypeFrame:
trackerImpl = s.createStreamTrackerFrame(layer)
}
if trackerImpl == nil {
return nil
}
tracker := streamtracker.NewStreamTracker(streamtracker.StreamTrackerParams{
StreamTrackerImpl: trackerImpl,
BitrateReportInterval: bitrateInterval,
Logger: s.logger.WithValues("layer", layer),
})
s.logger.Debugw("StreamTrackerManager add track", "layer", layer)
tracker.OnStatusChanged(func(status StreamStatus) {
tracker.OnStatusChanged(func(status streamtracker.StreamStatus) {
s.logger.Debugw("StreamTrackerManager OnStatusChanged", "layer", layer, "status", status)
if status == StreamStatusStopped {
if status == streamtracker.StreamStatusStopped {
s.removeAvailableLayer(layer)
} else {
s.addAvailableLayer(layer)
@@ -119,9 +150,11 @@ func (s *StreamTrackerManager) AddTracker(layer int32) *StreamTracker {
})
s.lock.Lock()
paused := s.paused
s.trackers[layer] = tracker
s.lock.Unlock()
tracker.SetPaused(paused)
tracker.Start()
return tracker
}
@@ -156,7 +189,7 @@ func (s *StreamTrackerManager) RemoveAllTrackers() {
}
}
func (s *StreamTrackerManager) GetTracker(layer int32) *StreamTracker {
func (s *StreamTrackerManager) GetTracker(layer int32) *streamtracker.StreamTracker {
s.lock.RLock()
defer s.lock.RUnlock()
@@ -176,6 +209,13 @@ func (s *StreamTrackerManager) SetPaused(paused bool) {
}
}
func (s *StreamTrackerManager) IsPaused() bool {
s.lock.RLock()
defer s.lock.RUnlock()
return s.paused
}
func (s *StreamTrackerManager) SetMaxExpectedSpatialLayer(layer int32) int32 {
s.lock.Lock()
prev := s.maxExpectedLayer
@@ -197,7 +237,7 @@ func (s *StreamTrackerManager) SetMaxExpectedSpatialLayer(layer int32) int32 {
// But, those conditions should be rare. In those cases, the restart will
// take longer.
//
var trackersToReset []*StreamTracker
var trackersToReset []*streamtracker.StreamTracker
for l := s.maxExpectedLayer + 1; l <= layer; l++ {
if s.hasSpatialLayerLocked(l) {
continue
@@ -393,13 +433,7 @@ func (s *StreamTrackerManager) removeAvailableLayer(layer int32) {
// remove from available if not exempt
//
exempt := false
var sourceExemptedLayers []int32
if s.trackInfo.Source == livekit.TrackSource_SCREEN_SHARE {
sourceExemptedLayers = s.exemptedLayersScreenshare
} else {
sourceExemptedLayers = s.exemptedLayersVideo
}
for _, l := range sourceExemptedLayers {
for _, l := range s.trackerConfig.ExemptedLayers {
if layer == l {
exempt = true
break
@@ -414,7 +448,8 @@ func (s *StreamTrackerManager) removeAvailableLayer(layer int32) {
newLayers := make([]int32, 0, DefaultMaxLayerSpatial+1)
for _, l := range s.availableLayers {
if exempt || l != layer {
// do not remove layers for non-simulcast
if exempt || l != layer || len(s.trackInfo.Layers) < 2 {
newLayers = append(newLayers, l)
}
}