From 2b031a511276d39da5e8644d549107d3284faca2 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Wed, 28 Dec 2022 13:00:21 +0530 Subject: [PATCH] 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 --- pkg/config/config.go | 128 ++++++++----- pkg/rtc/mediatrack.go | 1 + pkg/sfu/buffer/buffer.go | 10 +- pkg/sfu/receiver.go | 22 ++- pkg/sfu/streamtracker/interfaces.go | 42 +++++ pkg/sfu/{ => streamtracker}/streamtracker.go | 169 +++++++---------- pkg/sfu/streamtracker/streamtracker_frame.go | 177 ++++++++++++++++++ pkg/sfu/streamtracker/streamtracker_packet.go | 83 ++++++++ .../streamtracker_packet_test.go} | 90 +++++---- pkg/sfu/streamtrackermanager.go | 147 +++++++++------ 10 files changed, 630 insertions(+), 239 deletions(-) create mode 100644 pkg/sfu/streamtracker/interfaces.go rename pkg/sfu/{ => streamtracker}/streamtracker.go (60%) create mode 100644 pkg/sfu/streamtracker/streamtracker_frame.go create mode 100644 pkg/sfu/streamtracker/streamtracker_packet.go rename pkg/sfu/{streamtracker_test.go => streamtracker/streamtracker_packet_test.go} (66%) diff --git a/pkg/config/config.go b/pkg/config/config.go index eb3388cc0..a3ceaa6fb 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -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, + }, }, }, }, diff --git a/pkg/rtc/mediatrack.go b/pkg/rtc/mediatrack.go index eedfa6184..a370c3661 100644 --- a/pkg/rtc/mediatrack.go +++ b/pkg/rtc/mediatrack.go @@ -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), diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index 3ef846297..2e30a242d 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -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 diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index f8aec6a67..644f82329 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -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) { diff --git a/pkg/sfu/streamtracker/interfaces.go b/pkg/sfu/streamtracker/interfaces.go new file mode 100644 index 000000000..3837f8470 --- /dev/null +++ b/pkg/sfu/streamtracker/interfaces.go @@ -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 +} diff --git a/pkg/sfu/streamtracker.go b/pkg/sfu/streamtracker/streamtracker.go similarity index 60% rename from pkg/sfu/streamtracker.go rename to pkg/sfu/streamtracker/streamtracker.go index ad622514e..8f8715a8a 100644 --- a/pkg/sfu/streamtracker.go +++ b/pkg/sfu/streamtracker/streamtracker.go @@ -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() { diff --git a/pkg/sfu/streamtracker/streamtracker_frame.go b/pkg/sfu/streamtracker/streamtracker_frame.go new file mode 100644 index 000000000..b809cb3a6 --- /dev/null +++ b/pkg/sfu/streamtracker/streamtracker_frame.go @@ -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 + } + } +} diff --git a/pkg/sfu/streamtracker/streamtracker_packet.go b/pkg/sfu/streamtracker/streamtracker_packet.go new file mode 100644 index 000000000..78866e40e --- /dev/null +++ b/pkg/sfu/streamtracker/streamtracker_packet.go @@ -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 +} diff --git a/pkg/sfu/streamtracker_test.go b/pkg/sfu/streamtracker/streamtracker_packet_test.go similarity index 66% rename from pkg/sfu/streamtracker_test.go rename to pkg/sfu/streamtracker/streamtracker_packet_test.go index 82d9a9cc0..2aee13ecb 100644 --- a/pkg/sfu/streamtracker_test.go +++ b/pkg/sfu/streamtracker/streamtracker_packet_test.go @@ -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() }) diff --git a/pkg/sfu/streamtrackermanager.go b/pkg/sfu/streamtrackermanager.go index 28cd5f295..f4d278ab2 100644 --- a/pkg/sfu/streamtrackermanager.go +++ b/pkg/sfu/streamtrackermanager.go @@ -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) } }