diff --git a/pkg/rtc/mediatrack.go b/pkg/rtc/mediatrack.go index 3c167e0f2..be1436fce 100644 --- a/pkg/rtc/mediatrack.go +++ b/pkg/rtc/mediatrack.go @@ -95,6 +95,7 @@ type MediaTrackParams struct { EnableRTPStreamRestartDetection bool UpdateTrackInfoByVideoSizeChange bool ForceBackupCodecPolicySimulcast bool + MediaBatchIOEnabled bool OnSubscribedMaxQualityChange func( trackID livekit.TrackID, trackInfo *livekit.TrackInfo, @@ -272,6 +273,18 @@ func (t *MediaTrack) ToProto() *livekit.TrackInfo { return t.MediaTrackReceiver.TrackInfoClone() } +func (t *MediaTrack) getReceiverLBThreshold() int { + // With batch IO the send syscall moves to the BatchConn flush goroutine, so WriteRTP only + // translates the packet and enqueues it: ~3.3us per downtrack, against ~26us when the send runs + // inline. The last subscriber of 300 subscribers at a serial forwarder is enqueued ~1ms + // after the first. + if t.params.MediaBatchIOEnabled { + return 300 + } + + return 20 +} + // AddReceiver adds a new RTP receiver to the track, returns true when receiver represents a new codec // and if a receiver was added successfully func (t *MediaTrack) AddReceiver(receiver *webrtc.RTPReceiver, track sfu.TrackRemote, mid string) (bool, bool) { @@ -401,7 +414,7 @@ func (t *MediaTrack) AddReceiver(receiver *webrtc.RTPReceiver, track sfu.TrackRe t.params.VideoConfig.StreamTrackerManager, sfu.WithPliThrottleConfig(t.params.PLIThrottleConfig), sfu.WithAudioConfig(t.params.AudioConfig), - sfu.WithLoadBalanceThreshold(20), + sfu.WithLoadBalanceThreshold(t.getReceiverLBThreshold()), sfu.WithForwardStats(t.params.ForwardStats), sfu.WithEnableRTPStreamRestartDetection(t.params.EnableRTPStreamRestartDetection), ) diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index b0a663f27..a63529fb5 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -241,6 +241,7 @@ type ParticipantParams struct { MigrationWaitDuration time.Duration ExcludeIPv6LocalCandidates bool EnableWarp bool + MediaBatchIOEnabled bool } type ParticipantImpl struct { @@ -3573,6 +3574,7 @@ func (p *ParticipantImpl) addMediaTrack(signalCid string, ti *livekit.TrackInfo) EnableRTPStreamRestartDetection: p.params.EnableRTPStreamRestartDetection, UpdateTrackInfoByVideoSizeChange: p.params.UseOneShotSignallingMode, ForceBackupCodecPolicySimulcast: p.params.ForceBackupCodecPolicySimulcast, + MediaBatchIOEnabled: p.params.MediaBatchIOEnabled, OnSubscribedMaxQualityChange: p.onSubscribedMaxQualityChange, OnSubscribedAudioCodecChange: p.onSubscribedAudioCodecChange, }, ti) diff --git a/pkg/service/roommanager.go b/pkg/service/roommanager.go index 5aace9845..b667344f3 100644 --- a/pkg/service/roommanager.go +++ b/pkg/service/roommanager.go @@ -543,6 +543,9 @@ func (r *RoomManager) StartSession( EnableParticipantDataBlob: r.config.EnableParticipantDataBlob, EnableRTPStreamRestartDetection: r.config.RTC.EnableRTPStreamRestartDetection, EnableWarp: enableWarp, + MediaBatchIOEnabled: r.config.RTC.BatchIO.BatchSize > 0 && + (r.config.RTC.ICEPortRangeStart == 0 || r.config.RTC.ICEPortRangeEnd == 0) && + r.config.RTC.UDPPort.Valid(), }) if err != nil { prometheus.IncrementParticipantRtcCanceled(1, enableWarp, pi.Client.GetSdk())