diff --git a/README.md b/README.md index d1683918..cc23284d 100644 --- a/README.md +++ b/README.md @@ -1276,10 +1276,10 @@ rtsps://localhost:8322/mystream In some scenarios, when publishing or reading from the server with RTSP, frames can get corrupted. This can be caused by multiple reasons: -* the packet buffer of the server is too small and can't keep up with the stream throughput. A solution consists in increasing its size: +* the write queue of the server is too small and can't keep up with the stream throughput. A solution consists in increasing its size: ```yml - readBufferCount: 1024 + writeQueueSize: 1024 ``` * The stream throughput is too big and the stream can't be transmitted correctly with the UDP transport protocol. UDP is more performant, faster and more efficient than TCP, but doesn't have a retransmission mechanism, that is needed in case of streams that need a large bandwidth. A solution consists in switching to TCP: diff --git a/apidocs/openapi.yaml b/apidocs/openapi.yaml index da1f7a4f..43d6682c 100644 --- a/apidocs/openapi.yaml +++ b/apidocs/openapi.yaml @@ -31,7 +31,7 @@ components: type: string writeTimeout: type: string - readBufferCount: + writeQueueSize: type: integer udpMaxPayloadSize: type: integer diff --git a/internal/conf/conf.go b/internal/conf/conf.go index b81dd307..01ae0503 100644 --- a/internal/conf/conf.go +++ b/internal/conf/conf.go @@ -95,7 +95,8 @@ type Conf struct { LogFile string `json:"logFile"` ReadTimeout StringDuration `json:"readTimeout"` WriteTimeout StringDuration `json:"writeTimeout"` - ReadBufferCount int `json:"readBufferCount"` + ReadBufferCount int `json:"readBufferCount"` // deprecated + WriteQueueSize int `json:"writeQueueSize"` UDPMaxPayloadSize int `json:"udpMaxPayloadSize"` ExternalAuthenticationURL string `json:"externalAuthenticationURL"` API bool `json:"api"` @@ -218,8 +219,11 @@ func (conf Conf) Clone() *Conf { // Check checks the configuration for errors. func (conf *Conf) Check() error { // general - if (conf.ReadBufferCount & (conf.ReadBufferCount - 1)) != 0 { - return fmt.Errorf("'readBufferCount' must be a power of two") + if conf.ReadBufferCount != 0 { + conf.WriteQueueSize = conf.ReadBufferCount + } + if (conf.WriteQueueSize & (conf.WriteQueueSize - 1)) != 0 { + return fmt.Errorf("'writeQueueSize' must be a power of two") } if conf.UDPMaxPayloadSize > 1472 { return fmt.Errorf("'udpMaxPayloadSize' must be less than 1472") @@ -317,7 +321,7 @@ func (conf *Conf) UnmarshalJSON(b []byte) error { conf.LogFile = "mediamtx.log" conf.ReadTimeout = 10 * StringDuration(time.Second) conf.WriteTimeout = 10 * StringDuration(time.Second) - conf.ReadBufferCount = 512 + conf.WriteQueueSize = 512 conf.UDPMaxPayloadSize = 1472 conf.APIAddress = "127.0.0.1:9997" conf.MetricsAddress = "127.0.0.1:9998" diff --git a/internal/conf/conf_test.go b/internal/conf/conf_test.go index b623c396..ca939724 100644 --- a/internal/conf/conf_test.go +++ b/internal/conf/conf_test.go @@ -217,9 +217,9 @@ func TestConfErrors(t *testing.T) { "json: unknown field \"invalid\"", }, { - "invalid readBufferCount", - "readBufferCount: 1001\n", - "'readBufferCount' must be a power of two", + "invalid writeQueueSize", + "writeQueueSize: 1001\n", + "'writeQueueSize' must be a power of two", }, { "invalid udpMaxPayloadSize", diff --git a/internal/core/core.go b/internal/core/core.go index 6128bddc..175c9ada 100644 --- a/internal/core/core.go +++ b/internal/core/core.go @@ -246,7 +246,7 @@ func (p *Core) createResources(initial bool) error { p.conf.AuthMethods, p.conf.ReadTimeout, p.conf.WriteTimeout, - p.conf.ReadBufferCount, + p.conf.WriteQueueSize, p.conf.UDPMaxPayloadSize, p.conf.Paths, p.externalCmdPool, @@ -266,7 +266,7 @@ func (p *Core) createResources(initial bool) error { p.conf.AuthMethods, p.conf.ReadTimeout, p.conf.WriteTimeout, - p.conf.ReadBufferCount, + p.conf.WriteQueueSize, useUDP, useMulticast, p.conf.RTPAddress, @@ -301,7 +301,7 @@ func (p *Core) createResources(initial bool) error { p.conf.AuthMethods, p.conf.ReadTimeout, p.conf.WriteTimeout, - p.conf.ReadBufferCount, + p.conf.WriteQueueSize, false, false, "", @@ -335,7 +335,7 @@ func (p *Core) createResources(initial bool) error { p.conf.RTMPAddress, p.conf.ReadTimeout, p.conf.WriteTimeout, - p.conf.ReadBufferCount, + p.conf.WriteQueueSize, false, "", "", @@ -361,7 +361,7 @@ func (p *Core) createResources(initial bool) error { p.conf.RTMPSAddress, p.conf.ReadTimeout, p.conf.WriteTimeout, - p.conf.ReadBufferCount, + p.conf.WriteQueueSize, true, p.conf.RTMPServerCert, p.conf.RTMPServerKey, @@ -397,7 +397,7 @@ func (p *Core) createResources(initial bool) error { p.conf.HLSTrustedProxies, p.conf.HLSDirectory, p.conf.ReadTimeout, - p.conf.ReadBufferCount, + p.conf.WriteQueueSize, p.pathManager, p.metrics, p, @@ -419,7 +419,7 @@ func (p *Core) createResources(initial bool) error { p.conf.WebRTCTrustedProxies, p.conf.WebRTCICEServers2, p.conf.ReadTimeout, - p.conf.ReadBufferCount, + p.conf.WriteQueueSize, p.conf.WebRTCICEHostNAT1To1IPs, p.conf.WebRTCICEUDPMuxAddress, p.conf.WebRTCICETCPMuxAddress, @@ -439,7 +439,7 @@ func (p *Core) createResources(initial bool) error { p.conf.SRTAddress, p.conf.ReadTimeout, p.conf.WriteTimeout, - p.conf.ReadBufferCount, + p.conf.WriteQueueSize, p.conf.UDPMaxPayloadSize, p.externalCmdPool, p.pathManager, @@ -504,7 +504,7 @@ func (p *Core) closeResources(newConf *conf.Conf, calledByAPI bool) { !reflect.DeepEqual(newConf.AuthMethods, p.conf.AuthMethods) || newConf.ReadTimeout != p.conf.ReadTimeout || newConf.WriteTimeout != p.conf.WriteTimeout || - newConf.ReadBufferCount != p.conf.ReadBufferCount || + newConf.WriteQueueSize != p.conf.WriteQueueSize || newConf.UDPMaxPayloadSize != p.conf.UDPMaxPayloadSize || closeMetrics if !closePathManager && !reflect.DeepEqual(newConf.Paths, p.conf.Paths) { @@ -518,7 +518,7 @@ func (p *Core) closeResources(newConf *conf.Conf, calledByAPI bool) { !reflect.DeepEqual(newConf.AuthMethods, p.conf.AuthMethods) || newConf.ReadTimeout != p.conf.ReadTimeout || newConf.WriteTimeout != p.conf.WriteTimeout || - newConf.ReadBufferCount != p.conf.ReadBufferCount || + newConf.WriteQueueSize != p.conf.WriteQueueSize || !reflect.DeepEqual(newConf.Protocols, p.conf.Protocols) || newConf.RTPAddress != p.conf.RTPAddress || newConf.RTCPAddress != p.conf.RTCPAddress || @@ -539,7 +539,7 @@ func (p *Core) closeResources(newConf *conf.Conf, calledByAPI bool) { !reflect.DeepEqual(newConf.AuthMethods, p.conf.AuthMethods) || newConf.ReadTimeout != p.conf.ReadTimeout || newConf.WriteTimeout != p.conf.WriteTimeout || - newConf.ReadBufferCount != p.conf.ReadBufferCount || + newConf.WriteQueueSize != p.conf.WriteQueueSize || newConf.ServerCert != p.conf.ServerCert || newConf.ServerKey != p.conf.ServerKey || newConf.RTSPAddress != p.conf.RTSPAddress || @@ -555,7 +555,7 @@ func (p *Core) closeResources(newConf *conf.Conf, calledByAPI bool) { newConf.RTMPAddress != p.conf.RTMPAddress || newConf.ReadTimeout != p.conf.ReadTimeout || newConf.WriteTimeout != p.conf.WriteTimeout || - newConf.ReadBufferCount != p.conf.ReadBufferCount || + newConf.WriteQueueSize != p.conf.WriteQueueSize || newConf.RTSPAddress != p.conf.RTSPAddress || newConf.RunOnConnect != p.conf.RunOnConnect || newConf.RunOnConnectRestart != p.conf.RunOnConnectRestart || @@ -568,7 +568,7 @@ func (p *Core) closeResources(newConf *conf.Conf, calledByAPI bool) { newConf.RTMPSAddress != p.conf.RTMPSAddress || newConf.ReadTimeout != p.conf.ReadTimeout || newConf.WriteTimeout != p.conf.WriteTimeout || - newConf.ReadBufferCount != p.conf.ReadBufferCount || + newConf.WriteQueueSize != p.conf.WriteQueueSize || newConf.RTMPServerCert != p.conf.RTMPServerCert || newConf.RTMPServerKey != p.conf.RTMPServerKey || newConf.RTSPAddress != p.conf.RTSPAddress || @@ -594,7 +594,7 @@ func (p *Core) closeResources(newConf *conf.Conf, calledByAPI bool) { !reflect.DeepEqual(newConf.HLSTrustedProxies, p.conf.HLSTrustedProxies) || newConf.HLSDirectory != p.conf.HLSDirectory || newConf.ReadTimeout != p.conf.ReadTimeout || - newConf.ReadBufferCount != p.conf.ReadBufferCount || + newConf.WriteQueueSize != p.conf.WriteQueueSize || closePathManager || closeMetrics @@ -608,7 +608,7 @@ func (p *Core) closeResources(newConf *conf.Conf, calledByAPI bool) { !reflect.DeepEqual(newConf.WebRTCTrustedProxies, p.conf.WebRTCTrustedProxies) || !reflect.DeepEqual(newConf.WebRTCICEServers2, p.conf.WebRTCICEServers2) || newConf.ReadTimeout != p.conf.ReadTimeout || - newConf.ReadBufferCount != p.conf.ReadBufferCount || + newConf.WriteQueueSize != p.conf.WriteQueueSize || !reflect.DeepEqual(newConf.WebRTCICEHostNAT1To1IPs, p.conf.WebRTCICEHostNAT1To1IPs) || newConf.WebRTCICEUDPMuxAddress != p.conf.WebRTCICEUDPMuxAddress || newConf.WebRTCICETCPMuxAddress != p.conf.WebRTCICETCPMuxAddress || @@ -620,7 +620,7 @@ func (p *Core) closeResources(newConf *conf.Conf, calledByAPI bool) { newConf.SRTAddress != p.conf.SRTAddress || newConf.ReadTimeout != p.conf.ReadTimeout || newConf.WriteTimeout != p.conf.WriteTimeout || - newConf.ReadBufferCount != p.conf.ReadBufferCount || + newConf.WriteQueueSize != p.conf.WriteQueueSize || newConf.UDPMaxPayloadSize != p.conf.UDPMaxPayloadSize || closePathManager diff --git a/internal/core/hls_manager.go b/internal/core/hls_manager.go index a129387e..0f687b1a 100644 --- a/internal/core/hls_manager.go +++ b/internal/core/hls_manager.go @@ -42,7 +42,7 @@ type hlsManager struct { partDuration conf.StringDuration segmentMaxSize conf.StringSize directory string - readBufferCount int + writeQueueSize int pathManager *pathManager metrics *metrics parent hlsManagerParent @@ -78,7 +78,7 @@ func newHLSManager( trustedProxies conf.IPsOrCIDRs, directory string, readTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, pathManager *pathManager, metrics *metrics, parent hlsManagerParent, @@ -94,7 +94,7 @@ func newHLSManager( partDuration: partDuration, segmentMaxSize: segmentMaxSize, directory: directory, - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, pathManager: pathManager, parent: parent, metrics: metrics, @@ -241,7 +241,7 @@ func (m *hlsManager) createMuxer(pathName string, remoteAddr string) *hlsMuxer { m.partDuration, m.segmentMaxSize, m.directory, - m.readBufferCount, + m.writeQueueSize, &m.wg, pathName, m.pathManager, diff --git a/internal/core/hls_muxer.go b/internal/core/hls_muxer.go index dcaf2bd4..f1233b77 100644 --- a/internal/core/hls_muxer.go +++ b/internal/core/hls_muxer.go @@ -66,7 +66,7 @@ type hlsMuxer struct { partDuration conf.StringDuration segmentMaxSize conf.StringSize directory string - readBufferCount int + writeQueueSize int wg *sync.WaitGroup pathName string pathManager *pathManager @@ -96,7 +96,7 @@ func newHLSMuxer( partDuration conf.StringDuration, segmentMaxSize conf.StringSize, directory string, - readBufferCount int, + writeQueueSize int, wg *sync.WaitGroup, pathName string, pathManager *pathManager, @@ -113,7 +113,7 @@ func newHLSMuxer( partDuration: partDuration, segmentMaxSize: segmentMaxSize, directory: directory, - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, wg: wg, pathName: pathName, pathManager: pathManager, @@ -254,7 +254,7 @@ func (m *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{}) defer m.path.removeReader(pathRemoveReaderReq{author: m}) - m.ringBuffer, _ = ringbuffer.New(uint64(m.readBufferCount)) + m.ringBuffer, _ = ringbuffer.New(uint64(m.writeQueueSize)) var medias media.Medias diff --git a/internal/core/path.go b/internal/core/path.go index 104b6f8f..3f12b453 100644 --- a/internal/core/path.go +++ b/internal/core/path.go @@ -174,7 +174,7 @@ type path struct { rtspAddress string readTimeout conf.StringDuration writeTimeout conf.StringDuration - readBufferCount int + writeQueueSize int udpMaxPayloadSize int confName string conf *conf.PathConf @@ -225,7 +225,7 @@ func newPath( rtspAddress string, readTimeout conf.StringDuration, writeTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, udpMaxPayloadSize int, confName string, cnf *conf.PathConf, @@ -241,7 +241,7 @@ func newPath( rtspAddress: rtspAddress, readTimeout: readTimeout, writeTimeout: writeTimeout, - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, udpMaxPayloadSize: udpMaxPayloadSize, confName: confName, conf: cnf, @@ -310,7 +310,7 @@ func (pa *path) run() { pa.conf, pa.readTimeout, pa.writeTimeout, - pa.readBufferCount, + pa.writeQueueSize, pa) if !pa.conf.SourceOnDemand { diff --git a/internal/core/path_manager.go b/internal/core/path_manager.go index da1e13f1..1cf39ce1 100644 --- a/internal/core/path_manager.go +++ b/internal/core/path_manager.go @@ -69,7 +69,7 @@ type pathManager struct { authMethods conf.AuthMethods readTimeout conf.StringDuration writeTimeout conf.StringDuration - readBufferCount int + writeQueueSize int udpMaxPayloadSize int pathConfs map[string]*conf.PathConf externalCmdPool *externalcmd.Pool @@ -103,7 +103,7 @@ func newPathManager( authMethods conf.AuthMethods, readTimeout conf.StringDuration, writeTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, udpMaxPayloadSize int, pathConfs map[string]*conf.PathConf, externalCmdPool *externalcmd.Pool, @@ -118,7 +118,7 @@ func newPathManager( authMethods: authMethods, readTimeout: readTimeout, writeTimeout: writeTimeout, - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, udpMaxPayloadSize: udpMaxPayloadSize, pathConfs: pathConfs, externalCmdPool: externalCmdPool, @@ -352,7 +352,7 @@ func (pm *pathManager) createPath( pm.rtspAddress, pm.readTimeout, pm.writeTimeout, - pm.readBufferCount, + pm.writeQueueSize, pm.udpMaxPayloadSize, pathConfName, pathConf, diff --git a/internal/core/rtmp_conn.go b/internal/core/rtmp_conn.go index 220257ae..c214421b 100644 --- a/internal/core/rtmp_conn.go +++ b/internal/core/rtmp_conn.go @@ -60,7 +60,7 @@ type rtmpConn struct { rtspAddress string readTimeout conf.StringDuration writeTimeout conf.StringDuration - readBufferCount int + writeQueueSize int runOnConnect string runOnConnectRestart bool wg *sync.WaitGroup @@ -85,7 +85,7 @@ func newRTMPConn( rtspAddress string, readTimeout conf.StringDuration, writeTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, runOnConnect string, runOnConnectRestart bool, wg *sync.WaitGroup, @@ -101,7 +101,7 @@ func newRTMPConn( rtspAddress: rtspAddress, readTimeout: readTimeout, writeTimeout: writeTimeout, - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, runOnConnect: runOnConnect, runOnConnectRestart: runOnConnectRestart, wg: wg, @@ -241,7 +241,7 @@ func (c *rtmpConn) runRead(conn *rtmp.Conn, u *url.URL) error { c.pathName = pathName c.mutex.Unlock() - ringBuffer, _ := ringbuffer.New(uint64(c.readBufferCount)) + ringBuffer, _ := ringbuffer.New(uint64(c.writeQueueSize)) go func() { <-c.ctx.Done() ringBuffer.Close() diff --git a/internal/core/rtmp_server.go b/internal/core/rtmp_server.go index 35499c82..fbb66fc4 100644 --- a/internal/core/rtmp_server.go +++ b/internal/core/rtmp_server.go @@ -50,7 +50,7 @@ type rtmpServerParent interface { type rtmpServer struct { readTimeout conf.StringDuration writeTimeout conf.StringDuration - readBufferCount int + writeQueueSize int isTLS bool rtspAddress string runOnConnect string @@ -79,7 +79,7 @@ func newRTMPServer( address string, readTimeout conf.StringDuration, writeTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, isTLS bool, serverCert string, serverKey string, @@ -113,7 +113,7 @@ func newRTMPServer( s := &rtmpServer{ readTimeout: readTimeout, writeTimeout: writeTimeout, - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, rtspAddress: rtspAddress, runOnConnect: runOnConnect, runOnConnectRestart: runOnConnectRestart, @@ -185,7 +185,7 @@ outer: s.rtspAddress, s.readTimeout, s.writeTimeout, - s.readBufferCount, + s.writeQueueSize, s.runOnConnect, s.runOnConnectRestart, &s.wg, diff --git a/internal/core/rtsp_server.go b/internal/core/rtsp_server.go index 19eec72e..f90f4494 100644 --- a/internal/core/rtsp_server.go +++ b/internal/core/rtsp_server.go @@ -67,7 +67,7 @@ func newRTSPServer( authMethods []headers.AuthMethod, readTimeout conf.StringDuration, writeTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, useUDP bool, useMulticast bool, rtpAddress string, @@ -111,8 +111,8 @@ func newRTSPServer( Handler: s, ReadTimeout: time.Duration(readTimeout), WriteTimeout: time.Duration(writeTimeout), - ReadBufferCount: readBufferCount, - WriteBufferCount: readBufferCount, + ReadBufferCount: writeQueueSize, + WriteBufferCount: writeQueueSize, RTSPAddress: address, } diff --git a/internal/core/rtsp_source.go b/internal/core/rtsp_source.go index c2bce2fa..bf625500 100644 --- a/internal/core/rtsp_source.go +++ b/internal/core/rtsp_source.go @@ -66,23 +66,23 @@ type rtspSourceParent interface { } type rtspSource struct { - readTimeout conf.StringDuration - writeTimeout conf.StringDuration - readBufferCount int - parent rtspSourceParent + readTimeout conf.StringDuration + writeTimeout conf.StringDuration + writeQueueSize int + parent rtspSourceParent } func newRTSPSource( readTimeout conf.StringDuration, writeTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, parent rtspSourceParent, ) *rtspSource { return &rtspSource{ - readTimeout: readTimeout, - writeTimeout: writeTimeout, - readBufferCount: readBufferCount, - parent: parent, + readTimeout: readTimeout, + writeTimeout: writeTimeout, + writeQueueSize: writeQueueSize, + parent: parent, } } @@ -95,12 +95,13 @@ func (s *rtspSource) run(ctx context.Context, cnf *conf.PathConf, reloadConf cha s.Log(logger.Debug, "connecting") c := &gortsplib.Client{ - Transport: cnf.SourceProtocol.Transport, - TLSConfig: tlsConfigForFingerprint(cnf.SourceFingerprint), - ReadTimeout: time.Duration(s.readTimeout), - WriteTimeout: time.Duration(s.writeTimeout), - ReadBufferCount: s.readBufferCount, - AnyPortEnable: cnf.SourceAnyPortEnable, + Transport: cnf.SourceProtocol.Transport, + TLSConfig: tlsConfigForFingerprint(cnf.SourceFingerprint), + ReadTimeout: time.Duration(s.readTimeout), + WriteTimeout: time.Duration(s.writeTimeout), + ReadBufferCount: s.writeQueueSize, + WriteBufferCount: s.writeQueueSize, + AnyPortEnable: cnf.SourceAnyPortEnable, OnRequest: func(req *base.Request) { s.Log(logger.Debug, "c->s %v", req) }, diff --git a/internal/core/source_static.go b/internal/core/source_static.go index 8ce92af4..560350d7 100644 --- a/internal/core/source_static.go +++ b/internal/core/source_static.go @@ -49,7 +49,7 @@ func newSourceStatic( cnf *conf.PathConf, readTimeout conf.StringDuration, writeTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, parent sourceStaticParent, ) *sourceStatic { s := &sourceStatic{ @@ -66,7 +66,7 @@ func newSourceStatic( s.impl = newRTSPSource( readTimeout, writeTimeout, - readBufferCount, + writeQueueSize, s) case strings.HasPrefix(cnf.Source, "rtmp://") || diff --git a/internal/core/srt_conn.go b/internal/core/srt_conn.go index 4e8c1ee1..9dcd54eb 100644 --- a/internal/core/srt_conn.go +++ b/internal/core/srt_conn.go @@ -50,7 +50,7 @@ type srtConnParent interface { type srtConn struct { readTimeout conf.StringDuration writeTimeout conf.StringDuration - readBufferCount int + writeQueueSize int udpMaxPayloadSize int connReq srt.ConnRequest wg *sync.WaitGroup @@ -75,7 +75,7 @@ func newSRTConn( parentCtx context.Context, readTimeout conf.StringDuration, writeTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, udpMaxPayloadSize int, connReq srt.ConnRequest, wg *sync.WaitGroup, @@ -88,7 +88,7 @@ func newSRTConn( c := &srtConn{ readTimeout: readTimeout, writeTimeout: writeTimeout, - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, udpMaxPayloadSize: udpMaxPayloadSize, connReq: connReq, wg: wg, @@ -416,7 +416,7 @@ func (c *srtConn) runRead(req srtNewConnReq, pathName string, user string, pass c.conn = sconn c.mutex.Unlock() - ringBuffer, _ := ringbuffer.New(uint64(c.readBufferCount)) + ringBuffer, _ := ringbuffer.New(uint64(c.writeQueueSize)) go func() { <-c.ctx.Done() ringBuffer.Close() diff --git a/internal/core/srt_server.go b/internal/core/srt_server.go index 67edfb7a..ce10860b 100644 --- a/internal/core/srt_server.go +++ b/internal/core/srt_server.go @@ -59,7 +59,7 @@ type srtServerParent interface { type srtServer struct { readTimeout conf.StringDuration writeTimeout conf.StringDuration - readBufferCount int + writeQueueSize int udpMaxPayloadSize int externalCmdPool *externalcmd.Pool pathManager *pathManager @@ -84,7 +84,7 @@ func newSRTServer( address string, readTimeout conf.StringDuration, writeTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, udpMaxPayloadSize int, externalCmdPool *externalcmd.Pool, pathManager *pathManager, @@ -104,7 +104,7 @@ func newSRTServer( s := &srtServer{ readTimeout: readTimeout, writeTimeout: writeTimeout, - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, udpMaxPayloadSize: udpMaxPayloadSize, externalCmdPool: externalCmdPool, pathManager: pathManager, @@ -161,7 +161,7 @@ outer: s.ctx, s.readTimeout, s.writeTimeout, - s.readBufferCount, + s.writeQueueSize, s.udpMaxPayloadSize, req.connReq, &s.wg, diff --git a/internal/core/webrtc_manager.go b/internal/core/webrtc_manager.go index 7ebcc112..a9166233 100644 --- a/internal/core/webrtc_manager.go +++ b/internal/core/webrtc_manager.go @@ -276,13 +276,13 @@ type webRTCManagerParent interface { } type webRTCManager struct { - allowOrigin string - trustedProxies conf.IPsOrCIDRs - iceServers []conf.WebRTCICEServer - readBufferCount int - pathManager *pathManager - metrics *metrics - parent webRTCManagerParent + allowOrigin string + trustedProxies conf.IPsOrCIDRs + iceServers []conf.WebRTCICEServer + writeQueueSize int + pathManager *pathManager + metrics *metrics + parent webRTCManagerParent ctx context.Context ctxCancel func() @@ -314,7 +314,7 @@ func newWebRTCManager( trustedProxies conf.IPsOrCIDRs, iceServers []conf.WebRTCICEServer, readTimeout conf.StringDuration, - readBufferCount int, + writeQueueSize int, iceHostNAT1To1IPs []string, iceUDPMuxAddress string, iceTCPMuxAddress string, @@ -328,7 +328,7 @@ func newWebRTCManager( allowOrigin: allowOrigin, trustedProxies: trustedProxies, iceServers: iceServers, - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, pathManager: pathManager, metrics: metrics, parent: parent, @@ -436,7 +436,7 @@ outer: case req := <-m.chNewSession: sx := newWebRTCSession( m.ctx, - m.readBufferCount, + m.writeQueueSize, m.api, req, &wg, diff --git a/internal/core/webrtc_session.go b/internal/core/webrtc_session.go index 34d725ce..c7b68085 100644 --- a/internal/core/webrtc_session.go +++ b/internal/core/webrtc_session.go @@ -175,12 +175,12 @@ type webRTCSessionPathManager interface { } type webRTCSession struct { - readBufferCount int - api *webrtc.API - req webRTCNewSessionReq - wg *sync.WaitGroup - pathManager webRTCSessionPathManager - parent *webRTCManager + writeQueueSize int + api *webrtc.API + req webRTCNewSessionReq + wg *sync.WaitGroup + pathManager webRTCSessionPathManager + parent *webRTCManager ctx context.Context ctxCancel func() @@ -196,7 +196,7 @@ type webRTCSession struct { func newWebRTCSession( parentCtx context.Context, - readBufferCount int, + writeQueueSize int, api *webrtc.API, req webRTCNewSessionReq, wg *sync.WaitGroup, @@ -206,7 +206,7 @@ func newWebRTCSession( ctx, ctxCancel := context.WithCancel(parentCtx) s := &webRTCSession{ - readBufferCount: readBufferCount, + writeQueueSize: writeQueueSize, api: api, req: req, wg: wg, @@ -509,7 +509,7 @@ func (s *webRTCSession) runRead() (int, error) { s.pc = pc s.mutex.Unlock() - ringBuffer, _ := ringbuffer.New(uint64(s.readBufferCount)) + ringBuffer, _ := ringbuffer.New(uint64(s.writeQueueSize)) defer ringBuffer.Close() writeError := make(chan error) diff --git a/mediamtx.yml b/mediamtx.yml index 730d7f4e..91c98422 100644 --- a/mediamtx.yml +++ b/mediamtx.yml @@ -13,9 +13,9 @@ logFile: mediamtx.log readTimeout: 10s # Timeout of write operations. writeTimeout: 10s -# Number of read buffers. +# Size of the queue of outgoing packets. # A higher value allows to increase throughput, a lower value allows to save RAM. -readBufferCount: 512 +writeQueueSize: 512 # Maximum size of outgoing UDP packets. # This can be decreased to avoid fragmentation on networks with a low UDP MTU. udpMaxPayloadSize: 1472