From 534ea4d0c60627ddd4e01992c1331c1c8ce10370 Mon Sep 17 00:00:00 2001 From: Alessandro Ros Date: Tue, 29 Jul 2025 10:43:52 +0200 Subject: [PATCH] api: add additional WebRTC statistics (#4795) rtpPacketsReceived, rtpPacketsSent, rtpPacketsLost, rtpPacketsJitter, rtcpPacketsReceived, rtcpPacketsSent --- apidocs/openapi.yaml | 18 +++ internal/core/api_test.go | 6 + internal/defs/api.go | 6 + internal/protocols/webrtc/incoming_track.go | 11 +- internal/protocols/webrtc/outgoing_track.go | 12 +- internal/protocols/webrtc/peer_connection.go | 118 +++++++++++++----- internal/protocols/webrtc/stats.go | 13 ++ .../protocols/webrtc/stats_interceptor.go | 70 +++++++++++ internal/servers/webrtc/session.go | 20 +-- 9 files changed, 231 insertions(+), 43 deletions(-) create mode 100644 internal/protocols/webrtc/stats.go create mode 100644 internal/protocols/webrtc/stats_interceptor.go diff --git a/apidocs/openapi.yaml b/apidocs/openapi.yaml index df6298eb..21e01ced 100644 --- a/apidocs/openapi.yaml +++ b/apidocs/openapi.yaml @@ -1075,6 +1075,24 @@ components: bytesSent: type: integer format: int64 + rtpPacketsReceived: + type: integer + format: int64 + rtpPacketsSent: + type: integer + format: int64 + rtpPacketsLost: + type: integer + format: int64 + rtpPacketsJitter: + type: number + format: float64 + rtcpPacketsReceived: + type: integer + format: int64 + rtcpPacketsSent: + type: integer + format: int64 WebRTCSessionList: type: object diff --git a/internal/core/api_test.go b/internal/core/api_test.go index 685e9918..45161a5f 100644 --- a/internal/core/api_test.go +++ b/internal/core/api_test.go @@ -761,6 +761,12 @@ func TestAPIProtocolListGet(t *testing.T) { "remoteAddr": out1.(map[string]interface{})["items"].([]interface{})[0].(map[string]interface{})["remoteAddr"], "remoteCandidate": out1.(map[string]interface{})["items"].([]interface{})[0].(map[string]interface{})["remoteCandidate"], "state": "read", + "rtcpPacketsReceived": float64(0), + "rtcpPacketsSent": float64(2), + "rtpPacketsJitter": float64(0), + "rtpPacketsLost": float64(0), + "rtpPacketsReceived": float64(0), + "rtpPacketsSent": float64(1), }, }, }, out1) diff --git a/internal/defs/api.go b/internal/defs/api.go index 05830230..86aca7e8 100644 --- a/internal/defs/api.go +++ b/internal/defs/api.go @@ -357,6 +357,12 @@ type APIWebRTCSession struct { Query string `json:"query"` BytesReceived uint64 `json:"bytesReceived"` BytesSent uint64 `json:"bytesSent"` + RTPPacketsReceived uint64 `json:"rtpPacketsReceived"` + RTPPacketsSent uint64 `json:"rtpPacketsSent"` + RTPPacketsLost uint64 `json:"rtpPacketsLost"` + RTPPacketsJitter float64 `json:"rtpPacketsJitter"` + RTCPPacketsReceived uint64 `json:"rtcpPacketsReceived"` + RTCPPacketsSent uint64 `json:"rtcpPacketsSent"` } // APIWebRTCSessionList is a list of WebRTC sessions. diff --git a/internal/protocols/webrtc/incoming_track.go b/internal/protocols/webrtc/incoming_track.go index 71b3ca47..d391fb5f 100644 --- a/internal/protocols/webrtc/incoming_track.go +++ b/internal/protocols/webrtc/incoming_track.go @@ -1,6 +1,7 @@ package webrtc import ( + "sync/atomic" "time" "github.com/bluenviron/gortsplib/v4/pkg/rtcpreceiver" @@ -244,6 +245,8 @@ type IncomingTrack struct { receiver *webrtc.RTPReceiver writeRTCP func([]rtcp.Packet) error log logger.Writer + rtpPacketsReceived *uint64 + rtpPacketsLost *uint64 packetsLost *counterdumper.CounterDumper rtcpReceiver *rtcpreceiver.RTCPReceiver @@ -295,7 +298,8 @@ func (t *IncomingTrack) start() { panic(err) } - // incoming RTCP packets must always be read to make interceptors work + // read incoming RTCP packets. + // incoming RTCP packets must always be read to make interceptors work. go func() { buf := make([]byte, 1500) for { @@ -336,7 +340,7 @@ func (t *IncomingTrack) start() { }() } - // read incoming RTP packets + // read incoming RTP packets. go func() { reorderer := &rtpreorderer.Reorderer{} reorderer.Initialize() @@ -349,10 +353,13 @@ func (t *IncomingTrack) start() { packets, lost := reorderer.Process(pkt) if lost != 0 { + atomic.AddUint64(t.rtpPacketsLost, uint64(lost)) t.packetsLost.Add(uint64(lost)) // do not return } + atomic.AddUint64(t.rtpPacketsReceived, uint64(len(packets))) + err2 = t.rtcpReceiver.ProcessPacket(pkt, time.Now(), true) if err2 != nil { t.log.Log(logger.Warn, err2.Error()) diff --git a/internal/protocols/webrtc/outgoing_track.go b/internal/protocols/webrtc/outgoing_track.go index 354df894..42aeb8fb 100644 --- a/internal/protocols/webrtc/outgoing_track.go +++ b/internal/protocols/webrtc/outgoing_track.go @@ -2,6 +2,7 @@ package webrtc import ( "strings" + "sync/atomic" "time" "github.com/bluenviron/gortsplib/v4/pkg/rtcpsender" @@ -23,9 +24,10 @@ var multichannelOpusSDP = map[int]string{ type OutgoingTrack struct { Caps webrtc.RTPCodecCapability - track *webrtc.TrackLocalStaticRTP - ssrc uint32 - rtcpSender *rtcpsender.RTCPSender + track *webrtc.TrackLocalStaticRTP + ssrc uint32 + rtcpSender *rtcpsender.RTCPSender + rtpPacketsSent *uint64 } func (t *OutgoingTrack) isVideo() bool { @@ -67,7 +69,7 @@ func (t *OutgoingTrack) setup(p *PeerConnection) error { } t.rtcpSender.Initialize() - p.wr.GetSenders() + t.rtpPacketsSent = p.rtpPacketsSent // incoming RTCP packets must always be read to make interceptors work go func() { @@ -106,5 +108,7 @@ func (t *OutgoingTrack) WriteRTPWithNTP(pkt *rtp.Packet, ntp time.Time) error { t.rtcpSender.ProcessPacket(pkt, ntp, true) + atomic.AddUint64(t.rtpPacketsSent, 1) + return t.track.WriteRTP(pkt) } diff --git a/internal/protocols/webrtc/peer_connection.go b/internal/protocols/webrtc/peer_connection.go index c295235c..391e27b9 100644 --- a/internal/protocols/webrtc/peer_connection.go +++ b/internal/protocols/webrtc/peer_connection.go @@ -7,6 +7,7 @@ import ( "fmt" "strconv" "sync" + "sync/atomic" "time" "github.com/pion/ice/v4" @@ -31,17 +32,33 @@ func stringInSlice(a string, list []string) bool { return false } -// skip ConfigureRTCPReports -func registerInterceptors(mediaEngine *webrtc.MediaEngine, interceptorRegistry *interceptor.Registry) error { - if err := webrtc.ConfigureNack(mediaEngine, interceptorRegistry); err != nil { +// * skip ConfigureRTCPReports +// * add statsInterceptor +func registerInterceptors( + mediaEngine *webrtc.MediaEngine, + interceptorRegistry *interceptor.Registry, + onStatsInterceptor func(s *statsInterceptor), +) error { + err := webrtc.ConfigureNack(mediaEngine, interceptorRegistry) + if err != nil { return err } - if err := webrtc.ConfigureSimulcastExtensionHeaders(mediaEngine); err != nil { + err = webrtc.ConfigureSimulcastExtensionHeaders(mediaEngine) + if err != nil { return err } - return webrtc.ConfigureTWCCSender(mediaEngine, interceptorRegistry) + err = webrtc.ConfigureTWCCSender(mediaEngine, interceptorRegistry) + if err != nil { + return err + } + + interceptorRegistry.Add(&statsInterceptorFactory{ + onCreate: onStatsInterceptor, + }) + + return nil } // TracksAreValid checks whether tracks in the SDP are valid @@ -97,17 +114,24 @@ type PeerConnection struct { UseAbsoluteTimestamp bool Log logger.Writer - wr *webrtc.PeerConnection - stateChangeMutex sync.Mutex - newLocalCandidate chan *webrtc.ICECandidateInit - connected chan struct{} - failed chan struct{} - closed chan struct{} - gatheringDone chan struct{} - incomingTrack chan trackRecvPair - ctx context.Context - ctxCancel context.CancelFunc - incomingTracks []*IncomingTrack + wr *webrtc.PeerConnection + stateChangeMutex sync.Mutex + newLocalCandidate chan *webrtc.ICECandidateInit + connected chan struct{} + failed chan struct{} + closed chan struct{} + gatheringDone chan struct{} + incomingTrack chan trackRecvPair + ctx context.Context + ctxCancel context.CancelFunc + incomingTracks []*IncomingTrack + startedReading *int64 + rtpPacketsReceived *uint64 + rtpPacketsSent *uint64 + rtpPacketsLost *uint64 + statsInterceptor *statsInterceptor + // rtcpPacketsReceived *uint64 + // rtcpPacketsSent *uint64 } // Start starts the peer connection. @@ -216,7 +240,13 @@ func (co *PeerConnection) Start() error { interceptorRegistry := &interceptor.Registry{} - err := registerInterceptors(mediaEngine, interceptorRegistry) + err := registerInterceptors( + mediaEngine, + interceptorRegistry, + func(s *statsInterceptor) { + co.statsInterceptor = s + }, + ) if err != nil { return err } @@ -242,6 +272,11 @@ func (co *PeerConnection) Start() error { co.ctx, co.ctxCancel = context.WithCancel(context.Background()) + co.startedReading = new(int64) + co.rtpPacketsReceived = new(uint64) + co.rtpPacketsSent = new(uint64) + co.rtpPacketsLost = new(uint64) + if co.Publish { for _, tr := range co.OutgoingTracks { err = tr.setup(co) @@ -473,6 +508,8 @@ func (co *PeerConnection) GatherIncomingTracks(ctx context.Context) error { receiver: pair.receiver, writeRTCP: co.wr.WriteRTCP, log: co.Log, + rtpPacketsReceived: co.rtpPacketsReceived, + rtpPacketsLost: co.rtpPacketsLost, } t.initialize() co.incomingTracks = append(co.incomingTracks, t) @@ -542,6 +579,7 @@ func (co *PeerConnection) StartReading() { for _, track := range co.incomingTracks { track.start() } + atomic.StoreInt64(co.startedReading, 1) } // RemoteCandidate returns the remote candidate. @@ -566,26 +604,48 @@ func (co *PeerConnection) RemoteCandidate() string { return "" } -// BytesReceived returns received bytes. -func (co *PeerConnection) BytesReceived() uint64 { - for _, stats := range co.wr.GetStats() { +func bytesStats(wr *webrtc.PeerConnection) (uint64, uint64) { + for _, stats := range wr.GetStats() { if tstats, ok := stats.(webrtc.TransportStats); ok { if tstats.ID == "iceTransport" { - return tstats.BytesReceived + return tstats.BytesReceived, tstats.BytesSent } } } - return 0 + return 0, 0 } -// BytesSent returns sent bytes. -func (co *PeerConnection) BytesSent() uint64 { - for _, stats := range co.wr.GetStats() { - if tstats, ok := stats.(webrtc.TransportStats); ok { - if tstats.ID == "iceTransport" { - return tstats.BytesSent +// Stats returns statistics. +func (co *PeerConnection) Stats() *Stats { + bytesReceived, bytesSent := bytesStats(co.wr) + + v := float64(0) + n := float64(0) + + if atomic.LoadInt64(co.startedReading) == 1 { + for _, tr := range co.incomingTracks { + if recvStats := tr.rtcpReceiver.Stats(); recvStats != nil { + v += recvStats.Jitter + n++ } } } - return 0 + + var rtpPacketsJitter float64 + if n != 0 { + rtpPacketsJitter = v / n + } else { + rtpPacketsJitter = 0 + } + + return &Stats{ + BytesReceived: bytesReceived, + BytesSent: bytesSent, + RTPPacketsReceived: atomic.LoadUint64(co.rtpPacketsReceived), + RTPPacketsSent: atomic.LoadUint64(co.rtpPacketsSent), + RTPPacketsLost: atomic.LoadUint64(co.rtpPacketsLost), + RTPPacketsJitter: rtpPacketsJitter, + RTCPPacketsReceived: atomic.LoadUint64(co.statsInterceptor.rtcpPacketsReceived), + RTCPPacketsSent: atomic.LoadUint64(co.statsInterceptor.rtcpPacketsSent), + } } diff --git a/internal/protocols/webrtc/stats.go b/internal/protocols/webrtc/stats.go new file mode 100644 index 00000000..1cad863c --- /dev/null +++ b/internal/protocols/webrtc/stats.go @@ -0,0 +1,13 @@ +package webrtc + +// Stats are WebRTC statistics. +type Stats struct { + BytesReceived uint64 + BytesSent uint64 + RTPPacketsReceived uint64 + RTPPacketsSent uint64 + RTPPacketsLost uint64 + RTPPacketsJitter float64 + RTCPPacketsReceived uint64 + RTCPPacketsSent uint64 +} diff --git a/internal/protocols/webrtc/stats_interceptor.go b/internal/protocols/webrtc/stats_interceptor.go new file mode 100644 index 00000000..bd754ea4 --- /dev/null +++ b/internal/protocols/webrtc/stats_interceptor.go @@ -0,0 +1,70 @@ +package webrtc + +import ( + "sync/atomic" + + "github.com/pion/interceptor" + "github.com/pion/rtcp" +) + +type statsInterceptor struct { + rtcpPacketsSent *uint64 + rtcpPacketsReceived *uint64 +} + +func (*statsInterceptor) Close() error { + return nil +} + +func (s *statsInterceptor) BindRTCPReader(reader interceptor.RTCPReader) interceptor.RTCPReader { + return interceptor.RTCPReaderFunc(func(bytes []byte, + attributes interceptor.Attributes, + ) (int, interceptor.Attributes, error) { + n, attrs, err := reader.Read(bytes, attributes) + + pkts, err2 := attrs.GetRTCPPackets(bytes) + if err2 == nil { + atomic.AddUint64(s.rtcpPacketsReceived, uint64(len(pkts))) + } + + return n, attrs, err + }) +} + +func (s *statsInterceptor) BindRTCPWriter(writer interceptor.RTCPWriter) interceptor.RTCPWriter { + return interceptor.RTCPWriterFunc(func(pkts []rtcp.Packet, attributes interceptor.Attributes) (int, error) { + atomic.AddUint64(s.rtcpPacketsSent, uint64(len(pkts))) + return writer.Write(pkts, attributes) + }) +} + +func (s *statsInterceptor) BindLocalStream(_ *interceptor.StreamInfo, + writer interceptor.RTPWriter, +) interceptor.RTPWriter { + return writer +} + +func (*statsInterceptor) UnbindLocalStream(_ *interceptor.StreamInfo) {} + +func (s *statsInterceptor) BindRemoteStream(_ *interceptor.StreamInfo, + reader interceptor.RTPReader, +) interceptor.RTPReader { + return reader +} + +func (*statsInterceptor) UnbindRemoteStream(_ *interceptor.StreamInfo) {} + +type statsInterceptorFactory struct { + onCreate func(s *statsInterceptor) +} + +func (f *statsInterceptorFactory) NewInterceptor(_ string) (interceptor.Interceptor, error) { + s := &statsInterceptor{ + rtcpPacketsSent: new(uint64), + rtcpPacketsReceived: new(uint64), + } + + f.onCreate(s) + + return s, nil +} diff --git a/internal/servers/webrtc/session.go b/internal/servers/webrtc/session.go index 92de5121..cead6e84 100644 --- a/internal/servers/webrtc/session.go +++ b/internal/servers/webrtc/session.go @@ -433,15 +433,13 @@ func (s *session) apiItem() *defs.APIWebRTCSession { peerConnectionEstablished := false localCandidate := "" remoteCandidate := "" - bytesReceived := uint64(0) - bytesSent := uint64(0) + var stats *webrtc.Stats if s.pc != nil { peerConnectionEstablished = true localCandidate = s.pc.LocalCandidate() remoteCandidate = s.pc.RemoteCandidate() - bytesReceived = s.pc.BytesReceived() - bytesSent = s.pc.BytesSent() + stats = s.pc.Stats() } return &defs.APIWebRTCSession{ @@ -457,9 +455,15 @@ func (s *session) apiItem() *defs.APIWebRTCSession { } return defs.APIWebRTCSessionStateRead }(), - Path: s.req.pathName, - Query: s.req.httpRequest.URL.RawQuery, - BytesReceived: bytesReceived, - BytesSent: bytesSent, + Path: s.req.pathName, + Query: s.req.httpRequest.URL.RawQuery, + BytesReceived: stats.BytesReceived, + BytesSent: stats.BytesSent, + RTPPacketsReceived: stats.RTPPacketsReceived, + RTPPacketsSent: stats.RTPPacketsSent, + RTPPacketsLost: stats.RTPPacketsLost, + RTPPacketsJitter: stats.RTPPacketsJitter, + RTCPPacketsReceived: stats.RTCPPacketsReceived, + RTCPPacketsSent: stats.RTCPPacketsSent, } }