From 156114fcafc9b505fc2da0d9b6134720c34a0b25 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Mon, 2 Dec 2024 11:09:21 +0530 Subject: [PATCH] Clean up remote BWE a bit. (#3225) * Clean up remote BWE a bit. - Had forgotten to start worker, fix that - ensure correct type of channel observer (probe OR non-probe) based on probe state. - introduce congested hangover state to see better state transitions. Does not really affect operation, but state transitions are clearer. * prevent 0 ticker --- pkg/sfu/bwe/remotebwe/remote_bwe.go | 107 ++++++++++++++++++++-------- 1 file changed, 77 insertions(+), 30 deletions(-) diff --git a/pkg/sfu/bwe/remotebwe/remote_bwe.go b/pkg/sfu/bwe/remotebwe/remote_bwe.go index c934db456..a868afa9e 100644 --- a/pkg/sfu/bwe/remotebwe/remote_bwe.go +++ b/pkg/sfu/bwe/remotebwe/remote_bwe.go @@ -68,6 +68,7 @@ type RemoteBWE struct { lastExpectedBandwidthUsage int64 committedChannelCapacity int64 + isInProbe bool channelObserver *channelObserver congestionState bwe.CongestionState @@ -83,7 +84,10 @@ func NewRemoteBWE(params RemoteBWEParams) *RemoteBWE { r := &RemoteBWE{ params: params, } - r.channelObserver = r.newChannelObserverNonProbe() + + go r.worker() + + r.Reset() return r } @@ -105,8 +109,21 @@ func (r *RemoteBWE) Reset() { r.lock.Lock() defer r.lock.Unlock() - r.channelObserver = r.newChannelObserverNonProbe() - r.updateCongestionState(bwe.CongestionStateNone, channelCongestionReasonNone) + r.lastReceivedEstimate = 0 + r.lastExpectedBandwidthUsage = 0 + r.committedChannelCapacity = 100_000_000 + + r.isInProbe = false + r.newChannelObserver() + + r.congestionState = bwe.CongestionStateNone + r.congestionStateSwitchedAt = mono.Now() + + // notify worker for ticker interval management based on state + select { + case r.wake <- struct{}{}: + default: + } } func (r *RemoteBWE) Stop() { @@ -165,6 +182,15 @@ func (r *RemoteBWE) congestionDetectionStateMachine() (bool, bwe.CongestionState // update state as this needs to reset switch time to wait for congestion min duration again update = true } + } else { + newState = bwe.CongestionStateCongestedHangover + } + + case bwe.CongestionStateCongestedHangover: + if trend == channelTrendCongesting { + if r.estimateAvailableChannelCapacity(reason) { + newState = bwe.CongestionStateCongested + } } else if time.Since(r.congestionStateSwitchedAt) >= r.params.Config.CongestedMinDuration { newState = bwe.CongestionStateNone } @@ -218,7 +244,7 @@ func (r *RemoteBWE) estimateAvailableChannelCapacity(reason channelCongestionRea r.committedChannelCapacity = estimateToCommit // reset to get new set of samples for next trend - r.channelObserver = r.newChannelObserverNonProbe() + r.newChannelObserver() return true } @@ -243,14 +269,25 @@ func (r *RemoteBWE) updateCongestionState(state bwe.CongestionState, reason chan r.congestionStateSwitchedAt = mono.Now() } -func (r *RemoteBWE) newChannelObserverNonProbe() *channelObserver { - return newChannelObserver( - channelObserverParams{ - Name: "non-probe", - Config: r.params.Config.ChannelObserverNonProbe, - }, - r.params.Logger, - ) +func (r *RemoteBWE) newChannelObserver() { + if r.isInProbe { + r.channelObserver = newChannelObserver( + channelObserverParams{ + Name: "probe", + Config: r.params.Config.ChannelObserverProbe, + }, + r.params.Logger, + ) + r.channelObserver.SeedEstimate(r.lastReceivedEstimate) + } else { + r.channelObserver = newChannelObserver( + channelObserverParams{ + Name: "non-probe", + Config: r.params.Config.ChannelObserverNonProbe, + }, + r.params.Logger, + ) + } } func (r *RemoteBWE) ProbeClusterStarting(pci ccutils.ProbeClusterInfo) { @@ -266,14 +303,8 @@ func (r *RemoteBWE) ProbeClusterStarting(pci ccutils.ProbeClusterInfo) { "channel", r.channelObserver, ) - r.channelObserver = newChannelObserver( - channelObserverParams{ - Name: "probe", - Config: r.params.Config.ChannelObserverProbe, - }, - r.params.Logger, - ) - r.channelObserver.SeedEstimate(r.lastReceivedEstimate) + r.isInProbe = true + r.newChannelObserver() } func (r *RemoteBWE) ProbeClusterDone(_pci ccutils.ProbeClusterInfo) (bool, int64) { @@ -282,7 +313,8 @@ func (r *RemoteBWE) ProbeClusterDone(_pci ccutils.ProbeClusterInfo) (bool, int64 // switch to a non-probe channel observer on probe end pco := r.channelObserver - r.channelObserver = r.newChannelObserverNonProbe() + r.isInProbe = false + r.newChannelObserver() r.params.Logger.Debugw( "remote bwe: probe done", @@ -304,21 +336,36 @@ func (r *RemoteBWE) ProbeClusterDone(_pci ccutils.ProbeClusterInfo) (bool, int64 return trend == channelTrendClearing, r.committedChannelCapacity } +func (r *RemoteBWE) getCheckInterval() time.Duration { + r.lock.RLock() + state := r.congestionState + r.lock.RUnlock() + + switch state { + case bwe.CongestionStateCongested: + if r.params.Config.PeriodicCheckIntervalCongested != 0 { + return r.params.Config.PeriodicCheckIntervalCongested + } + + return DefaultRemoteBWEConfig.PeriodicCheckIntervalCongested + + default: + if r.params.Config.PeriodicCheckInterval != 0 { + return r.params.Config.PeriodicCheckInterval + } + + return DefaultRemoteBWEConfig.PeriodicCheckInterval + } +} + func (r *RemoteBWE) worker() { - ticker := time.NewTicker(r.params.Config.PeriodicCheckInterval) + ticker := time.NewTicker(r.getCheckInterval()) defer ticker.Stop() for { select { case <-r.wake: - r.lock.RLock() - state := r.congestionState - r.lock.RUnlock() - if state == bwe.CongestionStateCongested { - ticker.Reset(r.params.Config.PeriodicCheckIntervalCongested) - } else { - ticker.Reset(r.params.Config.PeriodicCheckInterval) - } + ticker.Reset(r.getCheckInterval()) case <-ticker.C: r.lock.Lock()